import { test } from "node:test"; import assert from "node:assert/strict"; import type { Worker } from "../lib/workers"; import { WorkerPool } from "./workerPool"; // Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/workerPool.test.ts // // Covers the foreground/background waiter ordering added so a manual // (foreground) transcription preempts queued auto-runner (background) units for // the next freed worker slot — the running unit is never interrupted, it drains. // Each test builds a fresh `new WorkerPool()` (not the global singleton) and // seeds it with reconfigure(workers, {applyEnabled:true}), which avoids reading // settings.json so the scheduler is fully isolated. function worker(id: string, priority = 0, tags?: string[]): Worker { return { id, name: id, kind: "local", enabled: true, priority, ...(tags ? { tags } : {}), }; } // A single-slot pool with its only worker already leased out (running), so every // subsequent acquire parks as a waiter. Returns the pool and the running lease. async function busyPool() { const pool = new WorkerPool(); pool.reconfigure([worker("w1")], { applyEnabled: true }); const lease = await pool.acquire(); return { pool, lease }; } test("a foreground acquire is served before an earlier-parked background one", async () => { const { pool, lease } = await busyPool(); const order: string[] = []; // Background parks FIRST, foreground SECOND — arrival order would serve B first. const pB = pool .acquire(undefined, { background: true }) .then((l) => (order.push("B"), l)); const pF = pool.acquire().then((l) => (order.push("F"), l)); lease.release(); // one slot frees → highest-priority waiter (foreground) wins const leaseF = await pF; assert.deepEqual(order, ["F"]); leaseF.release(); // next slot → the background waiter await pB; assert.deepEqual(order, ["F", "B"]); }); test("foreground waiters keep FIFO among themselves", async () => { const { pool, lease } = await busyPool(); const order: string[] = []; const p1 = pool.acquire().then((l) => (order.push("1"), l)); const p2 = pool.acquire().then((l) => (order.push("2"), l)); lease.release(); const l1 = await p1; assert.deepEqual(order, ["1"]); l1.release(); await p2; assert.deepEqual(order, ["1", "2"]); }); test("a foreground acquire slots ahead of several queued background ones", async () => { const { pool, lease } = await busyPool(); const order: string[] = []; const b1 = pool .acquire(undefined, { background: true }) .then((l) => (order.push("B1"), l)); const b2 = pool .acquire(undefined, { background: true }) .then((l) => (order.push("B2"), l)); const b3 = pool .acquire(undefined, { background: true }) .then((l) => (order.push("B3"), l)); // Arrives last but is foreground → must be served first, then B1,B2,B3 in order. const f = pool.acquire().then((l) => (order.push("F"), l)); let lease_ = lease; for (const p of [f, b1, b2, b3]) { lease_.release(); lease_ = await p; } assert.deepEqual(order, ["F", "B1", "B2", "B3"]); }); test("aborting a parked waiter removes it without disturbing the order", async () => { const { pool, lease } = await busyPool(); const order: string[] = []; const ac = new AbortController(); const p1 = pool.acquire().then((l) => (order.push("1"), l)); const p2 = pool .acquire(ac.signal) .then((l) => (order.push("2"), l)) .catch((e) => (order.push("2-aborted"), Promise.reject(e))); const p3 = pool.acquire().then((l) => (order.push("3"), l)); ac.abort(); // cancel the middle waiter while parked await assert.rejects(p2, /abort/i); lease.release(); const l1 = await p1; l1.release(); await p3; assert.deepEqual(order, ["2-aborted", "1", "3"]); }); test("all-background acquires fall back to plain FIFO (backward compatible)", async () => { const { pool, lease } = await busyPool(); const order: string[] = []; const b1 = pool .acquire(undefined, { background: true }) .then((l) => (order.push("B1"), l)); const b2 = pool .acquire(undefined, { background: true }) .then((l) => (order.push("B2"), l)); lease.release(); const l1 = await b1; assert.deepEqual(order, ["B1"]); l1.release(); await b2; assert.deepEqual(order, ["B1", "B2"]); }); test("tier ordering: urgent jumps ahead of foreground, which jumps ahead of background", async () => { const { pool, lease } = await busyPool(); const order: string[] = []; // Park in arrival order background, foreground, urgent — served in tier order. const b = pool .acquire(undefined, { tier: "background" }) .then((l) => (order.push("B"), l)); const f = pool .acquire(undefined, { tier: "foreground" }) .then((l) => (order.push("F"), l)); const u = pool .acquire(undefined, { tier: "urgent" }) .then((l) => (order.push("U"), l)); let lease_ = lease; for (const p of [u, f, b]) { lease_.release(); lease_ = await p; } assert.deepEqual(order, ["U", "F", "B"]); }); // --- Capability matching (Worker.tags + acquire requires) --- test("an untagged worker satisfies any requirement", async () => { const pool = new WorkerPool(); pool.reconfigure([worker("w1")], { applyEnabled: true }); const lease = await pool.acquire(undefined, { requires: ["diarization"] }); assert.equal(lease.worker.id, "w1"); lease.release(); }); test("a tagged worker matches when the requirement intersects its tags", async () => { const pool = new WorkerPool(); pool.reconfigure([worker("cpu-box", 0, ["cpu", "diarization"])], { applyEnabled: true, }); // Any one qualification suffices — the requirement lists acceptable ones. const lease = await pool.acquire(undefined, { requires: ["diarization", "gpu"], }); assert.equal(lease.worker.id, "cpu-box"); lease.release(); }); test("a requirement no configured worker matches rejects instead of parking forever", async () => { const pool = new WorkerPool(); pool.reconfigure([worker("cpu-box", 0, ["cpu"])], { applyEnabled: true }); await assert.rejects( pool.acquire(undefined, { requires: ["gpu"] }), /no configured worker matches/, ); }); test("a freed slot skips a head waiter it cannot serve (no head-of-line deadlock)", async () => { const pool = new WorkerPool(); pool.reconfigure( [worker("gpu-w", 0, ["gpu"]), worker("cpu-w", 1, ["cpu"])], { applyEnabled: true }, ); const gpuLease = await pool.acquire(undefined, { requires: ["gpu"] }); const cpuLease = await pool.acquire(undefined, { requires: ["cpu"] }); const order: string[] = []; // The GPU waiter parks FIRST — a blind head grant would hand it the freed // CPU slot (which it cannot use) or park the CPU waiter behind it forever. const pGpu = pool .acquire(undefined, { requires: ["gpu"] }) .then((l) => (order.push("gpu"), l)); const pCpu = pool .acquire(undefined, { requires: ["cpu"] }) .then((l) => (order.push("cpu"), l)); cpuLease.release(); // frees the CPU slot → the CPU waiter, not the head await pCpu; assert.deepEqual(order, ["cpu"]); gpuLease.release(); await pGpu; assert.deepEqual(order, ["cpu", "gpu"]); }); test("tier order is preserved within the waiters a slot can serve", async () => { const pool = new WorkerPool(); pool.reconfigure([worker("cpu-w", 0, ["cpu"])], { applyEnabled: true }); const lease = await pool.acquire(undefined, { requires: ["cpu"] }); const order: string[] = []; const b = pool .acquire(undefined, { tier: "background", requires: ["cpu"] }) .then((l) => (order.push("B"), l)); const f = pool .acquire(undefined, { tier: "foreground", requires: ["cpu"] }) .then((l) => (order.push("F"), l)); let lease_ = lease; for (const p of [f, b]) { lease_.release(); lease_ = await p; } assert.deepEqual(order, ["F", "B"]); }); test("tryAcquire claims a free matching slot and never parks", async () => { const pool = new WorkerPool(); pool.reconfigure([worker("cpu-w", 0, ["cpu"])], { applyEnabled: true }); const lease = pool.tryAcquire(["cpu"]); assert.ok(lease, "expected a lease on the free matching worker"); assert.equal(lease!.worker.id, "cpu-w"); // Busy now — a second claim yields null rather than parking. assert.equal(pool.tryAcquire(["cpu"]), null); // And a non-matching requirement never claims it. lease!.release(); assert.equal(pool.tryAcquire(["gpu"]), null); }); test("tryAcquire's kind filter keeps a fan-out off local workers", async () => { const pool = new WorkerPool(); pool.reconfigure([worker("w1")], { applyEnabled: true }); // The untagged local worker matches the requirement, but the claim is // restricted to remote capacity — the fan-out must never grab a local slot // out from under the transcription scheduler. assert.equal(pool.tryAcquire(["digest"], { kind: "remote" }), null); const lease = pool.tryAcquire(["digest"]); assert.ok(lease, "unrestricted claim still matches the untagged local"); lease!.release(); }); test("a slot freed while an unmatchable waiter parks stays claimable by others", async () => { const pool = new WorkerPool(); pool.reconfigure( [worker("gpu-w", 0, ["gpu"]), worker("cpu-w", 1, ["cpu"])], { applyEnabled: true }, ); const gpuLease = await pool.acquire(undefined, { requires: ["gpu"] }); const cpuLease = await pool.acquire(undefined, { requires: ["cpu"] }); const pGpu = pool.acquire(undefined, { requires: ["gpu"] }); // Freeing the CPU slot cannot serve the parked GPU waiter… cpuLease.release(); // …so it stays free, and a claim that does match it succeeds without // starving the waiter (which the freed slot could never serve). const claimed = pool.tryAcquire(["cpu"]); assert.ok(claimed, "the free CPU slot should be claimable"); claimed!.release(); gpuLease.release(); (await pGpu).release(); }); test("hasEligibleWorker honours a requirement", () => { const pool = new WorkerPool(); pool.reconfigure([worker("cpu-w", 0, ["cpu"])], { applyEnabled: true }); assert.equal(pool.hasEligibleWorker(), true); assert.equal(pool.hasEligibleWorker(["cpu"]), true); assert.equal(pool.hasEligibleWorker(["gpu"]), false); }); test("summary() carries tags and keeps its shape for untagged workers", () => { const pool = new WorkerPool(); pool.reconfigure([worker("plain"), worker("tagged", 1, ["cpu"])], { applyEnabled: true, }); const summary = pool.summary(); assert.equal(summary.find((w) => w.id === "plain")?.tags, undefined); assert.deepEqual(summary.find((w) => w.id === "tagged")?.tags, ["cpu"]); }); // --- Remote slot expansion (RemoteWorkerConfig.slots) --- function remoteWorker(id: string, slots?: number, priority = 0): Worker { return { id, name: id, kind: "remote", enabled: true, priority, // An explicit slot count means the capacity probe never fires, keeping // these tests fully offline. remote: { baseUrl: "http://127.0.0.1:9", ...(slots ? { slots } : {}) }, }; } test("a remote with slots: N expands to N independently-leasable entries", async () => { const pool = new WorkerPool(); pool.reconfigure([remoteWorker("box", 3)], { applyEnabled: true }); const ids = pool.summary().map((w) => w.id); // Slot 1 keeps the BASE id so saved defaults and degrade marks keep meaning // what they meant; extras are #2..#N. assert.deepEqual(ids.sort(), ["box", "box#2", "box#3"]); const l1 = await pool.acquire(); const l2 = await pool.acquire(); const l3 = await pool.acquire(); assert.equal(new Set([l1, l2, l3].map((l) => l.worker.id)).size, 3); // All three busy — a fourth claim yields nothing. assert.equal(pool.tryAcquire([]), null); for (const l of [l1, l2, l3]) l.release(); }); test("shrinking a remote's slots retires the surplus (drain if busy)", async () => { const pool = new WorkerPool(); pool.reconfigure([remoteWorker("box", 3)], { applyEnabled: true }); const l1 = await pool.acquire(); // box const l2 = await pool.acquire(); // box#2 pool.reconfigure([remoteWorker("box", 1)], { applyEnabled: true }); // The idle surplus slot is gone at once; the busy one drains. const byId = new Map(pool.summary().map((w) => [w.id, w])); assert.equal(byId.has("box#3"), false); assert.equal(byId.get("box#2")?.state, "draining"); l2.release(); assert.equal( pool.summary().some((w) => w.id === "box#2"), false, "the drained surplus slot is dropped when its lease releases", ); l1.release(); }); test("degrading one remote slot leaves its siblings eligible", async () => { const pool = new WorkerPool(); pool.reconfigure([remoteWorker("box", 2)], { applyEnabled: true }); pool.markDegraded("box"); const lease = await pool.acquire(); assert.equal(lease.worker.id, "box#2"); lease.release(); }); test("a slot added by expansion inherits the base slot's runtime state", async () => { const pool = new WorkerPool(); pool.reconfigure([remoteWorker("box", 1)], { applyEnabled: true }); // Operator disables the remote from the Workers page (runtime, not settings)… pool.disableWorker("box"); // …then a reconfigure grows the slot count (as a landed probe would). pool.reconfigure([remoteWorker("box", 2)]); assert.equal( pool.summary().find((w) => w.id === "box#2")?.state, "disabled", "expansion must not re-enable a remote the operator just disabled", ); }); // --- LLM endpoint workers --- function llmWorker(id: string, slots?: number, tags?: string[]): Worker { return { id, name: id, kind: "llm", enabled: true, priority: 0, ...(tags ? { tags } : {}), llm: { baseUrl: "http://127.0.0.1:9", ...(slots ? { slots } : {}) }, }; } test("an llm worker never satisfies a requirement-less claim", () => { const pool = new WorkerPool(); pool.reconfigure([llmWorker("mac")], { applyEnabled: true }); // Today's transcription path: no requirement. The endpoint must be invisible // to it — including to the "paused — waiting for a worker" probe. assert.equal(pool.tryAcquire([]), null); assert.equal(pool.hasEligibleWorker(), false); // But it serves its LLM operations. assert.equal(pool.hasEligibleWorker(["digest"]), true); const lease = pool.tryAcquire(["digest"], { kind: "llm" }); assert.ok(lease); lease!.release(); }); test("taggedOnly claims skip untagged workers (the unit-dispatch opt-in gate)", () => { const pool = new WorkerPool(); pool.reconfigure( [remoteWorker("plain-remote", 1, 0), worker("cpu-box", 1, ["cpu"])], { applyEnabled: true }, ); // The untagged remote matches everything under the universal rule, but an // untagged remote predates the unit protocol — units go only to workers the // operator explicitly tagged. assert.equal( pool.tryAcquire(["diarization", "cpu"], { kind: "remote", taggedOnly: true, }), null, ); assert.equal( pool.freeSlots(["diarization", "cpu"], { kind: "remote", taggedOnly: true }), 0, ); // Tag the remote and the same claim lands. pool.reconfigure( [ { ...remoteWorker("plain-remote", 1, 0), tags: ["cpu"] }, worker("cpu-box", 1, ["cpu"]), ], { applyEnabled: true }, ); const lease = pool.tryAcquire(["diarization", "cpu"], { kind: "remote", taggedOnly: true, }); assert.equal(lease?.worker.id, "plain-remote"); lease?.release(); }); test("an llm worker's slots expand like a remote's", async () => { const pool = new WorkerPool(); pool.reconfigure([llmWorker("mac", 2)], { applyEnabled: true }); assert.equal(pool.freeSlots(["digest"], { kind: "llm" }), 2); const l1 = pool.tryAcquire(["digest"], { kind: "llm" }); assert.ok(l1); assert.equal(pool.freeSlots(["digest"], { kind: "llm" }), 1); const l2 = pool.tryAcquire(["digest"], { kind: "llm" }); assert.ok(l2); assert.equal(pool.tryAcquire(["digest"], { kind: "llm" }), null); l1!.release(); l2!.release(); }); test("legacy { background: true } maps to the background tier", async () => { const { pool, lease } = await busyPool(); const order: string[] = []; // A background (legacy flag) waiter parks first; a tier:foreground waiter must // still preempt it, proving the boolean maps onto the same comparator. const b = pool .acquire(undefined, { background: true }) .then((l) => (order.push("B"), l)); const f = pool .acquire(undefined, { tier: "foreground" }) .then((l) => (order.push("F"), l)); lease.release(); const lf = await f; assert.deepEqual(order, ["F"]); lf.release(); await b; assert.deepEqual(order, ["F", "B"]); }); // opts.only (WorkerFilter): the one-off file transcription names a worker, or // keeps to local ones. The filter narrows the free-list and the parked grant; // a filter no configured worker passes is refused at once rather than parked. test("an acquire narrowed by `only` takes the named worker, even when a better one is free", async () => { const pool = new WorkerPool(); pool.reconfigure([worker("gpu", 0), worker("cpu", 5)], { applyEnabled: true }); const lease = await pool.acquire(undefined, { only: (w) => w.id === "cpu" }); assert.equal(lease.worker.id, "cpu"); lease.release(); }); test("an acquire narrowed by `only` parks for its worker and skips others that free", async () => { const pool = new WorkerPool(); pool.reconfigure([worker("gpu", 0), worker("cpu", 5)], { applyEnabled: true }); const gpu = await pool.acquire(); const cpu = await pool.acquire(); assert.equal(gpu.worker.id, "gpu"); let got: string | null = null; const p = pool .acquire(undefined, { only: (w) => w.id === "cpu" }) .then((l) => ((got = l.worker.id), l)); gpu.release(); // the wrong worker frees: the narrowed waiter stays parked await new Promise((r) => setImmediate(r)); assert.equal(got, null); cpu.release(); const l = await p; assert.equal(l.worker.id, "cpu"); l.release(); }); test("an `only` filter no configured worker passes is refused, not parked", async () => { const pool = new WorkerPool(); pool.reconfigure([worker("w1")], { applyEnabled: true }); await assert.rejects( pool.acquire(undefined, { only: (w) => w.id === "nope" }), /no configured worker can take this work/, ); });