// Preprocess /channels//data// into: // - LMDB cache at /index.mdb (for incremental rebuilds) // - /{manifest,page-NNNN}.json (paginated cross-channel // summaries index for browse/search) // - //{manifest,page-NNNN}.json // (per-channel paginated transcript detail + cues, oldest-first so // adding a newer video only dirties the last page) // // Per-channel config.json selects the transcript parser ("youtube" → VTT, // "transcribe" → whisper.cpp JSON). Short-circuits when mtimes already match. // // A CHANNEL WHOSE MEDIA IS NOT REACHABLE is HELD, not emptied: a relocated // `data/` on an unmounted drive, one mid-relocation, a link and a config that // disagree (inspectChannelMedia), or a data dir that fails to read. It is not // rescanned; its index records are kept as they are and its shared page trees // (transcripts, subs, digests) are left as they are, so the next site build // still publishes it. Until this hold the scan read such a channel as having // no videos, removed every record it had, and the site built next published // the channel as gone. The stats build has the same hold (buildStats.ts); the // words are shared (lib/channelMediaHold.ts). // // A FULL REBUILD (a schema change, or no index yet) with a channel held // REFUSES: it clears every channel's records, and a held channel cannot be // re-read, so it would come out empty. ARCHILYZER_INDEX_ALLOW_HELD=1 lets it // proceed; the held channel is then out of the index until its media is back // and the index is built again. import path from "node:path"; import { createHash } from "node:crypto"; import { lstat, mkdir, readdir, readFile, rename, rm, stat, writeFile, } from "node:fs/promises"; import { createWriteStream } from "node:fs"; import type { Dirent, WriteStream } from "node:fs"; import { once } from "node:events"; import { open } from "lmdb"; import { writeJsonAtomic as writeJsonAtomicShared } from "../lib/jsonFile-server"; import { parseVtt, type Cue } from "../lib/vtt"; import { parseTranscriptJson } from "../lib/whisper"; import { parseLiveChat } from "../lib/liveChat"; import { platformLabelStale, summarize, toDisplaySummary, type RawMetadata, } from "../lib/transcripts-server"; import { isCuesJsonFresh, readNormalizedTranscript, } from "./normalizeTranscript"; import { isLiveChatCuesFresh, readNormalizedLiveChat, } from "./normalizeLiveChat"; import type { TranscriptSummary, TranscriptDetail, DisplaySummary, } from "../lib/transcripts"; import type { StoredSubs, SubsDetail } from "../lib/subs"; import type { TrackFields } from "../lib/captionTracks"; import { readTrackFields } from "../lib/captionTracks-server"; import type { Manifest, ChannelEntry, ChannelTranscriptsManifest, ChannelSubsManifest, SubsChannelEntry, SubsManifest, } from "../lib/manifest"; import { MANIFEST_VERSION, SUMMARIES_PAGE_SIZE, TRANSCRIPTS_MANIFEST_VERSION, SUBS_CHANNEL_MANIFEST_VERSION, SUBS_MANIFEST_VERSION, pageFileName, } from "../lib/manifest"; import { getSettings } from "../lib/settings"; import { INDEX_SCANNED_AT_KEY } from "../lib/stats"; import { listSites, siteDigestsDir, siteIndexDir, sitePostsDir, siteSummariesDir, siteSubsDir, } from "../lib/site"; import { type ChannelConfig, type ChannelHandling, } from "../lib/channelConfig"; import { readChannelConfigFile } from "./channels"; import { inspectChannelMedia } from "../lib/channelMedia"; import { isTierable } from "../lib/mediaTier"; import { onDrive } from "../lib/storageHealth"; import { HELD_WAYS_OUT, describeHeld, heldReason, } from "../lib/channelMediaHold"; import { resolveChannelGroupId } from "../lib/channelGroups"; import type { Paths } from "../lib/paths"; import { CAPTION_TRACK_RULE_VERSION, captionInputs, englishVttsByPreference, pickIndexTranscript, readEnglishVttCues, readSubTracks, readVideoFiles, type IndexTranscript, type SubTrack, } from "../lib/videoStatus"; import { resolveEffectiveAvailability, resolveMaybeMissingState, } from "../lib/availability-server"; import { AVAILABILITY_FILENAME, stateFromAvailability, type Availability, type VideoState, } from "../lib/availability"; import { loadMaybeMissing } from "./quickAvailabilityCheck"; import { POSTS_MANIFEST_VERSION, isPostGone, SITE_POSTS_MANIFEST_VERSION, comparePostsNewestFirst, postsPageFileName, type ChannelPostsManifest, type Post, type PostPlatform, type PostsChannelEntry, type PostsManifest, } from "../lib/posts"; import { channelPostsDir, listPostShards, readPostShard, readPostAvailability, postsAvailabilityPath, } from "../lib/posts-server"; import { isSocialChannel } from "../lib/channelConfig"; import { publishedMemberSlugs } from "../lib/postsVisibility"; import { DIGESTS_MANIFEST_VERSION, SITE_DIGESTS_MANIFEST_VERSION, digestPageFileName, newestGeneratedAt, type ChannelDigestsManifest, type DigestsChannelEntry, type DigestsManifest, type VideoDigest, } from "../lib/digests"; import { DIGEST_FILENAME, DIGEST_OVERRIDES_FILENAME, effectiveDigest } from "../lib/digest"; import { loadDigest, loadDigestOverrides } from "../lib/digest-server"; import { TAG_COUNTS_FILENAME, applyCuratedTagsToSummary, clearCuratedPagesPending, createTagCounter, deriveCuratedTags, loadCuratedTagsRuntime, reapplyCuratedTags, } from "./curatedTagsIndex"; import { tagsForSite } from "../lib/curatedTags"; import { compileRecordedDateRule, deriveRecordedDate } from "../lib/recordedDate"; // v10: multi-site build. Shared per-channel transcript/subs pages are written // once; per-site summaries + subs manifests are filtered selections. Bumped to // force a clean rebuild into the new shared/ + sites/ staging layout. // v11: transcript pages now carry `tags` (for the "tags" search scope) — force // a re-extract so existing pages re-emit with the field. // v12: the social-post corpus. A `posts` sub-DB keyed [createdAt, channelSlug, // id] (ISO-8601 sorts correctly, unlike the [uploadDate, …] tuple videos use) // plus a shared /posts// page tree. // v13: the AI-digest corpus. A `digests` sub-DB holding the COMPOSED digest // (effectiveDigest of the machine sidecar + the human override file) plus a // shared /digests// page tree. Bumped so existing indexes populate the // new sub-DB — nothing re-derives it lazily. // NOT a v14 — curated tags (`curatedTags` on every summary, derived from // transcripts/tags.json) deliberately did NOT bump this, and the reason is // worth keeping: every prior bump existed because something would otherwise // never be derived at all (v13 populates a new sub-DB; nothing re-derives it // lazily). Curated tags DO re-derive lazily, by design. `curatedRulesHash` is // absent on an index built before them, so it can never equal the current hash // and `reapplyCuratedTags` re-derives every record on the first build after // this ships — measured at ~80 s across 30k videos, straight out of LMDB, and // zero page writes on an untagged corpus because an empty derivation leaves // each summary byte-identical. A bump would instead wipe the cache and re-read // 76k video directories to reach exactly the same state. No sub-DB was added, // so the clearAsync() enumeration above is unchanged too. const SCHEMA_VERSION = 13; // PLATFORM LABELS — also NOT a schema bump. When the app learns a platform // whose records may already be indexed under another label (BitChute's were // summarized "youtube", with no file to play), a bump would wipe the cache and // re-read every video directory to fix a handful. Instead, the first build // that sees a new PLATFORM_LABELS_VERSION reads every stored summary out of // LMDB once (no disk), queues each one platformLabelStale() flags as changed, // and records the version — only when no channel is held, so a held channel's // records are still looked at once its text is back. Bump this whenever // platformFromMetadata learns such a platform. const PLATFORM_LABELS_VERSION = 1; const PLATFORM_LABELS_KEY = "platformLabels"; // CAPTION TRACK — the same one-shot shape, for the caption-track rule // (CAPTION_TRACK_RULE_VERSION, lib/videoStatus.ts). When the rule changes, the // first build that sees the new version re-reads, from disk, every caption // record the change can reach — one with more than one English VTT (the rule // may now pick another), or whose stored cues are empty (the cue-block shape // used to parse to nothing) — and nothing else; then records the version, // again only when no channel is held. Each re-read is reported per channel: // how many now read different text, how many had none and now do. const CAPTION_TRACK_KEY = "captionTrackRule"; // RECORDED DATES — per channel, the same "re-derive what the change reaches" // shape: `recordedDateRule:` holds the `recordedDate.titlePattern` the // channel's records were last derived under (absent: none). A build that finds // another re-processes every record of that channel, then records the // pattern. const RECORDED_DATE_RULE_KEY_PREFIX = "recordedDateRule:"; // A cue list's text, for "did the words change" — timing alone is not a // different transcript. function cueText(list: readonly Cue[]): string { return list.map((c) => c.text).join("\n"); } // Whether the cue list stored under a key is empty or absent, WITHOUT decoding // it: a non-empty list of cues is far longer than the few bytes an empty // msgpack array takes, and decoding every caption record of the corpus to ask // this would cost the full read the pass exists to avoid. function storedCuesEmpty(db: { getBinaryFast(key: IndexKey): Buffer | undefined }, key: IndexKey): boolean { const raw = db.getBinaryFast(key); return raw === undefined || raw.length <= 4; } // ALTERNATE TRACKS — the same one-shot shape again (lib/captionTracks.ts). A // record's other English tracks, kept where their words differ from the // primary's, live in the `alts` sub-DB and ride on its transcript page // (`track` + `altTracks`). The first build that sees a new ALT_TRACKS_VERSION // re-reads every record that CAN hold one — a caption record with two or more // English VTTs, a transcribed record with any — and nothing else; then records // the version, only when no channel is held. Bump it when what an alternate is // changes. const ALT_TRACKS_VERSION = 1; const ALT_TRACKS_KEY = "altTracks"; // Whether a scanned record can hold an alternate track at all — no file read. function canHoldAltTracks(s: { transcriptKind: IndexTranscript["kind"] | null; englishVttCount: number; }): boolean { return ( (s.transcriptKind === "vtt" && s.englishVttCount >= 2) || (s.transcriptKind === "whisper" && s.englishVttCount >= 1) ); } export type CaptionTrackChannelReport = { reread: number; // Re-read records whose caption text is not what the index held. changed: number; // Of those, the ones the index held no text for. zeroToText: number; }; // Per-channel post stats, persisted so per-site aggregates survive a no-op // rebuild that doesn't re-encode the post pages. Mirrors ChannelSubsStat. type ChannelPostsStat = { name: string; postCount: number; platform: string; // Signature of the channel's on-disk posts dir, so an unchanged channel // skips re-encoding. Same idea as the mtime records the video scan keeps. signature: string; }; // Per-channel digest stats, persisted so the per-site digests manifest survives // a no-op rebuild that doesn't re-encode the digest pages. Mirrors // ChannelSubsStat / ChannelPostsStat. type ChannelDigestStat = { name: string; digestCount: number; }; // Per-channel subtitle stats, collected while writing the shared subs pages and // persisted to LMDB so per-site subs manifests can be assembled on a no-op // rebuild without re-encoding every channel's pages. type ChannelSubsStat = { name: string; videoCount: number; tracks: string[]; liveChatCount: number; }; type IndexKey = [string, string, string]; type ChannelKey = [string, string, string]; type PageHashKey = [string, number]; // [createdAt, channelSlug, postId] — see the `posts` sub-DB comment below. type PostIndexKey = [string, string, string]; type PathKey = [string, string]; type MtimeRecord = { metaMs: number; transcriptMs: number | null; subsMs: number | null; availabilityMs: number | null; // Newest mtime across BOTH digest sidecars — ai-digest.json (machine) and // ai-digest.overrides.json (human). Max-of-two, not a single stat, because a // human correction only ever touches the overrides file: keying off the // machine file alone would leave a corrected digest producing no index // mutation, so the correction would never ship. That is the entire reason // the overrides file exists. See the subsMs multi-file loop it follows. digestMs: number | null; // The effective availability at the last index of this video (availability // .json reconciled with download-outcome.json). Absent on records written // before the VideoState work, which is why the two booleans below stay as // the fallback (see stateFromRecord). Not worth a SCHEMA_VERSION bump // — that wipes the whole cache and re-pages the entire corpus, whereas // availabilityMs already participates in change detection, so any video whose // availability.json is rewritten from here on picks up the exact string. availability?: Availability; isDeleted: boolean; isUnlisted: boolean; indexKey: IndexKey; }; // The probe result an MtimeRecord carries, preferring the exact availability // string and falling back to the pre-VideoState booleans. function stateFromRecord(rec: MtimeRecord): VideoState { if (rec.availability) return stateFromAvailability(rec.availability); if (rec.isDeleted) return "deleted"; if (rec.isUnlisted) return "unlisted"; return "available"; } function indexKeyId(k: IndexKey): string { return `${k[0]}\x00${k[1]}\x00${k[2]}`; } type PageHashRecord = { hash: string; entryCount: number; sizeBytes: number; }; function indexToChannelKey(k: IndexKey): ChannelKey { return [k[1], k[0], k[2]]; } type LiveEntry = { channelSlug: string; handling: ChannelHandling; configName: string | undefined; videoDir: string; metaPath: string; metaMs: number; transcriptPath: string; transcriptMs: number | null; transcriptKind: IndexTranscript["kind"] | null; // How many English VTTs the dir holds (the caption-track pass's question). englishVttCount: number; subTracks: SubTrack[]; subsMs: number | null; availabilityMs: number | null; digestMs: number | null; }; async function exists(p: string): Promise { try { await stat(p); return true; } catch { return false; } } function errCode(err: unknown): string { const code = (err as NodeJS.ErrnoException | null)?.code; return typeof code === "string" ? code : String(err); } // `held` maps each channel whose media could not be read to why, in words with // no path in them. A held channel contributes no live entries; the caller keeps // its records and pages (see the file header). async function scanSource( channelsDir: string, log: (msg: string) => void, ): Promise<{ live: LiveEntry[]; channels: Map; held: Map; mediaAccess: Map; }> { const locations = getSettings().storage.locations; const channels = new Map(); const live: LiveEntry[] = []; const held = new Map(); const mediaAccess = new Map(); let channelEntries: Dirent[]; try { channelEntries = await readdir(channelsDir, { withFileTypes: true }); } catch { // Fresh transcripts dir with no channels yet. return { live, channels, held, mediaAccess }; } for (const ch of channelEntries) { if (!ch.isDirectory()) continue; const channelDir = path.join(channelsDir, ch.name); const cfg = await readChannelConfigFile(path.join(channelDir, "config.json")); if (!cfg) { log(`Skipping channel ${ch.name}: missing or invalid config.json`); continue; } if (cfg.excludeFromBuild) { log(`Skipping channel ${ch.name}: excluded from build`); continue; } channels.set(ch.name, cfg); // A social channel carries posts, not videos: it has no data/ dir to walk, // and routing it through the video scan would only ever produce noise. Its // posts tree is built from the JSONL shards further down. if (isSocialChannel(cfg)) continue; // An unreadable text tier is not an empty channel (lib/channelMedia.ts): // the readdir below would fail, and every record the channel has would be // removed as gone. // // THE TEXT GUARD, NOT THE MEDIA ONE (release 17): the index reads the text // tier only — presence of a media file is a name in the one readdir, and // every stat is a sidecar — so a moving, stalled or unmounted MEDIA drive // never holds it. Held: a `legacy` channel (its text is on the far drive), // a `data/` that is not a directory, a tier migration in flight. const media = await inspectChannelMedia({ channelsDir }, ch.name, cfg, { fresh: true, }); if (!media.text.readable) { held.set(ch.name, heldReason(media, cfg.mediaDir ?? cfg.dataDir, locations)); continue; } mediaAccess.set(ch.name, { readable: media.status === "ok" || media.status === "in-place", drive: cfg.mediaDir?.trim() || undefined, }); const dataDir = path.join(channelDir, "data"); let videoEntries: Dirent[]; try { videoEntries = await readdir(dataDir, { withFileTypes: true }); } catch (err) { const code = errCode(err); // No data/ at all: the text tier is on the corpus disk, so the channel // has downloaded nothing yet (or its files were deleted), and it IS empty. if (code === "ENOENT") { log(`Channel ${ch.name}: no data/ directory; indexed as a channel with no videos.`); continue; } // Anything else — a relocated drive gone between the check and the // read, a permission or I/O error — is a channel that could not be read. held.set(ch.name, `its data directory could not be read (${code})`); continue; } const channelLive: LiveEntry[] = []; let readFailure: string | null = null; for (const v of videoEntries) { if (!v.isDirectory()) continue; const videoDir = v.name; const fullVideoDir = path.join(dataDir, videoDir); const metaPath = path.join(fullVideoDir, "metadata.info.json"); let metaMs: number; try { metaMs = (await stat(metaPath)).mtimeMs; } catch (err) { // ENOENT is a video dir with no metadata yet (a download in flight, a // partial one): skipped, as always. Any other error is the channel's // media failing mid-scan; the channel is held below rather than read // as missing this video and every one after it. const code = errCode(err); if (code !== "ENOENT" && code !== "ENOTDIR") { readFailure = code; break; } continue; } // Hybrid: a single channel may contain both YouTube auto-subs (.vtt) // and whisper-generated transcript.json. Prefer whisper when both // exist; fall back to .vtt otherwise. const files = await readVideoFiles(fullVideoDir); const picked = pickIndexTranscript(files); const transcriptPath = picked ? path.join(fullVideoDir, picked.filename) : path.join(fullVideoDir, "transcript.json"); let transcriptMs: number | null = null; if (picked) { // Captions: the newest of every caption input (each English VTT and // the operator's pin), since the cues may come from any of them. A // transcribed record's captions are its alternate tracks // (lib/captionTracks.ts), so they count for it too. const inputs = picked.kind === "vtt" ? captionInputs(files.entries) : [picked.filename, ...captionInputs(files.entries)]; for (const name of inputs) { try { const ms = (await stat(path.join(fullVideoDir, name))).mtimeMs; if (transcriptMs === null || ms > transcriptMs) transcriptMs = ms; } catch { // ignore } } } const subTracks = await readSubTracks(fullVideoDir); let subsMs: number | null = null; for (const t of subTracks) { try { // A TIERED FILE IS `lstat`ED (release 17): `transcript.live_chat.json` // is a sub track AND media — on a tiered channel a link into // channels//media, possibly on another drive. The link carries // the file's times (the tier hook's `lutimes`), so this answers from // the corpus disk and never reaches the media drive. const p = path.join(fullVideoDir, t.filename); const ms = (await (isTierable(t.filename) ? lstat(p) : stat(p))).mtimeMs; if (subsMs === null || ms > subsMs) subsMs = ms; } catch { // ignore } } let availabilityMs: number | null = null; try { availabilityMs = ( await stat(path.join(fullVideoDir, AVAILABILITY_FILENAME)) ).mtimeMs; } catch { availabilityMs = null; } // MAX of both digest sidecars, following the subsMs loop above rather // than the single-stat availabilityMs below it. The overrides file is the // one a human writes; if it did not move this number, a hand-corrected // digest would produce no mutation and never reach a built site. let digestMs: number | null = null; for (const name of [DIGEST_FILENAME, DIGEST_OVERRIDES_FILENAME]) { try { const ms = (await stat(path.join(fullVideoDir, name))).mtimeMs; if (digestMs === null || ms > digestMs) digestMs = ms; } catch { // Sidecar absent — the common case (102 of ~76,000 videos have one). } } channelLive.push({ channelSlug: ch.name, handling: cfg.handling, configName: cfg.name, videoDir, metaPath, metaMs, transcriptPath, transcriptMs, transcriptKind: picked?.kind ?? null, englishVttCount: englishVttsByPreference(files.entries).length, subTracks, subsMs, availabilityMs, digestMs, }); } if (readFailure !== null) { // One video, not the directory: said apart from the readdir failure // above, so the log points at the right place. The whole channel is // held all the same. held.set(ch.name, `a video in its data directory could not be read (${readFailure})`); continue; } // Asked again after the walk: a text tier that became unreadable DURING it // (a tier migration that started) leaves the videos after that point // missing from this scan, which would remove them. A few syscalls a channel. const after = await inspectChannelMedia({ channelsDir }, ch.name, cfg, { fresh: true, }); if (!after.text.readable) { held.set(ch.name, heldReason(after, cfg.mediaDir ?? cfg.dataDir, locations)); continue; } for (const e of channelLive) live.push(e); } return { live, channels, held, mediaAccess }; } // Whether a channel's media tier may be read in this build, and through which // drive (scanSource, from the same inspect the text guard asked). type MediaAccess = { readable: boolean; drive?: string }; function pathKeyId(k: PathKey): string { return `${k[0]}\x00${k[1]}`; } function indexKeysEqual(a: IndexKey, b: IndexKey): boolean { return a[0] === b[0] && a[1] === b[1] && a[2] === b[2]; } // Compact, no trailing newline — the export pages' historical bytes. function writeJsonAtomic(filePath: string, value: unknown): Promise { return writeJsonAtomicShared(filePath, value, { indent: 0, newline: false }); } export type BuildIndexResult = { totalCount: number; pageCount: number; channelCount: number; durationMs: number; pagesWritten: number; pagesSkipped: number; pagesDeleted: number; added: number; changed: number; removed: number; shortCircuited: boolean; // Channels whose media could not be read, so they were not rescanned: their // index records and shared pages were kept as they were (see the header). heldChannels: string[]; // Per channel, what a caption-track pass re-read and changed; empty when no // pass ran (CAPTION_TRACK_KEY). captionTrack: Record; }; export type BuildIndexOptions = { paths: Paths; onLog?: (msg: string) => void; }; // Set to 1 (or true/yes/on) to let a FULL rebuild proceed with a channel held. // Declared in lib/envVars.ts. export const INDEX_ALLOW_HELD_ENV = "ARCHILYZER_INDEX_ALLOW_HELD"; const TRUTHY = new Set(["1", "true", "yes", "on"]); function allowsHeldFullRebuild( env: Record = process.env, ): boolean { const raw = env.ARCHILYZER_INDEX_ALLOW_HELD; return typeof raw === "string" && TRUTHY.has(raw.trim().toLowerCase()); } export async function buildIndex({ paths, onLog, }: BuildIndexOptions): Promise { const log = onLog ?? ((msg: string) => console.log(msg)); const t0 = Date.now(); const channelsDir = paths.channelsDir; const dbPath = paths.lmdbPath; // Shared, site-agnostic per-channel page trees (built once). Per-site // summaries + subs manifests are written under exportSitesIndexDir// // in the per-site loop below. const transcriptsOutDir = paths.exportSharedTranscriptsDir; const subsOutDir = paths.exportSharedSubsDir; const postsOutDir = paths.exportSharedPostsDir; const digestsOutDir = paths.exportSharedDigestsDir; await mkdir(path.dirname(dbPath), { recursive: true }); await mkdir(transcriptsOutDir, { recursive: true }); await mkdir(subsOutDir, { recursive: true }); await mkdir(postsOutDir, { recursive: true }); await mkdir(digestsOutDir, { recursive: true }); await mkdir(paths.exportSitesIndexDir, { recursive: true }); const root = open({ path: dbPath, maxDbs: 18, compression: true, }); const sums = root.openDB({ name: "sums", encoding: "msgpack", }); const cues = root.openDB({ name: "cues", encoding: "msgpack", }); const subs = root.openDB({ name: "subs", encoding: "msgpack", }); // A record's alternate English tracks (ALT_TRACKS_KEY), only where one // differs from the primary — sparse, like `subs`. const alts = root.openDB({ name: "alts", encoding: "msgpack", }); const mtimes = root.openDB({ name: "mtimes", encoding: "msgpack", }); const byChannel = root.openDB({ name: "byChannel", encoding: "msgpack", }); const pageHashes = root.openDB({ name: "pageHashes", encoding: "msgpack", }); const subPageHashes = root.openDB({ name: "subPageHashes", encoding: "msgpack", }); // slug -> ChannelSubsStat, persisted so per-site subs manifests survive a // no-op rebuild where the shared subs pages aren't re-encoded. const channelStatsDb = root.openDB({ name: "channelStats", encoding: "msgpack", }); // The social-post corpus. Keyed [createdAt, channelSlug, id]: createdAt is // ISO-8601 so a lexicographic key sort IS a chronological sort, which fixes // the intra-day ordering problem the videos' [uploadDate, …] tuple has. const posts = root.openDB({ name: "posts", encoding: "msgpack", }); const postPageHashes = root.openDB({ name: "postPageHashes", encoding: "msgpack", }); const channelPostsStatsDb = root.openDB({ name: "channelPostsStats", encoding: "msgpack", }); // The AI-digest corpus. Holds the COMPOSED digest (effectiveDigest of the // machine sidecar + the human overrides), keyed like sums/cues/subs, so the // page writer below can stream a channel's digests straight out of LMDB // without re-reading 76k video dirs on a no-op rebuild. const digests = root.openDB({ name: "digests", encoding: "msgpack", }); const digestPageHashes = root.openDB({ name: "digestPageHashes", encoding: "msgpack", }); const channelDigestStatsDb = root.openDB({ name: "channelDigestStats", encoding: "msgpack", }); // pathKeyId -> VideoState, sparse (only non-`available` videos). Written for // buildStats, which runs after us against this same LMDB file and whose own // per-video cache is keyed on metadata mtime alone — so it cannot see an // availability flip on its own. Keyed by path rather than indexKey because // that is what buildStats has to hand. const videoState = root.openDB({ name: "videoState", encoding: "msgpack", }); const meta = root.openDB({ name: "meta", encoding: "msgpack", }); const storedSchema = meta.get("schema") as number | undefined; const schemaBumped = storedSchema !== SCHEMA_VERSION; // Recorded as INDEX_SCANNED_AT_KEY only when this build completes, so // buildStats can tell apart a video with no `mtimes` record: metadata newer // than this is "not indexed yet"; older, and this build saw it and skipped it // (no upload_date, or processing failed) or its channel was held. // // The scan reads only the source tree, so it runs BEFORE a schema clear: a // full rebuild must know which channels it cannot read before it drops them. const scanStartedAt = Date.now(); const { live, channels: channelConfigs, held, mediaAccess, } = await scanSource(channelsDir, log); // A full rebuild clears every channel's records, and a held channel cannot be // re-read: it would come out of this build empty, and the site built next // would publish it as gone. Refuse, unless told to go on without it. const heldThroughClear = schemaBumped && held.size > 0; if (heldThroughClear && !allowsHeldFullRebuild()) { await root.close(); throw new Error( `The index must be rebuilt in full (index schema ${storedSchema ?? ""} -> ${SCHEMA_VERSION}), ` + `but ${held.size} channel(s) cannot be read: ${describeHeld(held)}. ` + `A full rebuild clears every channel's index records, so these would come out empty and the next site build would publish them as gone. ` + `${HELD_WAYS_OUT} ` + `To rebuild without them anyway (each stays out of the index until its media is back and the index is built again), set ${INDEX_ALLOW_HELD_ENV}=1 ` + `in the environment of the process that runs the build: for the command line, the command's own ` + `(\`${INDEX_ALLOW_HELD_ENV}=1 pnpm archilyzer index\`); for the editor's Build index job, or a site build started from the editor, ` + `the editor's own environment, which takes a restart of the editor.`, ); } if (schemaBumped) { log( `Schema change (${storedSchema ?? ""} -> ${SCHEMA_VERSION}); invalidating LMDB cache.`, ); await sums.clearAsync(); await cues.clearAsync(); await subs.clearAsync(); await alts.clearAsync(); await mtimes.clearAsync(); await byChannel.clearAsync(); await pageHashes.clearAsync(); await subPageHashes.clearAsync(); await channelStatsDb.clearAsync(); await posts.clearAsync(); await postPageHashes.clearAsync(); await channelPostsStatsDb.clearAsync(); await digests.clearAsync(); await digestPageHashes.clearAsync(); await channelDigestStatsDb.clearAsync(); await videoState.clearAsync(); await meta.put("schema", SCHEMA_VERSION); } const livePathIds = new Set(); const liveByPathId = new Map(); for (const s of live) { const id = pathKeyId([s.channelSlug, s.videoDir]); livePathIds.add(id); liveByPathId.set(id, s); } const added: LiveEntry[] = []; const changed: LiveEntry[] = []; const removed: { pathKey: PathKey; indexKey: IndexKey }[] = []; for (const s of live) { const pk: PathKey = [s.channelSlug, s.videoDir]; const prev = mtimes.get(pk); if (!prev) { added.push(s); } else if ( prev.metaMs !== s.metaMs || prev.transcriptMs !== s.transcriptMs || (prev.subsMs ?? null) !== s.subsMs || (prev.availabilityMs ?? null) !== s.availabilityMs || (prev.digestMs ?? null) !== s.digestMs ) { changed.push(s); } } // A held channel has no live entries, and its records are not "gone": they // are kept, and counted for the log. const keptHeld = new Map(); for (const { key, value } of mtimes.getRange()) { const k = key as PathKey; if (held.has(k[0])) { keptHeld.set(k[0], (keptHeld.get(k[0]) ?? 0) + 1); continue; } if (!livePathIds.has(pathKeyId(k))) { removed.push({ pathKey: k, indexKey: (value as MtimeRecord).indexKey }); } } let keptHeldTotal = 0; for (const [slug, why] of held) { const kept = keptHeld.get(slug) ?? 0; keptHeldTotal += kept; log( heldThroughClear ? `Channel ${slug}: ${why}; held under ${INDEX_ALLOW_HELD_ENV}: this full rebuild cleared its index records, so it is out of the index until its media is back and the index is built again. Its transcript, subtitle and digest pages are left as they are.` : `Channel ${slug}: ${why}; its ${kept} indexed video(s) are kept as they are, not rescanned, and its transcript, subtitle and digest pages are left as they are.`, ); } // Summaries labelled before their platform was known (PLATFORM_LABELS_ // VERSION): re-derived like a changed record. A schema bump has cleared // every summary, so it has none to look at. const relabelDue = meta.get(PLATFORM_LABELS_KEY) !== PLATFORM_LABELS_VERSION; if (relabelDue && !schemaBumped) { const queued = new Set( [...added, ...changed].map((s) => pathKeyId([s.channelSlug, s.videoDir])), ); let relabelled = 0; for (const s of live) { const pk: PathKey = [s.channelSlug, s.videoDir]; if (queued.has(pathKeyId(pk))) continue; const prev = mtimes.get(pk); if (!prev) continue; const sum = sums.get(prev.indexKey); if (sum && platformLabelStale(sum)) { changed.push(s); relabelled++; } } log( `Platform labels v${PLATFORM_LABELS_VERSION}: ${relabelled} record(s) labelled before their platform was known, re-derived.`, ); } // Caption records the caption-track rule may now read differently // (CAPTION_TRACK_KEY). A schema bump re-reads everything anyway. const captionTrackDue = meta.get(CAPTION_TRACK_KEY) !== CAPTION_TRACK_RULE_VERSION; // pathKeyIds of the records this pass re-reads. const captionReread = new Set(); if (captionTrackDue && !schemaBumped) { const queued = new Set( [...added, ...changed].map((s) => pathKeyId([s.channelSlug, s.videoDir])), ); for (const s of live) { if (s.transcriptKind !== "vtt") continue; const pk: PathKey = [s.channelSlug, s.videoDir]; const prev = mtimes.get(pk); if (!prev) continue; if (s.englishVttCount < 2 && !storedCuesEmpty(cues, prev.indexKey)) continue; captionReread.add(pathKeyId(pk)); if (!queued.has(pathKeyId(pk))) changed.push(s); } log( `Caption track v${CAPTION_TRACK_RULE_VERSION}: ${captionReread.size} record(s) re-read.`, ); } const captionReport = new Map(); // Records that can hold an alternate track (ALT_TRACKS_KEY). A schema bump // re-reads everything anyway. const altTracksDue = meta.get(ALT_TRACKS_KEY) !== ALT_TRACKS_VERSION; if (altTracksDue && !schemaBumped) { const queued = new Set( [...added, ...changed].map((s) => pathKeyId([s.channelSlug, s.videoDir])), ); let reread = 0; for (const s of live) { if (!canHoldAltTracks(s)) continue; const pk: PathKey = [s.channelSlug, s.videoDir]; if (!mtimes.get(pk)) continue; reread++; if (!queued.has(pathKeyId(pk))) changed.push(s); } log(`Alternate tracks v${ALT_TRACKS_VERSION}: ${reread} record(s) re-read.`); } let altTrackRecords = 0; // RECORDED DATES (release 20 D1, lib/recordedDate.ts): each channel's rule, // compiled once, and the channels whose rule differs from the one their // records were derived under (RECORDED_DATE_RULE_KEY_PREFIX) — every record // of those is re-processed like a changed one. A held channel is left for // the build that can read it. const recordedDateRes = new Map(); const recordedDateRuleChanged: string[] = []; for (const [slug, cfg] of channelConfigs) { const re = cfg.recordedDate ? compileRecordedDateRule(cfg.recordedDate) : null; if (re) recordedDateRes.set(slug, re); if (held.has(slug)) continue; const stored = meta.get(`${RECORDED_DATE_RULE_KEY_PREFIX}${slug}`); const now = re ? cfg.recordedDate!.titlePattern : undefined; if ((typeof stored === "string" ? stored : undefined) !== now) { recordedDateRuleChanged.push(slug); } } if (recordedDateRuleChanged.length > 0 && !schemaBumped) { const ruled = new Set(recordedDateRuleChanged); const queued = new Set( [...added, ...changed].map((s) => pathKeyId([s.channelSlug, s.videoDir])), ); const perChannel = new Map(); for (const s of live) { if (!ruled.has(s.channelSlug)) continue; const pk: PathKey = [s.channelSlug, s.videoDir]; if (!mtimes.get(pk)) continue; perChannel.set(s.channelSlug, (perChannel.get(s.channelSlug) ?? 0) + 1); if (!queued.has(pathKeyId(pk))) changed.push(s); } for (const [slug, n] of perChannel) { log(`Recorded dates: ${slug}: its rule changed; ${n} record(s) re-derived.`); } } const anyMutations = added.length > 0 || changed.length > 0 || removed.length > 0; const { maxTranscriptPageBytes: configuredMaxPageBytes } = getSettings(); // The shared per-channel page trees (transcripts + subs) only need rewriting // when source data changed, the schema bumped, or a previous shared build is // missing/incomplete. Per-site aggregates are always regenerated (cheaply) // from LMDB further below so they reflect current site config. const sharedManifestsPresent = async (): Promise => { for (const channelSlug of channelConfigs.keys()) { // A held channel's pages are not written this build either way. if (held.has(channelSlug)) continue; const mPath = path.join(transcriptsOutDir, channelSlug, "manifest.json"); const raw = await readFile(mPath, "utf8").catch(() => null); if (!raw) return false; try { const cm = JSON.parse(raw) as ChannelTranscriptsManifest; if ( cm.version !== TRANSCRIPTS_MANIFEST_VERSION || cm.maxPageBytes !== configuredMaxPageBytes ) { return false; } } catch { return false; } } return true; }; let sharedNeedsBuild = anyMutations || schemaBumped || !(await sharedManifestsPresent()); // The curated-tag vocabulary + assignments, loaded and compiled ONCE for the // whole build. Every per-video derivation below uses this; reapplyCuratedTags // (after the worker loop) decides whether a tag-only edit has to re-derive // records the mtime diff never looked at. const curated = loadCuratedTagsRuntime(paths); // indexKeyIds the worker derived fresh this build, so the re-apply pass does // not do them a second time. const curatedFresh = new Set(); log( `Diff: +${added.length} added, ~${changed.length} changed, -${removed.length} removed, ${live.length} total.` + (held.size > 0 ? ` Held: ${held.size} channel(s), ${keptHeldTotal} video(s) kept.` : ""), ); for (const channelSlug of channelConfigs.keys()) { await mkdir(path.join(transcriptsOutDir, channelSlug), { recursive: true }); } const toProcess = [...added, ...changed]; let processed = 0; const BATCH = 200; for (let i = 0; i < toProcess.length; i += BATCH) { const slice = toProcess.slice(i, i + BATCH); await Promise.all( slice.map(async (s) => { try { const videoFullDir = path.dirname(s.metaPath); const cuesFresh = await isCuesJsonFresh(videoFullDir); let summary: TranscriptSummary | undefined; let cueList: Cue[] | undefined; if (cuesFresh.fresh) { const normalized = await readNormalizedTranscript(cuesFresh.cuesPath); if (normalized) { const { version, source, cues: cuesField, ...rest } = normalized; void version; void source; summary = { ...rest, channel: s.configName ?? rest.channel, }; cueList = cuesField; // Normalized before its platform was known: the summary is // re-derived from the metadata below, the cues are kept. if (platformLabelStale(summary)) summary = undefined; } } if (!summary) { const metaRaw = await readFile(s.metaPath, "utf8"); const parsedMeta = JSON.parse(metaRaw) as RawMetadata; summary = summarize( s.channelSlug, s.videoDir, parsedMeta, s.configName, ); if (cueList === undefined && s.transcriptMs !== null && s.transcriptKind) { try { cueList = s.transcriptKind === "vtt" ? // The caption-track rule (lib/videoStatus.ts). (await readEnglishVttCues(videoFullDir))?.cues : parseTranscriptJson(await readFile(s.transcriptPath, "utf8")); } catch { cueList = undefined; } } } if (!summary.uploadDate) { log(`Skipping ${s.channelSlug}/${s.videoDir}: no upload_date`); return; } // The channel's recorded-date rule, applied to the summary as it // stands (a normalized one included): set when it reads a date, and // never carried over from an older derivation. delete summary.recordedDate; const recordedRe = recordedDateRes.get(s.channelSlug); if (recordedRe) { const recorded = deriveRecordedDate(recordedRe, summary.title, summary.uploadDate); if (recorded) summary.recordedDate = recorded; } const indexKey: IndexKey = [ summary.uploadDate, s.channelSlug, summary.id, ]; const pk: PathKey = [s.channelSlug, s.videoDir]; const prev = mtimes.get(pk); // What the last build held for this video's captions, for the // caption-track report — read before the put below replaces it. const heldBefore = prev && captionReread.has(pathKeyId(pk)) ? cueText(cues.get(prev.indexKey) ?? []) : undefined; // What the last build held for this video's sub tracks, read before // a re-keyed record is removed: a tiered live chat whose raw cannot // be read now keeps the cues it had (below). const prevSubs = prev ? subs.get(prev.indexKey) : undefined; if (prev && !indexKeysEqual(prev.indexKey, indexKey)) { sums.remove(prev.indexKey); cues.remove(prev.indexKey); subs.remove(prev.indexKey); alts.remove(prev.indexKey); digests.remove(prev.indexKey); byChannel.remove(indexToChannelKey(prev.indexKey)); } if (cueList) cues.put(indexKey, cueList); else cues.remove(indexKey); // The other English tracks, where their words differ from the // primary's (lib/captionTracks.ts). Only a record that can hold one // reads anything. const trackFields: TrackFields = canHoldAltTracks(s) && s.transcriptKind ? await readTrackFields(videoFullDir, s.transcriptKind, cueList) : {}; if (trackFields.altTracks) { alts.put(indexKey, trackFields); altTrackRecords++; } else alts.remove(indexKey); if (heldBefore !== undefined) { const r = captionReport.get(s.channelSlug) ?? { reread: 0, changed: 0, zeroToText: 0 }; r.reread++; const now = cueText(cueList ?? []); if (heldBefore !== now) { r.changed++; if (heldBefore === "" && now !== "") r.zeroToText++; } captionReport.set(s.channelSlug, r); } const parsedSubs: StoredSubs = []; // A tiered track that could be neither read nor kept: `subsMs` is // stored null so the next build's scan sees a change and retries. let retrySubs = false; for (const t of s.subTracks) { try { let trackCues: Cue[] | null = null; if (t.track === "live_chat") { const liveChatFresh = await isLiveChatCuesFresh(videoFullDir); if (liveChatFresh.fresh) { const normalized = await readNormalizedLiveChat( liveChatFresh.cuesPath, ); if (normalized) trackCues = normalized.cues ?? []; } } if (!trackCues) { const rawPath = path.join(path.dirname(s.metaPath), t.filename); // THE RAW REPLAY IS MEDIA (release 17). Tiered — a link into // channels//media — it is read only while the channel's // media is reachable: the drive is PROBED through the watchdog // (one stat, which is what its budget is sized for), then the // file is read directly — a raw replay of hundreds of MB on a // slow platter would outlast the budget, be refused, and, where // the disk's counters are unknown, mark the location stalled. // Not read (unreachable, not answering, unreadable): the cues // the last build held are kept; with none to keep, the record // is written with no sub-track time, so the next build retries. const tiered = t.track === "live_chat" && (await lstat(rawPath).then((l) => l.isSymbolicLink(), () => false)); if (tiered) { const access = mediaAccess.get(s.channelSlug); let raw: string | null = null; if (access?.readable) { try { if (access.drive) await onDrive(access.drive, () => stat(rawPath)); raw = await readFile(rawPath, "utf8"); } catch { raw = null; } } if (raw === null) { const kept = prevSubs?.find((x) => x.track === t.track); if (kept) parsedSubs.push(kept); else retrySubs = true; continue; } trackCues = parseLiveChat(raw); } else { const raw = await readFile(rawPath, "utf8"); trackCues = t.track === "live_chat" ? parseLiveChat(raw) : parseVtt(raw); } } if (trackCues.length > 0) { parsedSubs.push({ track: t.track, cues: trackCues }); } } catch { // Skip unreadable / malformed sub tracks; the rest of the // video's data is still useful. } } if (parsedSubs.length > 0) subs.put(indexKey, parsedSubs); else subs.remove(indexKey); // Curated tags. Derived here, with the caption cues and the live-chat // track already in hand, and stored ON the summary — which is why the // sums.put waited for the subs parse. Omitted when empty, so an // untagged corpus's pages stay byte-identical. Costs nothing at all // when the corpus has no tags (curated.isNoop short-circuits). if (!curated.isNoop) { applyCuratedTagsToSummary( summary, deriveCuratedTags(curated, { channelSlug: s.channelSlug, id: summary.id, title: summary.title, description: summary.description, tags: summary.tags, uploadDate: summary.uploadDate, captionCues: cueList, chatCues: parsedSubs.find((t) => t.track === "live_chat")?.cues, }), ); } curatedFresh.add(indexKeyId(indexKey)); sums.put(indexKey, summary); // The derived layer. Only opened when the scan saw a sidecar, so the // ~76k videos without one cost zero extra reads. What is stored is // effectiveDigest(machine, overrides) — human corrections applied, // `enabled: false` items dropped, chapters sorted by start — so the // corpus ships what a human approved rather than raw model output. if (s.digestMs !== null) { const [machine, overrides] = await Promise.all([ loadDigest(videoFullDir), loadDigestOverrides(videoFullDir), ]); const eff = effectiveDigest(machine, overrides); if (eff.chapters.length > 0 || eff.tags.length > 0) { const provenance = { ...(eff.sections.chapters ? { chapters: eff.sections.chapters.provenance } : {}), ...(eff.sections.tags ? { tags: eff.sections.tags.provenance } : {}), }; digests.put(indexKey, { slug: `${s.channelSlug}/${summary.id}`, id: summary.id, chapters: eff.chapters.map((c) => ({ id: c.id, start: c.start, clock: c.clock, title: c.title, decidedBy: c.decidedBy, })), tags: eff.tags.map((t) => ({ id: t.id, tag: t.tag, decidedBy: t.decidedBy, })), generatedAt: newestGeneratedAt(provenance), provenance, // Carried straight through: a shared digest must stay // identifiable as borrowed all the way to the viewer. ...(eff.derivedFrom ? { derivedFrom: eff.derivedFrom } : {}), }); } else { digests.remove(indexKey); } } else { digests.remove(indexKey); } // Effective, not just availability.json's top-level field: a download // attempt that fails with "deleted" records the class in // download-outcome.json and deliberately leaves the top-level fields // alone (see downloadOneManaged), so the file alone would still read // "public". It does append an availability *history* entry, which // bumps availability.json's mtime — which is what brings us back here // to notice. Same source buildStats' status chart used to read // directly, so the two can't disagree. let availability: Availability | undefined; if (s.availabilityMs !== null) { availability = (await resolveEffectiveAvailability(videoFullDir)) ?? undefined; } byChannel.put(indexToChannelKey(indexKey), 1); mtimes.put(pk, { metaMs: s.metaMs, transcriptMs: s.transcriptMs, subsMs: retrySubs ? null : s.subsMs, availabilityMs: s.availabilityMs, digestMs: s.digestMs, ...(availability ? { availability } : {}), isDeleted: availability === "deleted", isUnlisted: availability === "unlisted", indexKey, }); } catch (err) { log(`Failed to process ${s.channelSlug}/${s.videoDir}: ${String(err)}`); } }), ); processed += slice.length; if (toProcess.length > BATCH) { log(` processed ${processed}/${toProcess.length}`); } } // What the caption-track pass changed, per channel (CAPTION_TRACK_KEY). for (const [slug, r] of [...captionReport].sort(([a], [b]) => (a < b ? -1 : a > b ? 1 : 0))) { log( `Caption track v${CAPTION_TRACK_RULE_VERSION}: ${slug}: ${r.reread} re-read, ${r.changed} now read different text, ${r.zeroToText} had no text and now do.`, ); } if (altTrackRecords > 0) { log(`Alternate tracks: ${altTrackRecords} record(s) read this build hold a track whose words differ from the primary's.`); } for (const { pathKey, indexKey } of removed) { sums.remove(indexKey); cues.remove(indexKey); subs.remove(indexKey); alts.remove(indexKey); digests.remove(indexKey); byChannel.remove(indexToChannelKey(indexKey)); mtimes.remove(pathKey); } // Curated tags can change with no video on disk moving at all — a rule // edited, a video pinned — and nothing in the mtime diff above can see that. // The two meta hashes are therefore checked on EVERY build, and the videos // whose tags moved are re-derived straight out of LMDB (no disk re-read). // Anything that changed has to be re-paged, which is what the // sharedNeedsBuild flip below buys: the page writer's sha1 skip then keeps // the untouched pages untouched, so compose still copies only real changes. // It is awaited because it yields to the event loop as it walks: a rule edit // makes this a whole-corpus pass, and when the editor shares the process a // synchronous one blocks every page render for minutes. const curatedReapply = await reapplyCuratedTags({ runtime: curated, sums, cues, subs, byChannel, meta, alreadyFresh: curatedFresh, allFresh: schemaBumped, log, }); // `pagesPending` also covers a PREVIOUS build that re-derived records and was // then interrupted before writing the pages: its hashes are already stored, // so without the flag nothing would ever ask for those shards again and they // would ship stale curatedTags for ever. if (curatedReapply.pagesPending) sharedNeedsBuild = true; // Commit the flag (and the re-derived summaries) BEFORE the page build, so an // interrupt during it finds the debt recorded rather than lost in a pending // write batch. await meta.flushed; await sums.flushed; await cues.flushed; await subs.flushed; await alts.flushed; await digests.flushed; await byChannel.flushed; await mtimes.flushed; const maxTranscriptPageBytes = configuredMaxPageBytes; const generatedAt = new Date().toISOString(); let pagesWritten = 0; let pagesSkipped = 0; let pagesDeleted = 0; // Streaming page writer: writes each entry directly to the tmp output as it // arrives instead of buffering the full page in memory. Hash is computed // incrementally; on close, the hash is compared to the previously stored one // — equal means we delete the tmp and skip, different means we atomically // rename in place. Avoids the `[...].map(e => e.encoded).join(",")` pattern // that briefly held the entire concatenated page body in V8, which OOMed // when a single live_chat track produced a multi-hundred-MB encoded entry. type PageWriterOpts = { outDir: string; fileName: (idx: number) => string; maxPageBytes: number; ensureDir: boolean; getPrevHash: (idx: number) => string | undefined; setHash: (idx: number, record: PageHashRecord) => void; onLog: (msg: string) => void; }; type PageWriterResult = { pageCount: number; pagesWritten: number; pagesSkipped: number; slugToPage: Record; // Content hash of each emitted page, by page index. Recorded for BOTH the // written and the skipped branch (a skipped page's hash is by definition // the one already on disk), so this is the true current content hash // regardless of whether the page was rewritten this build. The digests tree // publishes these in its manifest as a client cache version. pageHashes: string[]; }; const createPageWriter = (opts: PageWriterOpts) => { let pageIdx = 0; let stream: WriteStream | null = null; let tmpPath = ""; let outPath = ""; let hash = createHash("sha1"); let payloadBytes = 0; let totalBytes = 0; let entryCount = 0; let firstInPage = true; let pagesWrittenLocal = 0; let pagesSkippedLocal = 0; let dirEnsured = !opts.ensureDir; const slugToPage: Record = {}; const pageHashes: string[] = []; const writeChunk = async (chunk: string): Promise => { const s = stream; if (!s) throw new Error("page writer: stream not open"); hash.update(chunk); totalBytes += Buffer.byteLength(chunk, "utf8"); if (!s.write(chunk)) { await once(s, "drain"); } }; const openPage = async (): Promise => { if (!dirEnsured) { await mkdir(opts.outDir, { recursive: true }); dirEnsured = true; } outPath = path.join(opts.outDir, opts.fileName(pageIdx)); tmpPath = `${outPath}.tmp-${process.pid}`; stream = createWriteStream(tmpPath); hash = createHash("sha1"); payloadBytes = 0; totalBytes = 0; entryCount = 0; firstInPage = true; await writeChunk("["); }; const closePage = async (): Promise => { const s = stream; if (!s) return; await writeChunk("]"); await new Promise((resolve, reject) => { s.end((err?: NodeJS.ErrnoException | null) => { if (err) reject(err); else resolve(); }); }); const digest = hash.digest("hex"); pageHashes[pageIdx] = digest; const prev = opts.getPrevHash(pageIdx); if (prev === digest && (await exists(outPath))) { await rm(tmpPath, { force: true }); pagesSkippedLocal++; } else { await rename(tmpPath, outPath); opts.setHash(pageIdx, { hash: digest, entryCount, sizeBytes: totalBytes, }); pagesWrittenLocal++; } stream = null; }; const push = async (encoded: string, id: string): Promise => { const entryBytes = Buffer.byteLength(encoded, "utf8"); const commaBytes = firstInPage ? 0 : 1; const delta = entryBytes + commaBytes; if (stream && !firstInPage && payloadBytes + delta + 2 > opts.maxPageBytes) { await closePage(); pageIdx++; } if (!stream) await openPage(); if (entryBytes + 2 > opts.maxPageBytes) { opts.onLog( ` ${id}: entry (${entryBytes} bytes) exceeds page budget; emitting solo page.`, ); } slugToPage[id] = pageIdx; if (!firstInPage) await writeChunk(","); await writeChunk(encoded); firstInPage = false; payloadBytes += delta; entryCount++; }; const finish = async (): Promise => { if (stream) { await closePage(); pageIdx++; } return { pageCount: pageIdx, pagesWritten: pagesWrittenLocal, pagesSkipped: pagesSkippedLocal, slugToPage, pageHashes, }; }; return { push, finish }; }; if (sharedNeedsBuild) { for (const channelSlug of Array.from(channelConfigs.keys()).sort()) { // A held channel's pages are left exactly as the last build wrote them: // not rewritten, not pruned, and (below) not removed. After a full rebuild // its records are gone, and a rewrite would publish it empty. if (held.has(channelSlug)) continue; const channelDir = path.join(transcriptsOutDir, channelSlug); await mkdir(channelDir, { recursive: true }); const writer = createPageWriter({ outDir: channelDir, fileName: pageFileName, maxPageBytes: maxTranscriptPageBytes, ensureDir: false, getPrevHash: (idx) => pageHashes.get([channelSlug, idx])?.hash, setHash: (idx, record) => { pageHashes.put([channelSlug, idx], record); }, onLog: log, }); for (const { key } of byChannel.getRange({ start: [channelSlug], end: [channelSlug, "￿"], })) { const ck = key as ChannelKey; if (ck[0] !== channelSlug) continue; const indexKey: IndexKey = [ck[1], ck[0], ck[2]]; const summary = sums.get(indexKey); if (!summary) continue; const cueList = cues.get(indexKey); // `track` + `altTracks` only on a record that has an alternate, so every // other record's page bytes are what they were. const trackFields = alts.get(indexKey); const detail: TranscriptDetail = { ...summary, cues: cueList, ...(trackFields?.altTracks?.length ? { track: trackFields.track, altTracks: trackFields.altTracks } : {}), }; const encoded = JSON.stringify(detail); await writer.push(encoded, summary.id); } const { pageCount, pagesWritten: chPagesWritten, pagesSkipped: chPagesSkipped, slugToPage } = await writer.finish(); pagesWritten += chPagesWritten; pagesSkipped += chPagesSkipped; const keep = new Set(["manifest.json"]); for (let i = 0; i < pageCount; i++) keep.add(pageFileName(i)); const existing = await readdir(channelDir).catch(() => [] as string[]); for (const name of existing) { if (keep.has(name)) continue; await rm(path.join(channelDir, name), { force: true }); pagesDeleted++; } for (const { key } of pageHashes.getRange({ start: [channelSlug, pageCount], end: [channelSlug, Number.MAX_SAFE_INTEGER], })) { pageHashes.remove(key as PageHashKey); } const channelManifest: ChannelTranscriptsManifest = { version: TRANSCRIPTS_MANIFEST_VERSION, channelSlug, pageCount, maxPageBytes: maxTranscriptPageBytes, generatedAt, slugToPage, }; await writeJsonAtomic( path.join(channelDir, "manifest.json"), channelManifest, ); } await pageHashes.flushed; const topEntries = await readdir(transcriptsOutDir, { withFileTypes: true, }).catch(() => [] as Dirent[]); for (const e of topEntries) { if (e.isDirectory()) { if (!channelConfigs.has(e.name)) { await rm(path.join(transcriptsOutDir, e.name), { recursive: true, force: true, }); } } else if (e.isFile()) { await rm(path.join(transcriptsOutDir, e.name), { force: true }); } } log( `Transcript pages: ${pagesWritten} written, ${pagesSkipped} unchanged, ${pagesDeleted} stale removed.`, ); } else { log(`Shared transcript pages up to date; skipping rewrite.`); } // Build lookups from indexKey → VideoState out of MtimeRecord, which caches // the availability check done at mutation-processing time. Previously this // loop read availability.json for every video on every build (~28k sequential // disk reads), which stalled GC right before the heaviest phase. // // Both maps are SPARSE: only non-`available` videos get an entry, so the // common case costs nothing. stateByPath is the same data keyed the way // buildStats needs it (see the videoState sub-DB written below). const stateByIndexKey = new Map(); const stateByPath = new Map(); for (const { key, value } of mtimes.getRange()) { const rec = value as MtimeRecord; const st = stateFromRecord(rec); if (st === "available") continue; stateByIndexKey.set(indexKeyId(rec.indexKey), st); stateByPath.set(pathKeyId(key as PathKey), st); } // Overlay the channel-level quick availability check. Computed here rather // than in the per-video loop above on purpose: maybe-missing.json is a // *channel*-level file, but the per-video `mtimes` cache keys off per-video // mtimes, so a channel-level change would never invalidate it and the flag // would simply never ship. Doing it here sidesteps invalidation entirely, and // since the maybe-missing set is small by construction, reading // availability.json for just those ids is cheap and exact. let maybeMissingCount = 0; for (const slug of channelConfigs.keys()) { // A held channel's availability.json files are on the media that cannot be // read, and a missing one reads as "maybe missing": its states are carried // over from the last build instead (below). if (held.has(slug)) continue; const record = await loadMaybeMissing(paths, slug); if (!record?.ids.length) continue; const scannedAtMs = Date.parse(record.checkedAt); if (!Number.isFinite(scannedAtMs)) continue; for (const id of record.ids) { // The quick check's ids ARE data-dir names, which is exactly the second // half of a PathKey (see quickAvailabilityCheck's own comment). const rec = mtimes.get([slug, id]); if (!rec) continue; // not indexed → nothing to show const ikId = indexKeyId(rec.indexKey); if (stateByIndexKey.has(ikId)) continue; // already a confirmed state const state = await resolveMaybeMissingState( path.join(channelsDir, slug, "data", id), scannedAtMs, ); if (state === "available") continue; if (state === "maybe_missing") maybeMissingCount++; stateByIndexKey.set(ikId, state); stateByPath.set(pathKeyId([slug, id]), state); } } if (maybeMissingCount > 0) { log( `Availability: ${maybeMissingCount} video(s) unconfirmed-missing from a channel listing.`, ); } // A held channel keeps the states the last build published for it (the // confirmed ones above come from its kept records; this adds the overlay's), // read back before the wholesale rewrite below. Nothing to carry after a full // rebuild: the clear took them, with the records they described. if (held.size > 0) { for (const { key, value } of videoState.getRange()) { const id = key as string; const cut = id.indexOf("\x00"); if (cut < 0) continue; const slug = id.slice(0, cut); if (!held.has(slug)) continue; const rec = mtimes.get([slug, id.slice(cut + 1)]); if (!rec) continue; const st = value as VideoState; stateByIndexKey.set(indexKeyId(rec.indexKey), st); stateByPath.set(id, st); } } // Publish the sparse map for buildStats, which runs after us against the same // LMDB file and would otherwise have to re-read ~76k availability.json files // to build the status chart. Rewritten wholesale each build: the map is small // and maybe_missing is derived at build time, so reconciling would cost more // than it saves. await videoState.clearAsync(); for (const [id, st] of stateByPath) videoState.put(id, st); await videoState.flushed; let subsPagesWritten = 0; let subsPagesSkipped = 0; let subsPagesDeleted = 0; let subsTotalCount = 0; // slug -> per-channel subtitle stats (site-agnostic). Populated when the // shared subs pages are (re)written, otherwise loaded from LMDB so per-site // subs manifests survive a no-op rebuild. const channelStats = new Map(); if (sharedNeedsBuild) { for (const channelSlug of Array.from(channelConfigs.keys()).sort()) { // Held: its subs dir is left as it is, and its stats carried over so the // site manifests still list it (none to carry after a full rebuild). if (held.has(channelSlug)) { const prev = channelStatsDb.get(channelSlug); if (prev) { channelStats.set(channelSlug, prev); subsTotalCount += prev.videoCount; } continue; } const cfg = channelConfigs.get(channelSlug)!; const subsChannelDir = path.join(subsOutDir, channelSlug); const tracksInChannel = new Set(); let videoCount = 0; let liveChatCount = 0; const subWriter = createPageWriter({ outDir: subsChannelDir, fileName: pageFileName, maxPageBytes: maxTranscriptPageBytes, ensureDir: true, getPrevHash: (idx) => subPageHashes.get([channelSlug, idx])?.hash, setHash: (idx, record) => { subPageHashes.put([channelSlug, idx], record); }, onLog: log, }); for (const { key } of byChannel.getRange({ start: [channelSlug], end: [channelSlug, "￿"], })) { const ck = key as ChannelKey; if (ck[0] !== channelSlug) continue; const indexKey: IndexKey = [ck[1], ck[0], ck[2]]; const stored = subs.get(indexKey); if (!stored || stored.length === 0) continue; const summary = sums.get(indexKey); if (!summary) continue; const display = toDisplaySummary(summary, { state: stateByIndexKey.get(indexKeyId(indexKey)), }); const tracks: Record = {}; let hasLiveChat = false; for (const t of stored) { tracks[t.track] = t.cues; tracksInChannel.add(t.track); if (t.track === "live_chat") hasLiveChat = true; } if (hasLiveChat) liveChatCount++; const detail: SubsDetail = { ...display, tracks }; const encoded = JSON.stringify(detail); await subWriter.push(encoded, summary.id); videoCount++; } const { pageCount: subPageCount, pagesWritten: chSubsWritten, pagesSkipped: chSubsSkipped, slugToPage: subSlugToPage } = await subWriter.finish(); subsPagesWritten += chSubsWritten; subsPagesSkipped += chSubsSkipped; if (videoCount === 0) { // No subs for this channel — remove any stale dir + hashes. await rm(subsChannelDir, { recursive: true, force: true }); for (const { key } of subPageHashes.getRange({ start: [channelSlug], end: [channelSlug, Number.MAX_SAFE_INTEGER], })) { subPageHashes.remove(key as PageHashKey); } continue; } const subKeep = new Set(["manifest.json"]); for (let i = 0; i < subPageCount; i++) subKeep.add(pageFileName(i)); const subExisting = await readdir(subsChannelDir).catch( () => [] as string[], ); for (const name of subExisting) { if (subKeep.has(name)) continue; await rm(path.join(subsChannelDir, name), { force: true }); subsPagesDeleted++; } for (const { key } of subPageHashes.getRange({ start: [channelSlug, subPageCount], end: [channelSlug, Number.MAX_SAFE_INTEGER], })) { subPageHashes.remove(key as PageHashKey); } const tracksList = Array.from(tracksInChannel).sort(); const channelSubsManifest: ChannelSubsManifest = { version: SUBS_CHANNEL_MANIFEST_VERSION, channelSlug, pageCount: subPageCount, maxPageBytes: maxTranscriptPageBytes, generatedAt, tracks: tracksList, slugToPage: subSlugToPage, }; await writeJsonAtomic( path.join(subsChannelDir, "manifest.json"), channelSubsManifest, ); channelStats.set(channelSlug, { name: cfg.name ?? channelSlug, videoCount, tracks: tracksList, liveChatCount, }); subsTotalCount += videoCount; } await subPageHashes.flushed; // Top-level cleanup: drop shared sub dirs for channels that no longer exist. const topSubsEntries = await readdir(subsOutDir, { withFileTypes: true, }).catch(() => [] as Dirent[]); const subsChannelSlugSet = new Set(channelStats.keys()); for (const e of topSubsEntries) { if (e.isDirectory()) { if (!subsChannelSlugSet.has(e.name) && !held.has(e.name)) { await rm(path.join(subsOutDir, e.name), { recursive: true, force: true, }); } } else if (e.isFile()) { // The site-level subs manifest is now per-site; the shared subs root // holds only per-channel dirs. await rm(path.join(subsOutDir, e.name), { force: true }); } } // Persist per-channel subs stats for no-op rebuilds. await channelStatsDb.clearAsync(); for (const [slug, statRec] of channelStats) { channelStatsDb.put(slug, statRec); } await channelStatsDb.flushed; log( `Sub pages: ${subsPagesWritten} written, ${subsPagesSkipped} unchanged, ${subsPagesDeleted} stale removed across ${channelStats.size} channels (${subsTotalCount} sub videos).`, ); } else { for (const { key, value } of channelStatsDb.getRange()) { channelStats.set(key as string, value as ChannelSubsStat); } log(`Shared sub pages up to date; ${channelStats.size} channels with subs.`); } // Both shared trees are now written (or were already current), so the // curated-tag page debt is settled. Deliberately AFTER the page build and not // beside the hashes: an interrupt anywhere above must leave the flag standing // so the next build rewrites the shards. // // Except for a held channel's pages, which were not written: the re-apply // pass re-derives its records in LMDB like any other, and its pages owe them. // The flag stays, so the first build with its media back rewrites them. if (held.size > 0 && curatedReapply.pagesPending) { log( `curated tags: the pages of ${held.size} held channel(s) are not rewritten while held; the re-derived tags reach them on the first build with their media back.`, ); } else { clearCuratedPagesPending(meta); } await meta.flushed; // --------------------------------------------------------------------------- // The social-post corpus: a parallel per-channel page tree beside transcripts // and subs. Social channels have no `data/` dir and so never appear in the // video scan's mtime bookkeeping; incrementality here keys off a signature of // each channel's on-disk posts shards instead. // --------------------------------------------------------------------------- const channelPostsStats = new Map(); let postsPagesWritten = 0; let postsPagesSkipped = 0; let postsTotalCount = 0; const socialSlugs = Array.from(channelConfigs.entries()) .filter(([, cfg]) => isSocialChannel(cfg)) .map(([slug]) => slug) .sort(); for (const channelSlug of socialSlugs) { const cfg = channelConfigs.get(channelSlug)!; const channelRoot = path.join(channelsDir, channelSlug); const postsChannelDir = path.join(postsOutDir, channelSlug); // Signature over the channel's month shards (name + size + mtime). Cheap // to compute, and it changes exactly when a fetch appended something. const shards = await listPostShards(channelRoot); const sigParts: string[] = []; for (const shard of shards) { try { const st = await stat( path.join(channelPostsDir(channelRoot), `${shard}.jsonl`), ); sigParts.push(`${shard}:${st.size}:${st.mtimeMs}`); } catch { /* shard vanished mid-scan — treat as absent */ } } // Include the availability sidecar: a deletion sweep changes what the // served pages must say, even though no shard changed. try { const st = await stat(postsAvailabilityPath(channelRoot)); sigParts.push(`avail:${st.size}:${st.mtimeMs}`); } catch { /* never checked */ } const signature = sigParts.join("|"); const prevStat = channelPostsStatsDb.get(channelSlug); const manifestPresent = await readFile( path.join(postsChannelDir, "manifest.json"), "utf8", ) .then(() => true) .catch(() => false); if ( !schemaBumped && prevStat && prevStat.signature === signature && manifestPresent ) { channelPostsStats.set(channelSlug, prevStat); postsTotalCount += prevStat.postCount; continue; } // Load this channel's posts and refresh its slice of the posts DB. Clearing // by range first means a deleted/rewritten shard can't leave orphans. for (const { key } of posts.getRange({})) { const pk = key as PostIndexKey; if (pk[1] === channelSlug) posts.remove(pk); } // Merge the deleted-post sidecar in, exactly as video availability is // merged into a summary: the sidecar is the source of truth, the flag on // the served record is derived. The JSONL is never rewritten. const postAvailability = await readPostAvailability(channelRoot); const channelPosts: Post[] = []; for (const shard of shards) { channelPosts.push(...(await readPostShard(channelRoot, shard))); } // Dedupe by id (an interrupted archive write can duplicate across shards) // then order newest-first, the order the page tree and the UI present. const byId = new Map(); for (const post of channelPosts) byId.set(post.id, post); const ordered = [...byId.values()].sort(comparePostsNewestFirst); for (const post of ordered) { const rec = postAvailability[post.id]; if (rec) { post.availability = rec.availability; post.availabilityCheckedAt = rec.checkedAt; if (isPostGone(rec.availability)) post.isDeleted = true; } posts.put([post.createdAt, channelSlug, post.id], post); } const postWriter = createPageWriter({ outDir: postsChannelDir, fileName: postsPageFileName, maxPageBytes: maxTranscriptPageBytes, ensureDir: true, getPrevHash: (idx) => postPageHashes.get([channelSlug, idx])?.hash, setHash: (idx, record) => { postPageHashes.put([channelSlug, idx], record); }, onLog: log, }); for (const post of ordered) { await postWriter.push(JSON.stringify(post), post.id); } const { pageCount: postPageCount, pagesWritten: chPostsWritten, pagesSkipped: chPostsSkipped, slugToPage: postSlugToPage, } = await postWriter.finish(); postsPagesWritten += chPostsWritten; postsPagesSkipped += chPostsSkipped; if (ordered.length === 0) { await rm(postsChannelDir, { recursive: true, force: true }); for (const { key } of postPageHashes.getRange({ start: [channelSlug], end: [channelSlug, Number.MAX_SAFE_INTEGER], })) { postPageHashes.remove(key as PageHashKey); } channelPostsStatsDb.remove(channelSlug); continue; } const postKeep = new Set(["manifest.json"]); for (let i = 0; i < postPageCount; i++) postKeep.add(postsPageFileName(i)); for (const name of await readdir(postsChannelDir).catch(() => [] as string[])) { if (postKeep.has(name)) continue; await rm(path.join(postsChannelDir, name), { force: true }); } for (const { key } of postPageHashes.getRange({ start: [channelSlug, postPageCount], end: [channelSlug, Number.MAX_SAFE_INTEGER], })) { postPageHashes.remove(key as PageHashKey); } const postsManifest: ChannelPostsManifest = { version: POSTS_MANIFEST_VERSION, channelSlug, pageCount: postPageCount, maxPageBytes: maxTranscriptPageBytes, generatedAt, slugToPage: postSlugToPage, }; await writeJsonAtomic( path.join(postsChannelDir, "manifest.json"), postsManifest, ); const statRecord: ChannelPostsStat = { name: cfg.name ?? channelSlug, postCount: ordered.length, platform: cfg.platform ?? "", signature, }; channelPostsStatsDb.put(channelSlug, statRecord); channelPostsStats.set(channelSlug, statRecord); postsTotalCount += ordered.length; } // Load stats for channels skipped above (unchanged) that we didn't touch, and // drop page trees for channels that are no longer social / no longer exist. for (const { key, value } of channelPostsStatsDb.getRange()) { const slug = key as string; if (!channelPostsStats.has(slug)) { if (!socialSlugs.includes(slug)) { await rm(path.join(postsOutDir, slug), { recursive: true, force: true }); channelPostsStatsDb.remove(slug); continue; } channelPostsStats.set(slug, value as ChannelPostsStat); } } for (const e of await readdir(postsOutDir, { withFileTypes: true }).catch( () => [] as Dirent[], )) { if (e.isDirectory() && !channelPostsStats.has(e.name)) { await rm(path.join(postsOutDir, e.name), { recursive: true, force: true }); } } await posts.flushed; await postPageHashes.flushed; await channelPostsStatsDb.flushed; if (socialSlugs.length > 0) { log( `Post pages: ${postsPagesWritten} written, ${postsPagesSkipped} unchanged ` + `across ${channelPostsStats.size} social channel(s) (${postsTotalCount} posts).`, ); } // --------------------------------------------------------------------------- // The AI-digest corpus: a third shared per-channel page tree, holding the // chapters + topic tags derived from each transcript. Unlike transcripts and // subs it is SPARSE — a channel emits a manifest only if at least one of its // videos has been digested, and a manifest's slugToPage lists only those // videos. That sparsity is load-bearing downstream: it is what lets the // viewer decide whether to offer a Digest control without a per-video fetch. // --------------------------------------------------------------------------- const channelDigestStats = new Map(); let digestPagesWritten = 0; let digestPagesSkipped = 0; let digestTotalCount = 0; if (sharedNeedsBuild) { for (const channelSlug of Array.from(channelConfigs.keys()).sort()) { // Held: as with subs, its digest dir is left as it is and its stats // carried over. if (held.has(channelSlug)) { const prev = channelDigestStatsDb.get(channelSlug); if (prev) { channelDigestStats.set(channelSlug, prev); digestTotalCount += prev.digestCount; } continue; } const cfg = channelConfigs.get(channelSlug)!; const digestChannelDir = path.join(digestsOutDir, channelSlug); let digestCount = 0; const digestWriter = createPageWriter({ outDir: digestChannelDir, fileName: digestPageFileName, maxPageBytes: maxTranscriptPageBytes, ensureDir: true, getPrevHash: (idx) => digestPageHashes.get([channelSlug, idx])?.hash, setHash: (idx, record) => { digestPageHashes.put([channelSlug, idx], record); }, onLog: log, }); // Oldest-first (the byChannel key order), matching the transcript tree so // adding a newer digest only dirties the last page. for (const { key } of byChannel.getRange({ start: [channelSlug], end: [channelSlug, "￿"], })) { const ck = key as ChannelKey; if (ck[0] !== channelSlug) continue; const indexKey: IndexKey = [ck[1], ck[0], ck[2]]; const digest = digests.get(indexKey); if (!digest) continue; await digestWriter.push(JSON.stringify(digest), digest.id); digestCount++; } const { pageCount: digestPageCount, pagesWritten: chDigestsWritten, pagesSkipped: chDigestsSkipped, slugToPage: digestSlugToPage, pageHashes: digestHashes, } = await digestWriter.finish(); digestPagesWritten += chDigestsWritten; digestPagesSkipped += chDigestsSkipped; if (digestCount === 0) { // No digests in this channel — drop any stale tree + hashes so a // channel whose digests were deleted stops advertising them. await rm(digestChannelDir, { recursive: true, force: true }); for (const { key } of digestPageHashes.getRange({ start: [channelSlug], end: [channelSlug, Number.MAX_SAFE_INTEGER], })) { digestPageHashes.remove(key as PageHashKey); } continue; } const digestKeep = new Set(["manifest.json"]); for (let i = 0; i < digestPageCount; i++) { digestKeep.add(digestPageFileName(i)); } for (const name of await readdir(digestChannelDir).catch( () => [] as string[], )) { if (digestKeep.has(name)) continue; await rm(path.join(digestChannelDir, name), { force: true }); } for (const { key } of digestPageHashes.getRange({ start: [channelSlug, digestPageCount], end: [channelSlug, Number.MAX_SAFE_INTEGER], })) { digestPageHashes.remove(key as PageHashKey); } const channelDigestsManifest: ChannelDigestsManifest = { version: DIGESTS_MANIFEST_VERSION, channelSlug, pageCount: digestPageCount, maxPageBytes: maxTranscriptPageBytes, generatedAt, slugToPage: digestSlugToPage, pageHashes: digestHashes, }; await writeJsonAtomic( path.join(digestChannelDir, "manifest.json"), channelDigestsManifest, ); channelDigestStats.set(channelSlug, { name: cfg.name ?? channelSlug, digestCount, }); digestTotalCount += digestCount; } await digestPageHashes.flushed; // Drop shared digest dirs for channels that no longer have any. const topDigestEntries = await readdir(digestsOutDir, { withFileTypes: true, }).catch(() => [] as Dirent[]); for (const e of topDigestEntries) { if (e.isDirectory()) { if (!channelDigestStats.has(e.name) && !held.has(e.name)) { await rm(path.join(digestsOutDir, e.name), { recursive: true, force: true, }); } } else if (e.isFile()) { // The site-level digests manifest is per-site; the shared root holds // only per-channel dirs. await rm(path.join(digestsOutDir, e.name), { force: true }); } } await channelDigestStatsDb.clearAsync(); for (const [slug, statRec] of channelDigestStats) { channelDigestStatsDb.put(slug, statRec); } await channelDigestStatsDb.flushed; if (digestTotalCount > 0 || digestPagesWritten > 0) { log( `Digest pages: ${digestPagesWritten} written, ${digestPagesSkipped} unchanged ` + `across ${channelDigestStats.size} channel(s) (${digestTotalCount} digests).`, ); } } else { for (const { key, value } of channelDigestStatsDb.getRange()) { channelDigestStats.set(key as string, value as ChannelDigestStat); } if (channelDigestStats.size > 0) { log( `Shared digest pages up to date; ${channelDigestStats.size} channel(s) with digests.`, ); } } // --------------------------------------------------------------------------- // Per-site aggregates: a filtered summaries index (pages + manifest) and a // site-level subs manifest, one bundle per configured site. The heavy // per-channel page trees above are shared; here we only select + regroup. // --------------------------------------------------------------------------- const pageSize = SUMMARIES_PAGE_SIZE; // Generation counter: bumped whenever the underlying video data changed, so a // site's cached fingerprint invalidates when shared data moved even if the // site's own config is untouched. let generation = (meta.get("generation") as number | undefined) ?? 0; if (anyMutations || schemaBumped) { generation++; await meta.put("generation", generation); } const sites = listSites(paths); const siteIds = new Set(sites.map((s) => s.siteId)); let sitesBuilt = 0; let sitesSkipped = 0; let aggregateSummaryPages = 0; let aggregateSummaries = 0; let representativeChannelCount = 0; // Which members a site may publish: an X channel is built only into private // sites while `social.x.visibility` is "private" (lib/postsVisibility.ts — // compose-site narrows its trees by the same rule). const visibilitySettings = getSettings(); for (const site of sites) { const memberSlugs = publishedMemberSlugs( site, (slug) => channelConfigs.get(slug), visibilitySettings, ); // Members the rule leaves out of THIS build, named in the fingerprint below // so flipping the setting (or the site's audience) rebuilds the site. const withheld = site.channels .map((c) => c.slug) .filter((slug) => !memberSlugs.includes(slug)); const slugSet = new Set(memberSlugs); const slugGroup = new Map(); for (const m of site.channels) { slugGroup.set( m.slug, resolveChannelGroupId(m.groupId, site.groups, site.defaultGroupId), ); } const summariesOut = siteSummariesDir(paths, site.siteId); const subsOut = siteSubsDir(paths, site.siteId); const postsOut = sitePostsDir(paths, site.siteId); const digestsOut = siteDigestsDir(paths, site.siteId); const summariesManifestPath = path.join(summariesOut, "manifest.json"); const subsManifestPath = path.join(subsOut, "manifest.json"); const postsManifestPath = path.join(postsOut, "manifest.json"); const digestsManifestPath = path.join(digestsOut, "manifest.json"); const fingerprint = JSON.stringify({ gen: generation, max: configuredMaxPageBytes, groups: site.groups, defaultGroupId: site.defaultGroupId, channels: site.channels, // A tag-only edit moves no site config and no video mtime, but it does // change this site's tag-counts.json (and the curatedTags on its summary // pages). Without these the site would report "up to date" and ship // yesterday's counts. curatedRules: curated.rulesHash, curatedAssign: curated.assignHash, // Only when the visibility rule withholds a member, so a site it does // not touch keeps the fingerprint it had. ...(withheld.length > 0 ? { withheld } : {}), }); const fpKey = `siteFp:${site.siteId}`; if ( !schemaBumped && (meta.get(fpKey) as string | undefined) === fingerprint && (await exists(summariesManifestPath)) && (await exists(subsManifestPath)) ) { sitesSkipped++; try { const m = JSON.parse( await readFile(summariesManifestPath, "utf8"), ) as Manifest; aggregateSummaryPages += m.pageCount; aggregateSummaries += m.totalCount; representativeChannelCount = Math.max( representativeChannelCount, m.channels.length, ); } catch { // ignore } log(`Site ${site.siteId}: up to date; skipping.`); continue; } await mkdir(summariesOut, { recursive: true }); await mkdir(subsOut, { recursive: true }); // --- summaries pages (filtered to member channels) --- // Keyed by slug so we can attach the per-site groupId; member channels with // no data yet still appear (count 0). const chan = new Map< string, { name: string; count: number; groupId: string } >(); for (const slug of memberSlugs) { const name = channelConfigs.get(slug)?.name; if (!name) continue; chan.set(slug, { name, count: 0, groupId: slugGroup.get(slug) ?? site.defaultGroupId, }); } let pageIndex = 0; let buffer: DisplaySummary[] = []; let total = 0; const flushPage = async () => { if (buffer.length === 0) return; await writeJsonAtomic( path.join(summariesOut, pageFileName(pageIndex)), buffer, ); buffer = []; pageIndex++; }; // Per-tag totals for THIS site, accumulated in the pass that is already // streaming every member record — no second walk of the corpus. compose // joins these to the site's effective vocabulary to write /tags.json. const tagCounter = createTagCounter(); for (const { key, value } of sums.getRange({ reverse: true })) { const ik = key as IndexKey; if (!slugSet.has(ik[1])) continue; const stored = value as TranscriptSummary; // A tag scoped to other sites (curated `sites`) is not this site's: // dropped from its summaries and its counts. A copy only when it moves — // the stored value is the LMDB cursor's. const siteTags = tagsForSite(stored.curatedTags, site.siteId, curated.defs); let s: TranscriptSummary = stored; if (siteTags !== stored.curatedTags) { s = { ...stored, curatedTags: siteTags }; if (!siteTags || siteTags.length === 0) delete s.curatedTags; } const acc = chan.get(ik[1]); if (acc) acc.count++; tagCounter.add(ik[1], s.curatedTags); buffer.push( toDisplaySummary(s, { state: stateByIndexKey.get(indexKeyId(ik)) }), ); total++; if (buffer.length >= pageSize) await flushPage(); } await flushPage(); // Written even when empty (a tiny `{tags:{}}`), so a site that loses its // last tagged video does not keep serving a stale count file. compose // decides from this whether to publish /tags.json at all. await writeJsonAtomic( path.join(siteIndexDir(paths, site.siteId), TAG_COUNTS_FILENAME), tagCounter.file(new Date().toISOString()), ); const expectedPages = new Set(); for (let i = 0; i < pageIndex; i++) expectedPages.add(pageFileName(i)); for (const name of await readdir(summariesOut).catch(() => [] as string[])) { if (!name.startsWith("page-")) continue; if (!expectedPages.has(name)) { await rm(path.join(summariesOut, name), { force: true }); } } const channelList: ChannelEntry[] = Array.from(chan.entries()) .map(([slug, { name, count, groupId }]) => ({ slug, name, count, groupId })) .sort((a, b) => a.name.localeCompare(b.name)); const summariesManifest: Manifest = { version: MANIFEST_VERSION, totalCount: total, pageSize, pageCount: pageIndex, generatedAt: new Date().toISOString(), channels: channelList, groups: site.groups, defaultGroupId: site.defaultGroupId, siteId: site.siteId, }; await writeJsonAtomic(summariesManifestPath, summariesManifest); // --- site-level subs manifest (filtered to member channels) --- const subsEntries: SubsChannelEntry[] = []; let subsTotalForSite = 0; let liveChatTotalForSite = 0; for (const slug of memberSlugs) { const statRec = channelStats.get(slug); if (!statRec) continue; subsEntries.push({ name: statRec.name, slug, videoCount: statRec.videoCount, tracks: statRec.tracks, liveChatCount: statRec.liveChatCount, groupId: slugGroup.get(slug) ?? site.defaultGroupId, }); subsTotalForSite += statRec.videoCount; liveChatTotalForSite += statRec.liveChatCount; } const siteSubsManifest: SubsManifest = { version: SUBS_MANIFEST_VERSION, channels: subsEntries.sort((a, b) => a.name.localeCompare(b.name)), totalCount: subsTotalForSite, liveChatTotalCount: liveChatTotalForSite, generatedAt: new Date().toISOString(), groups: site.groups, defaultGroupId: site.defaultGroupId, siteId: site.siteId, }; await writeJsonAtomic(subsManifestPath, siteSubsManifest); // --- site-level posts manifest (filtered to member channels) --- // Only social members contribute; a site with none still gets a manifest // with an empty channel list, so the client's fetch is unconditional. const postsEntries: PostsChannelEntry[] = []; let postsTotalForSite = 0; for (const slug of memberSlugs) { const statRec = channelPostsStats.get(slug); if (!statRec || statRec.postCount === 0) continue; postsEntries.push({ name: statRec.name, slug, postCount: statRec.postCount, platform: (statRec.platform || "bluesky") as PostPlatform, groupId: slugGroup.get(slug) ?? site.defaultGroupId, }); postsTotalForSite += statRec.postCount; } const sitePostsManifest: PostsManifest = { version: SITE_POSTS_MANIFEST_VERSION, channels: postsEntries.sort((a, b) => a.name.localeCompare(b.name)), totalCount: postsTotalForSite, generatedAt: new Date().toISOString(), siteId: site.siteId, }; await mkdir(postsOut, { recursive: true }); await writeJsonAtomic(postsManifestPath, sitePostsManifest); // --- site-level digests manifest (filtered to member channels) --- // Only channels that actually carry digests contribute, so a site with none // ships a manifest with an empty channel list rather than no manifest — // keeping the client's fetch unconditional, as with posts. const digestEntries: DigestsChannelEntry[] = []; let digestsTotalForSite = 0; for (const slug of memberSlugs) { const statRec = channelDigestStats.get(slug); if (!statRec || statRec.digestCount === 0) continue; digestEntries.push({ name: statRec.name, slug, digestCount: statRec.digestCount, groupId: slugGroup.get(slug) ?? site.defaultGroupId, }); digestsTotalForSite += statRec.digestCount; } const siteDigestsManifest: DigestsManifest = { version: SITE_DIGESTS_MANIFEST_VERSION, channels: digestEntries.sort((a, b) => a.name.localeCompare(b.name)), totalCount: digestsTotalForSite, generatedAt: new Date().toISOString(), siteId: site.siteId, }; await mkdir(digestsOut, { recursive: true }); await writeJsonAtomic(digestsManifestPath, siteDigestsManifest); await meta.put(fpKey, fingerprint); sitesBuilt++; aggregateSummaryPages += pageIndex; aggregateSummaries += total; representativeChannelCount = Math.max( representativeChannelCount, channelList.length, ); log( `Site ${site.siteId}: ${total} summaries across ${pageIndex} pages; ${channelList.length} channels.`, ); } // Drop staging dirs + fingerprints for sites that no longer exist. const stagedSiteDirs = await readdir(paths.exportSitesIndexDir, { withFileTypes: true, }).catch(() => [] as Dirent[]); for (const e of stagedSiteDirs) { if (e.isDirectory() && !siteIds.has(e.name)) { await rm(path.join(paths.exportSitesIndexDir, e.name), { recursive: true, force: true, }); } } const staleFpKeys: string[] = []; for (const { key } of meta.getRange()) { const k = key as string; if ( typeof k === "string" && k.startsWith("siteFp:") && !siteIds.has(k.slice("siteFp:".length)) ) { staleFpKeys.push(k); } } for (const k of staleFpKeys) meta.remove(k); await meta.put(INDEX_SCANNED_AT_KEY, scanStartedAt); if (relabelDue && held.size === 0) { await meta.put(PLATFORM_LABELS_KEY, PLATFORM_LABELS_VERSION); } if (captionTrackDue && held.size === 0) { await meta.put(CAPTION_TRACK_KEY, CAPTION_TRACK_RULE_VERSION); } if (altTracksDue && held.size === 0) { await meta.put(ALT_TRACKS_KEY, ALT_TRACKS_VERSION); } for (const slug of recordedDateRuleChanged) { const rule = channelConfigs.get(slug)?.recordedDate; const key = `${RECORDED_DATE_RULE_KEY_PREFIX}${slug}`; if (rule && recordedDateRes.has(slug)) await meta.put(key, rule.titlePattern); else await meta.remove(key); } await meta.flushed; await root.close(); const durationMs = Date.now() - t0; log( `Done in ${(durationMs / 1000).toFixed(2)}s. ${sites.length} site(s): ${sitesBuilt} built, ${sitesSkipped} up to date; ${live.length} transcripts in pool.` + (held.size > 0 ? ` Held, their media not readable: ${[...held.keys()].join(", ")}.` : ""), ); return { totalCount: aggregateSummaries, pageCount: aggregateSummaryPages, channelCount: representativeChannelCount, durationMs, pagesWritten, pagesSkipped, pagesDeleted, added: added.length, changed: changed.length, removed: removed.length, shortCircuited: !sharedNeedsBuild && sitesBuilt === 0, heldChannels: [...held.keys()], captionTrack: Object.fromEntries(captionReport), }; }