// Curated tags, as the index build sees them. // // A video's `curatedTags` is derived, never stored on disk beside the video: // it is (rule hits ∪ operator pins) − suppressions, folded fresh at every // build. Rule hits in particular are NOT persisted anywhere — edit a rule, // rebuild, done — which is why this module exists: nothing in the per-video // mtime diff can notice that a rule changed, so the build has to ask. // // Two hashes in the existing `meta` sub-DB answer that question: // curatedRulesHash the derivation-relevant shape of the vocabulary — every // def's id and order, and every ENABLED rule's kind, // pattern, channel scope and date range // curatedAssignHash every pin and suppression // Both are checked unconditionally on every build (a skipped check would make // a rule edit silently inert until some unrelated mtime moved), and alongside // them `curatedAssignSigs` keeps a per-video signature so an assignment-only // change can re-derive just the videos whose assignment actually moved instead // of the whole corpus. // // **Only the CORPUS vocabulary is evaluated here.** A site-only rule // (sites//tags.json) contributes its definition to that site's published // /tags.json, but it does not tag records: records are shared by every site // that carries the channel, so a per-site rule hit would have to live in a // per-site record. Promote a rule to transcripts/tags.json to make it bind. import { createHash } from "node:crypto"; import type { Cue } from "../lib/vtt"; import type { Paths } from "../lib/paths"; import type { TranscriptSummary } from "../lib/transcripts"; import { compileTagRules, effectiveTagsFor, evaluateCompiledRules, ruleApplies, scopeTagsToChannel, tagChannelScope, type CompiledTagRules, type CuratedTagAssignment, type CuratedTagDef, type CuratedTagsConfig, type PublishedTags, type TagRuleInput, } from "../lib/curatedTags"; import { readGlobalTags } from "../lib/curatedTagsStore"; import { listSites } from "../lib/site"; export const META_RULES_HASH = "curatedRulesHash"; export const META_ASSIGN_HASH = "curatedAssignHash"; export const META_ASSIGN_SIGS = "curatedAssignSigs"; // "Records were re-derived; the shared pages do not reflect them yet." Set when // the hashes are recorded with changes, cleared only once the shared page build // has finished. Without it, a Ctrl-C between the two would leave the hashes // stored and the shards stale for ever: the next build would compare equal, // re-derive nothing, and never dirty the pages. export const META_PAGES_PENDING = "curatedPagesPending"; export const TAG_COUNTS_FILENAME = "tag-counts.json"; export const TAG_COUNTS_VERSION = 1; type IndexKey = [string, string, string]; type ChannelKey = [string, string, string]; function sha1(s: string): string { return createHash("sha1").update(s).digest("hex"); } // The shape of the vocabulary that can change a record. Presentation fields // (label, colour, group, hidden) are deliberately absent: relabelling a tag // must not re-page 30,000 videos. // // ONE-TIME COST WHEN THIS SHAPE CHANGES: the hash is stored per build, so // adding a field to it (channelsExclude, 2026-09-22) makes every stored // `curatedRulesHash` stale and the first build-index after the change // re-derives the whole corpus once — ~80 s, read back out of LMDB with no cue // re-read. The log line below says `rules (changed)` when that happens; // that sentence is the operator's signal that the pass was the migration and // not a rule they edited. // // A SITE-SCOPED tag (`sites`, release 18) also changes records when the // membership of its sites changes, so its resolved channel set is part of the // hash — but only when some tag is scoped (`scope` non-empty): a corpus with // none keeps the hash it had, and pays no re-derivation for this field. export function hashCuratedRules( defs: CuratedTagDef[], scope: ReadonlyMap> = new Map(), ): string { const base = hashCuratedRulesShape(defs); if (scope.size === 0) return base; const scoped = [...scope.entries()] .map(([id, slugs]) => [id, [...slugs].sort()] as const) .sort((a, b) => (a[0] < b[0] ? -1 : a[0] > b[0] ? 1 : 0)); return sha1(JSON.stringify([base, scoped])); } function hashCuratedRulesShape(defs: CuratedTagDef[]): string { const shape = defs.map((def, i) => [ def.id, def.order ?? i, (def.rules ?? []) .filter((r) => r.enabled !== false) .map((r) => [ r.id, r.kind, r.pattern, [...(r.channels ?? [])].sort(), [...(r.channelsExclude ?? [])].sort(), r.dateFrom ?? null, r.dateTo ?? null, ]), ]); return sha1(JSON.stringify(shape)); } // One video's assignment, canonically. Also the per-key signature stored in // meta, so the diff is a string compare. export function assignmentSignature(a: CuratedTagAssignment): string { return `${[...(a.manual ?? [])].sort().join(",")}|${[...(a.suppressed ?? [])].sort().join(",")}`; } export function assignmentSignatures( assignments: Record, ): Record { const out: Record = {}; for (const key of Object.keys(assignments).sort()) { out[key] = assignmentSignature(assignments[key]); } return out; } export function hashCuratedAssignments( assignments: Record, ): string { return sha1(JSON.stringify(assignmentSignatures(assignments))); } // Everything the build needs, loaded and compiled ONCE per build. export type CuratedTagsRuntime = { config: CuratedTagsConfig; defs: CuratedTagDef[]; compiled: CompiledTagRules; rulesHash: string; assignHash: string; sigs: Record; // Site-scoped tags (`sites`): tag id -> the channels it may be derived for. // Empty when no tag is scoped. scope: Map>; // True when the corpus has no rules and no assignments at all — the state of // every untagged install, and the cue to do no work whatsoever. isNoop: boolean; }; export function loadCuratedTagsRuntime(paths: Paths): CuratedTagsRuntime { const members: Record = {}; for (const site of listSites(paths)) members[site.siteId] = site.channels.map((c) => c.slug); return curatedTagsRuntime(readGlobalTags(paths), members); } // `members` (site id -> its channel slugs, site.json `channels`) resolves the // site-scoped tags; without it a scoped tag's sites have no members and the // tag is derived for nothing. export function curatedTagsRuntime( config: CuratedTagsConfig, members: Readonly> = {}, ): CuratedTagsRuntime { const compiled = compileTagRules(config.tags); const scope = tagChannelScope(config.tags, members); return { config, defs: config.tags, compiled, scope, rulesHash: hashCuratedRules(config.tags, scope), assignHash: hashCuratedAssignments(config.assignments), sigs: assignmentSignatures(config.assignments), isNoop: compiled.isEmpty && Object.keys(config.assignments).length === 0, }; } // Cheap pre-checks, so a video on a channel no rule is scoped to never pays for // a cue read. They ask `ruleApplies` — the SAME predicate // evaluateCompiledRules uses — so this can never answer "no cues needed" for a // rule that would then have fired. export function needsCaptionCues( runtime: CuratedTagsRuntime, channelSlug: string, uploadDate?: string, ): boolean { return runtime.compiled.caption.some( (r) => inScope(runtime, r.tagId, channelSlug) && ruleApplies(r, { channelSlug, uploadDate }), ); } // A rule of a site-scoped tag never fires outside its sites' channels: no cue // read is paid for it there either. function inScope(runtime: CuratedTagsRuntime, tagId: string, channelSlug: string): boolean { return runtime.scope.get(tagId)?.has(channelSlug) ?? true; } export function needsChatCues( runtime: CuratedTagsRuntime, channelSlug: string, uploadDate?: string, ): boolean { return runtime.compiled.chatAuthor.some( (r) => inScope(runtime, r.tagId, channelSlug) && ruleApplies(r, { channelSlug, uploadDate }), ); } // The whole fold for one video: rule hits, then the operator's pins and // suppressions over the top. export function deriveCuratedTags( runtime: CuratedTagsRuntime, input: TagRuleInput, ): string[] { const hits = evaluateCompiledRules(input, runtime.compiled); const assignment = runtime.config.assignments[`${input.channelSlug}/${input.id}`]; if (hits.length === 0 && !assignment) return []; // A site-scoped tag exists only on its sites' channels: neither a rule hit // nor a pin carries it to a video on any other channel. return scopeTagsToChannel(effectiveTagsFor(hits, assignment, runtime.defs), input.channelSlug, runtime.scope); } // Set/clear `curatedTags` on a summary in place. Returns true when the stored // value actually changed — omitted-when-empty is the contract, so an untagged // record must come out byte-identical to the one already on disk. export function applyCuratedTagsToSummary( summary: TranscriptSummary, tags: string[], ): boolean { const before = summary.curatedTags; if (tags.length === 0) { if (before === undefined) return false; delete summary.curatedTags; return true; } if ( before && before.length === tags.length && before.every((t, i) => t === tags[i]) ) { return false; } summary.curatedTags = tags; return true; } // ─── re-derivation ─── type SumsDb = { get(key: IndexKey): TranscriptSummary | undefined; put(key: IndexKey, value: TranscriptSummary): unknown; getRange(options?: unknown): Iterable<{ key: unknown; value: unknown }>; }; type CuesDb = { get(key: IndexKey): Cue[] | undefined }; type SubsDb = { get(key: IndexKey): { track: string; cues: Cue[] }[] | undefined }; type ByChannelDb = { getRange(options: unknown): Iterable<{ key: unknown }>; }; type MetaDb = { get(key: string): unknown; put(key: string, value: unknown): unknown; }; // How many records the re-apply pass walks between yields to the event loop. // // This pass is the one place in the build that touches EVERY record in a single // stretch, and for a caption or chat-author rule it decodes each video's cues // out of LMDB while doing it — measured at roughly 500 videos/second, so an // unscoped caption rule over a 77k-record corpus is two to three minutes of // uninterrupted synchronous work. Run inside the editor that is minutes with no // event loop: every page render times out while only the cheapest API routes // answer in the gaps. The rest of the build does not behave this way because // its per-video worker awaits. // // 200 keeps the longest uninterrupted chunk well under half a second even on // the cue-decoding path, while the setImmediate turns themselves stay noise // against the decode cost (~385 of them across the whole corpus). export const CURATED_REAPPLY_YIELD_EVERY = 200; function yieldToEventLoop(): Promise { return new Promise((resolve) => setImmediate(resolve)); } export type ReapplyOptions = { runtime: CuratedTagsRuntime; sums: SumsDb; cues: CuesDb; subs: SubsDb; byChannel: ByChannelDb; meta: MetaDb; // indexKeyIds the per-video worker already derived this build. Skipped here // so a rule edit landing in the same build as new videos doesn't do them // twice. alreadyFresh?: Set; // A schema bump cleared the cache and the worker re-derived everything, so // there is nothing to re-apply — just record the hashes. allFresh?: boolean; // Return the changed index keys as well as the count. Off by default: the // build only needs the boolean, and a first build over 30k records would // otherwise hold 30k three-element tuples for no reason. Tests and any future // per-channel dirtying ask for them. collectChanged?: boolean; log?: (msg: string) => void; }; export type ReapplyResult = { rulesChanged: boolean; assignmentsChanged: boolean; // How many videos' stored curatedTags moved. Non-zero means the shared pages // must be rewritten. changedCount: number; // Which ones — only when `collectChanged` was asked for. changed?: IndexKey[]; // Videos examined (the cost). examined: number; // How many times the pass handed the event loop back. Zero is the normal // answer on a small or untouched corpus; it is non-zero exactly when the pass // would otherwise have blocked. yields: number; // True when the shared pages still owe a rewrite for curated tags: either // this call changed records, or an earlier build did and was interrupted // before its pages were written. The caller MUST OR this into whatever // decides to rebuild the shared trees, and clear it (clearCuratedPagesPending) // only once that build has finished. pagesPending: boolean; }; function indexKeyId(k: IndexKey): string { return `${k[0]}\x00${k[1]}\x00${k[2]}`; } function chatCuesOf(stored: { track: string; cues: Cue[] }[] | undefined): Cue[] | undefined { if (!stored) return undefined; for (const t of stored) if (t.track === "live_chat") return t.cues; return undefined; } // Re-derive `curatedTags` for the videos a tag change can reach, WITHOUT // re-reading a single video directory: cues and live chat come back out of the // LMDB sub-DBs the build already populated. // // rules changed every video is a candidate (a new pattern can match // anything), but cues are read only for the videos a // caption/chat rule is actually scoped to // assignments changed only the videos whose assignment signature moved — // found by a key-only range scan of `byChannel` per // affected channel, which is the cheap idiom // // Returns how many videos changed (and which, on request) so the caller can // force those pages to be rewritten: nothing in the mtime diff knows this // happened. It also sets `curatedPagesPending` — see the flag's own note. // // It is async ONLY to yield: every CURATED_REAPPLY_YIELD_EVERY records it hands // the event loop back, because the whole-corpus walk is otherwise minutes of // unbroken synchronous decoding and the editor serving the same process goes // dark for the duration. The LMDB reads are unchanged; the write ordering is // unchanged; `record()` still runs after the loops and before the log line. export async function reapplyCuratedTags( opts: ReapplyOptions, ): Promise { const { runtime, sums, cues, subs, byChannel, meta } = opts; const log = opts.log ?? (() => {}); const prevRules = meta.get(META_RULES_HASH) as string | undefined; const prevAssign = meta.get(META_ASSIGN_HASH) as string | undefined; const prevSigs = (meta.get(META_ASSIGN_SIGS) as Record | undefined) ?? {}; const rulesChanged = prevRules !== runtime.rulesHash; const assignmentsChanged = prevAssign !== runtime.assignHash; // A previous build re-derived records and was then interrupted before it // finished writing the shared pages. The hashes are already stored, so this // build would otherwise see "nothing changed" and ship stale shards forever. const pendingBefore = meta.get(META_PAGES_PENDING) === true; const record = (pages: boolean) => { meta.put(META_RULES_HASH, runtime.rulesHash); meta.put(META_ASSIGN_HASH, runtime.assignHash); meta.put(META_ASSIGN_SIGS, runtime.sigs); // Set BEFORE the pages are written and cleared only after they are; never // un-set here, or an interrupted build's debt would be forgotten. if (pages) meta.put(META_PAGES_PENDING, true); }; const result = (over: Partial & { changedCount: number }): ReapplyResult => ({ rulesChanged, assignmentsChanged, examined: 0, yields: 0, pagesPending: pendingBefore || over.changedCount > 0, ...over, }); if (!rulesChanged && !assignmentsChanged) { return result({ changedCount: 0, ...(opts.collectChanged ? { changed: [] } : {}) }); } if (opts.allFresh) { // A schema bump: the worker re-derived every record, and the pages are // being rebuilt anyway — but the debt is recorded all the same, so an // interrupt mid-rebuild is not forgotten either. record(true); log( `curated tags: rules ${runtime.rulesHash.slice(0, 8)}, assignments ${runtime.assignHash.slice(0, 8)} (full rebuild, nothing to re-apply)`, ); return { rulesChanged, assignmentsChanged, changedCount: 0, ...(opts.collectChanged ? { changed: [] } : {}), examined: 0, yields: 0, pagesPending: true, }; } const changed: IndexKey[] = []; let changedCount = 0; let examined = 0; let yields = 0; let sinceYield = 0; const skip = opts.alreadyFresh; // Called once per record the loops walk — INCLUDING the ones `alreadyFresh` // skips, because a 77k-record corpus where the worker already did every video // still walks 77k keys. const maybeYield = async (): Promise => { if (++sinceYield < CURATED_REAPPLY_YIELD_EVERY) return; sinceYield = 0; yields++; await yieldToEventLoop(); }; // `known` is the value the cursor already decoded, when there is one — a // second sums.get() per video would double the msgpack decodes on the // whole-corpus path. const rederive = (indexKey: IndexKey, known?: TranscriptSummary): void => { if (skip?.has(indexKeyId(indexKey))) return; const summary = known ?? sums.get(indexKey); if (!summary) return; examined++; const channelSlug = indexKey[1]; const tags = deriveCuratedTags(runtime, { channelSlug, id: summary.id, title: summary.title, description: summary.description, tags: summary.tags, uploadDate: summary.uploadDate, // Only read what a scoped rule can actually use. On an untagged corpus, // and on every channel no rule names, this reads nothing. captionCues: needsCaptionCues(runtime, channelSlug, summary.uploadDate) ? cues.get(indexKey) : undefined, chatCues: needsChatCues(runtime, channelSlug, summary.uploadDate) ? chatCuesOf(subs.get(indexKey)) : undefined, }); if (applyCuratedTagsToSummary(summary, tags)) { sums.put(indexKey, summary); changedCount++; if (opts.collectChanged) changed.push(indexKey); } }; if (rulesChanged) { // Every video is a candidate. This walks `sums` keys only; the value comes // from the same cursor, and cue reads are gated per video above. // // `snapshot: false` because this walk now awaits: lmdb-js otherwise keeps // ONE read transaction open for the whole iteration, and a read txn held // open across minutes of writes stops LMDB reusing the pages those writes // free. With it off the cursor renews its read transaction and repositions // on the key it last returned — the library's own long-duration-iterator // idiom. Nothing written here can confuse it: the only record this loop // puts is the one under the cursor, and the cursor only moves forward. for (const { key, value } of sums.getRange({ snapshot: false })) { rederive(key as IndexKey, value as TranscriptSummary); await maybeYield(); } } else { // Assignment-only change: the symmetric difference of the signatures. const movedKeys = new Set(); for (const [key, sig] of Object.entries(runtime.sigs)) { if (prevSigs[key] !== sig) movedKeys.add(key); } for (const key of Object.keys(prevSigs)) { if (!(key in runtime.sigs)) movedKeys.add(key); } // Group by channel so each channel costs one key-only range scan. const byChannelSlug = new Map>(); for (const key of movedKeys) { const at = key.indexOf("/"); if (at <= 0) continue; const slug = key.slice(0, at); const id = key.slice(at + 1); const set = byChannelSlug.get(slug) ?? new Set(); set.add(id); byChannelSlug.set(slug, set); } for (const [slug, ids] of byChannelSlug) { // Matched inside the scan, re-derived after it. The scan is key-only and // cheap; the re-derivation is what costs, and keeping the await out of // the cursor is free here because the match set is bounded by the // assignments that moved, not by the channel's size. const hits: IndexKey[] = []; // Key-only by convention, NOT by `values: false`: with that option // lmdb-js yields bare keys, not {key, value}, and this destructuring // reads undefined (it crashed the first assignment-only build on the // real corpus, 2026-09-22). recencyIndex.ts scans the same way. for (const { key } of byChannel.getRange({ start: [slug], end: [slug, "￿"], })) { const ck = key as ChannelKey; if (ck[0] !== slug) continue; if (!ids.has(ck[2])) continue; hits.push([ck[1], ck[0], ck[2]]); } for (const indexKey of hits) { rederive(indexKey); await maybeYield(); } } } record(changedCount > 0); const what = [ rulesChanged ? `rules ${runtime.rulesHash.slice(0, 8)} (changed)` : null, assignmentsChanged ? `assignments ${runtime.assignHash.slice(0, 8)} (changed)` : null, ] .filter(Boolean) .join(", "); log( `curated tags: ${what}, examined ${examined}, re-derived ${changedCount}` + (yields > 0 ? `, yielded ${yields}x` : ""), ); return result({ changedCount, ...(opts.collectChanged ? { changed } : {}), examined, yields, }); } // Clear the "the shared pages do not yet reflect the current tags" debt. Call // ONLY after the shared transcript/subs page build has completed: until then an // interrupt must leave the flag standing, which is the whole point of it. export function clearCuratedPagesPending(meta: MetaDb): void { meta.put(META_PAGES_PENDING, false); } // ─── per-site counts ─── export type TagCountsFile = { version: number; generatedAt: string; // tag id -> total on this site + the per-channel breakdown. tags: Record }>; }; export function emptyTagCounts(): TagCountsFile { return { version: TAG_COUNTS_VERSION, generatedAt: "", tags: {} }; } // Accumulator for the per-site summaries stream: one add() per record, no // second pass over the corpus. export function createTagCounter() { const tags = new Map }>(); return { add(channelSlug: string, curatedTags: string[] | undefined): void { if (!curatedTags || curatedTags.length === 0) return; for (const id of curatedTags) { let entry = tags.get(id); if (!entry) { entry = { count: 0, channels: new Map() }; tags.set(id, entry); } entry.count++; entry.channels.set(channelSlug, (entry.channels.get(channelSlug) ?? 0) + 1); } }, get size(): number { return tags.size; }, file(generatedAt: string): TagCountsFile { const out: TagCountsFile["tags"] = {}; for (const id of Array.from(tags.keys()).sort()) { const entry = tags.get(id)!; const channels: Record = {}; for (const slug of Array.from(entry.channels.keys()).sort()) { channels[slug] = entry.channels.get(slug)!; } out[id] = { count: entry.count, channels }; } return { version: TAG_COUNTS_VERSION, generatedAt, tags: out }; }, }; } // The published /tags.json for one site: its effective vocabulary joined to its // counts. Hidden tags and tags with no video ON THIS SITE are dropped — a // global tag that matched nothing here does not get a chip here — and a site // with nothing left ships no file at all (compose skips the write; a 404 is a // legitimate empty state). export function publishedTagsFrom( defs: CuratedTagDef[], counts: TagCountsFile | null, ): PublishedTags { const tags = defs .filter((def) => def.hidden !== true) .map((def, i) => ({ def, i, entry: counts?.tags[def.id] })) .filter((x) => (x.entry?.count ?? 0) > 0) .sort((a, b) => { const oa = a.def.order ?? a.i; const ob = b.def.order ?? b.i; if (oa !== ob) return oa - ob; return a.def.id < b.def.id ? -1 : a.def.id > b.def.id ? 1 : 0; }) .map(({ def, entry }) => ({ id: def.id, // The published document always carries a label: the wire format is read, // not layered, so an un-overridden tag falls back to its id exactly once — // here — and never by a sanitize pass inventing one on the site file. label: def.label ?? def.id, ...(def.group ? { group: def.group } : {}), ...(def.groupLabel ? { groupLabel: def.groupLabel } : {}), ...(def.color ? { color: def.color } : {}), ...(def.order !== undefined ? { order: def.order } : {}), count: entry!.count, channels: entry!.channels, })); return { version: 1, tags }; }