import { test } from "node:test"; import assert from "node:assert/strict"; import path from "node:path"; import os from "node:os"; import { mkdtemp, mkdir, writeFile, rm, utimes } from "node:fs/promises"; import { getOperation, allOperations, backfillLaneOperations, resolveBackfillLaneOperations, orderByDependencies, operationCatalog, operationsForLane, operationCostBasis, operationGroup, operationLabel, operationsActionLabel, operationsGroupLabel, digestLaneFor, diarizationLaneFor, laneYieldsToTranscription, addOperationState, emptyOperationCounts, reachableOperationWork, type OperationClassification, type Operation, backfillLaneEntriesOf, backfillLaneOperationEntriesOf, bucketLaneOperationId, presentOperationWork, } from "./operations"; import { pauseLaneFor } from "./pauseGates"; import { readVideoFiles, CUES_JSON_FILENAME, META_FILENAME, SOURCE_MEDIA_BASENAME, } from "./videoStatus"; import { CUES_FILE_VERSION } from "../controller/normalizeTranscript"; import { attributeOneVideo } from "../controller/attributeOne"; import { DIGEST_FILENAME, OLLAMA_DIGEST_APP_ID } from "./digest"; import { BACKFILL_QUEUE, DIGEST_LOCAL_QUEUE, DIGEST_REMOTE_QUEUE, } from "./queueKeys"; import { ATTRIBUTION_FILENAME, ATTRIBUTION_PROMPT_VERSION, type AttributionRecord, } from "./attribution"; import { DEFAULT_DIARIZATION_THRESHOLD, DIARIZATION_FILENAME, SORTFORMER_DIARIZATION_ENGINE, diarizationTarget, isDiarizationFresh, type DiarizationRecord, } from "./diarization"; import { defaultSiteSettings, type SiteSettings } from "./settings"; import { SAVED_VIDEO_POINTER_FILENAME } from "./savedVideo"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test common/lib/operations.test.ts // // Every case here is a FIXTURE DIRECTORY, on purpose: `state()` is defined as // "read from disk, never stored", and a test that hands it a hand-built record // would not be testing the thing the indicators actually call. const diarization = getOperation("diarization")!; // resolveTarget is channel-scoped since the digest entry (a digest's identity // includes the hash of the channel's context note). None of the kinds tested // here reads paths/channelSlug, so a stub is honest as well as convenient. function targetCtx(settings: SiteSettings) { return { settings, paths: { channelsDir: "/nonexistent" } as unknown as Parameters< typeof diarization.resolveTarget >[0]["paths"], channelSlug: "test-channel", }; } function settingsWithDiarization( over: Partial = {}, ): SiteSettings { const s = defaultSiteSettings(); return { ...s, diarization: { ...s.diarization, enabled: true, segModel: "/opt/models/seg-1.onnx", embModel: "/opt/models/emb-1.onnx", threshold: DEFAULT_DIARIZATION_THRESHOLD, ...over, }, }; } // The duration cap is OFF by default now that windowing exists, so a test of the // cap has to ask for one. Stated here once rather than inline, so it is obvious // that every OTHER test in this file runs with the cap disabled — which is the // shipped configuration. function settingsWithCap(hours: number): SiteSettings { return settingsWithDiarization({ maxAudioHours: hours }); } // A video dir built from a description of what is on disk. `sidecar` is written // verbatim so a MALFORMED file can be tested — that case is not hypothetical, // it is what a crash mid-write leaves behind. async function fixture(opts: { transcript?: boolean; audio?: boolean; container?: boolean; savedPointer?: boolean; sidecar?: string; // Written into metadata.info.json. Absent means NO metadata file at all, // which is the "duration unknown" case the cap has to get right. durationSec?: number; }): Promise<{ dir: string; cleanup: () => Promise }> { const root = await mkdtemp(path.join(os.tmpdir(), "backfill-kinds-")); const dir = path.join(root, "vid1"); await mkdir(dir, { recursive: true }); if (opts.transcript !== false) { await writeFile( path.join(dir, "transcript.json"), JSON.stringify({ transcription: [{ text: "hi" }] }), ); } if (opts.audio) await writeFile(path.join(dir, "audio.mp3"), "x"); // `source-media.` specifically — isSourceMediaFile keys off the basename, // and an arbitrarily-named mp4 in a video dir is not a persisted container. if (opts.container) { await writeFile(path.join(dir, `${SOURCE_MEDIA_BASENAME}.mp4`), "x"); } if (opts.savedPointer) { await writeFile( path.join(dir, SAVED_VIDEO_POINTER_FILENAME), JSON.stringify({ storedAt: "now", dir: "/nowhere", file: "video.mp4" }), ); } if (opts.sidecar !== undefined) { await writeFile(path.join(dir, DIARIZATION_FILENAME), opts.sidecar); } if (opts.durationSec !== undefined) { await writeFile( path.join(dir, META_FILENAME), JSON.stringify({ id: "vid1", duration: opts.durationSec }), ); } return { dir, cleanup: () => rm(root, { recursive: true, force: true }) }; } function sidecar(engine: DiarizationRecord["engine"]): string { return JSON.stringify({ videoId: "vid1", generatedAt: "2026-08-07T00:00:00.000Z", speakers: 2, turns: [{ start: 0, end: 4, speaker: 0 }], engine, } satisfies DiarizationRecord); } // The current identity, as scripts/diarize.mjs would record it: BASENAMES, not // the configured full paths. const CURRENT = { engine: "sherpa-onnx", segmentationModel: "seg-1.onnx", embeddingModel: "emb-1.onnx", threshold: DEFAULT_DIARIZATION_THRESHOLD, }; async function classify( dirOpts: Parameters[0], settings: SiteSettings = settingsWithDiarization(), ): Promise { const { dir, cleanup } = await fixture(dirOpts); try { const files = await readVideoFiles(dir, { checkUntranscribable: true }); return await diarization.state({ videoDir: dir, videoId: "vid1", files, target: await diarization.resolveTarget(targetCtx(settings)), settings, }); } finally { await cleanup(); } } test("present: the sidecar matches the identity we would produce now", async () => { assert.equal( await classify({ audio: true, sidecar: sidecar(CURRENT) }), "present", ); // Still present with the audio already cleaned away — the whole point of // capturing while the audio exists is that the result outlives it. assert.equal(await classify({ sidecar: sidecar(CURRENT) }), "present"); }); test("stale: a different threshold or model is work, not coverage", async () => { // The threshold is the single most consequential knob (it decides how many // speakers come out), so a change to it MUST show as work. Before the // comparator, this read as done. assert.equal( await classify({ audio: true, sidecar: sidecar({ ...CURRENT, threshold: 0.5 }), }), "stale", ); assert.equal( await classify({ audio: true, sidecar: sidecar({ ...CURRENT, segmentationModel: "seg-OLD.onnx" }), }), "stale", ); assert.equal( await classify({ audio: true, sidecar: sidecar({ ...CURRENT, embeddingModel: "emb-OLD.onnx" }), }), "stale", ); }); // THE ENGINE IS NOT COMPARED, and that is a fix rather than an omission. // scripts/diarize.mjs records whichever binary actually ran, and nothing in // DiarizationSettings can predict that — so a hardcoded engine in the target // would mark every sidecar from any other wrapper permanently stale, which at // ~500-680 s/audio-hour is an infinite regeneration loop. (The e2e fake engine // records "fake-diarize" and would have tripped it on the first run.) test("a different engine binary is not, by itself, stale", async () => { assert.equal( await classify({ audio: true, sidecar: sidecar({ ...CURRENT, engine: "some-other-engine" }), }), "present", ); // The comparison is still WRITTEN, so adding an engine setting later needs no // new logic: a target that does declare one still rejects a mismatch. assert.equal( isDiarizationFresh( { videoId: "v", generatedAt: "now", speakers: 1, turns: [], engine: { ...CURRENT, engine: "some-other-engine" }, }, { ...CURRENT, engine: "sherpa-onnx" }, ), false, ); }); // The engine-as-a-setting case the test above says needs "no new logic". It is // asserted for sortformer and NOT for the default, and that asymmetry is the // whole design: the default engine still cannot predict which binary a // `--engine` override runs, so it must keep asserting nothing. test("selecting sortformer asserts the engine; selecting the default still does not", () => { const rec = (engine: DiarizationRecord["engine"]): DiarizationRecord => ({ videoId: "v", generatedAt: "now", speakers: 1, turns: [], engine, }); const sherpaSidecar = rec({ engine: "sherpa-onnx", segmentationModel: "seg-1.onnx", embeddingModel: "emb-1.onnx", threshold: DEFAULT_DIARIZATION_THRESHOLD, }); const sortformerSidecar = rec({ engine: SORTFORMER_DIARIZATION_ENGINE, model: "sortformer-4spk.gguf", }); const sortformerTarget = diarizationTarget({ engine: SORTFORMER_DIARIZATION_ENGINE, sortformerModel: "/abs/path/sortformer-4spk.gguf", // Still configured, because settings carry one set of fields for both // engines. None of it may leak into the sortformer identity. segModel: "/abs/seg-1.onnx", embModel: "/abs/emb-1.onnx", threshold: DEFAULT_DIARIZATION_THRESHOLD, }); // Basename only, as everywhere else: the corpus is rsynced between shards. assert.equal(sortformerTarget.model, "sortformer-4spk.gguf"); assert.equal(sortformerTarget.segmentationModel, undefined); assert.equal(sortformerTarget.embeddingModel, undefined); assert.equal(isDiarizationFresh(sortformerSidecar, sortformerTarget), true); // The point of switching: every sherpa sidecar becomes work the backfill lane // will offer to redo, rather than silently staying half a corpus. assert.equal(isDiarizationFresh(sherpaSidecar, sortformerTarget), false); // And back the other way, with no special case needed — a sortformer record // carries neither segmentation nor embedding model, so it fails the sherpa // comparison on the models alone. const sherpaTarget = diarizationTarget({ engine: "sherpa-onnx", segModel: "/abs/seg-1.onnx", embModel: "/abs/emb-1.onnx", threshold: DEFAULT_DIARIZATION_THRESHOLD, }); assert.equal(sherpaTarget.engine, undefined); assert.equal(isDiarizationFresh(sherpaSidecar, sherpaTarget), true); assert.equal(isDiarizationFresh(sortformerSidecar, sherpaTarget), false); }); // Sortformer has no clustering step, so the threshold cannot have changed any of // its turns. Letting it into the identity would regenerate the whole corpus for // an edit that provably could not affect it. test("the clustering threshold does not stale a sortformer sidecar", () => { const sidecar: DiarizationRecord = { videoId: "v", generatedAt: "now", speakers: 1, turns: [], engine: { engine: SORTFORMER_DIARIZATION_ENGINE, model: "m.gguf" }, }; const at = (threshold: number) => diarizationTarget({ engine: SORTFORMER_DIARIZATION_ENGINE, sortformerModel: "m.gguf", threshold, }); assert.equal(isDiarizationFresh(sidecar, at(DEFAULT_DIARIZATION_THRESHOLD)), true); assert.equal(isDiarizationFresh(sidecar, at(0.4)), true); // The model itself IS the identity, though — a different one is a redo. assert.equal( isDiarizationFresh( sidecar, diarizationTarget({ engine: SORTFORMER_DIARIZATION_ENGINE, sortformerModel: "other.gguf", }), ), false, ); }); // Which resource diarization competes for is a CONFIGURATION outcome, and the // backfill lane reads it to decide whether a guaranteed share is even available. // Getting this wrong in the permissive direction puts ~4.4 GB of sortformer next // to parakeet on an 8 GB card. test("diarization contends for the GPU only as sortformer on vulkan", () => { const lane = (engine: "sherpa-onnx" | "sortformer", backend: "vulkan" | "cpu") => diarizationLaneFor({ engine, backend }); assert.equal(lane("sortformer", "vulkan").contendsFor, "gpu"); assert.equal(laneYieldsToTranscription(lane("sortformer", "vulkan")), true); // The same engine on the CPU backend is only after cores. assert.equal(lane("sortformer", "cpu").contendsFor, "cpu"); assert.equal(laneYieldsToTranscription(lane("sortformer", "cpu")), false); // sherpa-onnx is ONNX/CPU, so a stale `backend: vulkan` left in settings must // NOT make it claim the card — that would park the lane behind transcription // for work using no shaders at all, which is the bug digestYield.ts already // records having hit once. assert.equal(lane("sherpa-onnx", "vulkan").contendsFor, "cpu"); assert.equal(laneYieldsToTranscription(lane("sherpa-onnx", "vulkan")), false); // The queue key never changes: two diarizations must not run at once whichever // engine is selected. assert.equal(lane("sortformer", "vulkan").queueKey, BACKFILL_QUEUE); assert.equal(lane("sherpa-onnx", "cpu").queueKey, BACKFILL_QUEUE); }); // THE COMPATIBILITY RULE, and the reason it is written down: without it, adding // a field to the provenance would mark all 77,000 videos stale at once. test("an absent recorded field compares equal to today's default", async () => { assert.equal( await classify({ audio: true, // A sidecar from before the threshold was recorded at all. sidecar: sidecar({ engine: "sherpa-onnx", segmentationModel: "seg-1.onnx", embeddingModel: "emb-1.onnx", }), }), "present", ); // ...and it is equal to the DEFAULT specifically, not to anything: with a // non-default threshold configured, the same record is stale. assert.equal( await classify( { audio: true, sidecar: sidecar({ engine: "sherpa-onnx", segmentationModel: "seg-1.onnx", embeddingModel: "emb-1.onnx", }), }, settingsWithDiarization({ threshold: 0.5 }), ), "stale", ); }); test("a version bump alone does not invalidate captured work", async () => { // ~500-680 s/audio-hour of CPU says an engine point-release is not a reason to // redo everything when the models and the threshold are unchanged. assert.equal( await classify({ audio: true, sidecar: sidecar({ ...CURRENT, version: "9.9.9" }), }), "present", ); }); test("missing: no sidecar, and the input is still here", async () => { assert.equal(await classify({ audio: true }), "missing"); // A persisted source container counts — ffmpeg reads it directly, which is // what resolveDiarizableMedia does. assert.equal(await classify({ container: true }), "missing"); }); test("missing-input: no sidecar and nothing to diarize from", async () => { assert.equal(await classify({}), "missing-input"); // A saved-video POINTER whose stored file has gone (an unmounted backup disk) // is not an input either — the pointer is not the media. assert.equal(await classify({ savedPointer: true }), "missing-input"); }); // --------------------------------------------------------------------------- // The duration cap. A stopgap for an OOM that kills 6 of 10 videos over 6 hours // on this box, burning ~40 minutes each and producing nothing. test("deferred: over the cap, with the input right there", async () => { // 5 hours against the 4-hour default. The input EXISTS — that is the whole // point of a third state: this is not missing-input (nothing to work from) and // not missing (work to do); it is work deliberately not attempted. assert.equal( await classify({ audio: true, durationSec: 5 * 3600 }, settingsWithCap(4)), "deferred", ); }); test("under the cap is ordinary missing work", async () => { assert.equal( await classify({ audio: true, durationSec: 3 * 3600 }, settingsWithCap(4)), "missing", ); // Exactly at the cap is not over it. assert.equal( await classify({ audio: true, durationSec: 4 * 3600 }, settingsWithCap(4)), "missing", ); }); test("unknown duration is NOT deferred", async () => { // No metadata.info.json at all. Deferring here would quietly remove a video // from the work list on the strength of a file that could not be read, and it // would break every fixture in this file that predates the cap. assert.equal(await classify({ audio: true }, settingsWithCap(4)), "missing"); // Present but useless — the same answer, for the same reason. assert.equal( await classify({ audio: true, durationSec: 0 }, settingsWithCap(4)), "missing", ); }); test("maxAudioHours 0 turns the cap off — and 0 is the shipped default", async () => { // Where this setting went once windowed diarization landed. assert.equal( await classify( { audio: true, durationSec: 12 * 3600 }, settingsWithDiarization({ maxAudioHours: 0 }), ), "missing", ); }); test("the cap never overrules a sidecar that is already there", async () => { // A long video ALREADY diarized stays `present`: the cap decides what to // attempt, not what counts as done. Otherwise raising the cap would look like // work appearing and lowering it would look like work being undone. assert.equal( await classify( { audio: true, durationSec: 12 * 3600, sidecar: sidecar(CURRENT) }, settingsWithCap(4), ), "present", ); }); test("the cap does not resurrect a video whose input is gone", async () => { // missing-input is checked FIRST. A 12-hour video with nothing to diarize from // is unreachable, not deferred — deferred promises "we could do this if you // raised the cap", and that would be a lie here. // // With re-download ARMED it stops being a lie, and the answer changes; the // three cases below are that branch. assert.equal( await classify({ durationSec: 12 * 3600 }, settingsWithCap(4)), "missing-input", ); }); // The cap in the missing-input branch. `missing-input` is not inert once // settings.backfill.allowRedownload is on: it DISPATCHES A DOWNLOAD, and // diarizeOneVideo does not enforce the cap itself — so without this a 10-hour // VOD over the cap would be fetched in full and then diarized anyway. One // metadata read is nothing next to a download. function settingsWithRedownload(hours: number): SiteSettings { const s = settingsWithCap(hours); return { ...s, backfill: { ...s.backfill, allowRedownload: true } }; } test("re-download armed: over the cap is deferred BEFORE a download is spent", async () => { assert.equal( await classify({ durationSec: 12 * 3600 }, settingsWithRedownload(4)), "deferred", ); }); test("re-download armed: under the cap still asks for its media", async () => { assert.equal( await classify({ durationSec: 3 * 3600 }, settingsWithRedownload(4)), "missing-input", ); // Unknown duration gets the benefit of the doubt here too — the same rule as // the `missing` branch, and for the same reason. assert.equal( await classify({}, settingsWithRedownload(4)), "missing-input", ); }); test("the cap read stays off the hot path when re-download is off", async () => { // countBackfillWork calls state() for every video on every job start, so the // metadata read is bought only by the branch that can spend a download. With // the flag off the answer is missing-input whatever the duration says. assert.equal( await classify({ durationSec: 12 * 3600 }, settingsWithCap(4)), "missing-input", ); }); test("a malformed sidecar reads as absent, never as done", async () => { // The same rule diarization-server.ts's hasDiarization() encodes: a // half-written file must not be what convinces anything the work is captured // — that is what would let the cleanup sweep delete the only copy of the audio. assert.equal( await classify({ audio: true, sidecar: "{ not json" }), "missing", ); assert.equal( await classify({ audio: true, sidecar: JSON.stringify({ videoId: "x" }) }), "missing", ); }); test("not-applicable: an untranscribed video is not this backfill's business", async () => { assert.equal( await classify({ transcript: false, audio: true }), "not-applicable", ); }); test("a disabled or unconfigured feature reports no backfill at all", () => { const s = defaultSiteSettings(); // Off by default, so nothing advertises catch-up work for it. assert.equal(diarization.enabled(s), false); assert.equal(backfillLaneOperations(s).length, 0); // Enabled but with no models is "not set up", which must also report nothing // rather than a corpus-sized work list nobody can act on. assert.equal( diarization.enabled(settingsWithDiarization({ segModel: "" })), false, ); assert.equal(diarization.enabled(settingsWithDiarization()), true); assert.equal(backfillLaneOperations(settingsWithDiarization()).length, 1); }); test("an unknown or stale kind id in the sweep scope is dropped, not fatal", () => { const s = settingsWithDiarization(); // Empty scope = every enabled lane kind. assert.deepEqual( resolveBackfillLaneOperations(s, []).map((k) => k.id), ["diarization"], ); assert.deepEqual( resolveBackfillLaneOperations(s, undefined).map((k) => k.id), ["diarization"], ); // A settings file naming a kind from another build must not wedge the lane. assert.deepEqual(resolveBackfillLaneOperations(s, ["from-the-future"]), []); assert.deepEqual( resolveBackfillLaneOperations(s, ["diarization", "from-the-future"]).map((k) => k.id), ["diarization"], ); }); // --------------------------------------------------------------------------- // Attribution — the second and third kinds, and the pair that shares one file // --------------------------------------------------------------------------- const attrText = getOperation("attribution-text")!; const attrDiarized = getOperation("attribution-diarized")!; function settingsWithAttribution( over: Partial = {}, ): SiteSettings { const s = defaultSiteSettings(); return { ...s, attribution: { ...s.attribution, enabled: true, diarizedEnabled: true, textOnlyEnabled: true, appId: OLLAMA_DIGEST_APP_ID, model: "qwen2.5:7b", ...over, }, }; } function attrSidecar(over: Partial = {}): string { return JSON.stringify({ videoId: "vid1", generatedAt: "2026-08-07T00:00:00.000Z", speakers: [{ index: 0, label: "Host" }], segments: [{ start: 0, end: 10, speaker: 0 }], provenance: { method: "diarized", appId: OLLAMA_DIGEST_APP_ID, model: "qwen2.5:7b", modelRequested: "qwen2.5:7b", promptVersion: ATTRIBUTION_PROMPT_VERSION, generatedAt: "2026-08-07T00:00:00.000Z", ...over, }, } satisfies AttributionRecord); } // Same fixture discipline as above — a real directory, because state() is // defined as "read from disk, never stored". async function attrFixture(opts: { transcript?: boolean; diarization?: string; attribution?: string; }): Promise<{ dir: string; cleanup: () => Promise }> { const root = await mkdtemp(path.join(os.tmpdir(), "backfill-attr-")); const dir = path.join(root, "vid1"); await mkdir(dir, { recursive: true }); if (opts.transcript !== false) { await writeFile( path.join(dir, "transcript.json"), JSON.stringify({ transcription: [{ text: "hi" }] }), ); } if (opts.diarization !== undefined) { await writeFile(path.join(dir, DIARIZATION_FILENAME), opts.diarization); } if (opts.attribution !== undefined) { await writeFile(path.join(dir, ATTRIBUTION_FILENAME), opts.attribution); } return { dir, cleanup: () => rm(root, { recursive: true, force: true }) }; } async function classifyAttr( kind: typeof attrText, dirOpts: Parameters[0], settings: SiteSettings = settingsWithAttribution(), ): Promise { const { dir, cleanup } = await attrFixture(dirOpts); try { const files = await readVideoFiles(dir, { checkUntranscribable: true }); return await kind.state({ videoDir: dir, videoId: "vid1", files, target: await kind.resolveTarget(targetCtx(settings)), settings, }); } finally { await cleanup(); } } // A diarization sidecar at a known generatedAt, so the "re-diarizing invalidates // the naming" case has something to compare against. const DIARIZED_AT = "2026-08-06T00:00:00.000Z"; function diarizationSidecar(generatedAt = DIARIZED_AT): string { return JSON.stringify({ videoId: "vid1", generatedAt, speakers: 2, turns: [{ start: 0, end: 10, speaker: 0 }], engine: { engine: "sherpa-onnx" }, } satisfies DiarizationRecord); } test("attribution-text reaches every transcribed video and never reports missing-input", async () => { // The claim that prices this lane: its input is the cue stream, which every // transcribed video has. That is why it can reach the whole corpus, and why // running it over the whole corpus costs ~194,000 model calls. assert.equal(await classifyAttr(attrText, {}), "missing"); assert.equal( await classifyAttr(attrText, { transcript: false }), "not-applicable", ); }); test("attribution-text: fresh, stale, and a malformed file that reads as absent", async () => { assert.equal( await classifyAttr(attrText, { attribution: attrSidecar({ method: "text-only" }), }), "present", ); assert.equal( await classifyAttr(attrText, { attribution: attrSidecar({ method: "text-only", model: "llama3:8b", modelRequested: "llama3:8b" }), }), "stale", ); // A half-written sidecar must read as absent — and here that matters twice // over, because a malformed file that read as "present" could masquerade as a // diarized record and block this lane forever. assert.equal( await classifyAttr(attrText, { attribution: "{ not json" }), "missing", ); }); // THE DOWNGRADE RULE, in the counter. Reporting a diarized record as work would // put this lane in a loop of "attempt, refuse, still outstanding" on every pass // of a multi-day sweep. test("attribution-text has nothing to do where a diarized record exists", async () => { assert.equal( await classifyAttr(attrText, { attribution: attrSidecar() }), "present", ); // Not even when that diarized record is itself stale — it is not this lane's // record to redo. assert.equal( await classifyAttr(attrText, { attribution: attrSidecar({ model: "llama3:8b", modelRequested: "llama3:8b" }), }), "present", ); }); test("attribution-diarized: no diarization.json is BLOCKED, not missing-input", async () => { // ~73,000 videos on this corpus, against a handful reachable — so this state // has to be right or every surface is wrong. // // THIS ASSERTION USED TO SAY missing-input, AND THAT WAS THE DEFECT. // missing-input means one thing to the rest of the system: the media is gone, // re-acquire it. So these videos were counted as needing media re-fetched, // and with allowRedownload on the lane would have spent a download per video // fetching AUDIO — which cannot satisfy a wait for diarization.json — and // then deleted it again. What they are waiting for is the `diarization` kind, // which this table produces. assert.equal(await classifyAttr(attrDiarized, {}), "blocked"); assert.equal( await classifyAttr(attrDiarized, { attribution: attrSidecar() }), "blocked", ); }); // THE DISPATCH HALF OF THIS PROPERTY LIVES IN // controller/operationBatch.test.ts ("blocked is never dispatched, and // re-download cannot change that"). It used to live here and import // candidateAction across the layer boundary — a lib/ test reaching into // controller/, carried on architecture.test.ts's allow-list as a debt. The // classification is this file's business; what a dispatcher does with it is // the dispatcher's. test("blocked is counted, and is NOT reachable work", () => { const counts = emptyOperationCounts(); addOperationState(counts, "blocked"); addOperationState(counts, "blocked"); addOperationState(counts, "missing"); assert.equal(counts.blocked, 2); assert.equal(counts.missing, 1); // The load-bearing line. A corpus with one diarization and 73,000 waiting // attributions must not report 73,000 jobs ready to run. assert.equal(reachableOperationWork(counts), 1); // And it is its own number, not folded into the re-acquire population. assert.equal(counts.missingInput, 0); assert.equal(counts.deferred, 0); }); test("a prerequisite is ordered before the kind that declares it", () => { const base = settingsWithDiarization(); const kinds = resolveBackfillLaneOperations( { ...base, attribution: { ...base.attribution, enabled: true, diarizedEnabled: true, textOnlyEnabled: true, }, }, undefined, ); const ids = kinds.map((k) => k.id); const diarizationAt = ids.indexOf("diarization"); const diarizedAttrAt = ids.indexOf("attribution-diarized"); assert.ok(diarizationAt >= 0 && diarizedAttrAt >= 0, ids.join(",")); // Without this, a video diarized during a pass only becomes attributable on // whatever LATER pass happens to find the sidecar on disk. assert.ok( diarizationAt < diarizedAttrAt, `diarization must precede attribution-diarized, got ${ids.join(", ")}`, ); // AND the sort is STABLE: attribution-text stays last, which is a separate // deliberate decision (the better lane must reach a diarized video first, or // the text lane spends ~30 model calls to produce a record it must not write). assert.equal(ids[ids.length - 1], "attribution-text", ids.join(",")); }); test("ordering degrades safely when a prerequisite is absent or cyclic", () => { const a = { id: "a", dependsOn: ["b"] } as unknown as Operation; const b = { id: "b", dependsOn: ["a"] } as unknown as Operation; // A cycle must not wedge the lane or silently drop a kind: every entry comes // back, in declaration order. assert.deepEqual( orderByDependencies([a, b]).map((k) => k.id), ["a", "b"], ); // A dependency on something not in the selection imposes no ordering, and // does not remove the dependant from the run. const lonely = { id: "lonely", dependsOn: ["not-here"] } as unknown as Operation; assert.deepEqual( orderByDependencies([lonely]).map((k) => k.id), ["lonely"], ); }); test("attribution-diarized treats a text-only record as the upgrade queue", async () => { // `missing`, not `stale`: the diarized record this kind is responsible for // genuinely was never made. This is the whole of PLAN.md's bespoke "upgrade // job", and it falls out of the registry rather than needing new machinery. assert.equal( await classifyAttr(attrDiarized, { diarization: diarizationSidecar(), attribution: attrSidecar({ method: "text-only" }), }), "missing", ); assert.equal( await classifyAttr(attrDiarized, { diarization: diarizationSidecar() }), "missing", ); }); test("attribution-diarized: present, and stale when the clusters underneath change", async () => { assert.equal( await classifyAttr(attrDiarized, { diarization: diarizationSidecar(), attribution: attrSidecar({ diarizationGeneratedAt: DIARIZED_AT }), }), "present", ); // Re-diarized since. Cluster 3 is now a different person, or nobody, so the // names that pointed at it are work again. assert.equal( await classifyAttr(attrDiarized, { diarization: diarizationSidecar("2026-08-09T00:00:00.000Z"), attribution: attrSidecar({ diarizationGeneratedAt: DIARIZED_AT }), }), "stale", ); // A model change is stale the ordinary way too. assert.equal( await classifyAttr(attrDiarized, { diarization: diarizationSidecar(), attribution: attrSidecar({ diarizationGeneratedAt: DIARIZED_AT, model: "llama3:8b", modelRequested: "llama3:8b", }), }), "stale", ); }); // Capture may legitimately be switched off after a run: a sidecar on disk is a // perfectly good input, and refusing to name it would strand exactly the work // the capture lane exists to protect. test("attribution-diarized does not require diarization CAPTURE to still be on", async () => { const s = settingsWithAttribution(); assert.equal(s.diarization.enabled, false); assert.equal(attrDiarized.enabled(s), true); assert.equal( await classifyAttr( attrDiarized, { diarization: diarizationSidecar() }, s, ), "missing", ); }); test("each attribution lane is gated separately, under one master switch", () => { const off = defaultSiteSettings(); assert.equal(attrText.enabled(off), false); assert.equal(attrDiarized.enabled(off), false); // Turning the feature on must not by itself arm a ~194,000-call sweep. const onlyMaster = settingsWithAttribution({ diarizedEnabled: false, textOnlyEnabled: false, }); assert.equal(attrText.enabled(onlyMaster), false); assert.equal(attrDiarized.enabled(onlyMaster), false); assert.equal( attrText.enabled(settingsWithAttribution({ textOnlyEnabled: false })), false, ); assert.equal( attrDiarized.enabled(settingsWithAttribution({ diarizedEnabled: false })), false, ); // Both lanes on. Diarization CAPTURE is still off in this fixture, so two — // which is also the point: an attribution backfill does not need the capture // lane armed to have work. assert.equal(backfillLaneOperations(settingsWithAttribution()).length, 2); }); // The cheap, better lane must get to a video first: on a video that has // diarization, running the text lane first would spend ~30 model calls producing // a record the diarized lane then replaces. test("the diarized lane is ordered ahead of the text-only lane", () => { assert.deepEqual( resolveBackfillLaneOperations(settingsWithAttribution(), []).map((k) => k.id), ["attribution-diarized", "attribution-text"], ); }); // THE RUNNER'S OWN GUARD, and it is not the same test as the counter's. state() // runs at pull time; the pool can hold a candidate for minutes afterwards, and // the diarized lane can land a better record in that window. This asserts the // refusal happens with NO engine call at all — no settings resolved, no model // chosen, nothing that could fail for an unrelated reason. test("the text lane refuses to downgrade a diarized record, without calling an engine", async () => { const { dir, cleanup } = await attrFixture({ attribution: attrSidecar(), }); try { // A full transcript fixture, so the run gets PAST the cues guard and the // refusal is genuinely the downgrade rule rather than a missing transcript. await writeTranscriptFixture(dir); const outcome = await attributeOneVideo({ // A bogus engine: reaching it at all is the failure this test is looking // for. The guard fires before anything resolves an app. paths: { channelsDir: dir } as never, videoDir: dir, videoId: "vid1", channelSlug: "chan", method: "text-only", settings: { ...settingsWithAttribution().attribution, appId: "no-such-engine", }, // Even FORCED. Forcing a regeneration is not the same as asking for a // worse record, and nothing in the UI should be able to request the second // by accident. force: true, onLog: () => {}, }); assert.equal(outcome, "outranked"); } finally { await cleanup(); } }); // mtimes have to ascend: isCuesJsonFresh compares cues.json against the metadata // and the raw transcript, and a cues.json older than either means the transcript // changed underneath and must not be attributed. async function writeTranscriptFixture(dir: string): Promise { await writeFile( path.join(dir, META_FILENAME), JSON.stringify({ id: "vid1", title: "A video" }), ); await writeFile( path.join(dir, CUES_JSON_FILENAME), JSON.stringify({ version: CUES_FILE_VERSION, id: "vid1", title: "A video", cues: [{ start: 0, end: 5, text: "hello" }], }), ); const base = Date.now() / 1000; await utimes(path.join(dir, META_FILENAME), base, base); await utimes(path.join(dir, "transcript.json"), base, base); await utimes(path.join(dir, CUES_JSON_FILENAME), base + 10, base + 10); } // --------------------------------------------------------------------------- // The digest entry. Registered but NOT dispatched by the backfill lane — see // backfillLaneOperations — so these tests pin the classification, which is the half // that is now shared, and the lane declaration, which is what keeps the GPU and // CPU lanes from serializing. // --------------------------------------------------------------------------- const digestKind = getOperation("digest")!; // The digest freshness target, hand-built rather than resolved: resolveTarget // reads the channel's context note through paths, and these cases are about // what state() does with a record, not about how the identity is derived. const DIGEST_TARGET = { target: { appId: OLLAMA_DIGEST_APP_ID, model: "llama3:8b", promptVersion: 2, contextHash: "ctx-1", }, sections: ["chapters"] as const, }; function digestSidecar(over: Record = {}): string { return JSON.stringify({ videoId: "vid1", digestSchemaVersion: 1, sections: { chapters: { items: [{ start: 0, title: "Intro" }], provenance: { appId: OLLAMA_DIGEST_APP_ID, model: "llama3:8b", modelRequested: "llama3:8b", promptVersion: 2, contextHash: "ctx-1", }, }, }, ...over, }); } async function classifyDigest(opts: { transcript?: boolean; cues?: boolean; staleCues?: boolean; sidecar?: string; // Which sections must ALL be fresh. Defaults to the single-section shape the // live corpus is configured with; the two-section form is what makes a // part-done digest possible at all. sections?: readonly string[]; }): Promise { const { dir, cleanup } = await fixture({ transcript: opts.transcript }); try { if (opts.cues !== false && opts.transcript !== false) { await writeTranscriptFixture(dir); if (opts.staleCues) { // cues.json OLDER than the raw transcript: the transcript changed // underneath and the normalize pass owes this video a rewrite. const base = Date.now() / 1000; await utimes(path.join(dir, CUES_JSON_FILENAME), base - 100, base - 100); } } if (opts.sidecar) { await writeFile(path.join(dir, DIGEST_FILENAME), opts.sidecar); } const files = await readVideoFiles(dir, { checkUntranscribable: true }); return await digestKind.state({ videoDir: dir, videoId: "vid1", files, target: opts.sections ? { ...DIGEST_TARGET, sections: opts.sections } : DIGEST_TARGET, settings: defaultSiteSettings(), }); } finally { await cleanup(); } } test("digest: a video with no transcript is BLOCKED on transcription", async () => { // PLAN.md states "a video cannot be digested until it has a transcript" in // prose and nothing enforced it. Now it is declared (dependsOn) and reported. assert.equal(await classifyDigest({ transcript: false }), "blocked"); assert.deepEqual(digestKind.dependsOn, ["transcription"]); // And the dependency resolves to something nameable, which is the whole // reason the externally-dispatched operations are in the catalog at all. assert.equal(operationLabel("transcription"), "Transcription"); }); test("digest: a fresh sidecar is present, a mismatched one is stale", async () => { assert.equal( await classifyDigest({ sidecar: digestSidecar() }), "present", ); // A model change is work. This is the case the snapshot's old // "does an ai-digest.json exist?" test got wrong, reading "All digested" // while the batch reported the whole channel as stale. assert.equal( await classifyDigest({ sidecar: digestSidecar({ sections: { chapters: { items: [{ start: 0, title: "Intro" }], provenance: { appId: OLLAMA_DIGEST_APP_ID, model: "llama3:70b", modelRequested: "llama3:70b", promptVersion: 2, contextHash: "ctx-1", }, }, }, }), }), "stale", ); }); test("digest: no sidecar at all is missing, not stale", async () => { assert.equal(await classifyDigest({}), "missing"); }); test("digest: a stale cues.json defers rather than digesting superseded text", async () => { // Digesting now would describe text that is about to be rewritten, and would // then look fresh forever. Not a failure — the normalize pass fixes it — so // it must not be counted as reachable work or as a broken engine. assert.equal(await classifyDigest({ staleCues: true }), "deferred"); }); test("digest: a transcript with NO cues.json defers too — the case that is 97.6% of them", async () => { // The population this branch actually catches. Measured over 79,219 video // dirs: 1,942 have a raw transcript and no cues.json at all, against 47 with // a superseded one — and 1,683 of the 1,942 are a single `handling: "youtube"` // channel, which downloads subtitles with --skip-download and so never runs // transcribeOne, the only automatic caller of normalizeTranscript. // // Same classification as the stale case on purpose: one normalize pass fixes // both, so the distinction is about COPY (see CuesFreshReason), not dispatch. assert.equal( await classifyDigest({ transcript: true, cues: false }), "deferred", ); }); test("a kind that can defer says WHY, and digest's reason is not 'it clears itself'", () => { // SpeakersStage used to hardcode one sentence about the diarization duration // cap for every kind's deferred videos at once. Correct only while diarization // was the sole kind that could defer. assert.match( getOperation("diarization")!.deferredHint!, /Max audio hours/, ); // The claim this whole change exists to retract: the old copy said the // normalize pass clears these on its own. Nothing runs it on its own, so the // hint has to name the action. const digestHint = digestKind.deferredHint!; assert.match(digestHint, /Normalize/); assert.doesNotMatch(digestHint, /clears? (itself|them|it)/i); }); test("digest: a digest shared from a duplicate cluster counts as done", async () => { // Worth ~11% of the sweep. This entry IS the snapshot's digest work list, so // getting the share rule wrong here regenerates every mirror. assert.equal( await classifyDigest({ sidecar: digestSidecar({ derivedFrom: { videoId: "canonical" } }), }), "present", ); }); test("digest declares its own lane, and it is NOT the backfill queue", () => { // THE PROPERTY THAT KEEPS THEM CONCURRENT. registry.ts runs every non-empty // queueKey at concurrency 1, so sharing a key with diarization would make the // GPU lane wait on the CPU lane and vice versa. assert.notEqual(digestKind.lane.queueKey, BACKFILL_QUEUE); assert.equal(digestKind.lane.queueKey, DIGEST_LOCAL_QUEUE); assert.equal(getOperation("diarization")!.lane.queueKey, BACKFILL_QUEUE); // The two digest lanes are distinct from each other for the same reason. assert.notEqual( digestLaneFor("local-gpu").queueKey, digestLaneFor("remote-api").queueKey, ); // Only the GPU lane stands aside for transcription. Yielding the metered lane // would park something that costs nothing to keep running. assert.equal(laneYieldsToTranscription(digestLaneFor("local-gpu")), true); assert.equal(laneYieldsToTranscription(digestLaneFor("remote-api")), false); }); test("the backfill lane never dispatches digest", () => { // backfillLaneOperations feeds the backfill lane's executor, the channel // Backfill card and the dashboard instrument. Digest must be in the CATALOG // and out of THAT list, or it both serializes behind diarization and gets // double-counted. const settings = settingsWithDiarization(); const laneIds = backfillLaneOperations(settings).map((k) => k.id); const allIds = allOperations(settings).map((k) => k.id); assert.ok(allIds.includes("digest"), allIds.join(",")); assert.ok(!laneIds.includes("digest"), laneIds.join(",")); // Which means the backfill lane's operation resolution cannot reach it // either, even when it is asked for by name. assert.deepEqual(resolveBackfillLaneOperations(settings, ["digest"]), []); }); test("the catalog covers every operation, dispatched here or not", () => { const ids = operationCatalog().map((o) => o.id); for (const id of [ "sync", "download", "transcription", "diarization", "attribution-diarized", "attribution-text", "digest", ]) { assert.ok(ids.includes(id), `${id} missing from ${ids.join(",")}`); } // Every declared dependency resolves to a catalogued operation. This is the // check that would have caught digest.dependsOn naming something that did not // exist. for (const op of operationCatalog()) { for (const dep of op.dependsOn ?? []) { assert.ok(ids.includes(dep), `${op.id} depends on unknown ${dep}`); } } }); test("`runner` names the auto-queue runner, and only for the two that have one", () => { // The console reads this instead of asking whether the id happens to be // "download" or "transcription". `dispatch` cannot answer it: both of those // are `external` WITH a runner, and an external operation sharing a runner's // queue with no runner of its own is a legal entry — so an id-shaped guess // would hand such an entry that runner's controls. const runners = new Map(operationCatalog().map((o) => [o.id, o.runner])); assert.equal(runners.get("download"), "download"); assert.equal(runners.get("transcription"), "transcription"); // THE THIRD ONE IS CHANNEL-SCOPED, and that is the point of it being here // rather than inferred: the download runner dispatches the metadata scan as a // channel-scoped unit before it picks any video, so the scan's console — and // its pause — are the download lane's. assert.equal(runners.get("metadata-scan"), "download"); // Everything the sweep dispatches must leave it unset — a backfill kind with // a runner would render a runner console over a lane no runner feeds. const dispatched = new Set(["download", "transcription", "metadata-scan"]); for (const op of operationCatalog()) { if (dispatched.has(op.id)) continue; assert.equal(op.runner, undefined, `${op.id} declares a runner`); } // Sync included, and it is the interesting one: the sync scheduler's // heartbeat IS a runner in the IA doc's sense, but `runner` is typed to the // AutoQueueKinds and the heartbeat is not one of them. Its console is chosen // off `trigger` instead — see the scope/trigger test below. assert.equal(runners.get("sync"), undefined); }); test("`settingsBlock` names the settings.json block, off the descriptor", () => { // The operation page renders an operation's settings form by switching on // THIS, not on a table keyed by operation id — the rule slice 2 set for // `runner` and the lane. Both attribution operations name the same block // because there is one `settings.attribution` and one form drawn on both // pages; the backfill LANE's block is not an operation's and is not here. const expected: Record = { // The one external operation with a block: the sync scheduler's settings // are the sync operation's, and /operations/sync is where they are edited. sync: "syncScheduler", digest: "digest", diarization: "diarization", "attribution-text": "attribution", "attribution-diarized": "attribution", }; for (const op of operationCatalog()) { assert.equal( op.settingsBlock, expected[op.id], `${op.id} declares settingsBlock ${String(op.settingsBlock)}`, ); } // Every id the map names is actually in the catalog — otherwise the loop // above passes by never visiting a key that was renamed out from under it. const ids = new Set(operationCatalog().map((o) => o.id)); for (const id of Object.keys(expected)) assert.ok(ids.has(id), `${id} is not catalogued`); const attribution = operationCatalog().filter((o) => o.id.startsWith("attribution-")); assert.equal(attribution.length, 2); assert.equal(attribution[0].settingsBlock, attribution[1].settingsBlock); }); test("every catalogued operation declares a group and a cost basis", () => { // Both are read unconditionally by the UI — the transit line groups stations // by `group`, and every armed operation prints `costBasis` beside its // backlog. An entry missing either renders a blank where a fact should be, // which is the failure mode that let attribution-text sit armed at ~194,000 // model calls while every screen called it "Backfill". for (const op of operationCatalog()) { assert.ok(op.group, `${op.id} has no group`); assert.ok(op.shortLabel.length > 0, `${op.id} has no shortLabel`); assert.ok( op.shortLabel.length <= 12, `${op.id} shortLabel "${op.shortLabel}" is too long for a column header`, ); assert.ok(op.costBasis.length > 0, `${op.id} has no costBasis`); } }); test("the three speaker operations share one group; digest does not", () => { assert.equal(operationGroup("diarization"), "speakers"); assert.equal(operationGroup("attribution-diarized"), "speakers"); assert.equal(operationGroup("attribution-text"), "speakers"); assert.equal(operationGroup("digest"), "digest"); assert.equal(operationGroup("download"), "media"); assert.equal(operationGroup("transcription"), "transcript"); // Sync is its own group, not `media`: that group is download's /channels // column group, and sync draws no column. assert.equal(operationGroup("sync"), "sync"); // An id the catalog does not know is null, NOT filed under the first group. assert.equal(operationGroup("no-such-operation"), null); }); test("`scope` and `trigger`: sync is the one cadence-triggered operation", () => { // The two fields exist so a per-video surface can EXCLUDE a channel-scoped // entry by a declared fact rather than by its id. // // TWO entries are channel-scoped now: sync and the metadata scan. They differ // on `trigger`, and that difference is the point — sync runs on a cadence // (the scheduler's heartbeat decides), the scan runs off a BACKLOG (listed // videos nobody has read the title of). /operations/sync is chosen off // `trigger`, which is why only one entry may claim "cadence". const channelScoped = operationCatalog() .filter((op) => op.scope === "channel") .map((op) => op.id) .sort(); assert.deepEqual(channelScoped, ["metadata-scan", "sync"]); for (const op of operationCatalog()) { assert.equal(op.trigger === "cadence", op.id === "sync", `${op.id} trigger ${op.trigger}`); if (op.scope !== "channel") assert.equal(op.scope, "video", op.id); } assert.equal(operationCatalog()[0].id, "sync"); // upstream first }); test("a set's label is derived, so a mixed lane cannot claim one member's name", () => { // The whole reason this is derived: the backfill lane holds three speaker // operations today and its station can honestly say "Speakers". Add a kind // from another group and it degrades to the generic name rather than // continuing to advertise a label that now describes two thirds of it. assert.equal( operationsGroupLabel([ "diarization", "attribution-diarized", "attribution-text", ]), "Speakers", ); assert.equal( operationsGroupLabel(["diarization", "digest"]), "Derived data", ); assert.equal(operationsGroupLabel([]), "Derived data"); // Unknown ids contribute nothing rather than poisoning a single-group set. assert.equal(operationsGroupLabel(["digest", "no-such-op"]), "Digest"); }); test("a group has a station name AND a name you can put a verb in front of", () => { // "Run speakers work" is why these are two declarations rather than one // lower-cased derivation. The station eyebrow needs a noun; the button needs // an object. assert.equal( operationsActionLabel([ "diarization", "attribution-diarized", "attribution-text", ]), "speaker work", ); assert.equal(operationsActionLabel(["digest"]), "digests"); assert.equal(operationsActionLabel([]), "derived data"); assert.equal(operationsActionLabel(["diarization", "digest"]), "derived data"); }); test("attribution-text's cost basis states the CHUNK unit, not the video", () => { // The measured fact that hid behind the word "backfill": this lane's unit is // the transcript chunk, so its 11,337 reachable videos are on the order of // 194,000 model calls. Every other lane in the table is per-video, and a // reader who assumes that of this one is wrong by ~17x. assert.match(operationCostBasis("attribution-text"), /chunk/); assert.match(operationCostBasis("attribution-diarized"), /per video/); assert.match(operationCostBasis("diarization"), /per video/); assert.equal(operationCostBasis("no-such-operation"), ""); }); // The accounting rule the whole feature turns on: reachable work and // needs-re-acquiring are never added together. test("counts keep reachable work and needs-re-acquiring apart", () => { const counts = emptyOperationCounts(); for (const state of [ "missing", "missing", "stale", "missing-input", "missing-input", "missing-input", "present", "not-applicable", ] as OperationClassification[]) { addOperationState(counts, state); } assert.deepEqual(counts, { missing: 2, stale: 1, partial: 0, missingInput: 3, deferred: 0, blocked: 0, }); // 3, not 6. Measured on the real corpus the difference is 835 vs 77,105, and // reporting the larger number is what would make every surface useless. assert.equal(reachableOperationWork(counts), 3); }); test("deferred is counted, and is NOT reachable work", async () => { const counts = emptyOperationCounts(); for (const state of [ "missing", "deferred", "deferred", "stale", ] as OperationClassification[]) { addOperationState(counts, state); } assert.deepEqual(counts, { missing: 1, stale: 1, partial: 0, missingInput: 0, deferred: 2, blocked: 0, }); // 2, not 4. This is the assertion that keeps a capped corpus from ever reading // as finished, and the one that fails if someone "tidies up" by folding // deferred into the total. assert.equal(reachableOperationWork(counts), 2); }); // --------------------------------------------------------------------------- // PARTIAL: some sections at the current identity, some not. test("digest: some sections fresh and some not is PARTIAL, not stale", async () => { // THE CASE THIS EXISTS FOR. `sections` is a setting and the pre-sweep decision // is to turn tags on; the moment that happens every already-digested video in // the corpus is part-done at once. Without the split all ~77,000 would read as // `stale` — indistinguishable on a stage card from a PROMPT_VERSION bump that // really did invalidate everything, when in fact digestVideo would regenerate // only the tags. assert.equal( await classifyDigest({ sidecar: digestSidecar(), sections: ["chapters", "tags"], }), "partial", ); // The same sidecar against the section list it was made for is DONE, which is // what makes the line above about the sections and not about the sidecar. assert.equal( await classifyDigest({ sidecar: digestSidecar(), sections: ["chapters"] }), "present", ); }); test("digest: partial is REACHABLE work, unlike deferred and blocked", () => { // The distinction is about what the work COSTS, never about whether the lane // can do it. A part-done video is dispatched exactly like a stale one, so the // sum must be unchanged by splitting them — otherwise a corpus mid-tags- // backfill reports less work than it has. const split = emptyOperationCounts(); for (const st of ["missing", "partial", "stale"] as OperationClassification[]) { addOperationState(split, st); } assert.equal(reachableOperationWork(split), 3); const blockedAndDeferred = emptyOperationCounts(); for (const st of ["blocked", "deferred"] as OperationClassification[]) { addOperationState(blockedAndDeferred, st); } assert.equal(reachableOperationWork(blockedAndDeferred), 0); }); test("reachableOperationWork survives a snapshot written before `partial`", () => { // Every snapshot on disk predates the field. Reading it as undefined and // adding it would produce NaN, which renders as "NaN" and sorts unpredictably // — strictly worse than under-reporting. const old = { missing: 2, stale: 1, missingInput: 5, deferred: 0, blocked: 0, } as unknown as ReturnType; assert.equal(reachableOperationWork(old), 3); }); // --------------------------------------------------------------------------- // backfillLaneEntriesOf: the read-side twin of backfillLaneOperations. test("backfillLaneEntriesOf keeps a digest entry OUT of the lane's sums", () => { // THE REGRESSION THIS WHOLE DESIGN EXISTS TO PREVENT. Four surfaces summed // Object.values(snapshot.backfill) on the assumption that the map WAS the // backfill lane. Now that the snapshot carries an entry per catalog operation, // that assumption would fold ~75,000 digest videos into the dashboard's // backfill instrument, the /channels speaker columns and the widget. const backfill = { diarization: { ...emptyOperationCounts(), missing: 3, ids: [], eligible: 3 }, digest: { ...emptyOperationCounts(), missing: 75_000, ids: [], eligible: 75_000, }, }; const lane = backfillLaneEntriesOf(backfill); assert.equal(lane.length, 1); assert.equal( lane.reduce((n, e) => n + reachableOperationWork(e), 0), 3, ); }); test("backfillLaneEntriesOf filters by the DECLARATION, not by a hardcoded id", () => { // `key !== "digest"` would pass the test above and leave the identical trap // armed for the next operation registered on a lane of its own. The rule is // the queue key, asked of the registry — so every lane kind is in, and an id // the catalog does not know is out. const laneIds = backfillLaneOperations(settingsWithDiarization()).map((k) => k.id); const backfill: Record & { ids: string[] }> = {}; for (const id of [...laneIds, "digest", "some-kind-from-a-newer-build"]) { backfill[id] = { ...emptyOperationCounts(), missing: 1, ids: [] }; } assert.equal(backfillLaneEntriesOf(backfill).length, laneIds.length); assert.ok(laneIds.length > 0, "expected at least one lane kind enabled"); // Empty and absent are both simply nothing, never a throw: a snapshot may // predate the field entirely. assert.deepEqual(backfillLaneEntriesOf(undefined), []); assert.deepEqual(backfillLaneEntriesOf({}), []); }); test("backfillLaneOperationEntriesOf applies the SAME filter, keyed by kind", () => { // The keyed form exists so the corpus-wide backfill card can say WHICH kind a // number came from: summed, this lane reads "77,952 reachable · 77,134 need // media", where 99.5% of the first is attribution-text at ~1 model call per // transcript CHUNK and all of the second is diarization at a few hundred audio // passes. A breakdown that re-derived its own filter is how a "Digest" row // ends up on the backfill card contradicting the figure above it — so the two // are one function, and this asserts they cannot drift. const laneIds = backfillLaneOperations(settingsWithDiarization()).map((k) => k.id); const backfill: Record< string, ReturnType & { ids: string[] } > = {}; for (const id of [...laneIds, "digest", "some-kind-from-a-newer-build"]) { backfill[id] = { ...emptyOperationCounts(), missing: 1, ids: [] }; } assert.deepEqual( backfillLaneOperationEntriesOf(backfill).map(([id]) => id).sort(), [...laneIds].sort(), ); // Exactly the entries backfillLaneEntriesOf returns, in the same order. assert.deepEqual( backfillLaneOperationEntriesOf(backfill).map(([, e]) => e), backfillLaneEntriesOf(backfill), ); assert.deepEqual(backfillLaneOperationEntriesOf(undefined), []); assert.deepEqual(backfillLaneOperationEntriesOf({}), []); }); test("presentOperationWork says UNKNOWN rather than zero on an old snapshot", () => { // A 0 here would render as "nothing digested" on a fully digested channel, // which is the most dangerous direction for a coverage number to be wrong. assert.equal( presentOperationWork({ ...emptyOperationCounts(), ids: [] }), null, ); assert.equal( presentOperationWork({ ...emptyOperationCounts(), missing: 2, blocked: 1, ids: [], eligible: 10, }), 7, ); }); // --- operationsForLane ------------------------------------------------------ // // The registry's half of the four-lane model: which operations a LANE draws. // Off the QUEUE KEY and off laneFor() before lane, because digest's key follows // the engine it is configured with. function settingsWithBackfillLane(): SiteSettings { const d = settingsWithDiarization(); const a = settingsWithAttribution(); return { ...d, attribution: a.attribution }; } test("bucketLaneOperationId is pauseLaneFor read backwards", () => { // The lane -> operation direction, off the SAME declaration // (ExternalOperation.runner) that pauseLaneFor reads the other way. Neither // spells "download" or "transcription" as a literal, so the two cannot drift. assert.equal(bucketLaneOperationId("download"), "download"); assert.equal(bucketLaneOperationId("transcription"), "transcription"); // The operation lanes draw from the registry, not from a single external // entry, so there is nothing for this to answer. assert.equal(bucketLaneOperationId("digest"), null); assert.equal(bucketLaneOperationId("backfill"), null); // Round-trips through the other direction for every lane that has one. for (const lane of ["download", "transcription"] as const) { assert.equal(pauseLaneFor(bucketLaneOperationId(lane)!), lane); } }); test("operationsForLane: the bucket lanes draw no registry operations", () => { // Slice 1.5 gave download and transcription a snapshot entry, but not a // registry one: they still have no state() and no run(), so this function — // which returns Operations — still has nothing to hand back. The lane asks // bucketLaneOperationId for its work list instead. const s = settingsWithBackfillLane(); assert.deepEqual(operationsForLane("transcription", s), []); assert.deepEqual(operationsForLane("download", s), []); }); test("operationsForLane: the backfill lane is its three kinds, in dependency order", () => { const s = settingsWithBackfillLane(); assert.deepEqual( operationsForLane("backfill", s).map((o) => o.id), ["diarization", "attribution-diarized", "attribution-text"], ); // And it is EMPTY when the features are off — a lane with nothing enabled // draws nothing rather than drawing a disabled kind's backlog. assert.deepEqual(operationsForLane("backfill", defaultSiteSettings()), []); }); test("operationsForLane: digest stays on the digest lane on BOTH its queues", () => { // THE laneFor CASE. `digest.remoteEnabled` moves the operation from // DIGEST_LOCAL_QUEUE to DIGEST_REMOTE_QUEUE. Both are the digest lane, so the // answer must not change — a queue-key map naming only the local key would // empty the lane the moment the metered engine was switched on. const local = defaultSiteSettings(); assert.equal(digestKind.laneFor!(local).queueKey, DIGEST_LOCAL_QUEUE); assert.deepEqual( operationsForLane("digest", local).map((o) => o.id), ["digest"], ); const remote: SiteSettings = { ...local, digest: { ...local.digest, remoteEnabled: true }, }; assert.equal(digestKind.laneFor!(remote).queueKey, DIGEST_REMOTE_QUEUE); assert.deepEqual( operationsForLane("digest", remote).map((o) => o.id), ["digest"], ); // And the backfill lane never picks it up on either engine. assert.deepEqual(operationsForLane("backfill", remote).map((o) => o.id), []); });