commit dec95a926b9ed146db505fd22a41b751b20df572
parent 065042ab3d891035385b7d1c579829ef103bb219
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 26 Jun 2026 15:44:06 -0400
Phase 6: run the auto-runner and whisperBatch on the shared runPool
Replace both hand-rolled concurrency loops with the single audited runPool
primitive, deleting the bespoke wake/pendingWake/waiter/waitNext machinery and
the terminal `while (inFlight) await waitNext` drain-wait. This eliminates the
event-loop-starvation bug class ("Drain all hangs the app") by construction:
the one no-spin wait lives in runPool and is unit-tested.
- autoRunner.ts: runLoop's fill-to-capacity loop becomes a runPool call. `next`
selects via the unchanged policy (+ download per-platform backoff/in-flight
gating) and RESERVES the slot so the next selection sees it; `run` executes
launchUnit and releases with the exact former `.finally` accounting; `limit`
is eligibleSlots()/maxWorkers. A graceful stop (job gone after e2e reset, or
policy disabled) is modeled as a soft drain via a stopController combined into
the drain signal. Hard-cancel still actively cancels in-flight child download
jobs via a signal abort listener.
- whisperBatch.ts: Promise.all over all videos becomes a finite runPool whose
limit is the pool's eligible slot count — fills the pool without dumping a
backlog of parked acquires; drain stops pulling and lets in-flight finish; 0
workers ⇒ paused (not failed).
- concurrentRunner.ts: re-read the live abort signals after the async next() so
a stop arriving mid-next() takes effect immediately; consume a rejecting run()
via then(settle, settle) instead of void finally() (no leaked unhandledReject).
Verification: 33 jobs unit tests green; e2e auto-queue 11/11 (both drain
regressions) + jobs-batch-tasks-drain 6/6, zero unhandled rejections in the
server log. Typecheck (common+editor) clean.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Diffstat:
3 files changed, 180 insertions(+), 178 deletions(-)
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -6,6 +6,7 @@ import { detectPlatform } from "../lib/platform";
import { getWorkerPool } from "../jobs/workerPool";
import { getRegistry } from "../jobs/registry";
import { runManagedFunction } from "../jobs/streamCommand";
+import { runPool } from "../jobs/concurrentRunner";
import { requestChannelSnapshot } from "../jobs/snapshotScheduler";
import { makeTaskTracker } from "../jobs/taskHooks";
import type { JobRunContext } from "../jobs/streamCommand";
@@ -220,6 +221,10 @@ function eligibleSlots(): number {
// --- The fill-to-capacity loop ---------------------------------------------
+// One unit of work selected by the policy, carried from the runPool `next`
+// (selection + reservation) into `run` (execution + release).
+type Picked = { pick: WorkPick; channelSlug: string; unitPlatform: string | null };
+
async function runLoop(
kind: AutoQueueKind,
paths: Paths,
@@ -239,56 +244,7 @@ async function runLoop(
pruneExpired(kindState.platformBackoff, Date.now());
const completed = new Set<string>();
const completedOrder: string[] = [];
- let counters = { transcribed: 0, failed: 0, skipped: 0 };
-
- // Wake coordination: a finishing unit (or an abort) wakes the loop out of its
- // wait. pendingWake guards against a wake that fires between waits.
- let waiter: (() => void) | null = null;
- let pendingWake = false;
- const wake = () => {
- if (waiter) {
- const w = waiter;
- waiter = null;
- w();
- } else {
- pendingWake = true;
- }
- };
- // Wait up to maxMs, woken early by wake() (a finishing unit) or by a hard
- // cancel / drain firing during the wait (the abort listeners below). 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 the loops that wait for in-flight units to clear
- // (`while (inFlight) await waitNext()`, and the at-capacity wait once
- // target == 0 while draining) would become a timer-less microtask spin that
- // starves the event loop — a still-running unit's child-process exit (a
- // macrotask) would then never be delivered, inFlight would never reach zero,
- // and a CPU core would peg forever (the "Drain all hangs the app" bug). So we
- // only fast-path a real pending wake here; an already-aborted signal falls
- // through to a real timer, and a finishing unit still wakes us promptly.
- const waitNext = (maxMs: number): Promise<void> =>
- new Promise<void>((resolve) => {
- if (pendingWake) {
- pendingWake = false;
- resolve();
- return;
- }
- let done = false;
- const finish = () => {
- if (done) return;
- done = true;
- clearTimeout(timer);
- signal.removeEventListener("abort", finish);
- ctx.drainSignal.removeEventListener("abort", finish);
- if (waiter === finish) waiter = null;
- resolve();
- };
- const timer = setTimeout(finish, maxMs);
- timer.unref?.();
- waiter = finish;
- signal.addEventListener("abort", finish, { once: true });
- ctx.drainSignal.addEventListener("abort", finish, { once: true });
- });
+ const counters = { transcribed: 0, failed: 0, skipped: 0 };
const markCompleted = (videoId: string) => {
if (completed.has(videoId)) return;
@@ -322,6 +278,12 @@ async function runLoop(
let metaCache: ChannelMeta[] = [];
let metaAt = 0;
+ // A graceful stop (the job record is gone after an e2e reset, or the policy was
+ // disabled) is modeled as a soft drain: stop picking, let in-flight finish.
+ // Combined with the job's real drain signal so either one ends the run.
+ const stopController = new AbortController();
+ const drainSignal = AbortSignal.any([ctx.drainSignal, stopController.signal]);
+
// NOTE: do NOT call getWorkerPool().reconfigure() here. At server boot the
// runner is the first thing to touch the pool, and reconfigure() sets the
// pool's `initialized` flag WITHOUT applying the saved "Set as default"
@@ -332,36 +294,45 @@ async function runLoop(
// exactly once. Live settings changes are pushed by saveSettingsAction's own
// reconfigure(applyEnabled:true), so nothing is lost by not refreshing here.
- onLog(`Auto-${kind} runner started.`);
+ // On a HARD cancel, actively cancel each in-flight/queued child download job: a
+ // running one's yt-dlp is SIGTERM'd (via downloadOneManaged's signal), a
+ // still-queued one is dropped — either way launchUnit's `await res.done`
+ // resolves and the unit's run() settles so runPool's in-flight clears. (Drain
+ // is soft: it leaves them to finish. Transcription units settle via the hard
+ // signal threaded into the unit, so they need no explicit cancel here.)
+ const onHardAbort = () => {
+ for (const jid of childJobIds.values()) getRegistry().cancel(jid);
+ };
+ if (signal.aborted) onHardAbort();
+ else signal.addEventListener("abort", onHardAbort, { once: true });
+
+ // Transcription is throttled by the worker pool's slot count; download has no
+ // pool, so its real limiter is the per-platform gate in next() — leave the
+ // numeric cap at the policy max (or unbounded). runPool drops the target to 0
+ // itself while draining, so this never needs to check the drain signal.
+ const limit = (): number => {
+ const policy = getSettings().autoQueue[kind];
+ return kind === "transcription"
+ ? Math.min(eligibleSlots(), policy.maxWorkers ?? Number.POSITIVE_INFINITY)
+ : (policy.maxWorkers ?? Number.POSITIVE_INFINITY);
+ };
- while (!signal.aborted) {
+ // Select the next unit and RESERVE its slot (so the next selection sees it),
+ // or return null when there's nothing to pick right now / the runner should
+ // stop. runPool calls this only while below the limit and not draining.
+ const next = async (): Promise<Picked | null> => {
// Stop if the job record is gone (e2e registry reset) or no longer running
- // (hard cancel flips it to "cancelled"). Drain keeps it "running" and is
- // handled via ctx.drainSignal below.
+ // (a hard cancel flips it to "cancelled" and also fires `signal`).
const rec = getRegistry().get(live.jobId);
- if (!rec || rec.status !== "running") break;
-
+ if (!rec || rec.status !== "running") {
+ stopController.abort();
+ return null;
+ }
const policy: AutoQueuePolicy = getSettings().autoQueue[kind];
if (!policy.enabled) {
onLog(`Auto-${kind} disabled — stopping runner.`);
- break;
- }
- const drained = ctx.drainSignal.aborted;
- // Transcription is throttled by the worker pool's slot count; download has no
- // pool, so its real limiter is the per-platform gate below — leave the
- // numeric cap at the policy max (or unbounded).
- const target = drained
- ? 0
- : kind === "transcription"
- ? Math.min(eligibleSlots(), policy.maxWorkers ?? Number.POSITIVE_INFINITY)
- : (policy.maxWorkers ?? Number.POSITIVE_INFINITY);
-
- if (live.inFlight.size >= target) {
- // At capacity (or paused/drained). When draining/cancelling and nothing is
- // left in flight, the run is complete.
- if (live.inFlight.size === 0 && (drained || signal.aborted)) break;
- await waitNext(IDLE_POLL_MS);
- continue;
+ stopController.abort();
+ return null;
}
// Refresh the (rarely-changing) channel list/platforms on a TTL.
@@ -413,93 +384,19 @@ async function runLoop(
}
const pick = selectNextWork(policy.root, pending, runtime, live.active);
- if (!pick) {
- if (live.inFlight.size === 0 && (drained || signal.aborted)) break;
- await waitNext(IDLE_POLL_MS);
- continue;
- }
+ if (!pick) return null;
const channelSlug = owner.get(pick.videoId);
if (!channelSlug) {
// Shouldn't happen (every pending id came from a channel) — guard anyway.
markCompleted(pick.videoId);
- continue;
+ return null;
}
const unitPlatform =
kind === "download" ? platformKey(channelSlug, slugToPlatform) : null;
- let result: UnitResult = { outcome: "failed" };
- launchUnit({
- kind,
- paths,
- pick,
- channelSlug,
- tracker,
- onLog,
- signal,
- drainSignal: ctx.drainSignal,
- onChildJob: (jid) => childJobIds.set(pick.videoId, jid),
- })
- .then((r) => {
- result = r;
- })
- .catch((err) => {
- result = { outcome: "failed" };
- onLog(`Auto-${kind} unit error on ${pick.videoId}: ${String(err)}`);
- })
- .finally(() => {
- childJobIds.delete(pick.videoId);
- live.inFlight.delete(pick.videoId);
- for (const id of pick.path) {
- live.active[id] = Math.max(0, (live.active[id] ?? 0) - 1);
- }
- if (unitPlatform) {
- platformInFlight.set(
- unitPlatform,
- Math.max(0, (platformInFlight.get(unitPlatform) ?? 0) - 1),
- );
- }
-
- if (result.outcome === "transcribed") counters.transcribed++;
- else if (result.outcome === "failed") counters.failed++;
- else counters.skipped++;
-
- // Per-platform backoff: a rate-limit / network failure pauses the whole
- // platform with an exponential cooldown, and the video is deliberately
- // NOT marked completed so it's retried once the cooldown lapses. Any
- // other outcome marks the video done-for-session; a success also clears
- // the platform's cooldown.
- const backoffHit =
- unitPlatform !== null &&
- (result.failureClass === "rate_limit" ||
- result.failureClass === "network");
- if (backoffHit && unitPlatform) {
- const entry = nextBackoff(
- kindState.platformBackoff[unitPlatform],
- Date.now(),
- );
- kindState.platformBackoff[unitPlatform] = entry;
- const secs = Math.round((entry.until - Date.now()) / 1000);
- onLog(
- `Auto-download: ${unitPlatform} ${result.failureClass} — backing off ${secs}s (attempt ${entry.fails}). ${pick.videoId} will retry after cooldown.`,
- );
- } else {
- if (unitPlatform && result.outcome === "transcribed") {
- clearBackoff(kindState.platformBackoff, unitPlatform);
- }
- markCompleted(pick.videoId);
- }
- // Persist so the cooldown (and reset) survive a restart.
- persist();
- // The runner job has no channelSlug, so the registry's own
- // snapshot-on-finish hook doesn't fire — request a regen here so the
- // channel's pending counts refresh and a restart sees fresh buckets.
- requestChannelSnapshot(paths, channelSlug);
- wake();
- });
-
- // Reserve the slot and account active counts up the path BEFORE looping, so
- // the next pick sees the reservation (and ancestor caps reflect it).
+ // Reserve the slot and account active counts up the path BEFORE returning,
+ // so the next selection sees the reservation (and ancestor caps reflect it).
live.inFlight.set(pick.videoId, {
videoId: pick.videoId,
leafId: pick.leafId,
@@ -521,17 +418,95 @@ async function runLoop(
});
persist();
onLog(`Auto-${kind}: ${channelSlug}/${pick.videoId} (leaf ${pick.leafId})`);
- }
+ return { pick, channelSlug, unitPlatform };
+ };
+
+ // Run one reserved unit, then release its reservation and fold its outcome
+ // into counters/backoff. Mirrors the old `.finally` accounting exactly.
+ const run = async (picked: Picked, runSignal: AbortSignal): Promise<void> => {
+ const { pick, channelSlug, unitPlatform } = picked;
+ let result: UnitResult = { outcome: "failed" };
+ try {
+ result = await launchUnit({
+ kind,
+ paths,
+ pick,
+ channelSlug,
+ tracker,
+ onLog,
+ signal: runSignal,
+ drainSignal: ctx.drainSignal,
+ onChildJob: (jid) => childJobIds.set(pick.videoId, jid),
+ });
+ } catch (err) {
+ result = { outcome: "failed" };
+ onLog(`Auto-${kind} unit error on ${pick.videoId}: ${String(err)}`);
+ } finally {
+ childJobIds.delete(pick.videoId);
+ live.inFlight.delete(pick.videoId);
+ for (const id of pick.path) {
+ live.active[id] = Math.max(0, (live.active[id] ?? 0) - 1);
+ }
+ if (unitPlatform) {
+ platformInFlight.set(
+ unitPlatform,
+ Math.max(0, (platformInFlight.get(unitPlatform) ?? 0) - 1),
+ );
+ }
+
+ if (result.outcome === "transcribed") counters.transcribed++;
+ else if (result.outcome === "failed") counters.failed++;
+ else counters.skipped++;
+
+ // Per-platform backoff: a rate-limit / network failure pauses the whole
+ // platform with an exponential cooldown, and the video is deliberately NOT
+ // marked completed so it's retried once the cooldown lapses. Any other
+ // outcome marks the video done-for-session; a success also clears the
+ // platform's cooldown.
+ const backoffHit =
+ unitPlatform !== null &&
+ (result.failureClass === "rate_limit" ||
+ result.failureClass === "network");
+ if (backoffHit && unitPlatform) {
+ const entry = nextBackoff(
+ kindState.platformBackoff[unitPlatform],
+ Date.now(),
+ );
+ kindState.platformBackoff[unitPlatform] = entry;
+ const secs = Math.round((entry.until - Date.now()) / 1000);
+ onLog(
+ `Auto-download: ${unitPlatform} ${result.failureClass} — backing off ${secs}s (attempt ${entry.fails}). ${pick.videoId} will retry after cooldown.`,
+ );
+ } else {
+ if (unitPlatform && result.outcome === "transcribed") {
+ clearBackoff(kindState.platformBackoff, unitPlatform);
+ }
+ markCompleted(pick.videoId);
+ }
+ // Persist so the cooldown (and reset) survive a restart.
+ persist();
+ // The runner job has no channelSlug, so the registry's own
+ // snapshot-on-finish hook doesn't fire — request a regen here so the
+ // channel's pending counts refresh and a restart sees fresh buckets.
+ requestChannelSnapshot(paths, channelSlug);
+ }
+ };
+
+ onLog(`Auto-${kind} runner started.`);
+
+ // One audited fill-to-capacity primitive replaces the old hand-written
+ // wake/waitNext loop AND its terminal drain-wait. runPool stops pulling on a
+ // drain (or stop), lets in-flight units finish, and never spins the event
+ // loop on an already-aborted signal — the "Drain all hangs the app" bug class.
+ await runPool<Picked>({
+ next,
+ run,
+ limit,
+ signal,
+ drainSignal,
+ idlePollMs: IDLE_POLL_MS,
+ });
- // Drain remaining in-flight units (graceful stop / drain / cancel). On a HARD
- // cancel, actively cancel each in-flight/queued child download job: a running
- // one's yt-dlp is SIGTERM'd (via downloadOneManaged's signal), a still-queued
- // one is dropped — either way launchUnit's `await res.done` resolves and the
- // unit's .finally clears the slot. Drain (soft) leaves them to finish.
- if (signal.aborted) {
- for (const jid of childJobIds.values()) getRegistry().cancel(jid);
- }
- while (live.inFlight.size > 0) await waitNext(IDLE_POLL_MS);
persist();
onLog(
`Auto-${kind} runner finished: ${counters.transcribed} done, ${counters.failed} failed, ${counters.skipped} skipped.`,
diff --git a/common/controller/whisperBatch.ts b/common/controller/whisperBatch.ts
@@ -8,6 +8,7 @@ import {
readVideoFiles,
} from "../lib/videoStatus";
import { transcribeOneFromQueue } from "./transcribeOneFromQueue";
+import { runPool } from "../jobs/concurrentRunner";
import { pruneFailedTranscriptions } from "./failedTranscriptions";
import { resolveShardItems } from "./shard";
import { countNotYetTranscribed } from "./channels";
@@ -190,13 +191,17 @@ export async function runWhisperBatch({
let failed = 0;
let skipped = 0;
- // No pLimit: every video starts concurrently and parks in pool.acquire() until
- // a worker slot frees. The pool's total free slots throttle the batch, and a
- // freed high-priority slot is handed to the oldest waiter. When all workers are
- // disabled, acquire() blocks indefinitely — the batch pauses (it does not fail)
- // and resumes when a worker is re-enabled. The per-video gating + transcribe is
- // shared with the cross-channel auto-runner via transcribeOneFromQueue.
- const runOne = async (videoDir: string): Promise<void> => {
+ // The same runPool primitive the cross-channel auto-runner uses, instead of
+ // Promise.all over every video at once. It keeps up to `limit()` units in
+ // flight (the pool's eligible slot count), pulling the next as each finishes —
+ // so the batch fills the pool without dumping a backlog of parked acquires, and
+ // a drain stops pulling while letting in-flight units finish (no event-loop
+ // spin). When all workers are disabled the limit is 0 and the batch PAUSES (it
+ // does not fail), resuming when a worker is re-enabled. Per-video gating +
+ // transcribe is shared with the auto-runner via transcribeOneFromQueue.
+ const NEVER = new AbortController().signal;
+ let idx = 0;
+ const runOne = async (videoDir: string, runSignal: AbortSignal): Promise<void> => {
const res = await transcribeOneFromQueue({
paths,
channelSlug,
@@ -206,7 +211,7 @@ export async function runWhisperBatch({
failedSet,
tracker,
onLog: log,
- signal,
+ signal: runSignal,
drainSignal,
});
if (res.attempted) attempted++;
@@ -215,7 +220,18 @@ export async function runWhisperBatch({
else skipped++;
};
- await Promise.all(videoDirs.map(runOne));
+ await runPool<string>({
+ next: async () => (idx < videoDirs.length ? videoDirs[idx++] : null),
+ run: runOne,
+ // Eligible (enabled, non-degraded) worker slots — the real throughput
+ // ceiling. 0 ⇒ paused until a worker is re-enabled (runPool idle-waits).
+ limit: () =>
+ pool.summary().filter((s) => s.state === "enabled" && !s.degraded).length,
+ signal: signal ?? NEVER,
+ drainSignal: drainSignal ?? NEVER,
+ finite: true,
+ idlePollMs: 3000,
+ });
return { attempted, succeeded: succeededCount, failed, skipped };
}
diff --git a/common/jobs/concurrentRunner.ts b/common/jobs/concurrentRunner.ts
@@ -89,16 +89,24 @@ export async function runPool<T>(opts: RunPoolOptions<T>): Promise<void> {
// 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 && (drained || signal.aborted)) break;
+ 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 && (drained || signal.aborted || opts.finite)) break;
+ if (
+ inFlight.size === 0 &&
+ (drainSignal.aborted || signal.aborted || opts.finite)
+ ) {
+ break;
+ }
await waitNext(idle);
continue;
}
@@ -107,13 +115,16 @@ export async function runPool<T>(opts: RunPoolOptions<T>): Promise<void> {
// 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 (success or failure) and
- // wake the loop so it can pull more / observe drain completion. `p` is a
- // const assigned before this runs, so the closure reference is safe.
- void p.finally(() => {
+ // 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