commit 065042ab3d891035385b7d1c579829ef103bb219
parent 33cc197977ccaeb0323512e9228f90ae16e41055
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 26 Jun 2026 15:30:06 -0400
Phase 5: unify the worker-pool waiter ordering with compareTier
Delete the second hand-rolled copy of the foreground-before-background priority
logic. The pool's Waiter now carries a SchedulerTier instead of a background
boolean, and enqueueWaiter inserts via the shared compareTier comparator — the
exact ordering the registry's scheduler uses. One priority model now governs
both job queueing and worker-slot waiting; an urgent tier (used by Phase 8
retry) jumps ahead of both for free.
acquire(signal, { tier }) is the new form; the legacy { background: true } is
still accepted and maps to the "background" tier, so transcribeOne/whisperBatch
callers are untouched. The pool keeps all its own semantics (enable/disable/
drain/degrade/pause-all) — the scheduler only shares the comparator, it does
not subsume the pool.
workerPool.test.ts: existing 5 pass unchanged (back-compat); +2 new — tier
ordering (urgent > foreground > background) and the legacy-flag mapping.
Typecheck clean; auto-queue e2e 11/11 (both drain regressions green).
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Diffstat:
2 files changed, 63 insertions(+), 15 deletions(-)
diff --git a/common/jobs/workerPool.test.ts b/common/jobs/workerPool.test.ts
@@ -118,3 +118,45 @@ test("all-background acquires fall back to plain FIFO (backward compatible)", as
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"]);
+});
+
+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"]);
+});
diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts
@@ -28,6 +28,8 @@ import type { Worker } from "../lib/workers";
import { getSettings } from "../lib/settings";
import { getPaths } from "../lib/paths";
import { readWorkerDefaults } from "./workerDefaults";
+import { compareTier } from "./scheduler";
+import type { SchedulerTier } from "./jobKinds";
// Consecutive transport/exec failures before a worker is auto-marked degraded
// (skipped by the scheduler until re-enabled). See markFailure / Phase 5.
@@ -73,9 +75,10 @@ type Waiter = {
resolve: (lease: Lease) => void;
reject: (err: Error) => void;
onAbort?: () => void;
- // Background waiters queue behind foreground ones (see enqueueWaiter). Auto
- // (background) transcriptions yield to manual (foreground) ones for free slots.
- background?: boolean;
+ // Priority tier (see enqueueWaiter). Background waiters (auto transcriptions)
+ // queue behind foreground (manual) ones for free slots; an urgent waiter would
+ // jump ahead of both. Same compareTier ordering the registry's scheduler uses.
+ tier: SchedulerTier;
};
export class WorkerPool {
@@ -277,19 +280,15 @@ export class WorkerPool {
}
}
- // Insert a parked waiter respecting priority: a foreground waiter slots in
- // BEFORE the first background waiter (so manual work jumps ahead of queued auto
- // work), preserving FIFO within each class; a background waiter always appends.
- // Mirrors the registry queue's enqueue (registry.ts) so one priority model
- // governs both schedulers.
+ // Insert a parked waiter respecting priority: a waiter slots in BEFORE the
+ // first waiter of strictly lower priority (so manual/foreground work jumps
+ // ahead of queued auto/background work), preserving FIFO within a tier. Uses
+ // the SAME compareTier comparator as the registry's scheduler — one priority
+ // model governs both, instead of two hand-rolled copies.
private enqueueWaiter(waiter: Waiter): void {
- if (waiter.background) {
- this.waiters.push(waiter);
- return;
- }
let insertAt = this.waiters.length;
for (let i = 0; i < this.waiters.length; i++) {
- if (this.waiters[i].background) {
+ if (compareTier(waiter.tier, this.waiters[i].tier) < 0) {
insertAt = i;
break;
}
@@ -303,7 +302,12 @@ export class WorkerPool {
// foreground (manual) waiter — see enqueueWaiter. If `signal` aborts while
// parked (or before), rejects with an AbortError so the caller treats it like a
// cancel.
- acquire(signal?: AbortSignal, opts?: { background?: boolean }): Promise<Lease> {
+ // opts.tier sets the priority tier; the legacy opts.background is still
+ // accepted and maps to the "background" tier (default "foreground").
+ acquire(
+ signal?: AbortSignal,
+ opts?: { tier?: SchedulerTier; background?: boolean },
+ ): Promise<Lease> {
this.ensureInit();
if (signal?.aborted) {
return Promise.reject(abortError());
@@ -312,8 +316,10 @@ export class WorkerPool {
if (entry && this.waiters.length === 0) {
return Promise.resolve(this.grant(entry));
}
+ const tier: SchedulerTier =
+ opts?.tier ?? (opts?.background ? "background" : "foreground");
return new Promise<Lease>((resolve, reject) => {
- const waiter: Waiter = { resolve, reject, background: opts?.background };
+ const waiter: Waiter = { resolve, reject, tier };
if (signal) {
const onAbort = () => {
const idx = this.waiters.indexOf(waiter);