import { mkdir, readdir, readFile, rm, stat, writeFile } from "node:fs/promises"; import { test, expect } from "@playwright/test"; import type { APIRequestContext } from "@playwright/test"; import { baseUrl } from "./baseUrl"; import { channelStage, generateReport, pathExists, readJson, resetData, resolvePath, writeDigestVideo, writeSettings, } from "./helpers"; // Active jobs, straight off the JSON API. Reloading /jobs costs hundreds of // milliseconds per sample, which is wider than the overlap window some of these // assertions are trying to catch. async function activeJobs( request: APIRequestContext, ): Promise<{ kind: string; status: string; queueKey: string }[]> { const res = await request.get(`${baseUrl}/api/jobs/active`); const body = await res.json(); return Array.isArray(body) ? body : (body.jobs ?? []); } // The backfill lane: a declared kind, a lane with a fair share, and a visible // remainder. // // What these specs are for, in order of how much they would cost to get wrong: // // 1. THE RE-DOWNLOAD MODE MUST NOT LEAK. It writes media onto a disk with // ~45 GB free at 97% full, for a population of ~76,000 videos. A file left // behind after a failed item fills the disk. Three specs pin it: gone after // success, gone after FAILURE, and kept only when the operator said so. // 2. THE TWO NUMBERS MUST STAY APART. Reachable work and needs-re-acquiring // differ by ~91x on the real corpus; a surface that adds them is useless // the day it ships. // 3. THE LANE MUST BE CONCURRENT with transcription — that is the entire // reason it has its own queueKey — and must survive a restart. // // The diarization engine is faked via DIARIZE_BIN (e2e/fixtures/bin/ // fake-diarize.mjs), as in diarization.spec.ts. const SLUG = "test-transcribe"; function dataRel(videoId: string, file: string): string { return `test-transcripts/channels/${SLUG}/data/${videoId}/${file}`; } const BASE_SETTINGS = { adminTitle: "Test Admin", maxTranscriptPageBytes: 8388608, sleepBetweenDownloadsSeconds: 0, minFreeDiskGB: 0, verifyAvailabilityBeforeClean: false, syncScheduler: { fullSweepIntervalMinutes: 0 }, }; function backfillSettings(over: { diarization?: Record; backfill?: Record; minFreeDiskGB?: number; } = {}) { return { ...BASE_SETTINGS, ...(over.minFreeDiskGB !== undefined ? { minFreeDiskGB: over.minFreeDiskGB } : {}), diarization: { enabled: true, // OFF: the whole point of a backfill lane is that the capture did not // happen inline, so these specs must not have it quietly done for them. inlineAfterTranscribe: false, threshold: 0.5, threads: 1, python: "python3", // Any non-empty pair — an empty one reports "not configured" and the lane // never runs. segModel: "/dev/null", embModel: "/dev/null", concurrency: 1, ...over.diarization, }, backfill: { concurrency: 1, allowRedownload: false, ...over.backfill, }, // THE LANE'S GATE, SPELLED. It used to be spelled by `backfill.enabled: // true` in the block above — the inverted retired field, where `true` meant // NOT held. S0-pause deleted it, and the lane's gate defaults SHUT // (`defaultHeldFor`), so a fixture that wants the lane to run says so on // the lane. autoQueue: { backfill: { held: false } }, }; } async function seedTranscript(videoId: string): Promise { await writeFile( resolvePath(dataRel(videoId, "transcript.json")), '{"transcription":[{"text":"hi"}]}\n', ); } // Real audio files in a video dir — what the re-download specs measure. async function audioFiles(videoId: string): Promise { const dir = resolvePath( `test-transcripts/channels/${SLUG}/data/${videoId}`, ); const entries = await readdir(dir).catch(() => [] as string[]); return entries.filter( (e) => e.startsWith("audio.") && !e.endsWith(".info.json"), ); } // The backfill lane as settings.json spells it — the persisted intent the boot // hook reads. It was `backfill.sweepEnabled` until slice 1.3; a lane is armed by // its policy now, and its SCOPE is the tree beside the switch, so both are read // together here rather than one being watched while the other silently moves. type LanePolicy = { enabled?: boolean; order?: string; root?: { children?: { id?: string; match?: Record }[] }; }; async function backfillLane(): Promise { const s = await readJson<{ autoQueue?: { backfill?: LanePolicy } }>( "test-settings.json", ).catch(() => ({}) as { autoQueue?: { backfill?: LanePolicy } }); return s.autoQueue?.backfill ?? {}; } async function laneArmed(): Promise { return (await backfillLane()).enabled ?? false; } const SLOW = 120_000; // (1) The stage card reports the two populations SEPARATELY. test("the stage card separates reachable work from what needs its media back", async ({ page, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); // vidA keeps its audio (reachable now); vidB is transcribed with its audio // already cleaned away (needs re-acquiring). That 1-vs-1 split is the corpus's // 835-vs-76,270 in miniature, and the card must never show "2". await seedTranscript("vidA"); await seedTranscript("vidB"); for (const f of await audioFiles("vidB")) { await rm(resolvePath(dataRel("vidB", f)), { force: true }); } await rm(resolvePath(`test-transcripts/channels/${SLUG}/snapshot.json`), { force: true, }); await generateReport(page, SLUG); await page.goto(channelStage(SLUG, "speakers")); const section = page.getByLabel("speakers section"); await expect(section.getByLabel("speakers reachable")).toContainText( "1 video can be worked on now", ); await expect(section.getByLabel("speakers needs re-acquiring")).toContainText( "1", ); // The id list is the reachable half only. await expect( section.getByLabel("video needing speaker work vidA"), ).toBeVisible(); await expect( section.getByLabel("video needing speaker work vidB"), ).toHaveCount(0); }); // (2) The lane runs, and only over what it can reach. test("running the lane captures the reachable video and skips the one with no media", async ({ page, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); await seedTranscript("vidA"); await seedTranscript("vidB"); for (const f of await audioFiles("vidB")) { await rm(resolvePath(dataRel("vidB", f)), { force: true }); } await page.goto(channelStage(SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); await expect .poll(async () => pathExists(dataRel("vidA", "diarization.json")), { timeout: 60_000, }) .toBe(true); // Never touched: re-download is off, so a video whose input is gone is // COUNTED, not attempted. expect(await pathExists(dataRel("vidB", "diarization.json"))).toBe(false); await expect(page.getByLabel("Run speaker work output")).toContainText( "need their media re-acquired", { timeout: 30_000 }, ); }); // (3) A stale sidecar is redone. This is the freshness comparator end-to-end: // before it, diarizeOne short-circuited on mere existence and a threshold change // left the whole corpus looking done. test("a sidecar from a different threshold is regenerated", async ({ page }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); await seedTranscript("vidA"); // Written at a threshold the current settings no longer use. await writeFile( resolvePath(dataRel("vidA", "diarization.json")), JSON.stringify({ videoId: "vidA", generatedAt: "2020-01-01T00:00:00.000Z", speakers: 1, turns: [{ start: 0, end: 1, speaker: 0 }], engine: { engine: "fake-diarize", segmentationModel: "null", embeddingModel: "null", threshold: 0.1, }, }) + "\n", ); await page.goto(channelStage(SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); await expect .poll( async () => { const r = await readJson<{ engine: { threshold: number } }>( dataRel("vidA", "diarization.json"), ).catch(() => null); return r?.engine.threshold ?? null; }, { timeout: 60_000 }, ) .toBe(0.5); }); // (4) THE RE-DOWNLOAD, THE SAFE PATH: fetched, used, and gone again. test("re-acquired media is deleted after a successful backfill", async ({ page, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings( backfillSettings({ backfill: { allowRedownload: true } }), ); await seedTranscript("vidB"); for (const f of await audioFiles("vidB")) { await rm(resolvePath(dataRel("vidB", f)), { force: true }); } expect(await audioFiles("vidB")).toEqual([]); await page.goto(channelStage(SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); // The work landed… await expect .poll(async () => pathExists(dataRel("vidB", "diarization.json")), { timeout: 60_000, }) .toBe(true); // …and the media it needed did NOT survive it. This is the property the whole // mode is built around; without it the corpus-wide run fills the disk. await expect.poll(async () => audioFiles("vidB"), { timeout: 30_000 }).toEqual( [], ); }); // (5) THE RE-DOWNLOAD, THE FAILING PATH — the one a `finally` exists for. vidC // is the fake diarizer's designated failure, so the backfill throws AFTER the // media has been fetched. The file must still be gone. test("re-acquired media is deleted even when the backfill fails", async ({ page, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings( backfillSettings({ backfill: { allowRedownload: true } }), ); await seedTranscript("vidC"); for (const f of await audioFiles("vidC")) { await rm(resolvePath(dataRel("vidC", f)), { force: true }); } await page.goto(channelStage(SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); await expect(page.getByLabel("Run speaker work output")).toContainText( /failed|Removed re-acquired/, { timeout: 60_000 }, ); // No sidecar (the engine failed) AND no leftover media. Both halves matter: a // failure that leaves the file behind is exactly the leak. expect(await pathExists(dataRel("vidC", "diarization.json"))).toBe(false); await expect.poll(async () => audioFiles("vidC"), { timeout: 30_000 }).toEqual( [], ); }); // (6) …unless the operator said to keep it. do-not-clean is the one exception, // honoured here for the same reason every cleanup controller honours it. test("re-acquired media is KEPT when the video is marked do-not-clean", async ({ page, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings( backfillSettings({ backfill: { allowRedownload: true } }), ); await seedTranscript("vidB"); for (const f of await audioFiles("vidB")) { await rm(resolvePath(dataRel("vidB", f)), { force: true }); } await writeFile( resolvePath(dataRel("vidB", "do-not-clean.json")), JSON.stringify({ setAt: new Date().toISOString(), note: "e2e" }) + "\n", ); await page.goto(channelStage(SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); await expect .poll(async () => pathExists(dataRel("vidB", "diarization.json")), { timeout: 60_000, }) .toBe(true); // Kept, and SAID SO in the log — an unexplained file on a full disk is how a // leak gets discovered the hard way. await expect(page.getByLabel("Run speaker work output")).toContainText( "do not clean", { timeout: 30_000 }, ); expect((await audioFiles("vidB")).length).toBeGreaterThan(0); }); // (7) The disk floor refuses to START. Not "stops partway" — a preflight, per // item, so a long run's twentieth video cannot inherit the first one's headroom. test("the disk floor refuses to re-acquire anything", async ({ page }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings( backfillSettings({ backfill: { allowRedownload: true }, // Nothing has this much free. minFreeDiskGB: 100_000_000, }), ); await seedTranscript("vidB"); for (const f of await audioFiles("vidB")) { await rm(resolvePath(dataRel("vidB", f)), { force: true }); } await page.goto(channelStage(SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); await expect(page.getByLabel("Run speaker work output")).toContainText( "below the", { timeout: 60_000 }, ); // Nothing fetched, nothing produced. expect(await audioFiles("vidB")).toEqual([]); expect(await pathExists(dataRel("vidB", "diarization.json"))).toBe(false); }); // (8) THE LANE ARMS FROM THE DASHBOARD, AND DISARMING LEAVES ITS RULES ALONE. // // This was the sweep's arm/disarm pair. The two bugs it pinned were "the flag is // persisted without the scope" (a bounded run comes back corpus-wide after a // restart) and "stopping leaves the scope behind" (the next run is silently // narrowed). Both are gone by construction now — the scope IS the tree, written // in the same settings write as the switch — so what has to be pinned instead is // the other half of that: disarming must NOT clear the tree, because an operator // authored it and switching a lane off is not throwing it away. test("the lane arms from the dashboard and disarms without losing its rules", async ({ page, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); await page.goto("/"); const arm = page.getByRole("button", { name: "arm backfill lane" }); await expect(arm).toBeEnabled({ timeout: 30_000 }); await arm.click(); await expect.poll(laneArmed, { timeout: 30_000 }).toBe(true); // An unscoped arm is one catch-all leaf — what an empty `sweepChannels` meant. expect( (await backfillLane()).root?.children?.map((c) => c.match), ).toEqual([{ type: "all" }]); const disarm = page.getByRole("button", { name: "disarm backfill lane" }); await expect(disarm).toBeEnabled({ timeout: 30_000 }); await disarm.click(); await expect .poll( async () => { const lane = await backfillLane(); return { enabled: lane.enabled ?? false, rules: lane.root?.children?.length ?? 0, }; }, { timeout: 30_000 }, ) .toEqual({ enabled: false, rules: 1 }); // The sweep's own controls are gone from this surface, not merely relabelled // — BOTH halves, and both aria-labels, because the deck drew whichever half // matched the state and a test that only names the arm half would miss a // disarm button left over an action that no longer exists. for (const name of [ "start backfill sweep", "stop backfill sweep", "start digest sweep", "stop digest sweep", ]) { await expect(page.getByRole("button", { name })).toHaveCount(0); } await expect(page.getByRole("button", { name: /Stop sweeping/i })).toHaveCount( 0, ); }); // (8c) THE SCOPE IS SET FROM A SCREEN, AND IT IS A TREE. // // `backfill.sweepKinds` was honoured by the run since the sweep was written and // settable by nothing but a hand-edit of settings.json. That mattered rather // than merely being untidy: an unscoped run does EVERY enabled operation, and on // the live corpus that includes speaker-names-from-the-transcript at ~1 model // call per transcript chunk — on the order of 194,000 calls. "Diarization only" // had no expression in the product. // // It has one now, and it is the same control every lane uses: a rule that names // an operation. The scope is persisted as the lane's ROOT, in the same write as // the switch — so a scope that is not recorded is not a thing that can happen. test("a lane can be scoped to one operation from its console, and the scope is a tree", async ({ page, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings({ ...backfillSettings(), // Three operations, so there is a choice to make — and so the ladder draws // an operation select at all. With diarization alone it draws none, because // one operation is not a choice and a leaf naming none claims it anyway. attribution: { enabled: true, diarizedEnabled: true, textOnlyEnabled: true, }, }); // ONE OPERATION, ONE PAGE. Diarization's page renders the lane it is // dispatched by — the SHARED backfill lane, which is also attribution's — so // the rules below are the same rules, reached by an address instead of a // dropdown. await page.goto("/operations/diarization"); const lane = page.locator('section[data-lane="backfill"]'); await expect(lane).toHaveAttribute("data-hydrated", "true", { timeout: 30_000, }); // The retired panel took its attribute with it. await expect(page.locator("section[data-sweep-lane]")).toHaveCount(0); // WHAT THE LANE WOULD DO NEXT is drawn before anything is armed. That is the // sweep plan's job, done by the thing that actually dispatches. await expect(lane.getByText("Next up")).toBeVisible(); await lane.getByLabel("rule operation").selectOption("diarization"); await lane.getByLabel("Enable auto-backfill").check(); await lane.getByRole("button", { name: "Save policy" }).click(); await expect .poll( async () => { const b = await backfillLane(); return { enabled: b.enabled ?? false, operations: (b.root?.children ?? []).map( (c) => (c.match ?? {}).operation ?? null, ), }; }, { timeout: 30_000 }, ) .toEqual({ enabled: true, operations: ["diarization"] }); // Switch it back off so an armed lane does not bleed into the next spec. await lane.getByLabel("Enable auto-backfill").uncheck(); await lane.getByRole("button", { name: "Save policy" }).click(); await expect.poll(laneArmed, { timeout: 30_000 }).toBe(false); }); // (8b) …AND SURVIVES A RESTART. This is the reason the switch is persisted at // all: a corpus pass is days of work and will outlive several restarts by // construction. The suite cannot restart the dev server mid-run, so // /api/test/resume-lane does the two halves a restart does — drop the live // runner, then call the boot hook — and the lane must come back from the // persisted policy alone. test("an armed lane comes back after a restart", async ({ page, request }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); await page.goto("/"); const arm = page.getByRole("button", { name: "arm backfill lane" }); await expect(arm).toBeEnabled({ timeout: 30_000 }); await arm.click(); await expect.poll(laneArmed, { timeout: 30_000 }).toBe(true); const res = await request.get( `${baseUrl}/api/test/resume-lane?lane=backfill`, ); const body = (await res.json()) as { cancelled: string | null; resumed: string | null; }; // A NEW job id: the old one is gone (as a restart would leave it) and the hook // started a fresh runner rather than the test merely observing the old one. expect(body.resumed).toBeTruthy(); expect(body.resumed).not.toBe(body.cancelled); // Clean up, so an armed lane does not bleed into the next spec. await page.goto("/"); const disarm = page.getByRole("button", { name: "disarm backfill lane" }); await expect(disarm).toBeEnabled({ timeout: 30_000 }); await disarm.click(); await expect.poll(laneArmed, { timeout: 30_000 }).toBe(false); }); // (9) THE WHOLE POINT OF THE SEPARATE QUEUE KEY: a backfill job and a // transcription job run AT THE SAME TIME. The registry submits every named queue // at concurrency 1, so if these shared a key one would sit queued behind the // other and the "lane" would be a lie. // // The transcription is deliberately a SLOWOP video (fake-whisper runs ~7s of // real wall-time for a dir whose name contains "slowop"), because with the // instant fixtures both jobs finish before anything can observe the overlap — // the test would be measuring the poller, not the queue. Polled over the JSON // API rather than by reloading /jobs for the same reason. test("a backfill runs concurrently with a transcription", async ({ page, request, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); // Backfill work: a transcribed video that still has its audio. await seedTranscript("vidA"); // Transcription work that LASTS: a fresh dir with audio and no transcript. const slowDir = resolvePath( `test-transcripts/channels/${SLUG}/data/slowop1`, ); await mkdir(slowDir, { recursive: true }); await writeFile(`${slowDir}/audio.mp3`, "fake audio slowop1\n"); // Transcription first — it is the one that has to still be running when the // backfill lands. await page.goto(channelStage(SLUG, "transcribe")); await page .getByRole("button", { name: "Transcribe missing", exact: true }) .click(); await expect .poll(async () => activeJobs(request).then((j) => j.some((x) => x.kind === "whisper-all" && x.status === "running"), ), { timeout: 30_000 }) .toBe(true); await page.goto(channelStage(SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); // Both running at once. This is the assertion the separate queueKey exists for. const seen = await expect .poll( async () => { const jobs = await activeJobs(request); const backfill = jobs.find((j) => j.kind === "backfill-channel"); const transcribe = jobs.find((j) => j.kind === "whisper-all"); // A backfill that is QUEUED while a transcription runs is the exact // failure mode a shared key would produce, so record it distinctly. if (backfill?.status === "running" && transcribe?.status === "running") { // And the reason it can be: two different keys. expect(backfill.queueKey).not.toBe(transcribe.queueKey); return "concurrent"; } return backfill?.status ?? "absent"; }, { timeout: 30_000, intervals: [100] }, ) .toBe("concurrent"); // SETTLE BEFORE LEAVING. Ending here with a ~7s slowop whisper still running // leaks its child into the next spec — the teardown reaper catches it, but // "the anti-orphan layers did not hold" is a failure of this spec's hygiene, // not a safety net doing its job. await expect .poll( async () => { const jobs = await activeJobs(request); return jobs.filter( (j) => j.kind === "whisper-all" || j.kind === "backfill-channel", ).length; }, { timeout: 60_000, intervals: [250] }, ) .toBe(0); return seen; }); // THE PROPERTY REGISTERING DIGEST MUST NOT BREAK. // // Digest is now a registered operation in the same table as diarization, and // the obvious way to "unify" them would have been to put both on BACKFILL_QUEUE. // That would be a regression, not a simplification: registry.ts runs every // non-empty queueKey at concurrency 1, so one shared key makes the GPU digest // lane wait on CPU diarization and vice versa — across a sweep measured in // weeks. The lane is declared per operation precisely so this keeps working. test("a digest runs concurrently with a backfill, not behind it", async ({ page, request, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings({ ...backfillSettings(), digest: { localAppId: "ollama-direct", remoteAppId: "claude-code", sections: ["chapters"], }, }); // ENOUGH backfill work that the lane is provably still running when the // digest is submitted. The fake diarizer takes ~0.2s per video, so the three // fixture videos finish in well under a second — faster than a click — and // the test would measure nothing. 40 dirs buys ~8s of overlap window. for (const id of ["vidA", "vidB"]) await seedTranscript(id); for (let i = 0; i < 40; i++) { const id = `slowdiar${String(i).padStart(2, "0")}`; const dir = resolvePath(`test-transcripts/channels/${SLUG}/data/${id}`); await mkdir(dir, { recursive: true }); await writeFile(`${dir}/audio.m4a`, `fake audio ${id}\n`); await writeFile( `${dir}/transcript.json`, '{"transcription":[{"text":"hi"}]}\n', ); } // Digest work: a video with a real normalized transcript. await writeDigestVideo({ channelSlug: SLUG, videoId: "digestme" }); await generateReport(page, SLUG); await page.goto(channelStage(SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); await expect .poll( async () => activeJobs(request).then((j) => j.some((x) => x.kind === "backfill-channel" && x.status === "running"), ), { timeout: 30_000 }, ) .toBe(true); await page.goto(channelStage(SLUG, "digest")); await page.getByRole("button", { name: "Digest channel" }).click(); const seen = await expect .poll( async () => { const jobs = await activeJobs(request); const backfill = jobs.find((j) => j.kind === "backfill-channel"); const digest = jobs.find((j) => j.kind.startsWith("digest-channel")); // A digest sitting QUEUED behind a running backfill is exactly the // failure a shared queue key produces, so record that state distinctly // rather than letting it time out as a generic "never ran". if (digest?.status === "running" && backfill?.status === "running") { expect(digest.queueKey).not.toBe(backfill.queueKey); return "concurrent"; } return digest?.status ?? "absent"; }, { timeout: 30_000, intervals: [100] }, ) .toBe("concurrent"); // Settle, so neither lane's children leak into the next spec. await expect .poll( async () => { const jobs = await activeJobs(request); return jobs.filter( (j) => j.kind.startsWith("digest-channel") || j.kind === "backfill-channel", ).length; }, { timeout: 90_000, intervals: [250] }, ) .toBe(0); return seen; }); // (10) The operation's own page lists the channel in the sweep plan and keeps // the two populations apart — reachable now, and the ones whose media has to // come back first. // // This used to assert /actionable's speakers section, which also carried the // property that the page was not hidden behind a page-wide "nothing pending" // gate. That gate went with the page: an operation page has no such gate to // hide behind, so there is nothing left to assert there. test("the diarization page shows the channel and both populations", async ({ page, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); await seedTranscript("vidA"); await seedTranscript("vidB"); for (const f of await audioFiles("vidB")) { await rm(resolvePath(dataRel("vidB", f)), { force: true }); } await rm(resolvePath(`test-transcripts/channels/${SLUG}/snapshot.json`), { force: true, }); await generateReport(page, SLUG); await page.goto("/operations/diarization"); const lane = page.locator('section[data-lane="backfill"]'); await expect(lane).toBeVisible(); // The channel is what the lane would dispatch next — the plan the sweep's // itinerary used to draw, answered by the thing that actually dispatches. await expect(lane.getByText(new RegExp(`${SLUG}/`))).toBeVisible({ timeout: 30_000, }); // vidB's audio was removed above, so the re-acquire population is non-zero — // and it is STATED SEPARATELY on the rail's row, never added to the count // beside it. On the measured corpus these two are 91x apart. const row = page.locator('li[data-operation="diarization"]'); await expect(row).toContainText(/reachable/); await expect(row).toContainText(/no media/); }); // PAUSE IS A HOLD, NOT A STOP, and it is reachable from where the work is // watched. The hold has always existed — the backfill lane's limit() re-reads // the lane's gate at dispatch and returns 0, so the pool idle-waits and the job // keeps its place — but the only way to set it was the Settings page, which is // a strange place to look for a control over a job you are watching run on the // dashboard. // // The button writes THE SAME KEY as that checkbox rather than a second // `backfillPaused` flag, so the two cannot drift; this asserts the key, // not just the label, for exactly that reason. It also asserts the button // renders WITHOUT a sweep armed — an operator pauses a hand-clicked // backfill-channel job too, not only an armed lane. test("the dashboard pauses and resumes the backfill lane", async ({ page }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); // THE LANE'S GATE, on the lane. Slice 1.4 moved it off the inverted // `backfill.enabled` (deleted by S0-pause) and onto // `autoQueue.backfill.held`, so the polarity here is the plain one: held // means held. const laneHeld = async () => ( await readJson<{ autoQueue?: { backfill?: { held?: boolean } } }>( "test-settings.json", ).catch(() => ({}) as { autoQueue?: { backfill?: { held?: boolean } } }) ).autoQueue?.backfill?.held ?? false; await page.goto("/"); const pause = page.getByRole("button", { name: "pause backfill" }); await expect(pause).toBeEnabled({ timeout: 30_000 }); await pause.click(); await expect.poll(laneHeld, { timeout: 30_000 }).toBe(true); // THE LANE'S ARM SWITCH IS UNTOUCHED: pausing must not disarm the lane, or // "pause" would quietly become "stop" and days of queued work would need // re-arming. (The fixture leaves it disarmed, so `false` is what "unchanged" // looks like here.) expect(await laneArmed()).toBe(false); const resume = page.getByRole("button", { name: "resume backfill" }); await expect(resume).toBeEnabled({ timeout: 30_000 }); await resume.click(); await expect.poll(laneHeld, { timeout: 30_000 }).toBe(false); }); // ONE GATE, TWO SURFACES. /operations/diarization and the dashboard draw the // SAME control over the SAME lane — three speaker operations share one backfill // queue and therefore one pause, so an operation page is not a per-operation // switch. Before slice 7 these were two components with different aria-labels // and opposite emphasis, and only one of them was tested. // // The assertion is deliberately the FIELD ON DISK plus the other surface's // label, not two clicks on one page: what would break silently is the two // surfaces disagreeing about which lane they hold. test("an operation page holds the same lane the dashboard does", async ({ page, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); // THE LANE'S GATE, on the lane. Slice 1.4 moved it off the inverted // `backfill.enabled` (deleted by S0-pause) and onto // `autoQueue.backfill.held`, so the polarity here is the plain one: held // means held. const laneHeld = async () => ( await readJson<{ autoQueue?: { backfill?: { held?: boolean } } }>( "test-settings.json", ).catch(() => ({}) as { autoQueue?: { backfill?: { held?: boolean } } }) ).autoQueue?.backfill?.held ?? false; await page.goto("/operations/diarization"); // Scoped to the lane's own section so this cannot accidentally match a // control the rail draws elsewhere on the page. const lane = page.locator('section[data-lane="backfill"]'); const hold = lane.getByRole("button", { name: "pause backfill" }); // The hydration wait: LaneActionButton is disabled until React attaches, and // a click before that fires nothing at all. await expect(hold).toBeEnabled({ timeout: 30_000 }); await hold.click(); await expect.poll(laneHeld, { timeout: 30_000 }).toBe(true); // A hold is not a stop: the lane's arm switch is untouched. expect(await laneArmed()).toBe(false); // The dashboard shows the same lane held, and releases it. await page.goto("/"); const resumeOnDashboard = page.getByRole("button", { name: "resume backfill", }); await expect(resumeOnDashboard).toBeEnabled({ timeout: 30_000 }); await resumeOnDashboard.click(); await expect.poll(laneHeld, { timeout: 30_000 }).toBe(false); }); // (11) THE RE-DOWNLOAD ON A SUBTITLE CHANNEL: audio, not captions. // // A `handling: "youtube"` channel downloads with --skip-download --write-subs // --write-auto-subs. Run a re-acquire with that config and yt-dlp re-fetches the // captions the video already has, rewrites metadata.info.json, and lands nothing // a diarizer can read — which is what happened to ~16,000 videos on eight // channels between 2026-08-22 and 08-26: zero diarizations, and every touched // video left reading `deferred` to the digest lane because the metadata rewrite // had made transcript.cues.json stale. // // Both halves are pinned here: the transcribe override that makes there be audio // at all, and the re-normalize that keeps the cues fresh afterwards. const YT_SLUG = "test-yt-subs"; const YT_ROOT = `test-transcripts/channels/${YT_SLUG}`; function ytRel(videoId: string, file: string): string { return `${YT_ROOT}/data/${videoId}/${file}`; } // A youtube-handling channel with captions and NO audio — the shape every video // on a subtitle channel has, and the shape that makes diarization missing-input. // audioFormat is pinned so the forced transcribe-handling download lands on // audio.mp3, exactly as auto-subs-replace.spec.ts seeds it. const SEEDED_VTT = "WEBVTT\n\n00:00:00.000 --> 00:00:05.000\nSeeded caption line one.\n\n" + "00:00:05.000 --> 00:00:10.000\nSeeded caption line two.\n"; // Shaped after real yt-dlp --write-auto-subs output (a copy of // auto-subs-replace.spec.ts's). It matters WHICH vtt a case seeds: the hand-off // asks isAutoSubsOnly, whose 4 KB sniff is what decides ASR vs manual, and // SEEDED_VTT above carries no ASR fingerprints at all — provenance then falls // back to metadata.info.json, which the fake yt-dlp's re-acquire rewrites // WITHOUT automatic_captions. const ASR_VTT = `WEBVTT Kind: captions Language: en 00:00:00.030 --> 00:00:03.919 align:start position:0% so<00:00:00.719> today<00:00:01.199> we're<00:00:01.439> going<00:00:01.680> to 00:00:03.919 --> 00:00:03.929 align:start position:0% so today we're going to 00:00:03.929 --> 00:00:07.070 align:start position:0% so today we're going to talk<00:00:04.320> about<00:00:04.639> the<00:00:04.879> whole<00:00:05.199> thing `; // Audio this channel's video actually has on disk right now. async function ytAudioFiles(videoId: string): Promise { const entries = await readdir( resolvePath(`${YT_ROOT}/data/${videoId}`), ).catch(() => [] as string[]); return entries.filter( (e) => e.startsWith("audio.") && !e.endsWith(".info.json"), ); } async function seedSubtitleChannel( videoId: string, vtt: string = SEEDED_VTT, ): Promise { const dir = resolvePath(`${YT_ROOT}/data/${videoId}`); await mkdir(dir, { recursive: true }); await writeFile( resolvePath(`${YT_ROOT}/config.json`), JSON.stringify({ handling: "youtube", name: "Subtitle test channel", url: "https://www.youtube.com/@subs/videos", audioFormat: "mp3", }) + "\n", ); // findVideoSourceUrl resolves the id's URL from the stored playlist. await writeFile( resolvePath(`${YT_ROOT}/playlist`), `https://www.youtube.com/watch?v=${videoId}\n`, ); await writeFile(`${dir}/transcript.en.vtt`, vtt); await writeFile( `${dir}/metadata.info.json`, JSON.stringify({ id: videoId, title: `Synthetic ${videoId}`, upload_date: "20240101", duration: 60, extractor_key: "Youtube", webpage_url: `https://www.youtube.com/watch?v=${videoId}`, subtitles: {}, automatic_captions: { en: [{ ext: "vtt", url: "fake://subs" }] }, }) + "\n", ); await fetch(`${baseUrl}/api/test/invalidate-cache`).catch(() => {}); } test("re-acquiring on a subtitle channel downloads audio and re-normalizes the cues", async ({ page, }) => { test.setTimeout(SLOW); const VID = "ytsubs000001"; await resetData(); await writeSettings( backfillSettings({ backfill: { allowRedownload: true } }), ); await seedSubtitleChannel(VID); expect(await pathExists(ytRel(VID, "audio.mp3"))).toBe(false); await generateReport(page, YT_SLUG); await page.goto(channelStage(YT_SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); // The lane got what it needed — which it can only have done by downloading // AUDIO. With the channel's own config this assertion is what failed silently // 16,000 times as "nothing usable landed". await expect .poll(async () => pathExists(ytRel(VID, "diarization.json")), { timeout: 60_000, }) .toBe(true); // …and the audio still does not survive the item that fetched it (case (4)'s // property, on a channel that never had audio to begin with). await expect .poll(async () => ytAudioFiles(VID), { timeout: 30_000 }) .toEqual([]); // NO SUBTITLE RE-FETCH. The captions on disk are byte-identical to what was // seeded; the fake yt-dlp's youtube-handling mode writes its own transcript, // so any drift here means the override did not apply. expect(await readFile(resolvePath(ytRel(VID, "transcript.en.vtt")), "utf8")) .toBe(SEEDED_VTT); // The override, as the download itself recorded it. const outcome = await readJson<{ attempts: { kind: string; handling: string }[]; }>(ytRel(VID, "download-outcome.json")); const attempts = outcome.attempts ?? []; expect(attempts.length).toBeGreaterThan(0); expect(attempts[attempts.length - 1]?.handling).toBe("transcribe"); // THE SECOND HALF: the fetch rewrote metadata.info.json, and the cues were // re-normalized after it. Without this the next lane in the same sweep skips // the video as no-transcript and the digest lane defers it — 16,081 videos // corpus-wide, until an operator runs Normalize. expect(await pathExists(ytRel(VID, "transcript.cues.json"))).toBe(true); const [cuesStat, metaStat] = await Promise.all([ stat(resolvePath(ytRel(VID, "transcript.cues.json"))), stat(resolvePath(ytRel(VID, "metadata.info.json"))), ]); expect(cuesStat.mtimeMs).toBeGreaterThanOrEqual(metaStat.mtimeMs); }); // (12) THE HAND-OFF. The audio case (11) deletes is exactly what // autoQueue.transcription draws from `downloadedAutoSubsOnly` (ASR VTT AND audio // present), and the snapshot regenerates ~1 s after any unit on the channel // finishes — with no per-video lock anywhere. Deleting it under a running // whisper costs that video's run and blacklists it in failed-transcriptions. So // when the policy WOULD take it, the backfill hands it over instead. // // The runner is deliberately NOT started here: startAutoRunnersIfEnabled runs // only at boot (editor/instrumentation.ts), and the invalidate-cache route // starts nothing, so the policy is READ but nothing races the assertions. Once // the audio is in that bucket the rest of the lane is what // auto-subs-replace.spec.ts already walks end to end. test("re-acquired audio is handed to auto-transcribe when the policy would replace the auto-captions", async ({ page, }) => { test.setTimeout(SLOW); const VID = "ytsubs000002"; await resetData(); await writeSettings({ ...backfillSettings({ backfill: { allowRedownload: true } }), autoQueue: { // Named here because this literal REPLACES backfillSettings' autoQueue, // and the lane's gate defaults shut. backfill: { held: false }, transcription: { enabled: true, maxWorkers: 1, replaceAutoSubs: true, root: { id: "root", mode: "strict", children: [{ id: "leaf-all", match: { type: "all" } }], }, }, download: {}, }, }); // ASR-shaped captions: the 4 KB sniff is what makes this video ASR-only, and // therefore a candidate for the bucket. await seedSubtitleChannel(VID, ASR_VTT); await generateReport(page, YT_SLUG); await page.goto(channelStage(YT_SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); // The diarization still happens — the hand-off changes what happens to the // audio AFTERWARDS, not whether the backfill does its work. await expect .poll(async () => pathExists(ytRel(VID, "diarization.json")), { timeout: 60_000, }) .toBe(true); // Said so, per video and in the run's one-line summary. An unexplained file on // a full disk is how a leak gets discovered the hard way, so "kept" must never // be silent. await expect(page.getByLabel("Run speaker work output")).toContainText( "handed to auto-transcribe", { timeout: 30_000 }, ); await expect(page.getByLabel("Run speaker work output")).toContainText( "0 cleaned up, 1 handed to auto-transcribe", { timeout: 30_000 }, ); // The file itself survives the item that fetched it — the ONE case where that // is correct. expect((await ytAudioFiles(VID)).length).toBeGreaterThan(0); // …and it lands where the runner will find it. The snapshot regenerates at // job end (operationJobs.ts), so this is the hand-off actually completing // rather than a file left lying around. await expect .poll( async () => { const snap = await readJson<{ buckets?: { downloadedAutoSubsOnly?: string[] }; }>(`${YT_ROOT}/snapshot.json`).catch(() => null); return snap?.buckets?.downloadedAutoSubsOnly ?? []; }, { timeout: 60_000 }, ) .toContain(VID); }); // (12b) …and with no such policy, today's behaviour, through the same new code // path. The keep decision has one "yes" and six "no"s; this is the "no" that // every non-opted-in corpus gets. test("re-acquired audio is still removed when no auto-transcribe policy would take it", async ({ page, }) => { test.setTimeout(SLOW); const VID = "ytsubs000003"; await resetData(); await writeSettings( backfillSettings({ backfill: { allowRedownload: true } }), ); await seedSubtitleChannel(VID, ASR_VTT); await generateReport(page, YT_SLUG); await page.goto(channelStage(YT_SLUG, "speakers")); await page .getByRole("button", { name: "Run speaker work", exact: true }) .click(); await expect .poll(async () => pathExists(ytRel(VID, "diarization.json")), { timeout: 60_000, }) .toBe(true); await expect(page.getByLabel("Run speaker work output")).toContainText( "Removed re-acquired media", { timeout: 30_000 }, ); expect( await page.getByLabel("Run speaker work output").textContent(), ).not.toContain("handed to"); await expect .poll(async () => ytAudioFiles(VID), { timeout: 30_000 }) .toEqual([]); }); // (N) THE REGRESSION THIS STEP'S DESIGN EXISTS TO PREVENT. // // The channel snapshot now carries a work-list entry for EVERY catalog // operation, not just the ones on the backfill lane — that is what makes the // digest kind's state() reachable by anything at all. Four surfaces used to sum // `Object.values(snapshot.backfill)` on the assumption that the map WAS the // lane: the channel Backfill card, the /channels speaker columns and the // widget sync payload. On the measured corpus the digest entry is ~75,000 // videos against the lane's 9, so getting this wrong does not shift a number // slightly — it replaces it. // // The spec asserts BOTH halves, and the first half is what stops it passing // vacuously: if the digest entry were simply never written, every "unchanged" // assertion below would still pass while the feature was missing entirely. test("a digest entry in the snapshot does not move the backfill instrument", async ({ page, request, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); // Both videos get a real, digestable transcript (fresh cues.json against a // backdated vtt — a stale one would classify `deferred` and prove nothing). // vidA keeps its audio, so diarization is reachable; vidB's is removed, so // diarization needs it re-acquired. Digest applies to both. await writeDigestVideo({ channelSlug: SLUG, videoId: "vidA" }); await writeDigestVideo({ channelSlug: SLUG, videoId: "vidB" }); for (const f of await audioFiles("vidB")) { await rm(resolvePath(dataRel("vidB", f)), { force: true }); } await rm(resolvePath(`test-transcripts/channels/${SLUG}/snapshot.json`), { force: true, }); await generateReport(page, SLUG); // HALF ONE: the digest entry is really there, with the videos in it. const snapshot = await readJson<{ backfill?: Record< string, { missing?: number; ids?: string[]; eligible?: number; blocked?: number } >; buckets?: Record; }>(`test-transcripts/channels/${SLUG}/snapshot.json`); // ONE definition of the digest work list. The writer used to keep a second // one, `buckets.noDigest`, whose answer disagreed with this entry by 11,777 // videos corpus-wide (it had no cues-staleness gate and no transcript gate). // This is the only place the writer's output is asserted end-to-end, so // without this line the deletion is pinned by nothing. expect(snapshot.buckets).not.toHaveProperty("noDigest"); const digestEntry = snapshot.backfill?.digest; expect(digestEntry, "snapshot.backfill.digest must exist").toBeTruthy(); expect(digestEntry?.missing).toBe(2); expect(digestEntry?.ids?.slice().sort()).toEqual(["vidA", "vidB"]); // vidC is in the fixture with audio and NO transcript, so it is waiting on // transcription: counted, named, and never in the work list. Before the digest // operation declared that dependency it was in no bucket at all, which is how // a channel of untranscribed videos read as fully digested. expect(digestEntry?.blocked).toBe(1); // The denominator the dashboard's coverage percentage divides by: every video // the operation has an opinion about, which includes the blocked one (it // becomes digestable the moment it is transcribed) and excludes nothing else // here. 3 = 2 missing + 1 blocked. expect(digestEntry?.eligible).toBe(3); // And the lane kind is still counted separately, on its own entry — vidC is // NOT-APPLICABLE to diarization (untranscribed), so the two operations // legitimately disagree about the same video. That is the whole reason these // are per-operation entries rather than one shared number. expect(snapshot.backfill?.diarization?.missing).toBe(1); // HALF TWO: every surface that sums the lane ignores it. await page.goto(channelStage(SLUG, "speakers")); const backfillSection = page.getByLabel("speakers section"); // 1 — vidA's diarization. NOT 3, which is what folding the digest entry in // would produce here. await expect(backfillSection.getByLabel("speakers reachable")).toContainText( "1 video can be worked on now", ); // The digest card reports its own work, off the same snapshot entry. await page.goto(channelStage(SLUG, "digest")); await expect(page.getByLabel("digest section").getByRole("heading")).toContainText( "Generate digests (2)", ); // The widget payload the dashboard instrument renders: reachable is the lane // only, and stays 1. const sync = await request.get(`${baseUrl}/api/widget/sync`); const body = await sync.json(); expect(body.backfill.reachable).toBe(1); expect(body.backfill.needsMedia).toBe(1); // And the PER-KIND breakdown the lane card renders is filtered by the same // rule as the sums it splits. This is the half a sum cannot guard: a // breakdown built off Object.entries would list a "Digest" row on the backfill // card while `reachable` above stayed correct, so the card would contradict // its own figure. const kindIds = (body.backfill.kinds as { id: string }[]).map((k) => k.id); expect(kindIds).not.toContain("digest"); expect(kindIds).toContain("diarization"); // Off the same entries as the sums, so the two can never disagree. const kinds = body.backfill.kinds as { id: string; reachable: number; needsMedia: number; }[]; expect(kinds.reduce((n, k) => n + k.reachable, 0)).toBe( body.backfill.reachable, ); expect(kinds.reduce((n, k) => n + k.needsMedia, 0)).toBe( body.backfill.needsMedia, ); }); // (O) THE SAME REGRESSION, ONE SLICE LATER — and the same two halves. // // Slice 1.5 gave the two BUCKET lanes a snapshot entry too: // `backfill.download` and `backfill.transcription` are the fold of each lane's // default buckets, so `snapshot.backfill[op].ids` is where every lane's work // list lives. Six operations now write into a map four surfaces once summed // generically. // // The /channels pipeline columns are the surface with the most to lose, because // they draw a Download and a Transcribe band and could plausibly be "fixed" to // read the new entries. They must not: the band's transcription `reachable` is // downloadedNoTranscript ALONE while the lane's work list also carries the // retry bucket (1,873 videos against 881 on the live corpus), and the band's // `blocked` is what the entry calls `missingInput`. // // So this asserts invariance directly rather than against numbers copied out of // an older build: render the row with the entries present, strip them back out // of the snapshot exactly as every un-regenerated live snapshot has them, and // require the two cells to read the same both times. test("the bucket lanes' work lists land in the snapshot and move no band", async ({ page, }) => { test.setTimeout(SLOW); await resetData("one-transcribe-channel-with-audio"); await writeSettings(backfillSettings()); await rm(resolvePath(`test-transcripts/channels/${SLUG}/snapshot.json`), { force: true, }); await generateReport(page, SLUG); // HALF ONE: the entries are really there. Without this the invariance below // would pass vacuously over a feature that never shipped. const snapRel = `test-transcripts/channels/${SLUG}/snapshot.json`; type Snap = { backfill?: Record< string, { missing?: number; ids?: string[]; missingInput?: number } >; buckets?: Record; undownloadedIds?: string[]; }; const snapshot = await readJson(snapRel); const download = snapshot.backfill?.download; const transcription = snapshot.backfill?.transcription; expect(download, "snapshot.backfill.download must exist").toBeTruthy(); expect(transcription, "snapshot.backfill.transcription must exist").toBeTruthy(); // The fold, restated from the buckets in the same file: partialDownloads then // undownloadedIds, downloadedNoTranscript then failedListed, deduped, in that // order. Asserting the LIST and not just a count is the point — the download // lane's order is playlist order and sorting it would reorder the queue. const union = (names: string[]): string[] => { const out: string[] = []; for (const name of names) { const ids = name === "undownloadedIds" ? (snapshot.undownloadedIds ?? []) : (snapshot.buckets?.[name] ?? []); for (const id of ids) if (!out.includes(id)) out.push(id); } return out; }; expect(download?.ids).toEqual(union(["partialDownloads", "undownloadedIds"])); expect(transcription?.ids).toEqual( union(["downloadedNoTranscript", "failedListed"]), ); expect(download?.missing).toBe(download?.ids?.length); expect(transcription?.missing).toBe(transcription?.ids?.length); // Download's input is the channel listing, which is never missing. A video // with no audio is the TRANSCRIPTION lane's missing input. expect(download?.missingInput).toBe(0); expect(transcription?.missingInput).toBe( snapshot.buckets?.noTranscript?.length ?? 0, ); // The opt-in auto-captions bucket is a POLICY switch, never folded into a // corpus fact. for (const id of snapshot.buckets?.downloadedAutoSubsOnly ?? []) { expect(transcription?.ids).not.toContain(id); } // HALF TWO: the two bands read the same with the entries and without them. const cells = async (): Promise => { await fetch(`${baseUrl}/api/test/invalidate-cache`).catch(() => {}); await page.goto("/channels"); const out: string[] = []; for (const label of ["downloads count", "transcripts count"]) { const cell = page.getByLabel(`${label} for ${SLUG}`); await expect(cell).toBeVisible(); out.push( `${label}=${(await cell.textContent()) ?? ""}|${await cell.getAttribute("title")}`, ); } return out; }; const withEntries = await cells(); // Every live snapshot looks like this today: the slice regenerates nothing, // so the bands and the runner both have to read a file with no entry. const stripped = JSON.parse( await readFile(resolvePath(snapRel), "utf8"), ) as Snap; delete stripped.backfill?.download; delete stripped.backfill?.transcription; await writeFile( resolvePath(snapRel), JSON.stringify(stripped, null, 2) + "\n", ); expect(await cells()).toEqual(withEntries); });