import { writeFile } from "node:fs/promises"; import { join } from "node:path"; import { test, expect } from "@playwright/test"; import type { APIRequestContext, Page } from "@playwright/test"; import { baseUrl } from "./baseUrl"; import { generateReport, pathExists, readJson, resetData, resolvePath, writeChannelConfig, writeDigestVideo, writeSettings, } from "./helpers"; // THE AUTO-QUEUE RUNNER, DRIVING THE OPERATION LANES. // // Until slice 1.2 the digest and backfill lanes were dispatched by their own // sweeps and the runner knew how to transcribe and how to download and nothing // else. Now one loop drives all four, and these specs are what say so: // // 1. A LANE THAT STARTS AND ACTUALLY DISPATCHES. The digest lane picks a // video from its snapshot work list, runs one unit through the ollama // stub, and leaves an ai-digest.json — end to end, with nobody clicking a // per-channel button. // 2. A HOLD IS NOT A STOP. Holding the lane's gate mid-run must leave the // runner UP and idling at zero, not end it. Every pause in this repo // depends on that distinction (runPool idle-waits at limit 0; returning // null from next() would end the job), and it is the one property a // screenshot cannot tell apart from "wedged". // 3. IT COMES BACK AFTER A RESTART. The runner is process state; the policy // is on disk. /api/test/resume-lane does the two halves a restart does. // 4. DEPENDENCY ORDER WITHIN THE BACKFILL LANE. A video gets diarized and // then attributed FROM ITS DIARIZATION — the second operation consumes the // first's output, and the runner walks that chain one unit at a time. // 5. THE TWO LANES OVERLAP. This is the `backfill.spec.ts` "a digest runs // concurrently with a backfill, not behind it" invariant, at the runner: // there it is two per-channel JOBS holding distinct queue keys, here it is // two runner loops, both on queueKey "", dispatching in-process units that // hold no key at all. The mechanism differs; the property must not. It is // measured from the two timestamped PICK LOGS afterwards — do their // intervals overlap — rather than by catching both lanes in flight at one // sampled instant. The second is the same property measured by a coin // toss, and a loaded machine loses it to sampling rather than to a // regression. // 6. ONE GATE, THREE SURFACES. The dashboard deck, the widget rail and the // lane's own console all hold the SAME lane by writing the SAME key // (`autoQueue.digest.held`, since slice 1.4), and none of them touches // another lane's. Asserted from the file, because every one of those // buttons flips its own label optimistically. // // Engines: the digest and attribution lanes use the ollama HTTP stub // (e2e/fixtures/ollama-stub.mjs, wired via OLLAMA_URL), diarization uses // fake-diarize.mjs (DIARIZE_BIN) — the same fixtures digest.spec, attribution. // spec and backfill.spec already run on. const CHANNEL = "lane-runner-channel"; const SLOW = 180_000; type Leaf = { id: string; match: { type: string; value?: string; operation?: string }; weight?: number; maxWorkers?: number | null; }; type Group = { id: string; mode: string; children: (Group | Leaf)[]; maxWorkers?: number | null; }; // One catch-all rule. A leaf naming NO operation draws the lane's whole // dispatch set — the operation half of the default bucket union — which is the // projection change slice 1.2 made and the thing these specs exercise. const CATCH_ALL: Group = { id: "root", mode: "strict", children: [{ id: "all", match: { type: "all" }, weight: 1, maxWorkers: null }], }; function dataRel(videoId: string, file: string): string { return join("test-transcripts", "channels", CHANNEL, "data", videoId, file); } // Settings with both engines named and the lanes' own gates OPEN. `autoQueue` // is spelled per test, because which lane is enabled is the subject. function laneSettings(over: Record = {}) { return { adminTitle: "Test Admin", maxTranscriptPageBytes: 8388608, sleepBetweenDownloadsSeconds: 0, minFreeDiskGB: 0, verifyAvailabilityBeforeClean: false, syncScheduler: { fullSweepIntervalMinutes: 0 }, digest: { localAppId: "ollama-direct", remoteAppId: "claude-code", sections: ["chapters"], // The runner must not stand aside for a transcription that is not // happening; nothing transcribes in these specs, but the flag is stated // rather than inherited so a default change cannot silently park a lane. yieldToTranscription: false, ...((over.digest as Record) ?? {}), }, diarization: { enabled: true, inlineAfterTranscribe: false, threshold: 0.5, threads: 1, python: "python3", segModel: "/dev/null", embModel: "/dev/null", concurrency: 1, // CPU, not the vulkan default: a GPU-bound backfill operation makes the // whole lane idle-only, which is correct behaviour and would park this // lane the moment anything else looked busy. backend: "cpu", ...((over.diarization as Record) ?? {}), }, attribution: { enabled: true, appId: "ollama-direct", model: "qwen2.5:7b", diarizedEnabled: true, // The TEXT lane off: this spec is about the diarized one, which is the // half with a dependency to walk. textOnlyEnabled: false, promptVersion: 1, ...((over.attribution as Record) ?? {}), }, backfill: { concurrency: 1, allowRedownload: false, ...((over.backfill as Record) ?? {}), }, ...over, }; } type RunnerStatus = { runner: { running: boolean; jobId: string | null; idleReason: string | null; inFlight: unknown[]; }; policy: { enabled?: boolean; order?: string }; pendingByLeaf: Record; // The persisted pick log — one entry per unit the runner handed out, newest // first, each stamped with the moment it was dispatched. picks: { at: number; leafId: string; videoId: string }[]; }; async function status( request: APIRequestContext, lane: string, ): Promise { const res = await request.get(`${baseUrl}/api/auto-queue/status`); expect(res.ok()).toBeTruthy(); const body = (await res.json()) as Record; return body[lane]; } async function control( request: APIRequestContext, kind: string, action: "start" | "stop" | "drain", ): Promise<{ started?: boolean; blocked?: string }> { const res = await request.post(`${baseUrl}/api/auto-queue/control`, { // Behind the ops token since release 19 (A3): it starts and stops lanes. headers: { authorization: "Bearer test-worker-token" }, data: { kind, action }, }); expect(res.ok()).toBeTruthy(); return res.json(); } async function activeJobKinds(request: APIRequestContext): Promise { const res = await request.get(`${baseUrl}/api/jobs/active`); const body = await res.json(); const jobs = Array.isArray(body) ? body : (body.jobs ?? []); return jobs.map((j: { kind: string }) => j.kind); } // Stop every lane between tests, so one spec's loop cannot bleed into the next. test.afterEach(async ({ request }) => { for (const lane of ["digest", "backfill"]) { await control(request, lane, "stop").catch(() => ({})); } }); async function seedDigestChannel(page: Page, videoIds: string[]) { await resetData(null); await writeChannelConfig(CHANNEL); for (const id of videoIds) { await writeDigestVideo({ channelSlug: CHANNEL, videoId: id }); } } // --------------------------------------------------------------------------- // (1) The digest lane, driven by the runner. test("the digest lane runner dispatches a unit, holds on pause, and stops", async ({ page, request, }) => { test.setTimeout(SLOW); await seedDigestChannel(page, ["lanevid0001"]); await writeSettings( laneSettings({ autoQueue: { digest: { enabled: true, maxWorkers: 1, root: CATCH_ALL } }, }), ); // The snapshot is the lane's WORK LIST: `backfill.digest.ids` is what the // leaf draws from, so there is nothing to dispatch until it exists. await generateReport(page, CHANNEL); const started = await control(request, "digest", "start"); expect(started.started).toBe(true); // It appears as a job of its own kind — one long-lived runner, not one job // per video (the registry keeps 100 records; a job per video would evict the // history of the run that made them). await expect .poll(() => activeJobKinds(request), { timeout: 30_000 }) .toContain("auto-digest"); // AND IT ACTUALLY RAN ONE. The sidecar on disk is the assertion, not a // counter: it is what the operator gets. await expect .poll(() => pathExists(dataRel("lanevid0001", "ai-digest.json")), { timeout: 90_000, intervals: [250], }) .toBe(true); // A HOLD IS NOT A STOP. Holding the gate mid-run leaves the runner UP, // idling at zero — runPool idle-waits at limit 0, where returning null from // next() would end the job. await page.goto("/operations/digest"); const lane = page.locator('section[data-lane="digest"]'); await expect(lane).toHaveAttribute("data-hydrated", "true", { timeout: 30_000, }); await lane.getByRole("button", { name: "pause digests" }).click(); await expect .poll(async () => (await status(request, "digest")).runner.idleReason, { timeout: 30_000, intervals: [250], }) .toBe("lane-held"); expect((await status(request, "digest")).runner.running).toBe(true); await lane.getByRole("button", { name: "resume digests" }).click(); // Stop is the other thing entirely: the job ends. await control(request, "digest", "stop"); await expect .poll(async () => (await status(request, "digest")).runner.running, { timeout: 30_000, intervals: [250], }) .toBe(false); }); // --------------------------------------------------------------------------- // (2) A restart brings it back. test("the digest lane runner comes back after a restart", async ({ page, request, }) => { test.setTimeout(SLOW); await seedDigestChannel(page, ["lanevid0002"]); await writeSettings( laneSettings({ autoQueue: { digest: { enabled: true, maxWorkers: 1, root: CATCH_ALL } }, }), ); await generateReport(page, CHANNEL); await control(request, "digest", "start"); await expect .poll(async () => (await status(request, "digest")).runner.running, { timeout: 30_000, intervals: [250], }) .toBe(true); const before = (await status(request, "digest")).runner.jobId; // The two halves a server restart does: lose the in-process runner, then call // the boot hook. The POLICY is on disk, which is the whole point — nothing // re-arms it, it was never disarmed. const res = await request.get(`${baseUrl}/api/test/resume-lane?lane=digest`); expect(res.ok()).toBeTruthy(); const body = (await res.json()) as { resumed: string | null }; expect(body.resumed).not.toBeNull(); expect(body.resumed).not.toBe(before); expect((await status(request, "digest")).runner.running).toBe(true); }); // --------------------------------------------------------------------------- // (3) The backfill lane walks its dependency chain. test("the backfill lane runner diarizes, then attributes from that diarization", async ({ page, request, }) => { test.setTimeout(SLOW); const VIDEO = "lanevid0003"; await seedDigestChannel(page, [VIDEO]); // Diarization's INPUT is the audio; attribution-diarized's input is the // diarization sidecar the first one writes. So at the start exactly one of // the two is reachable, and the second becomes reachable only because the // first ran. await writeFile(resolvePath(dataRel(VIDEO, "audio.mp3")), "fake audio\n"); await writeSettings( laneSettings({ autoQueue: { // `held: false` because the backfill lane's gate defaults SHUT — the // reading its inverted `backfill.enabled` always gave a file that named // no gate, kept when S0-pause deleted the field. backfill: { enabled: true, held: false, maxWorkers: 1, root: CATCH_ALL }, }, }), ); await generateReport(page, CHANNEL); await control(request, "backfill", "start"); await expect .poll(() => activeJobKinds(request), { timeout: 30_000 }) .toContain("auto-backfill"); await expect .poll(() => pathExists(dataRel(VIDEO, "diarization.json")), { timeout: 90_000, intervals: [250], }) .toBe(true); // THE SECOND OPERATION, on the SAME video, from the first one's output. A // runner that retired the video after one unit — the shape the completed-set // had before this slice — would never reach here. await expect .poll(() => pathExists(dataRel(VIDEO, "attribution.json")), { timeout: 120_000, intervals: [250], }) .toBe(true); }); // --------------------------------------------------------------------------- // (4) The two lanes overlap. // // THE PROPERTY REGISTERING DIGEST MUST NOT BREAK, at the runner. backfill.spec // pins it for the per-channel JOBS, where it is bought with distinct queue keys // (registry.ts runs every non-empty key at concurrency 1, so one shared key // would make the GPU digest lane wait on CPU diarization across a sweep // measured in weeks). The runner's units hold no queue key at all — they run // in-process inside two loops, both on queueKey "" — so the mechanism is // different and the property has to be re-asserted rather than assumed. test("the digest and backfill lanes dispatch at the same time", async ({ page, request, }) => { test.setTimeout(SLOW); // ENOUGH WORK that each lane's dispatching spans seconds rather than one // tick: 12 digests on one side, 12 diarizations plus 12 attributions on the // other. Both stay inside the pick log's 50-entry cap. const ids = Array.from( { length: 12 }, (_, i) => `lanevid1${String(i).padStart(3, "0")}`, ); await seedDigestChannel(page, ids); for (const id of ids) { await writeFile(resolvePath(dataRel(id, "audio.mp3")), "fake audio\n"); } await writeSettings( laneSettings({ autoQueue: { digest: { enabled: true, maxWorkers: 1, root: CATCH_ALL }, backfill: { enabled: true, held: false, maxWorkers: 1, root: CATCH_ALL }, }, }), ); await generateReport(page, CHANNEL); await control(request, "digest", "start"); await control(request, "backfill", "start"); // Let both lanes do their work, then read the two PICK LOGS — which are // timestamped and persisted, so this measures what happened rather than what // a poll happened to catch. An earlier version of this test sampled for "both // in flight at one instant" and lost the race on a loaded machine: same // property, measured by a coin toss. await expect .poll( async () => { const digest = await status(request, "digest"); const backfill = await status(request, "backfill"); return Math.min(digest.picks.length, backfill.picks.length); }, { timeout: 120_000, intervals: [500] }, ) .toBeGreaterThanOrEqual(4); const digest = await status(request, "digest"); const backfill = await status(request, "backfill"); const span = (s: RunnerStatus): [number, number] => { const at = s.picks.map((p) => p.at); return [Math.min(...at), Math.max(...at)]; }; const [dFrom, dTo] = span(digest); const [bFrom, bTo] = span(backfill); // OVERLAPPING INTERVALS IS THE WHOLE ASSERTION. If either lane were waiting on // the other, every one of its picks would fall after the other's last — two // disjoint intervals, which is exactly what a shared queue key produces for // the per-channel jobs backfill.spec pins. expect( dFrom <= bTo && bFrom <= dTo, `digest picks ${dFrom}-${dTo} vs backfill picks ${bFrom}-${bTo} do not overlap`, ).toBe(true); }); // --------------------------------------------------------------------------- // (5) THREE SURFACES, ONE KEY. // // A lane's pause is drawn on the dashboard deck, in the monitor widget's rail // and on the lane's own runner console, and until slice 1.4 those three wrote a // settings FIELD chosen per lane — one of four, with three polarities. They // write `autoQueue[lane].held` now, and this is the spec that says all three // mean it: the file is read back after each click, the OTHER three lanes' keys // are compared before and after, and the runner is asked whether it is still // running each time. // // WHY THE FILE AND NOT THE LABEL. Every one of these buttons flips its own label // optimistically, so three green label assertions would pass over three // surfaces writing three different things. The key on disk is the thing they // have to agree about. // // A HOLD, NEVER A STOP, at every one of them: `running` stays true and the idle // reason is `lane-held` — the same pair test (1) asserts for the console, here // for the two surfaces that are not the console. type HeldByLane = Record; async function heldByLane(): Promise { const s = await readJson<{ autoQueue?: Record; }>("test-settings.json").catch(() => ({}) as { autoQueue?: undefined }); return Object.fromEntries( ["transcription", "download", "digest", "backfill"].map((lane) => [ lane, s.autoQueue?.[lane]?.held ?? null, ]), ); } test("a hold from the dashboard, the widget and the console writes one key", async ({ page, request, }) => { test.setTimeout(SLOW); await seedDigestChannel(page, ["lanevid0005"]); await writeSettings( laneSettings({ autoQueue: { digest: { enabled: true, maxWorkers: 1, root: CATCH_ALL } }, }), ); await generateReport(page, CHANNEL); await control(request, "digest", "start"); await expect .poll(async () => (await status(request, "digest")).runner.running, { timeout: 30_000, intervals: [250], }) .toBe(true); // The other three lanes, as the fixture leaves them. Nothing below may move // any of these — a surface that held "the lane" by holding the wrong one is // the failure this comparison exists for, and it is silent in every label. const settle = async (want: boolean | null) => { await expect .poll(async () => (await heldByLane()).digest, { timeout: 30_000, intervals: [250], }) .toBe(want); return heldByLane(); }; // The first write through the editor migrates every lane's gate onto its // policy, so the neighbours are compared from AFTER that write, not before. // // `scope` is not decoration on the console: an operation page draws the rail // as well as the lane's own section, and both carry this control (the same // reason backfill.spec scopes its hold to `section[data-lane]`). The other // two pages have one. const surfaces: { name: string; url: string; scope: string | null }[] = [ { name: "console", url: "/operations/digest", scope: 'section[data-lane="digest"]' }, { name: "dashboard", url: "/", scope: null }, { name: "widget", url: "/widget?controls=1", scope: null }, ]; let neighbours: HeldByLane | null = null; for (const surface of surfaces) { await page.goto(surface.url); const within = surface.scope ? page.locator(surface.scope) : page; const hold = within.getByRole("button", { name: "pause digests" }); // 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(); const afterHold = await settle(true); // HELD IS A HOLD. The runner is still up, idling at zero with the lane's // own reason — not `stopped`, and not the download lane's // `downloads-paused`, which names a different gate. await expect .poll(async () => (await status(request, "digest")).runner.idleReason, { timeout: 30_000, intervals: [250], }) .toBe("lane-held"); expect( (await status(request, "digest")).runner.running, `${surface.name} stopped the runner`, ).toBe(true); const release = within.getByRole("button", { name: "resume digests" }); await expect(release).toBeEnabled({ timeout: 30_000 }); await release.click(); const afterRelease = await settle(false); // ONE KEY MOVED, and only one. Held and released are compared separately so // a surface writing a second lane on either edge is caught. const others = (h: HeldByLane) => ({ ...h, digest: null }); if (neighbours === null) neighbours = others(afterHold); expect(others(afterHold), `${surface.name} hold moved another lane`).toEqual( neighbours, ); expect( others(afterRelease), `${surface.name} release moved another lane`, ).toEqual(neighbours); } // And it survived all three round trips as a runner, not as a restart. expect((await status(request, "digest")).runner.running).toBe(true); });