// Build archives of every channel's compact transcript.cues.json files. // Two output modes share the same staging logic: // // archiveTranscripts() — one . per channel // archiveCombinedTranscripts() — one all-transcripts. total // // Internal layout of .: // /channel.json // //transcript.cues.json // // Internal layout of all-transcripts.: // manifest.json (cross-channel index) // /channel.json (one per channel) // //transcript.cues.json // // Both modes stage hard-link trees under transcripts/export/.staging/ so // the compressor can stream in one pass with no extra disk usage, then // unlink the staging dir on success or failure. When includeMetadata or // prettyPrint is overridden, staging rewrites each cues file instead of // hard-linking it (slower, but bounded). import path from "node:path"; import { access, link, mkdir, readdir, readFile, rm, writeFile, stat, } from "node:fs/promises"; import { execa } from "execa"; import pLimit from "p-limit"; import { listChannelStatsFromDisk, type ChannelStat } from "./channels"; import { openChannelSigner, type ChannelSigner } from "../lib/channelSignature"; import { normalizeTranscript, type NormalizedTranscript, } from "./normalizeTranscript"; import { CUES_JSON_FILENAME } from "../lib/videoStatus"; import { detectPlatform } from "../lib/platform"; import type { Paths } from "../lib/paths"; import { ARCHIVE_FORMATS, DEFAULT_ARCHIVE_OPTIONS, type ArchiveBuildOptions, type ArchiveFormat, } from "../lib/archiveOptions"; export { ARCHIVE_FORMATS, DEFAULT_ARCHIVE_OPTIONS, type ArchiveBuildOptions, type ArchiveFormat, }; export type ArchiveTranscriptsOptions = { paths: Paths; onLog?: (msg: string) => void; signal?: AbortSignal; concurrency?: number; // How many channels' archives to build concurrently. Each channel shells out // to a `zip`/`tar` subprocess, so a few in parallel uses idle cores; kept // modest to avoid thrashing disk on a large corpus. Staging within a channel // stays bounded by `concurrency`. channelConcurrency?: number; channelSlugs?: string[]; build?: Partial; // Destination for the finished archives. Defaults to transcripts/export/ // archives (the standalone editor actions' location); compose-site passes the // served public/archives dir so the zips ship with the site build. outDir?: string; // Reuse an already-open per-channel signer (see channelSignature.ts) so the // caller can share one LMDB view across the transcripts + live-chat passes. // When omitted, one is opened for the duration of this call. signer?: ChannelSigner; // Materialize-only mode for the docker per-site compose: never stage, compress, // or write under the (read-only-mounted) cache. Report the already-cached // archive for each member channel — the host's build:archives pass warmed the // cache first. A member with no cached zip is a channel that legitimately had // nothing to archive (build:archives ran over a superset of these slugs), so it // is reported as skipped, exactly as a fresh generation would omit it. readOnly?: boolean; }; export type ArchiveTranscriptsResult = { archives: { slug: string; archivePath: string; transcriptCount: number }[]; totalTranscripts: number; skipped: { slug: string; reason: string }[]; }; export type ArchiveCombinedResult = { archivePath: string | null; channels: { slug: string; transcriptCount: number }[]; totalTranscripts: number; skipped: { slug: string; reason: string }[]; }; const COMBINED_BASENAME = "all-transcripts"; const COMBINED_STAGING_DIR = "all-transcripts"; function archivesDir(paths: Paths): string { return path.join(paths.transcriptsDir, "export", "archives"); } // The persistent, shared archive cache dir (the default archive outDir). It // lives outside any site's public/ tree and survives across builds, so // signature-gated archives can be reused build-to-build and across sites. // compose-site builds into this cache, then copies each site's member subset // into that site's public/archives. export function archiveCacheDir(paths: Paths): string { return archivesDir(paths); } function stagingDir(paths: Paths): string { return path.join(paths.transcriptsDir, "export", ".staging"); } export function resolveBuild( partial: Partial | undefined, ): ArchiveBuildOptions { const merged = { ...DEFAULT_ARCHIVE_OPTIONS, ...partial }; const lvl = Math.round(merged.compressionLevel); return { ...merged, compressionLevel: Math.min(9, Math.max(1, Number.isFinite(lvl) ? lvl : 6)), }; } function channelManifest(ch: ChannelStat, linkedCount: number) { return { slug: ch.slug, name: ch.config.name ?? null, handling: ch.config.handling, platform: ch.config.platform ?? detectPlatform(ch.config.url ?? null) ?? null, url: ch.config.url ?? null, transcriptCount: linkedCount, generatedAt: new Date().toISOString(), }; } // Hard-link or transform-write every video's transcript.cues.json from // //data/ into //transcript.cues.json. // Normalises on demand so the stage tree is always self-consistent. Returns // the count of staged entries (zero means the channel had no normalisable // transcripts). async function stageChannel( paths: Paths, ch: ChannelStat, destDir: string, build: ArchiveBuildOptions, log: (msg: string) => void, limit: ReturnType, ): Promise<{ linkedCount: number; normalizedCount: number; normalizeFailed: number }> { const dataDir = path.join(paths.channelsDir, ch.slug, "data"); await rm(destDir, { recursive: true, force: true }); await mkdir(destDir, { recursive: true }); let videoIds: string[]; try { videoIds = await readdir(dataDir); } catch { videoIds = []; } log(`Archive ${ch.slug}: scanning ${videoIds.length} videos`); let normalizedCount = 0; let linkedCount = 0; let normalizeFailed = 0; // Hard-linking is only safe when the on-disk file matches what the // archive should contain. If the user toggled metadata off or pretty- // print on, we have to re-encode each entry. const requiresRewrite = !build.includeMetadata || build.prettyPrint; await Promise.all( videoIds.map((id) => limit(async () => { const videoDir = path.join(dataDir, id); try { const s = await stat(videoDir); if (!s.isDirectory()) return; } catch { return; } let outcome; try { outcome = await normalizeTranscript({ videoDir, channelSlug: ch.slug, configName: ch.config.name, }); } catch (err) { normalizeFailed++; log(` ! normalize ${ch.slug}/${id}: ${(err as Error).message}`); return; } if (outcome.status === "skipped") return; if (outcome.status === "wrote") normalizedCount++; const dest = path.join(destDir, id, CUES_JSON_FILENAME); await mkdir(path.dirname(dest), { recursive: true }); if (requiresRewrite) { await writeTransformedCues(outcome.cuesPath, dest, build); } else { await link(outcome.cuesPath, dest); } linkedCount++; }), ), ); if (normalizedCount > 0) { log(` ${ch.slug}: normalized ${normalizedCount} cues files on demand`); } if (normalizeFailed > 0) { log(` ${ch.slug}: ${normalizeFailed} videos failed to normalize`); } if (linkedCount > 0) { await writeFile( path.join(destDir, "channel.json"), JSON.stringify(channelManifest(ch, linkedCount), null, 2) + "\n", ); } return { linkedCount, normalizedCount, normalizeFailed }; } export async function writeTransformedCues( sourcePath: string, destPath: string, build: ArchiveBuildOptions, ): Promise { const raw = await readFile(sourcePath, "utf8"); const parsed = JSON.parse(raw) as NormalizedTranscript; const payload = build.includeMetadata ? parsed : { version: parsed.version, source: parsed.source, cues: parsed.cues, }; const body = build.prettyPrint ? JSON.stringify(payload, null, 2) + "\n" : JSON.stringify(payload); await writeFile(destPath, body); } export function archiveExtension(format: ArchiveFormat): string { return format; } // --- Incremental cache: skip re-zipping a channel whose content is unchanged -- // Each produced archive gets a sibling `.sig` sidecar recording the // channel content signature it was built from (plus the entry count, which the // Downloads manifest needs and which we'd otherwise have to recompute on a // cache hit). A build reuses the existing zip when the current signature matches // and the zip is still present. export function archiveSidecarPath(archivePath: string): string { return `${archivePath}.sig`; } export async function readArchiveSidecar( archivePath: string, ): Promise<{ signature: string; count: number } | null> { try { const parsed = JSON.parse( await readFile(archiveSidecarPath(archivePath), "utf8"), ); if ( parsed && typeof parsed.signature === "string" && typeof parsed.count === "number" ) { return { signature: parsed.signature, count: parsed.count }; } return null; } catch { return null; } } export async function writeArchiveSidecar( archivePath: string, signature: string, count: number, ): Promise { await writeFile( archiveSidecarPath(archivePath), JSON.stringify({ signature, count }) + "\n", ); } async function fileExists(p: string): Promise { try { await access(p); return true; } catch { return false; } } // The non-mtime inputs that also determine a channel's archive bytes: the build // options and the config fields baked into channel.json. Folded into the // signature so a settings/config change forces a rebuild even when no video // mtime moved. export function archiveSignatureExtra( build: ArchiveBuildOptions, ch: ChannelStat, ): string { return JSON.stringify({ format: build.format, level: build.compressionLevel, meta: build.includeMetadata, pretty: build.prettyPrint, cfg: { name: ch.config.name ?? null, handling: ch.config.handling, platform: ch.config.platform ?? null, url: ch.config.url ?? null, }, }); } // Compress a single staged top-level entry (a slug dir for per-channel // mode, or one of multiple entries for combined mode) into archivePath. export async function compressArchive( archivePath: string, cwd: string, members: string[], build: ArchiveBuildOptions, ): Promise { const level = String(build.compressionLevel); if (build.format === "zip") { // `zip -r -` walks each member dir recursively. -X strips // uid/gid/extras so output is reproducible. Overwrite via -f? No — // we delete archivePath first so zip doesn't append to a stale file. await rm(archivePath, { force: true }); await execa( "zip", ["-r", `-${level}`, "-X", "-q", archivePath, ...members], { cwd, stdio: "ignore" }, ); return; } // tar.gz / tar.xz: route compression through --use-compress-program so we // can pass the level flag. const compressor = build.format === "tar.xz" ? `xz -${level} -T 0 -c` : `gzip -${level} -c`; await execa( "tar", [ "--use-compress-program", compressor, "-cf", archivePath, "--owner=0", "--group=0", "--mtime=@0", "-C", cwd, ...members, ], { stdio: "ignore" }, ); } export async function archiveTranscripts( opts: ArchiveTranscriptsOptions, ): Promise { const build = resolveBuild(opts.build); const log = opts.onLog ?? ((m: string) => console.log(m)); const limit = pLimit(opts.concurrency ?? 8); const archives: ArchiveTranscriptsResult["archives"] = []; const skipped: ArchiveTranscriptsResult["skipped"] = []; let totalTranscripts = 0; log( `Format: ${build.format} level=${build.compressionLevel} metadata=${build.includeMetadata} pretty=${build.prettyPrint}`, ); const outDir = opts.outDir ?? archivesDir(opts.paths); // In readOnly mode outDir is the shared cache on a read-only mount — never // create dirs or touch staging there. if (!opts.readOnly) { await mkdir(outDir, { recursive: true }); await mkdir(stagingDir(opts.paths), { recursive: true }); } const allChannels = await listChannelStatsFromDisk(opts.paths); const target = opts.channelSlugs ? allChannels.filter((c) => opts.channelSlugs!.includes(c.slug)) : allChannels; const signer = opts.readOnly ? null : (opts.signer ?? openChannelSigner(opts.paths)); const channelLimit = pLimit(opts.channelConcurrency ?? 4); let reused = 0; // Build (or reuse) one channel's archive. Returns exactly one of archive / // skipped so the caller can reduce results deterministically. const buildOne = async ( ch: ChannelStat, ): Promise<{ archive?: ArchiveTranscriptsResult["archives"][number]; skipped?: ArchiveTranscriptsResult["skipped"][number]; reused?: boolean; }> => { if (opts.signal?.aborted) return {}; const archivePath = path.join( outDir, `${ch.slug}.${archiveExtension(build.format)}`, ); // Materialize-only: consume the pre-warmed cache, never generate. if (opts.readOnly) { if (await fileExists(archivePath)) { const prev = await readArchiveSidecar(archivePath); return { archive: { slug: ch.slug, archivePath, transcriptCount: prev?.count ?? 0, }, reused: true, }; } return { skipped: { slug: ch.slug, reason: "no cached archive" } }; } const sig = signer!.signature(ch.slug, archiveSignatureExtra(build, ch)); if (sig) { const prev = await readArchiveSidecar(archivePath); if (prev && prev.signature === sig && (await fileExists(archivePath))) { log(` ${ch.slug}: unchanged — reusing cached archive`); return { archive: { slug: ch.slug, archivePath, transcriptCount: prev.count }, reused: true, }; } } const staging = path.join(stagingDir(opts.paths), ch.slug); try { const { linkedCount } = await stageChannel( opts.paths, ch, staging, build, log, limit, ); if (linkedCount === 0) { await rm(staging, { recursive: true, force: true }); // A channel that lost all transcripts must stop being served: drop any // stale cached zip + sidecar so it isn't reused next build. await rm(archivePath, { force: true }).catch(() => {}); await rm(archiveSidecarPath(archivePath), { force: true }).catch(() => {}); log(` ${ch.slug}: skipped (no transcripts to archive)`); return { skipped: { slug: ch.slug, reason: "no normalized transcripts" } }; } try { await compressArchive( archivePath, stagingDir(opts.paths), [ch.slug], build, ); } finally { await rm(staging, { recursive: true, force: true }); } if (sig) await writeArchiveSidecar(archivePath, sig, linkedCount); log( ` ${ch.slug}: ${linkedCount} transcripts → ${path.basename(archivePath)}`, ); return { archive: { slug: ch.slug, archivePath, transcriptCount: linkedCount }, }; } catch (err) { await rm(staging, { recursive: true, force: true }).catch(() => {}); log(` ! ${ch.slug}: ${(err as Error).message}`); return { skipped: { slug: ch.slug, reason: (err as Error).message } }; } }; try { const results = await Promise.all( target.map((ch) => channelLimit(() => buildOne(ch))), ); for (const r of results) { if (r.archive) { archives.push(r.archive); totalTranscripts += r.archive.transcriptCount; if (r.reused) reused++; } else if (r.skipped) { skipped.push(r.skipped); } } } finally { if (!opts.signer && signer) await signer.close(); } if (reused > 0) log(`Reused ${reused} unchanged cached archive(s).`); log(""); log( `Wrote ${archives.length} archive${archives.length === 1 ? "" : "s"} containing ${totalTranscripts} transcripts to:`, ); log(` ${outDir}`); for (const a of archives) { log(` - ${path.basename(a.archivePath)} (${a.transcriptCount} transcripts)`); } if (skipped.length > 0) { log(""); log(`Skipped ${skipped.length} channel(s):`); for (const s of skipped) log(` - ${s.slug}: ${s.reason}`); } return { archives, totalTranscripts, skipped }; } export async function archiveCombinedTranscripts( opts: ArchiveTranscriptsOptions, ): Promise { const build = resolveBuild(opts.build); const log = opts.onLog ?? ((m: string) => console.log(m)); const limit = pLimit(opts.concurrency ?? 8); const skipped: ArchiveCombinedResult["skipped"] = []; const channels: ArchiveCombinedResult["channels"] = []; let totalTranscripts = 0; log( `Format: ${build.format} level=${build.compressionLevel} metadata=${build.includeMetadata} pretty=${build.prettyPrint}`, ); const outDir = opts.outDir ?? archivesDir(opts.paths); await mkdir(outDir, { recursive: true }); await mkdir(stagingDir(opts.paths), { recursive: true }); const combinedStaging = path.join(stagingDir(opts.paths), COMBINED_STAGING_DIR); const archivePath = path.join( outDir, `${COMBINED_BASENAME}.${archiveExtension(build.format)}`, ); await rm(combinedStaging, { recursive: true, force: true }); await mkdir(combinedStaging, { recursive: true }); try { const allChannels = await listChannelStatsFromDisk(opts.paths); const target = opts.channelSlugs ? allChannels.filter((c) => opts.channelSlugs!.includes(c.slug)) : allChannels; for (const ch of target) { if (opts.signal?.aborted) { log("Aborted."); break; } const dest = path.join(combinedStaging, ch.slug); try { const { linkedCount } = await stageChannel( opts.paths, ch, dest, build, log, limit, ); if (linkedCount === 0) { await rm(dest, { recursive: true, force: true }); skipped.push({ slug: ch.slug, reason: "no normalized transcripts" }); log(` ${ch.slug}: skipped (no transcripts to archive)`); continue; } channels.push({ slug: ch.slug, transcriptCount: linkedCount }); totalTranscripts += linkedCount; } catch (err) { await rm(dest, { recursive: true, force: true }).catch(() => {}); skipped.push({ slug: ch.slug, reason: (err as Error).message }); log(` ! ${ch.slug}: ${(err as Error).message}`); } } if (channels.length === 0) { log("No channels had normalisable transcripts; nothing to archive."); return { archivePath: null, channels, totalTranscripts, skipped, }; } await writeFile( path.join(combinedStaging, "manifest.json"), JSON.stringify( { version: 1, generatedAt: new Date().toISOString(), channelCount: channels.length, transcriptCount: totalTranscripts, channels, }, null, 2, ) + "\n", ); const tarMembers = ["manifest.json", ...channels.map((c) => c.slug)]; log( `Compressing ${tarMembers.length - 1} channels into ${path.basename(archivePath)}…`, ); await compressArchive(archivePath, combinedStaging, tarMembers, build); log(""); log( `Wrote ${path.basename(archivePath)} containing ${totalTranscripts} transcripts across ${channels.length} channels to:`, ); log(` ${archivePath}`); if (skipped.length > 0) { log(""); log(`Skipped ${skipped.length} channel(s):`); for (const s of skipped) log(` - ${s.slug}: ${s.reason}`); } return { archivePath, channels, totalTranscripts, skipped, }; } finally { await rm(combinedStaging, { recursive: true, force: true }); } }