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 = { 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(); private byId = new Map(); 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__; }