// Build archives of every channel's parsed live_chat.cues.json files. // Mirrors archiveTranscripts but produces a separate archive so the // transcript archives stay scoped to spoken-content cues. // // archiveLiveChat() — one .live_chat. per channel // archiveCombinedLiveChat() — one all-live-chat. total // // Internal layout of .live_chat.: // /channel.json // //live_chat.cues.json // // Internal layout of all-live-chat.: // manifest.json // /channel.json // //live_chat.cues.json import path from "node:path"; import { access, link, lstat, mkdir, readdir, rm, writeFile, stat, } from "node:fs/promises"; import pLimit from "p-limit"; import { listChannelStatsFromDisk, type ChannelStat } from "./channels"; import { isLiveChatCuesFresh, normalizeLiveChat } from "./normalizeLiveChat"; import { inspectChannelMedia } from "../lib/channelMedia"; import { LIVE_CHAT_MEDIA_FILENAME } from "../lib/mediaTier"; import { archiveExtension, archiveSidecarPath, archiveSignatureExtra, compressArchive, readArchiveSidecar, resolveBuild, writeArchiveSidecar, writeTransformedCues, } from "./archiveTranscripts"; import { openChannelSigner, type ChannelSigner } from "../lib/channelSignature"; import { LIVE_CHAT_CUES_FILENAME } from "../lib/videoStatus"; import { detectPlatform } from "../lib/platform"; import type { Paths } from "../lib/paths"; import { type ArchiveBuildOptions, } from "../lib/archiveOptions"; async function fileExists(p: string): Promise { try { await access(p); return true; } catch { return false; } } export type ArchiveLiveChatOptions = { paths: Paths; onLog?: (msg: string) => void; signal?: AbortSignal; concurrency?: number; // Channels built concurrently (see ArchiveTranscriptsOptions.channelConcurrency). channelConcurrency?: number; channelSlugs?: string[]; build?: Partial; // Destination for the finished archives. Defaults to transcripts/export/ // archives; compose-site passes the served public/archives dir. Mirrors // ArchiveTranscriptsOptions.outDir. outDir?: string; // Reuse an already-open per-channel signer across passes (see // ArchiveTranscriptsOptions.signer). signer?: ChannelSigner; // Materialize-only mode for the docker per-site compose (see // ArchiveTranscriptsOptions.readOnly). readOnly?: boolean; }; export type ArchiveLiveChatResult = { archives: { slug: string; archivePath: string; liveChatCount: number }[]; totalLiveChats: number; skipped: { slug: string; reason: string }[]; }; export type ArchiveCombinedLiveChatResult = { archivePath: string | null; channels: { slug: string; liveChatCount: number }[]; totalLiveChats: number; skipped: { slug: string; reason: string }[]; }; const COMBINED_BASENAME = "all-live-chat"; const COMBINED_STAGING_DIR = "all-live-chat"; const ARCHIVE_INFIX = "live_chat"; function archivesDir(paths: Paths): string { return path.join(paths.transcriptsDir, "export", "archives"); } // Keep live-chat staging isolated from transcript staging so concurrent runs // (e.g. on separate queues) can't trample each other. function stagingDir(paths: Paths): string { return path.join(paths.transcriptsDir, "export", ".staging", "live-chat"); } function channelManifest(ch: ChannelStat, liveChatCount: 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, liveChatCount, generatedAt: new Date().toISOString(), }; } // Hard-link or transform-write every video's live_chat.cues.json from // //data/ into //live_chat.cues.json. // Normalizes on demand so the stage tree is always self-consistent. 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 live chat ${ch.slug}: scanning ${videoIds.length} videos`); // THE RAW REPLAY IS MEDIA (release 17): `transcript.live_chat.json` is a // tiered link onto the channel's media drive. A build re-normalizes a stale // cues file from it only while that drive is reachable; otherwise the cues // the corpus disk holds are staged as they are (the freshness check `lstat`s // the link, so asking it costs the drive nothing), and a video with no cues // yet is left out of this build. Never a read of a drive that is unmounted, // stalled or mid-move. const media = await inspectChannelMedia(paths, ch.slug, ch.config, { fresh: true, }); const rawReadable = media.status === "ok" || media.status === "in-place"; let normalizedCount = 0; let linkedCount = 0; let normalizeFailed = 0; let keptStale = 0; // Videos with a raw replay (a link onto the away drive, lstat'd) and no cues // on disk yet: left out of this build, and COUNTED — a drive outage must not // shrink a site silently (review L4). let leftOut = 0; 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 { if (rawReadable) { outcome = await normalizeLiveChat({ videoDir, channelSlug: ch.slug, configName: ch.config.name, // A build moves no media (release 17): the jobs tier, this reads. tier: false, }); } else { const { fresh, cuesPath } = await isLiveChatCuesFresh(videoDir); const have = await stat(cuesPath).then( (st) => st.isFile(), () => false, ); if (!have) { const hasRaw = await lstat(path.join(videoDir, LIVE_CHAT_MEDIA_FILENAME)).then( () => true, () => false, ); if (hasRaw) leftOut++; return; } if (!fresh) keptStale++; outcome = { status: "fresh" as const, cuesPath }; } } 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, LIVE_CHAT_CUES_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} live-chat cues files on demand`); } if (normalizeFailed > 0) { log(` ${ch.slug}: ${normalizeFailed} videos failed to normalize`); } if (!rawReadable) { log( ` ${ch.slug}: its media is not reachable (${media.detail ?? media.status}) — ` + `staged the live-chat cues on disk as they are` + (keptStale > 0 ? ` (${keptStale} older than their raw replay)` : "") + `; ${leftOut} video(s) with a raw replay and no cues left out of this build`, ); } 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 archiveLiveChat( opts: ArchiveLiveChatOptions, ): Promise { const build = resolveBuild(opts.build); const log = opts.onLog ?? ((m: string) => console.log(m)); const limit = pLimit(opts.concurrency ?? 8); const archives: ArchiveLiveChatResult["archives"] = []; const skipped: ArchiveLiveChatResult["skipped"] = []; let totalLiveChats = 0; log( `Format: ${build.format} level=${build.compressionLevel} metadata=${build.includeMetadata} pretty=${build.prettyPrint}`, ); const outDir = opts.outDir ?? archivesDir(opts.paths); 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; const buildOne = async ( ch: ChannelStat, ): Promise<{ archive?: ArchiveLiveChatResult["archives"][number]; skipped?: ArchiveLiveChatResult["skipped"][number]; reused?: boolean; }> => { if (opts.signal?.aborted) return {}; const archivePath = path.join( outDir, `${ch.slug}.${ARCHIVE_INFIX}.${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, liveChatCount: prev?.count ?? 0, }, reused: true, }; } return { skipped: { slug: ch.slug, reason: "no cached live-chat archive" } }; } // `live-chat:` marks this signature as the live-chat variant so it never // collides with the transcripts archive's sidecar (different files anyway, // but keeps the two kinds' signatures independent). const sig = signer!.signature( ch.slug, `live-chat:${archiveSignatureExtra(build, ch)}`, ); if (sig) { const prev = await readArchiveSidecar(archivePath); if (prev && prev.signature === sig && (await fileExists(archivePath))) { log(` ${ch.slug}: unchanged — reusing cached live-chat archive`); return { archive: { slug: ch.slug, archivePath, liveChatCount: 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 }); await rm(archivePath, { force: true }).catch(() => {}); await rm(archiveSidecarPath(archivePath), { force: true }).catch(() => {}); log(` ${ch.slug}: skipped (no live chat to archive)`); return { skipped: { slug: ch.slug, reason: "no normalized live chat" } }; } 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} live chats → ${path.basename(archivePath)}`, ); return { archive: { slug: ch.slug, archivePath, liveChatCount: 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); totalLiveChats += r.archive.liveChatCount; 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 live-chat archive(s).`); log(""); log( `Wrote ${archives.length} archive${archives.length === 1 ? "" : "s"} containing ${totalLiveChats} live chats to:`, ); log(` ${outDir}`); for (const a of archives) { log(` - ${path.basename(a.archivePath)} (${a.liveChatCount} live chats)`); } if (skipped.length > 0) { log(""); log(`Skipped ${skipped.length} channel(s):`); for (const s of skipped) log(` - ${s.slug}: ${s.reason}`); } return { archives, totalLiveChats, skipped }; } export async function archiveCombinedLiveChat( opts: ArchiveLiveChatOptions, ): Promise { const build = resolveBuild(opts.build); const log = opts.onLog ?? ((m: string) => console.log(m)); const limit = pLimit(opts.concurrency ?? 8); const skipped: ArchiveCombinedLiveChatResult["skipped"] = []; const channels: ArchiveCombinedLiveChatResult["channels"] = []; let totalLiveChats = 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 live chat" }); log(` ${ch.slug}: skipped (no live chat to archive)`); continue; } channels.push({ slug: ch.slug, liveChatCount: linkedCount }); totalLiveChats += 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 live chat; nothing to archive."); return { archivePath: null, channels, totalLiveChats, skipped, }; } await writeFile( path.join(combinedStaging, "manifest.json"), JSON.stringify( { version: 1, generatedAt: new Date().toISOString(), channelCount: channels.length, liveChatCount: totalLiveChats, 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 ${totalLiveChats} live chats 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, totalLiveChats, skipped, }; } finally { await rm(combinedStaging, { recursive: true, force: true }); } }