import { test } from "node:test"; import assert from "node:assert/strict"; import { runPool } from "./concurrentRunner"; // Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/concurrentRunner.test.ts // // runPool is the single audited fill-to-capacity primitive. The headline test is // the already-aborted-signal no-spin case: it reproduces the "Drain all hangs" // bug class, where a microtask-only wait starves the macrotask that delivers a // running item's completion. With the correct wait, these resolve promptly; a // regression would hang, which the timeout guard surfaces as a failure. const never = new AbortController().signal; const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); function withTimeout(p: Promise, ms: number, msg: string): Promise { return Promise.race([ p, sleep(ms).then(() => { throw new Error(`timeout: ${msg}`); }), ]) as Promise; } test("respects the dynamic limit(): never exceeds the ceiling", async () => { let i = 0; const items = Array.from({ length: 8 }, (_, n) => n); let inFlight = 0; let maxInFlight = 0; let ran = 0; await runPool({ next: async () => (i < items.length ? items[i++] : null), run: async () => { inFlight++; maxInFlight = Math.max(maxInFlight, inFlight); await sleep(5); ran++; inFlight--; }, limit: () => 2, signal: never, drainSignal: new AbortController().signal, idlePollMs: 50, finite: true, }); assert.equal(ran, 8, "all items ran"); assert.ok(maxInFlight <= 2, `max concurrent was ${maxInFlight}, expected <= 2`); }); test("finite mode stops once next() is exhausted and nothing is in flight", async () => { let i = 0; let ran = 0; await withTimeout( runPool({ next: async () => (i < 3 ? i++ : null), run: async () => { await sleep(1); ran++; }, limit: () => 4, signal: never, drainSignal: new AbortController().signal, finite: true, }), 2000, "finite pool should terminate", ); assert.equal(ran, 3, "all 3 items ran, then the pool terminated on exhaustion"); }); test("drain stops pulling new work but lets in-flight finish (no-spin)", async () => { const drainAc = new AbortController(); let pulled = 0; let ranAfterDrain = 0; let drained = false; await withTimeout( runPool({ // Daemon: would pull forever, but drain fires from inside the first run. next: async () => { if (drained) { // Should never be pulled after drain (target drops to 0). ranAfterDrain++; return null; } return pulled++; }, run: async () => { // Simulate "drain fires while a unit is actively running", then the unit // completes via a macrotask (setTimeout). If the post-drain wait spun on // microtasks, this timer would be starved and the pool would hang. drainAc.abort(); drained = true; await sleep(20); }, limit: () => 1, signal: never, drainSignal: drainAc.signal, idlePollMs: 1000, }), 2000, "drain with a running unit must complete, not hang", ); assert.equal(pulled, 1, "exactly one item was pulled before drain"); assert.equal(ranAfterDrain, 0, "no work pulled after drain"); }); test("already-aborted drain at submit time: in-flight unit still completes", async () => { // Pre-abort the drain, but stage one unit that is launched in the SAME tick it // observes the abort. The classic starvation case: signal already aborted, a // real timer (not a microtask) must pace the wait for the running unit. const drainAc = new AbortController(); let launched = false; let finished = false; await withTimeout( runPool({ next: async () => { if (launched) return null; launched = true; // Abort the drain right as we hand back the first (and only) item, so the // very next loop turn is already-drained with a unit about to run. drainAc.abort(); return 0; }, run: async () => { await sleep(20); // completes via a macrotask finished = true; }, limit: () => 1, signal: never, drainSignal: drainAc.signal, idlePollMs: 1000, }), 2000, "already-aborted drain must not starve the running unit", ); // The unit was pulled before the loop re-checked the (now aborted) drain, so it // runs to completion under soft drain. assert.ok(finished, "the in-flight unit completed under drain"); }); test("hard cancel aborts in-flight work and returns", async () => { const cancelAc = new AbortController(); let observedAbort = false; await withTimeout( runPool({ next: async () => 0, run: async (_item, signal) => { await new Promise((resolve) => { if (signal.aborted) { observedAbort = true; resolve(); return; } signal.addEventListener( "abort", () => { observedAbort = true; resolve(); }, { once: true }, ); // Fire the hard cancel shortly after the unit starts. setTimeout(() => cancelAc.abort(), 10); }); }, limit: () => 1, signal: cancelAc.signal, drainSignal: new AbortController().signal, idlePollMs: 1000, }), 2000, "hard cancel must terminate the pool", ); assert.ok(observedAbort, "in-flight run() observed the hard-cancel signal"); });