commit 8fed89768a6e1a58fd90a4b7418bdfd778f9885f
parent 956eb67bf450a99ec72f1fe6a868ae990a8e85f5
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 26 Jun 2026 15:02:12 -0400
Phase 3: scheduler core + runPool primitive (behind the API)
Land the two new core primitives as dead code (nothing imports them yet) so
the Phase 4-6 cutovers are pure swaps:
- common/jobs/scheduler.ts: the unified queue-ordering core. One compareTier
comparator (urgent < foreground < background) replaces the priority logic
duplicated in registry.enqueue and workerPool. Per-queue concurrency, plus
reorder/promote for granular control. queueKey "" stays parallel/untracked.
- common/jobs/concurrentRunner.ts: runPool(), one event-driven fill-to-capacity
primitive to replace both the auto-runner's hand-written waitNext/wake loop
and whisperBatch's Promise.all. Encodes the no-spin invariant ONCE: an
already-aborted signal never fast-paths to a resolved promise, so the wait
paces on a real timer and a running unit's child-exit macrotask is never
starved (the "Drain all hangs" bug class).
Tests: scheduler.test.ts (17) — tiered ordering, FIFO-in-tier, concurrency>1,
reorder/promote only touch queued, queueKey "" parallelism. concurrentRunner
.test.ts (5) — respects limit(), finite termination, drain-awaits-in-flight,
the already-aborted-signal no-spin case, hard-cancel aborts in-flight. All
green; typecheck (common+editor) clean.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Diffstat:
4 files changed, 693 insertions(+), 0 deletions(-)
diff --git a/common/jobs/concurrentRunner.test.ts b/common/jobs/concurrentRunner.test.ts
@@ -0,0 +1,175 @@
+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<void>((r) => setTimeout(r, ms));
+
+function withTimeout<T>(p: Promise<T>, ms: number, msg: string): Promise<T> {
+ return Promise.race([
+ p,
+ sleep(ms).then(() => {
+ throw new Error(`timeout: ${msg}`);
+ }),
+ ]) as Promise<T>;
+}
+
+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<number>({
+ 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<number>({
+ 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<number>({
+ // 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<number>({
+ 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<number>({
+ next: async () => 0,
+ run: async (_item, signal) => {
+ await new Promise<void>((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");
+});
diff --git a/common/jobs/concurrentRunner.ts b/common/jobs/concurrentRunner.ts
@@ -0,0 +1,123 @@
+// 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<T> = {
+ // 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<T | null | undefined>;
+ // 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<void>;
+ // 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<T>(opts: RunPoolOptions<T>): Promise<void> {
+ const { signal, drainSignal } = opts;
+ const idle = opts.idlePollMs ?? 3000;
+ const inFlight = new Set<Promise<void>>();
+
+ // 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<void> =>
+ new Promise<void>((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 && (drained || signal.aborted)) break;
+ await waitNext(idle);
+ continue;
+ }
+
+ const item = await opts.next();
+ 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;
+ 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 (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(() => {
+ inFlight.delete(p);
+ wake();
+ });
+ }
+
+ // 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]);
+}
diff --git a/common/jobs/scheduler.test.ts b/common/jobs/scheduler.test.ts
@@ -0,0 +1,182 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { createScheduler, compareTier, type ScheduledTask } from "./scheduler";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/scheduler.test.ts
+//
+// The scheduler is the single source of truth for queue ordering, generalizing
+// registry.enqueue's foreground-before-background insertion to priority tiers.
+// Each test uses a fresh createScheduler() so there is no shared global state.
+
+let seq = 0;
+function task(
+ queueKey: string,
+ tier: ScheduledTask["tier"],
+ concurrency = 1,
+): { task: ScheduledTask; started: () => boolean; cancelled: () => boolean } {
+ let started = false;
+ let cancelled = false;
+ const t: ScheduledTask = {
+ id: `t${seq++}`,
+ queueKey,
+ tier,
+ concurrency,
+ start: () => {
+ started = true;
+ },
+ onCancel: () => {
+ cancelled = true;
+ },
+ };
+ return { task: t, started: () => started, cancelled: () => cancelled };
+}
+
+test("compareTier orders urgent < foreground < background", () => {
+ assert.ok(compareTier("urgent", "foreground") < 0);
+ assert.ok(compareTier("foreground", "background") < 0);
+ assert.ok(compareTier("background", "urgent") > 0);
+ assert.equal(compareTier("foreground", "foreground"), 0);
+});
+
+test("first task on a queue runs immediately", () => {
+ const s = createScheduler();
+ const a = task("k", "foreground");
+ const r = s.submit(a.task);
+ assert.deepEqual(r, { willRunNow: true, position: 0 });
+ assert.ok(a.started());
+ assert.deepEqual(s.queues(), [{ name: "k", running: [a.task.id], queued: [] }]);
+});
+
+test("tiers order the queued section: urgent before foreground before background", () => {
+ const s = createScheduler();
+ const head = task("k", "foreground"); // runs
+ const bg = task("k", "background");
+ const fg = task("k", "foreground");
+ const urg = task("k", "urgent");
+ s.submit(head.task);
+ s.submit(bg.task);
+ s.submit(fg.task);
+ s.submit(urg.task);
+ const view = s.queues()[0];
+ assert.deepEqual(view.running, [head.task.id]);
+ assert.deepEqual(view.queued, [urg.task.id, fg.task.id, bg.task.id]);
+ // Only the running head actually started.
+ assert.ok(head.started());
+ assert.ok(!fg.started() && !bg.started() && !urg.started());
+});
+
+test("FIFO within a tier", () => {
+ const s = createScheduler();
+ const head = task("k", "foreground");
+ const a = task("k", "foreground");
+ const b = task("k", "foreground");
+ s.submit(head.task);
+ s.submit(a.task);
+ s.submit(b.task);
+ assert.deepEqual(s.queues()[0].queued, [a.task.id, b.task.id]);
+});
+
+test("concurrency > 1 runs that many at once, queues the rest", () => {
+ const s = createScheduler();
+ const ts = Array.from({ length: 5 }, () => task("k", "foreground", 3));
+ ts.forEach((t) => s.submit(t.task));
+ const view = s.queues()[0];
+ assert.equal(view.running.length, 3);
+ assert.equal(view.queued.length, 2);
+ assert.ok(ts.slice(0, 3).every((t) => t.started()));
+ assert.ok(ts.slice(3).every((t) => !t.started()));
+});
+
+test("complete promotes the next queued task and fires its start", () => {
+ const s = createScheduler();
+ const head = task("k", "foreground");
+ const next = task("k", "foreground");
+ s.submit(head.task);
+ s.submit(next.task);
+ assert.ok(!next.started());
+ s.complete(head.task.id);
+ assert.ok(next.started());
+ assert.deepEqual(s.queues(), [{ name: "k", running: [next.task.id], queued: [] }]);
+});
+
+test("completing the last task removes the queue", () => {
+ const s = createScheduler();
+ const a = task("k", "foreground");
+ s.submit(a.task);
+ s.complete(a.task.id);
+ assert.deepEqual(s.queues(), []);
+ assert.equal(s.positionInQueue(a.task.id), -1);
+});
+
+test("cancel removes a queued task and fires onCancel; running/unknown return false", () => {
+ const s = createScheduler();
+ const head = task("k", "foreground");
+ const q = task("k", "foreground");
+ s.submit(head.task);
+ s.submit(q.task);
+ assert.equal(s.cancel(q.task.id), true);
+ assert.ok(q.cancelled());
+ assert.deepEqual(s.queues()[0].queued, []);
+ // Running head can't be cancelled via the scheduler (caller hard-aborts).
+ assert.equal(s.cancel(head.task.id), false);
+ assert.equal(s.cancel("nope"), false);
+});
+
+test("promote moves a queued task to the front of the queued section", () => {
+ const s = createScheduler();
+ const head = task("k", "foreground");
+ const a = task("k", "background");
+ const b = task("k", "background");
+ s.submit(head.task);
+ s.submit(a.task);
+ s.submit(b.task);
+ assert.deepEqual(s.queues()[0].queued, [a.task.id, b.task.id]);
+ assert.equal(s.promote(b.task.id), true);
+ assert.deepEqual(s.queues()[0].queued, [b.task.id, a.task.id]);
+ // Running task can't be promoted.
+ assert.equal(s.promote(head.task.id), false);
+});
+
+test("reorder swaps adjacent queued tasks but never moves the running head", () => {
+ const s = createScheduler();
+ const head = task("k", "foreground");
+ const a = task("k", "foreground");
+ const b = task("k", "foreground");
+ s.submit(head.task);
+ s.submit(a.task);
+ s.submit(b.task);
+ // Move b up past a.
+ assert.equal(s.reorder(b.task.id, -1), true);
+ assert.deepEqual(s.queues()[0].queued, [b.task.id, a.task.id]);
+ // b is now at the front of the queued section: can't move up into running.
+ assert.equal(s.reorder(b.task.id, -1), false);
+ // a is last: can't move down past the end.
+ assert.equal(s.reorder(a.task.id, 1), false);
+ // Running head can't be reordered.
+ assert.equal(s.reorder(head.task.id, 1), false);
+});
+
+test('queueKey "" runs immediately in parallel and is not tracked', () => {
+ const s = createScheduler();
+ const a = task("", "foreground");
+ const b = task("", "foreground");
+ assert.deepEqual(s.submit(a.task), { willRunNow: true, position: 0 });
+ assert.deepEqual(s.submit(b.task), { willRunNow: true, position: 0 });
+ assert.ok(a.started() && b.started());
+ assert.deepEqual(s.queues(), []);
+ assert.equal(s.positionInQueue(a.task.id), -1);
+ s.complete(a.task.id); // no-op, no throw
+});
+
+test("positionInQueue reports 0 for running head, then queued offsets", () => {
+ const s = createScheduler();
+ const head = task("k", "foreground");
+ const a = task("k", "foreground");
+ const b = task("k", "foreground");
+ s.submit(head.task);
+ s.submit(a.task);
+ s.submit(b.task);
+ assert.equal(s.positionInQueue(head.task.id), 0);
+ assert.equal(s.positionInQueue(a.task.id), 1);
+ assert.equal(s.positionInQueue(b.task.id), 2);
+});
diff --git a/common/jobs/scheduler.ts b/common/jobs/scheduler.ts
@@ -0,0 +1,213 @@
+import type { SchedulerTier } from "./jobKinds";
+
+// The unified queue-ordering core. Before this, "foreground preempts background"
+// was implemented twice — once in registry.enqueue (job ordering) and again in
+// workerPool.acquire (slot waiters). Both are now expressible through ONE
+// comparator (compareTier) and one ordering policy here, generalized from a
+// background boolean to priority tiers.
+//
+// A scheduler owns N independent queues keyed by queueKey. Within a queue, up to
+// `concurrency` tasks run at once (the "running" prefix); the rest wait, ordered
+// by tier (urgent < foreground < background) with FIFO inside a tier. queueKey
+// === "" is special: such tasks bypass serialization and run immediately in
+// parallel (preserving the registry's long-standing behavior), and are not
+// tracked in any queue.
+
+// Lower number = higher priority = runs sooner. Negative when `a` outranks `b`.
+const TIER_RANK: Record<SchedulerTier, number> = {
+ urgent: 0,
+ foreground: 1,
+ background: 2,
+};
+
+export function compareTier(a: SchedulerTier, b: SchedulerTier): number {
+ return TIER_RANK[a] - TIER_RANK[b];
+}
+
+export type ScheduledTask = {
+ id: string;
+ queueKey: string;
+ tier: SchedulerTier;
+ // Max simultaneously-running tasks on this queueKey (1 = strict serial, today's
+ // behavior for every non-"" queue). The first submitter's value wins for a
+ // given live queue.
+ concurrency: number;
+ // Becomes-head-of-queue callback: invoked exactly once when the task starts
+ // running (immediately on submit if there's room, else later via complete()).
+ start: () => void;
+ // Invoked if the task is cancelled while still queued (never started).
+ onCancel: () => void;
+};
+
+export type QueueView = { name: string; running: string[]; queued: string[] };
+
+export interface Scheduler {
+ submit(t: ScheduledTask): { willRunNow: boolean; position: number };
+ // Mark a task done (normal completion OR post-cancel cleanup) and promote the
+ // next eligible queued task(s). Idempotent; no-op if the id is unknown.
+ complete(id: string): void;
+ // Remove a QUEUED task, firing its onCancel. Returns false for a running or
+ // unknown task (the caller hard-aborts a running one, then calls complete()).
+ cancel(id: string): boolean;
+ // Granular control: move a QUEUED task one slot toward (-1) or away from (+1)
+ // the head, within the queued section only. Never moves a running task.
+ reorder(id: string, dir: -1 | 1): boolean;
+ // Move a QUEUED task to the front of the queued section. Never preempts a
+ // running task.
+ promote(id: string): boolean;
+ // 0 = running head, 1.. = queued position. -1 if gone or untracked ("").
+ positionInQueue(id: string): number;
+ queues(): QueueView[];
+}
+
+type Entry = ScheduledTask & { status: "running" | "queued" };
+
+// Invariant for every queue array: all "running" entries form a contiguous
+// prefix, followed by all "queued" entries. submit/complete/promote/reorder all
+// preserve this.
+class InMemoryScheduler implements Scheduler {
+ private queuesByKey = new Map<string, Entry[]>();
+ private byId = new Map<string, Entry>();
+
+ submit(t: ScheduledTask): { willRunNow: boolean; position: number } {
+ // Unserialized: run immediately, never tracked.
+ if (t.queueKey === "") {
+ t.start();
+ return { willRunNow: true, position: 0 };
+ }
+ let q = this.queuesByKey.get(t.queueKey);
+ if (!q) {
+ q = [];
+ this.queuesByKey.set(t.queueKey, q);
+ }
+ const runningCount = q.filter((e) => e.status === "running").length;
+ if (runningCount < t.concurrency) {
+ // Room to run now: append to the end of the running prefix.
+ const entry: Entry = { ...t, status: "running" };
+ q.splice(runningCount, 0, entry);
+ this.byId.set(t.id, entry);
+ entry.start();
+ return { willRunNow: true, position: runningCount };
+ }
+ // At capacity: queue by tier, inserting before the first queued entry of a
+ // strictly lower priority (preserving FIFO within a tier).
+ const entry: Entry = { ...t, status: "queued" };
+ let insertAt = q.length;
+ for (let i = runningCount; i < q.length; i++) {
+ if (compareTier(t.tier, q[i].tier) < 0) {
+ insertAt = i;
+ break;
+ }
+ }
+ q.splice(insertAt, 0, entry);
+ this.byId.set(t.id, entry);
+ return { willRunNow: false, position: insertAt };
+ }
+
+ complete(id: string): void {
+ const entry = this.byId.get(id);
+ if (!entry) return;
+ this.byId.delete(id);
+ const q = this.queuesByKey.get(entry.queueKey);
+ if (!q) return;
+ const idx = q.indexOf(entry);
+ if (idx >= 0) q.splice(idx, 1);
+ this.fill(entry.queueKey, q);
+ }
+
+ cancel(id: string): boolean {
+ const entry = this.byId.get(id);
+ if (!entry || entry.status !== "queued") return false;
+ this.byId.delete(id);
+ const q = this.queuesByKey.get(entry.queueKey);
+ if (!q) return true;
+ const idx = q.indexOf(entry);
+ if (idx >= 0) q.splice(idx, 1);
+ // Removing a queued entry frees no slot, so no promotion is needed.
+ if (q.length === 0) this.queuesByKey.delete(entry.queueKey);
+ entry.onCancel();
+ return true;
+ }
+
+ reorder(id: string, dir: -1 | 1): boolean {
+ const entry = this.byId.get(id);
+ if (!entry || entry.status !== "queued") return false;
+ const q = this.queuesByKey.get(entry.queueKey);
+ if (!q) return false;
+ const i = q.indexOf(entry);
+ const j = i + dir;
+ // Stay within the queued section: never swap with a running entry or past
+ // the array bounds.
+ if (j < 0 || j >= q.length || q[j].status !== "queued") return false;
+ [q[i], q[j]] = [q[j], q[i]];
+ return true;
+ }
+
+ promote(id: string): boolean {
+ const entry = this.byId.get(id);
+ if (!entry || entry.status !== "queued") return false;
+ const q = this.queuesByKey.get(entry.queueKey);
+ if (!q) return false;
+ const i = q.indexOf(entry);
+ const runningCount = q.filter((e) => e.status === "running").length;
+ if (i <= runningCount) return false; // already at the front of the queue
+ q.splice(i, 1);
+ q.splice(runningCount, 0, entry);
+ return true;
+ }
+
+ positionInQueue(id: string): number {
+ const entry = this.byId.get(id);
+ if (!entry) return -1;
+ const q = this.queuesByKey.get(entry.queueKey);
+ if (!q) return -1;
+ return q.indexOf(entry);
+ }
+
+ queues(): QueueView[] {
+ const out: QueueView[] = [];
+ for (const [name, q] of this.queuesByKey.entries()) {
+ out.push({
+ name,
+ running: q.filter((e) => e.status === "running").map((e) => e.id),
+ queued: q.filter((e) => e.status === "queued").map((e) => e.id),
+ });
+ }
+ return out.sort((a, b) => a.name.localeCompare(b.name));
+ }
+
+ // Promote queued entries to running until the concurrency ceiling is reached.
+ // The queued entries are already in priority order, and sit immediately after
+ // the running prefix, so promoting the boundary entry preserves the invariant.
+ private fill(queueKey: string, q: Entry[]): void {
+ if (q.length === 0) {
+ this.queuesByKey.delete(queueKey);
+ return;
+ }
+ let runningCount = q.filter((e) => e.status === "running").length;
+ const cap = q[0]?.concurrency ?? 1;
+ while (runningCount < cap && runningCount < q.length) {
+ const next = q[runningCount];
+ if (next.status !== "queued") break;
+ next.status = "running";
+ runningCount++;
+ next.start();
+ }
+ }
+}
+
+declare global {
+ // eslint-disable-next-line no-var
+ var __yttScheduler__: Scheduler | undefined;
+}
+
+export function createScheduler(): Scheduler {
+ return new InMemoryScheduler();
+}
+
+export function getScheduler(): Scheduler {
+ if (!globalThis.__yttScheduler__) {
+ globalThis.__yttScheduler__ = new InMemoryScheduler();
+ }
+ return globalThis.__yttScheduler__;
+}