commit 33cc197977ccaeb0323512e9228f90ae16e41055
parent 8fed89768a6e1a58fd90a4b7418bdfd778f9885f
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 26 Jun 2026 15:26:59 -0400
Phase 4: delegate registry queue ordering to the scheduler
The registry now owns only job lifecycle/state; per-queueKey ordering and the
foreground-before-background priority are delegated to the shared Scheduler
(every job submitted with concurrency 1, preserving strict per-queue serial).
JobRecord, the meta sidecars, the /api/jobs/active payload, and cancel/drain/
task methods are unchanged. enqueue->submit, finalize->complete, queued
cancel->scheduler.cancel, listQueues/activeQueueNames/positionInQueue read the
scheduler's id views back to JobRecords.
- registry.ts: drop the queues Map + start/cancel WeakMaps; delegate to
getScheduler(). background -> tier mapping (background|foreground).
- api/test/invalidate-cache: also reset __yttScheduler__ in lockstep with the
registry, else a prior spec's still-"running" queue entry makes the next
spec's same-named queue see a phantom busy slot and queue forever.
Also fixes a PRE-EXISTING stale locator in jobs-batch-tasks-drain.spec.ts
(failing on main too): it filtered job rows by the raw kind "whisper-all", but
the table renders the label "Transcribe all". These 3 tests cover exactly the
queued-cancel/drain path this phase rewrote, so they now give real coverage.
Verification: registry.test.ts passes UNCHANGED (the insertion contract);
scheduler/runner/jobKinds unit tests green; e2e queues, jobs-active-order,
job-stream-cancel, jobs-batch-tasks-drain all green. Typecheck (common+editor)
clean.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Diffstat:
3 files changed, 71 insertions(+), 88 deletions(-)
diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts
@@ -1,5 +1,7 @@
import type { ChildProcess } from "node:child_process";
import type { JobSpec } from "./jobSpec";
+import { getScheduler } from "./scheduler";
+import { getJobKind, type SchedulerTier } from "./jobKinds";
export type JobStatus =
| "queued"
@@ -88,11 +90,14 @@ export type QueueSnapshot = {
type StartFn = () => void;
type CancelFn = () => void;
+// The registry owns job LIFECYCLE/STATE (the JobRecord, meta sidecars, tasks,
+// terminal transitions). Queue ORDERING — which job runs vs. waits, and the
+// foreground-before-background priority — is delegated to the shared Scheduler
+// (scheduler.ts), the single source of truth that the worker pool also shares.
+// Every non-"" job is submitted with concurrency 1 (strict per-queueKey serial),
+// preserving the long-standing one-running-job-per-queue behavior.
class JobRegistry {
private jobs = new Map<string, JobRecord>();
- private queues = new Map<string, JobRecord[]>();
- private starts = new WeakMap<JobRecord, StartFn>();
- private cancels = new WeakMap<JobRecord, CancelFn>();
register(record: JobRecord): void {
this.jobs.set(record.id, record);
@@ -118,65 +123,45 @@ class JobRegistry {
);
}
- // Add a job to its queue. If the queue is empty, start it immediately and
- // mark "running". Otherwise leave it "queued" and remember the start/cancel
- // callbacks for when it becomes head-of-queue or is cancelled.
+ // Submit a job for scheduling. The scheduler decides whether it runs now (its
+ // queueKey has a free slot) or waits, and invokes `start` when it becomes head
+ // of its queue — immediately, or later when an earlier job finalizes.
//
- // Priority: a foreground (default) job is inserted AHEAD of any queued
- // background jobs but never displaces the running head (q[0]) — so a manually
- // clicked Sync jumps ahead of the auto-download runner's queued units without
- // interrupting one that's already mid-download. Foreground jobs keep FIFO
- // order among themselves; background jobs always append. finalize() still just
- // promotes q[0], which this ordering keeps as the highest-priority job.
+ // Priority: a foreground (default) job runs ahead of any queued background
+ // jobs but never displaces a running one — so a manually clicked Sync jumps
+ // ahead of the auto-download runner's queued units without interrupting one
+ // that's already mid-download. This is the scheduler's compareTier ordering
+ // (background -> "background" tier, everything else -> "foreground").
//
- // queueKey === "" is special: the job bypasses queue serialization entirely
- // and runs immediately, parallel to any other work. Cancellation still
- // works (running → SIGTERM/abort) but there's no queue map entry to clean
- // up in finalize.
+ // queueKey === "" is special: the job bypasses serialization entirely and runs
+ // immediately, parallel to any other work, untracked by the scheduler.
+ // Cancellation still works (running -> SIGTERM/abort).
enqueue(
record: JobRecord,
callbacks: { start: StartFn; onCancel: CancelFn },
): { willRunNow: boolean; position: number } {
- this.starts.set(record, callbacks.start);
- this.cancels.set(record, callbacks.onCancel);
- if (record.queueKey === "") {
+ // Setting status/startedAt is coupled with actually launching, so wrap the
+ // caller's start: the scheduler runs this exactly when the job starts.
+ const start: StartFn = () => {
record.status = "running";
record.startedAt = Date.now();
callbacks.start();
- return { willRunNow: true, position: 0 };
- }
- let q = this.queues.get(record.queueKey);
- if (!q) {
- q = [];
- this.queues.set(record.queueKey, q);
- }
- if (q.length === 0) {
- q.push(record);
- record.status = "running";
- record.startedAt = Date.now();
- callbacks.start();
- return { willRunNow: true, position: 0 };
- }
- if (record.background) {
- // Background work always queues at the very back.
- q.push(record);
- } else {
- // Foreground work slots in before the first QUEUED background job (never
- // before the running head at index 0), preserving FIFO among foreground.
- let insertAt = q.length;
- for (let i = 1; i < q.length; i++) {
- if (q[i].background) {
- insertAt = i;
- break;
- }
- }
- q.splice(insertAt, 0, record);
- }
- return { willRunNow: false, position: q.indexOf(record) };
+ };
+ const tier: SchedulerTier = record.background
+ ? "background"
+ : getJobKind(record.kind)?.defaultTier ?? "foreground";
+ return getScheduler().submit({
+ id: record.id,
+ queueKey: record.queueKey,
+ tier,
+ concurrency: 1,
+ start,
+ onCancel: callbacks.onCancel,
+ });
}
- // Idempotent: marks the job terminal (if not already), splices it out of
- // its queue, and starts the next queued job in that queue.
+ // Idempotent: marks the job terminal (if not already), then tells the
+ // scheduler it is done so the next queued job on that queueKey is promoted.
finalize(
id: string,
status: "done" | "failed" | "cancelled",
@@ -192,29 +177,16 @@ class JobRegistry {
// Release per-task and drain references on terminal jobs.
job.tasks = [];
job.drainController = undefined;
- const q = this.queues.get(job.queueKey);
- if (!q) return;
- const idx = q.indexOf(job);
- if (idx >= 0) q.splice(idx, 1);
- if (q.length === 0) {
- this.queues.delete(job.queueKey);
- return;
- }
- const next = q[0];
- if (next.status === "queued") {
- next.status = "running";
- next.startedAt = Date.now();
- const startFn = this.starts.get(next);
- if (startFn) startFn();
- }
+ getScheduler().complete(id);
}
cancel(id: string): boolean {
const job = this.jobs.get(id);
if (!job) return false;
if (job.status === "queued") {
- const onCancel = this.cancels.get(job);
- if (onCancel) onCancel();
+ // Remove from the scheduler (fires the queued job's onCancel), then mark
+ // it terminal. finalize's scheduler.complete is a no-op by then.
+ getScheduler().cancel(id);
this.finalize(id, "cancelled");
return true;
}
@@ -307,31 +279,36 @@ class JobRegistry {
job.completedTaskMs = (job.completedTaskMs ?? 0) + durationMs;
}
- // Snapshot of queues for UI. Sorted by queue name.
+ // Snapshot of queues for UI, resolved from the scheduler's id views back to
+ // JobRecords. Already sorted by queue name (the scheduler sorts). Running/
+ // queued ids always resolve: the 100-job eviction never drops non-terminal
+ // jobs, but filter defensively anyway.
listQueues(): QueueSnapshot[] {
- const out: QueueSnapshot[] = [];
- for (const [name, q] of this.queues.entries()) {
- const running = q.find((j) => j.status === "running");
- const queued = q.filter((j) => j.status === "queued");
- out.push({ name, running, queued });
- }
- return out.sort((a, b) => a.name.localeCompare(b.name));
+ return getScheduler()
+ .queues()
+ .map((v) => ({
+ name: v.name,
+ running: v.running
+ .map((id) => this.jobs.get(id))
+ .find((j): j is JobRecord => Boolean(j)),
+ queued: v.queued
+ .map((id) => this.jobs.get(id))
+ .filter((j): j is JobRecord => Boolean(j)),
+ }));
}
// Names of queues that currently have any non-terminal jobs. Used by the
// UI to populate the QueuePicker dropdown.
activeQueueNames(): string[] {
- return Array.from(this.queues.keys()).sort();
+ return getScheduler()
+ .queues()
+ .map((v) => v.name);
}
- // Position of the job in its queue (0 = currently running). Returns -1 if
- // the job is no longer in any queue (terminal).
+ // Position of the job in its queue (0 = currently running). Returns -1 if the
+ // job is no longer in any queue (terminal or unserialized "").
positionInQueue(id: string): number {
- const job = this.jobs.get(id);
- if (!job) return -1;
- const q = this.queues.get(job.queueKey);
- if (!q) return -1;
- return q.indexOf(job);
+ return getScheduler().positionInQueue(id);
}
}
diff --git a/editor/app/api/test/invalidate-cache/route.ts b/editor/app/api/test/invalidate-cache/route.ts
@@ -30,6 +30,12 @@ function invalidate() {
// use. (In production the pool legitimately persists; this is e2e-only.)
// eslint-disable-next-line @typescript-eslint/no-explicit-any
(globalThis as any).__yttWorkerPool__ = undefined;
+ // Reset the shared scheduler: it owns the per-queueKey ordering for the
+ // (now-wiped) registry, so a prior spec's still-"running" queue entry would
+ // otherwise make the next spec's same-named queue see a phantom busy slot and
+ // queue forever. Cleared in lockstep with the registry it serves.
+ // eslint-disable-next-line @typescript-eslint/no-explicit-any
+ (globalThis as any).__yttScheduler__ = undefined;
// Reset the auto-queue runner singleton too: its live entries reference job
// ids in the (now-wiped) registry. A still-executing runner loop from a prior
// spec detects its job is gone (registry.get -> undefined) and exits on its
diff --git a/editor/e2e/jobs-batch-tasks-drain.spec.ts b/editor/e2e/jobs-batch-tasks-drain.spec.ts
@@ -204,7 +204,7 @@ test("'Drain all' drains the running batch and cancels the queued one", async ({
const bRow = page
.getByRole("row")
.filter({ hasText: "drainall-b" })
- .filter({ hasText: "whisper-all" })
+ .filter({ hasText: "Transcribe all" })
.first();
await expect(bRow).toContainText("cancelled", { timeout: 20_000 });
expect(await transcriptCount("drainall-b", bIds)).toBe(0);
@@ -212,7 +212,7 @@ test("'Drain all' drains the running batch and cancels the queued one", async ({
const aRow = page
.getByRole("row")
.filter({ hasText: "drainall-a" })
- .filter({ hasText: "whisper-all" })
+ .filter({ hasText: "Transcribe all" })
.first();
await expect(aRow).toContainText("done", { timeout: 20_000 });
const aDone = await transcriptCount("drainall-a", aIds);
@@ -252,7 +252,7 @@ test("a queued job can be cancelled directly from its row without opening the lo
const bRow = page
.getByRole("row")
.filter({ hasText: "qcancel-b" })
- .filter({ hasText: "whisper-all" })
+ .filter({ hasText: "Transcribe all" })
.first();
await expect(bRow).toContainText("cancelled", { timeout: 20_000 });
expect(await transcriptCount("qcancel-b", bIds)).toBe(0);
@@ -294,7 +294,7 @@ test("hard Cancel during a drain ends the job cancelled without stream errors",
const row = page
.getByRole("row")
.filter({ hasText: "cancel-drain" })
- .filter({ hasText: "whisper-all" })
+ .filter({ hasText: "Transcribe all" })
.first();
await expect(row).toContainText("cancelled", { timeout: 20_000 });