import { test } from "node:test"; import assert from "node:assert/strict"; import { backfillLimit, candidateAction, laneLimit, operationLaneLive, type OperationRun, } from "./operationBatch"; import { defaultBackfill, defaultDigest, defaultDiarization, defaultSiteSettings, sanitizeBackfill, sanitizeDiarization, type SiteSettings, } from "../lib/settings"; import type { Operation, OperationClassification } from "../lib/operations"; import { withGateHeld } from "../lib/pauseGates"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/operationBatch.test.ts // // The PURE function, not a live pool — same discipline as digestYield.test.ts, // and for the same reason: this decides whether a multi-day lane runs at all, // and a test that needs a GPU, a registry and a worker pool is a test nobody // runs before shipping. test("idle-only means full slots when the primary lane is quiet", () => { assert.equal(backfillLimit({ idleOnly: true, slots: 4, primaryBusy: false }), 4); assert.equal(backfillLimit({ idleOnly: true, slots: 1, primaryBusy: false }), 1); }); test("idle-only stands aside completely while the primary lane works", () => { // Zero is a HOLD, not a stop: runPool idle-waits at a zero limit rather than // finishing, so the lane resumes the moment transcription is free without // re-deriving anything. That is the whole reason this is a limit and not a // scheduler. assert.equal(backfillLimit({ idleOnly: true, slots: 4, primaryBusy: true }), 0); assert.equal(backfillLimit({ idleOnly: true, slots: 16, primaryBusy: true }), 0); }); test("a run that is not idle-only keeps its slots, busy or not", () => { // The `weight` scalar retired in slice 1.3 used to answer this with a share. // The question it was really being asked is "does this run contend with // transcription", which the operation's declared `contendsFor` answers — so a // network-bound attribution run keeps its slots rather than parking behind a // GPU it is not competing for. assert.equal(backfillLimit({ idleOnly: false, slots: 4, primaryBusy: true }), 4); assert.equal(backfillLimit({ idleOnly: false, slots: 4, primaryBusy: false }), 4); }); test("no slots means no work, idle-only or not", () => { assert.equal(backfillLimit({ idleOnly: true, slots: 0, primaryBusy: false }), 0); assert.equal(backfillLimit({ idleOnly: false, slots: 0, primaryBusy: false }), 0); }); test("the shipped default is off, and holds no disk", () => { // Stated as a test because these are promises the feature makes: nothing runs // until an operator says so, and nothing re-downloads media without an // explicit opt-in. // // "Off" is now TWO keys on the lane, and both ship off: the arm // (`autoQueue.backfill.enabled`) and the gate (`autoQueue.backfill.held`, // which `defaultHeldFor` shuts because the `backfill.enabled` this line used // to read was inverted and defaulted false). S0-pause deleted that field; it // did not change what a fresh install does. assert.equal(defaultSiteSettings().autoQueue.backfill.enabled, false); assert.equal(defaultSiteSettings().autoQueue.backfill.held, true); const d = defaultBackfill(); assert.equal(d.allowRedownload, false); assert.equal(d.concurrency, 1); // Concurrency shares clampPositiveInt's floor of 1, so a hand-edited 0 cannot // silently park the lane. assert.equal(sanitizeBackfill({ concurrency: 0 }).concurrency, 1); }); // --------------------------------------------------------------------------- // The candidate pull's dispatch decision. // // Untested until the duration cap needed it, and it is the single most // consequential branch in the file: it used to fall through to DISPATCHING any // state it did not recognize, so "add a state" and "run that state anyway" were // the same edit. const DISPATCH = { force: false, allowRedownload: false }; test("only reachable work is dispatched by default", () => { assert.equal(candidateAction("missing", DISPATCH), "dispatch"); assert.equal(candidateAction("stale", DISPATCH), "dispatch"); }); test("deferred is NEVER dispatched", () => { // The whole point of the cap. Without this the classification would be // computed, ignored, and the six-hour video handed to the engine — 40 minutes // of CPU and a kernel OOM kill, producing nothing. assert.equal(candidateAction("deferred", DISPATCH), "deferred"); assert.equal( candidateAction("deferred", { force: true, allowRedownload: true }), "deferred", ); }); test("force redoes present work but does not overrule the cap", () => { assert.equal(candidateAction("present", DISPATCH), "fresh"); assert.equal( candidateAction("present", { ...DISPATCH, force: true }), "dispatch", ); // `force` means "redo work that looks done", not "ignore the cap". assert.equal( candidateAction("deferred", { ...DISPATCH, force: true }), "deferred", ); }); test("missing-input is counted unless re-download is armed", () => { assert.equal(candidateAction("missing-input", DISPATCH), "missing-input"); assert.equal( candidateAction("missing-input", { ...DISPATCH, allowRedownload: true }), "dispatch", ); }); test("not-applicable is counted as nothing at all", () => { assert.equal(candidateAction("not-applicable", DISPATCH), "skip"); assert.equal( candidateAction("not-applicable", { force: true, allowRedownload: true }), "skip", ); }); test("every classification has an explicit decision", () => { // The list is written out rather than derived so that adding a state to the // union without deciding what the pull does with it fails HERE as well as at // the compile step — a runtime backstop for the `never` check, since the // hazard this replaces was precisely a silent fall-through. const ALL: OperationClassification[] = [ "present", "stale", // Part-done. Dispatched like stale — the runner regenerates only the // sections that are not fresh — and listed here so that stays a decision // rather than something a later edit can quietly drop. "partial", "missing", "missing-input", "deferred", // This list had already drifted: `blocked` was added to the union and to // candidateAction, and never to the backstop meant to catch exactly that. "blocked", "not-applicable", ]; for (const state of ALL) { assert.doesNotThrow(() => candidateAction(state, DISPATCH), state); } assert.throws( () => candidateAction("invented" as OperationClassification, DISPATCH), /unhandled backfill state/, ); }); test("the duration cap round-trips, and 0 means off", () => { // 0 is a MEANINGFUL value here, not an empty one — it is how the cap is turned // off once windowed diarization makes it unnecessary. clampPositiveInt would // have floored it to 1, i.e. a one-hour cap, which is why this knob does not // use it. assert.equal(sanitizeDiarization({ maxAudioHours: 0 }).maxAudioHours, 0); assert.equal(sanitizeDiarization({ maxAudioHours: 6.5 }).maxAudioHours, 6.5); // Junk falls back to the default rather than to "no cap" — a typo in a // hand-edited settings.json must not silently re-arm the OOM. assert.equal( sanitizeDiarization({ maxAudioHours: "four" }).maxAudioHours, defaultDiarization().maxAudioHours, ); assert.equal( sanitizeDiarization({ maxAudioHours: -2 }).maxAudioHours, defaultDiarization().maxAudioHours, ); // The shipped default is OFF: windowed diarization removed the OOM the cap // was a stopgap for. This is asserted rather than assumed because turning it // back on by accident would silently stop diarizing long videos. assert.equal(defaultDiarization().maxAudioHours, 0); }); test("an unknown engine or backend falls back to the default, never to nothing", () => { // The DEFAULT ENGINE IS LOAD-BEARING as a fallback, not just as a starting // point: it is the one every sidecar already on disk matches, so falling back // to it leaves the corpus fresh. Falling back to sortformer on a typo would // mark all of it stale and offer weeks of rework. assert.equal(defaultDiarization().engine, "sherpa-onnx"); assert.equal(sanitizeDiarization({ engine: "sortformer" }).engine, "sortformer"); assert.equal( sanitizeDiarization({ engine: "sortfromer" }).engine, defaultDiarization().engine, ); assert.equal(sanitizeDiarization({ engine: 7 }).engine, defaultDiarization().engine); // The backend defaults to the GPU, which is safe only because the lane yields // the card rather than sharing it — see diarizationLaneFor. assert.equal(defaultDiarization().backend, "vulkan"); assert.equal(sanitizeDiarization({ backend: "cpu" }).backend, "cpu"); assert.equal( sanitizeDiarization({ backend: "rocm" }).backend, defaultDiarization().backend, ); // Paths are plain strings and empty means "not configured", which diarizeOne // reports as a skip rather than a failure. assert.equal(sanitizeDiarization({}).sortformerBin, ""); assert.equal(sanitizeDiarization({}).sortformerModel, ""); assert.equal( sanitizeDiarization({ sortformerBin: " /opt/diarize-file " }).sortformerBin, "/opt/diarize-file", ); }); // --------------------------------------------------------------------------- // The layer boundary. This test used to live in lib/operations.test.ts and // import candidateAction across it — a lib/ test reaching into controller/, // carried on architecture.test.ts's allow-list as a debt. The classification is // the registry's business; what a DISPATCHER does with it is this file's, so // the assertion belongs here. test("blocked is never dispatched, and re-download cannot change that", () => { // The dispatch decision is the consequential one: `blocked` must not reach a // runner whatever the flags say. Note allowRedownload — the flag that DOES // turn missing-input into a dispatch — is deliberately inert here. for (const force of [false, true]) { for (const allowRedownload of [false, true]) { assert.equal( candidateAction("blocked", { force, allowRedownload }), "blocked", `force=${force} allowRedownload=${allowRedownload}`, ); } } // The contrast, so this test fails if the two ever get conflated again. assert.equal( candidateAction("missing-input", { force: false, allowRedownload: true }), "dispatch", ); }); // --------------------------------------------------------------------------- // laneLimit — the whole of a lane's dispatch gate, asked as one rule. // // No pool, no registry, no GPU: transcriptionActivity() fails OPEN in a bare // process (an unreadable pool and an unreadable registry both read as "nothing // there"), which is the property that makes this testable at all — and the same // property that means a broken read causes contention rather than deadlock. function settingsWith(over: Partial): SiteSettings { return { ...defaultSiteSettings(), ...over }; } test("the backfill lane holds when its gate is shut, and says so", () => { // THE GATE IS `autoQueue.backfill.held`. It was the inverted // `backfill.enabled` — the lane's "master switch" — until slice 1.4 moved it // onto the lane and S0-pause deleted the field; `withGateHeld` is the one // writer, so this is the shape every control produces. const verdict = laneLimit( withGateHeld(settingsWith({ backfill: defaultBackfill() }), "backfill", true), { lane: "backfill", operations: [], llmOps: [], remoteEligible: false, llmActive: 0, unitActive: 0, }, ); // A HOLD, not a stop: runPool idle-waits at zero, so the lane keeps its place // and resumes on the next poll after the switch comes back. assert.equal(verdict.limit, 0); assert.equal(verdict.hold?.reason, "paused"); }); test("an enabled backfill lane with no GPU-bound operation runs at its slots", () => { // Released explicitly: the backfill lane's gate DEFAULTS shut // (`defaultHeldFor`), which is the reading its inverted `backfill.enabled` // always gave a file that named no gate. const verdict = laneLimit( withGateHeld( settingsWith({ backfill: { ...defaultBackfill(), concurrency: 3 } }), "backfill", false, ), { lane: "backfill", operations: [], llmOps: [], remoteEligible: false, llmActive: 0, unitActive: 0, }, ); assert.equal(verdict.limit, 3); assert.equal(verdict.hold, null); }); test("the GPU carve-out is keyed on contendsFor, not on laneFor existing", () => { // THE RULE THIS SLICE CHANGED. It used to be `k.laneFor && ...` — the field's // PRESENCE as a proxy for "this one's resource depends on settings". That // made `laneFor` load-bearing by existing: giving every operation one (a // refactor that looks like a no-op) would have enrolled all of them and made // the lane idle-only whenever a statically GPU-bound operation was in the // run. Now the declaration itself answers, so an operation with a fixed // `lane: { contendsFor: "gpu" }` and NO laneFor is caught, and one with a // laneFor resolving to CPU is not. const settings = withGateHeld( settingsWith({ backfill: { ...defaultBackfill(), concurrency: 4 } }), "backfill", false, ); const cpuOnly = laneLimit(settings, { lane: "backfill", operations: [fakeOperation("cpu-op", "cpu")], llmOps: [], remoteEligible: false, llmActive: 0, unitActive: 0, }); // Sharing is a CPU concept, so a CPU-bound run keeps its slots. assert.equal(cpuOnly.limit, 4); const gpuBound = laneLimit(settings, { lane: "backfill", operations: [fakeOperation("cpu-op", "cpu"), fakeOperation("gpu-op", "gpu")], llmOps: [], remoteEligible: false, llmActive: 0, unitActive: 0, }); // One GPU-bound operation makes the WHOLE run idle-only, including its cheap // CPU siblings — deliberately blunt, because the alternative is an OOM when // two engines allocate the same 8 GB card. assert.equal(gpuBound.limit, 4, "idle: nothing is transcribing in this process"); const weightless = laneLimit( withGateHeld( settingsWith({ backfill: { ...defaultBackfill(), concurrency: 4 } }), "backfill", false, ), { lane: "backfill", operations: [fakeOperation("gpu-op", "gpu")], llmOps: [], remoteEligible: false, llmActive: 0, unitActive: 0, }, ); // Same shape as backfillLimit({idleOnly: true, primaryBusy: false}) — the // carve-out made it idle-only, and an idle primary still gives full slots. assert.equal(weightless.limit, 4); }); test("the digest lane holds while digests are paused", () => { const verdict = laneLimit( withGateHeld(settingsWith({ digest: defaultDigest() }), "digest", true), { lane: "digest", appLane: "local-gpu", metered: false, costUsd: 0, concurrency: 1, fanOutEligible: false, llmActive: 0, }, ); assert.equal(verdict.limit, 0); assert.equal(verdict.hold?.reason, "paused"); }); // HELD IS A HOLD, NEVER A STOP — at the runner, both lanes at once. // // The two tests above each hold ONE lane through `autoQueue[lane].held`; this // one walks both and pins that the other lane is untouched by it. There is no // second spelling to cross-check any more: slice 1.4 moved the gate onto the // lane and S0-pause deleted the four retired fields it had migrated from. // // WHAT "A HOLD" MEANS HERE is `limit: 0` with `hold.reason === "paused"`, which // is what makes runPool idle-wait and the runner report `lane-held` while // staying `running` (lane-runner.spec.ts test 1 asserts the other half, from // the browser). A stop would be `next()` aborting the controller, and nothing // on this path does that. // // Only the two OPERATION lanes hold through limit(): the download lane's gate // is asked in `next()` (idle reason `downloads-paused`) and transcription's is // the worker pool's slot count. Same gate, three shapes — see pauseGates.ts. function heldLane(lane: "digest" | "backfill", held: boolean): SiteSettings { const base = defaultSiteSettings(); return { ...base, backfill: { ...defaultBackfill(), concurrency: 2 }, autoQueue: { ...base.autoQueue, [lane]: { ...base.autoQueue[lane], held }, }, }; } test("an operation lane held through autoQueue[lane].held stops dispatching without stopping", () => { for (const lane of ["digest", "backfill"] as const) { const live = lane === "digest" ? ({ lane: "digest", appLane: "local-gpu", metered: false, costUsd: 0, concurrency: 1, fanOutEligible: false, llmActive: 0, } as const) : ({ lane: "backfill", operations: [], llmOps: [], remoteEligible: false, llmActive: 0, unitActive: 0, } as const); const held = laneLimit(heldLane(lane, true), live); assert.equal(held.limit, 0, `${lane} held should dispatch nothing`); assert.equal(held.hold?.reason, "paused", `${lane} held should say paused`); // And `held: false` releases the lane even while the retired field it // replaced still says "held" — the key wins, in both directions. const free = laneLimit(heldLane(lane, false), live); assert.ok(free.limit > 0, `${lane} free should dispatch`); assert.equal(free.hold?.reason, undefined, `${lane} free should not hold`); } }); test("the digest spend cap is per RUN, and parks the metered lane", () => { const settings = settingsWith({ digest: { ...defaultDigest(), spendCapUsd: 5 }, }); const under = laneLimit(settings, { lane: "digest", appLane: "remote-api", metered: true, costUsd: 4.99, concurrency: 2, fanOutEligible: false, llmActive: 0, }); assert.equal(under.limit, 2); const over = laneLimit(settings, { lane: "digest", appLane: "remote-api", metered: true, costUsd: 5, concurrency: 2, fanOutEligible: false, llmActive: 0, }); assert.equal(over.limit, 0); assert.equal(over.hold?.reason, "spend-cap"); }); test("the metered digest lane does not yield the GPU it never touches", () => { // laneYieldsToTranscription reads `contendsFor` off the engine's declared // lane: local-gpu contends, remote-api does not. Asserted through laneLimit // rather than through the declaration so the wiring is covered too. const settings = settingsWith({ digest: defaultDigest() }); for (const appLane of ["local-gpu", "remote-api"] as const) { const verdict = laneLimit(settings, { lane: "digest", appLane, metered: appLane === "remote-api", costUsd: 0, concurrency: 1, fanOutEligible: false, llmActive: 0, }); // Nothing is transcribing in this process, so neither holds — what is // asserted is that neither holds for the WRONG reason. assert.equal(verdict.hold, null, appLane); assert.equal(verdict.limit, 1, appLane); } }); // A minimal registered-shaped operation. Only `id` and `lane` are read by the // limit; laneForOperation returns null for an id the registry does not know, so // the declared `lane` below is what answers — which is exactly the fallback the // rule relies on. function fakeOperation(id: string, contendsFor: "cpu" | "gpu" | "network") { return { id, lane: { queueKey: "backfill", contendsFor }, } as unknown as Operation; } // --------------------------------------------------------------------------- // The ledger a refreshed run keeps. // // A long-lived runner re-opens its OperationRun when the operation settings // change or its TTL lapses, and the idle check that guards that swap cannot // close the window: it runs before the await that opens the new run, so one // more unit can be dispatched against the OLD one meanwhile. That unit settles // afterwards — releasing its lease and reporting its cost — against the run it // was dispatched with. So the refresh must carry the whole ledger OBJECT, not // copies of the fields that happened to look like state. function fakeDigestRun(live: OperationRun["live"]): OperationRun { return { lane: "digest", live, digest: { laneChoice: "local", engine: { app: { lane: "remote-api", metered: true, fields: { baseUrl: true } }, }, }, } as unknown as OperationRun; } test("a refreshed run keeps the lane's whole ledger, leases and spend alike", () => { const ledger = { llmActive: 0, unitActive: 0, costUsd: 0, diskFloorHit: false }; const before = fakeDigestRun(ledger); // A unit is dispatched on `before` and takes a remote lease. before.live.llmActive += 1; // The refresh installs a NEW run carrying the SAME ledger — what // autoRunner's refreshLaneRun does with `opened.live = laneRun.run.live`. const after = fakeDigestRun(before.live); assert.equal(after.live, before.live, "the ledger is shared, not copied"); // The in-flight unit now settles against the run it was dispatched with. before.live.llmActive -= 1; before.live.costUsd += 6; // And the limit, which reads the NEW run, sees both halves. Without the // shared object it would see llmActive stuck at 0 with a lease outstanding // (an under-count that never recovers) and costUsd at 0 (a spend cap the // metered lane silently stops enforcing). const liveNow = operationLaneLive(after, 1); assert.equal(liveNow.lane, "digest"); assert.equal(liveNow.lane === "digest" ? liveNow.llmActive : -1, 0); assert.equal(liveNow.lane === "digest" ? liveNow.costUsd : -1, 6); const verdict = laneLimit( settingsWith({ digest: { ...defaultDigest(), spendCapUsd: 5 } }), liveNow, ); assert.equal(verdict.limit, 0); assert.equal(verdict.hold?.reason, "spend-cap"); });