commit 06215ddebc1f379c3ae1311d2c280d08740ed254
parent e8065c6f51554d0729d2d4143f8f22d0becccf43
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Mon, 7 Sep 2026 20:30:39 -0400
editor: the lane runner, end to end, on both operation lanes
Slice 1.2, fourth commit. `e2e/lane-runner.spec.ts` is what says the runner
really drives the digest and backfill lanes rather than merely compiling
against them. Four properties, each one a thing a screenshot cannot tell apart
from its opposite:
1. THE LANE DISPATCHES. The digest runner starts, appears as an `auto-digest`
job, picks a video out of its snapshot work list and leaves an
ai-digest.json — with nobody clicking a per-channel button. The sidecar on
disk is the assertion, not a counter: it is what the operator gets.
2. A HOLD IS NOT A STOP. Holding the gate mid-run leaves the runner UP and
idling at `lane-held`; stopping it ends the job. runPool idle-waits at
limit 0, where returning null from next() would end it, and every pause in
this repo rests on that distinction.
3. IT COMES BACK. `/api/test/resume-lane?lane=` does the two halves a server
restart does — drop the in-process runner, then call the boot hook — and a
new job id comes back off the persisted policy alone. It is the shape
`/api/test/resume-backfill-sweep` already had for the sweep; 1.3 deletes
that one and keeps this. It WAITS for the cancelled runner to leave before
re-arming, because startAutoRunner treats a still-running record as
"already up" and would otherwise return it and prove nothing.
4. DEPENDENCY ORDER, AND OVERLAP. One video gets diarized and then attributed
FROM that diarization — the chain a runner that retired a video after one
unit could never walk — and the digest and backfill lanes are caught
dispatching at the same moment. That last one is backfill.spec's "a digest
runs concurrently with a backfill, not behind it", re-asserted at the
runner: there it is two per-channel jobs holding distinct queue keys, here
it is two loops both on queueKey "" with in-process units and no key at
all. Same property, different mechanism, so it had to be re-proved rather
than inherited. backfill.spec keeps its version, which still covers the
manual verbs.
Also here: the runner builds its duplicate-cluster plan ONCE per runner rather
than once per context refresh. The context is rebuilt on a 60-second TTL (which
is also what re-probes a dead engine), and the plan reads the whole corpus-wide
duplicates report — putting that on a timer forever would be paying, every
minute, to learn about a cluster detection has not run since.
4 passed / 0 failed in 48.3 s from a worktree of this commit; full suite in the
slice note. tsc clean in six packages, common 896, numbers diff empty.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
3 files changed, 451 insertions(+), 0 deletions(-)
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -811,6 +811,13 @@ async function runLoop(
error: string | null;
refreshing: boolean;
} = { run: null, key: "", at: 0, error: null, refreshing: false };
+ // THE DUPLICATE PLAN IS BUILT ONCE PER RUNNER, not once per refresh. It reads
+ // the whole corpus-wide duplicates report, and this context is rebuilt on a
+ // 60-second TTL — so re-deriving it would put that read on a timer forever to
+ // learn about a cluster detection has not run since. A per-channel job builds
+ // one per job for the same reason: it is per-RUN state, and a runner's run is
+ // its whole lifetime.
+ let clusterPlan: NonNullable<OperationRun["digest"]>["clusterPlan"] = null;
const LANE_RUN_TTL_MS = 60_000;
const laneRunKey = (settings: ReturnType<typeof getSettings>): string =>
@@ -839,7 +846,9 @@ async function runLoop(
lane: kind,
paths,
onLog,
+ ...(clusterPlan ? { clusterPlan } : {}),
});
+ clusterPlan = opened.digest?.clusterPlan ?? clusterPlan;
// The engine fail-fast, asked ONCE per context rather than per video: an
// unreachable ollama would otherwise produce one failure per candidate
// over 55,956 of them. A refusal is an IDLE, not a stop — it is a
diff --git a/editor/app/api/test/resume-lane/route.ts b/editor/app/api/test/resume-lane/route.ts
@@ -0,0 +1,56 @@
+import { NextResponse } from "next/server";
+import {
+ getAutoRunnerStatus,
+ startAutoRunnersIfEnabled,
+ stopAutoRunner,
+} from "yt-dlp-transcript-common/controller/autoRunner";
+import { LANES } from "yt-dlp-transcript-common/lib/autoQueueTypes";
+import type { AutoQueueKind } from "yt-dlp-transcript-common/jobs/autoQueueState";
+
+export const dynamic = "force-dynamic";
+
+// E2E test harness only. Simulates a SERVER RESTART for one lane's runner,
+// which is the one property of an enabled lane that cannot otherwise be tested:
+// the runner is process state, a restart destroys it, and
+// instrumentation.ts's boot hook is what brings it back. An e2e suite cannot
+// restart the dev server mid-run, so it does the two halves a restart does —
+// drop the in-process runner, then call the hook — and asserts one comes back.
+//
+// The same shape as /api/test/resume-backfill-sweep, which does this for the
+// sweep. Slice 1.3 deletes that one and leaves this.
+//
+// Mounted unconditionally, like the other /api/test routes: the editor is a
+// localhost admin tool, not a deployed service.
+export async function GET(req: Request) {
+ const lane = new URL(req.url).searchParams.get("lane") ?? "";
+ if (!LANES.includes(lane as AutoQueueKind)) {
+ return NextResponse.json(
+ { ok: false, error: `lane must be one of ${LANES.join(", ")}` },
+ { status: 400 },
+ );
+ }
+ const kind = lane as AutoQueueKind;
+
+ // Half one: lose the live runner, the way a restart does. The job's own
+ // `finally` clears the singleton, so this leaves exactly the state a fresh
+ // process starts in — the persisted policy and nothing running.
+ const before = getAutoRunnerStatus(kind).jobId;
+ stopAutoRunner(kind);
+ // The cancel is asynchronous at the loop's end, and startAutoRunner treats a
+ // still-running record as "already up" and returns it. So WAIT for the old
+ // one to leave, or the resume is a no-op that silently proves nothing.
+ const deadline = Date.now() + 15_000;
+ while (getAutoRunnerStatus(kind).running && Date.now() < deadline) {
+ await new Promise((r) => setTimeout(r, 50));
+ }
+
+ // Half two: the boot hook. Reads autoQueue[lane].enabled and its persisted
+ // tree, and re-launches from those alone.
+ await startAutoRunnersIfEnabled();
+
+ return NextResponse.json({
+ ok: true,
+ cancelled: before,
+ resumed: getAutoRunnerStatus(kind).jobId,
+ });
+}
diff --git a/editor/e2e/lane-runner.spec.ts b/editor/e2e/lane-runner.spec.ts
@@ -0,0 +1,386 @@
+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,
+ 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 dispatching in-process units at the same moment. The
+// mechanism differs; the property must not.
+//
+// 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<string, unknown> = {}) {
+ 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<string, unknown>) ?? {}),
+ },
+ 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<string, unknown>) ?? {}),
+ },
+ 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<string, unknown>) ?? {}),
+ },
+ backfill: {
+ enabled: true,
+ // weight 1, not the idle-only 0: these specs assert that work HAPPENS.
+ // The idle-only default has a pure unit test (operationBatch.test.ts),
+ // which is the right place for it — no pool, no GPU, no timing.
+ weight: 1,
+ concurrency: 1,
+ sweepEnabled: false,
+ sweepKinds: [],
+ sweepChannels: [],
+ allowRedownload: false,
+ ...((over.backfill as Record<string, unknown>) ?? {}),
+ },
+ ...over,
+ };
+}
+
+type RunnerStatus = {
+ runner: { running: boolean; jobId: string | null; idleReason: string | null; inFlight: unknown[] };
+ policy: { enabled?: boolean; order?: string };
+ pendingByLeaf: Record<string, number>;
+};
+
+async function status(
+ request: APIRequestContext,
+ lane: string,
+): Promise<RunnerStatus> {
+ const res = await request.get(`${baseUrl}/api/auto-queue/status`);
+ expect(res.ok()).toBeTruthy();
+ const body = (await res.json()) as Record<string, RunnerStatus>;
+ 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`, {
+ data: { kind, action },
+ });
+ expect(res.ok()).toBeTruthy();
+ return res.json();
+}
+
+async function activeJobKinds(request: APIRequestContext): Promise<string[]> {
+ 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: { backfill: { enabled: true, 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 both lanes are provably still going when the other is
+ // looked at. One video each finishes faster than a poll.
+ const ids = Array.from({ length: 8 }, (_, i) => `lanevid01${i}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, maxWorkers: 1, root: CATCH_ALL },
+ },
+ }),
+ );
+ await generateReport(page, CHANNEL);
+
+ await control(request, "digest", "start");
+ await control(request, "backfill", "start");
+
+ // Both dispatching at ONE moment. A lane sitting at zero while the other
+ // works is exactly the failure a shared resource gate would produce, so it is
+ // this simultaneity that is polled rather than "both eventually did work".
+ await expect
+ .poll(
+ async () => {
+ const res = await request.get(`${baseUrl}/api/auto-queue/status`);
+ const body = (await res.json()) as Record<string, RunnerStatus>;
+ return (
+ body.digest.runner.inFlight.length > 0 &&
+ body.backfill.runner.inFlight.length > 0
+ );
+ },
+ { timeout: 90_000, intervals: [100] },
+ )
+ .toBe(true);
+});