// One event-driven fill-to-capacity primitive, replacing both the auto-runner's // hand-written waitNext/wake loop and whisperBatch's Promise.all. Pulls work // from `next()` up to a dynamic `limit()`, runs items concurrently, and stops on // hard cancel (signal) or soft drain (drainSignal). Encoding the wait logic ONCE // here — with the no-spin invariant below — removes the whole class of // event-loop-starvation bugs (the "Drain all hangs the app" bug) that came from // hand-rolling this loop in two places. export type RunPoolOptions = { // Pull the next item to run, or null/undefined if nothing is available right // now. For a finite batch (`finite: true`) null means "exhausted". next: () => Promise; // Process one item. Receives the HARD-cancel signal so it can abort promptly; // a soft drain intentionally lets in-flight items finish. run: (item: T, signal: AbortSignal) => Promise; // Dynamic concurrency ceiling (e.g. eligible worker slots). Re-read each loop. limit: () => number; signal: AbortSignal; // hard cancel: stop pulling, abort in-flight, return drainSignal: AbortSignal; // soft drain: stop pulling, let in-flight finish, return // Poll interval when waiting for capacity / new work. The wait is woken early // by a finishing item or an abort, so this is just a safety ceiling. idlePollMs?: number; // true for a fixed work set that should stop once `next()` is exhausted and // nothing is in flight (batch). false (default) for a long-lived daemon that // idles waiting for new work and only stops on drain/cancel. finite?: boolean; }; export async function runPool(opts: RunPoolOptions): Promise { const { signal, drainSignal } = opts; const idle = opts.idlePollMs ?? 3000; const inFlight = new Set>(); // Wait up to maxMs, woken early by wake() (a finishing item) or by a hard // cancel / drain firing DURING the wait (the abort listeners below). // // The invariant that prevents event-loop starvation: an ALREADY-aborted signal // must NOT short-circuit to an immediately-resolved promise. After a drain (or // hard cancel) the signal stays aborted for the rest of the run, so a loop that // waits for in-flight items to clear would otherwise become a timer-less // microtask spin — starving the macrotask/timer/I/O phases, so a still-running // item's child-process exit (a macrotask) never arrives, inFlight never empties, // and a CPU core pegs forever. So we only fast-path a real pending wake; an // already-aborted signal falls through to a real timer, and a finishing item // still wakes us promptly. let pendingWake = false; let waiter: (() => void) | null = null; const wake = (): void => { if (waiter) { const w = waiter; waiter = null; w(); } else { pendingWake = true; } }; const waitNext = (maxMs: number): Promise => new Promise((resolve) => { if (pendingWake) { pendingWake = false; resolve(); return; } let done = false; const finish = (): void => { if (done) return; done = true; clearTimeout(timer); signal.removeEventListener("abort", finish); drainSignal.removeEventListener("abort", finish); if (waiter === finish) waiter = null; resolve(); }; const timer = setTimeout(finish, maxMs); timer.unref?.(); waiter = finish; // Registering on an ALREADY-aborted signal does not fire — the timer paces // the wait once draining, which is exactly what we want. signal.addEventListener("abort", finish, { once: true }); drainSignal.addEventListener("abort", finish, { once: true }); }); while (true) { if (signal.aborted) break; const drained = drainSignal.aborted; const target = drained ? 0 : Math.max(0, opts.limit()); if (inFlight.size >= target) { // At/over capacity, or draining (target 0). Only stop once nothing is in // flight AND we've been told to (drain/cancel). A transiently-zero limit // while neither draining nor cancelling just waits — it does NOT terminate. if (inFlight.size === 0 && (drainSignal.aborted || signal.aborted)) break; await waitNext(idle); continue; } const item = await opts.next(); // Re-read the abort state after the async next(): a drain/cancel may have // fired while it ran (using the live signals, not the stale top-of-loop // snapshot, makes a stop that arrives mid-next() take effect immediately). if (item === null || item === undefined) { // Nothing available. A finite batch with nothing in flight is done; a // daemon idles unless draining/cancelling. if ( inFlight.size === 0 && (drainSignal.aborted || signal.aborted || opts.finite) ) { break; } await waitNext(idle); continue; } // async thunk so a synchronous throw from run() becomes a contained // rejection (settled by Promise.allSettled) rather than escaping the pool. const p = (async () => opts.run(item, signal))(); inFlight.add(p); // Remove from the in-flight set once it settles and wake the loop so it can // pull more / observe drain completion. Use then(settle, settle) rather than // finally(): the rejection handler CONSUMES a rejecting run() (no leaked // unhandledRejection), while Promise.allSettled below still awaits `p`. `p` // is declared before `settle`, so the closure reference is safe. const settle = () => { inFlight.delete(p); wake(); }; void p.then(settle, settle); } // Let in-flight items settle: a soft drain lets them finish; a hard cancel // lets them observe `signal` and abort. Either way we don't return until the // event loop has delivered their completions. await Promise.allSettled([...inFlight]); }