commit bafc953c31e8555d3182c764418f8d4a2ae5577a
parent a3c5706ddb9a2fe46c5ff85663f71585b7d8114e
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Mon, 28 Sep 2026 12:46:59 -0400
jobs: a job cancelled while queued writes its cancelled sidecar
The queued-cancel path never rewrote <id>.meta.json: onCancel only closed
the stream and settled `done`, and registry.cancel() marked the record
terminal only after the scheduler had fired it. So the sidecar kept the
"queued" its enqueue wrote — /jobs read an evicted one back as queued, and
the boot pass could re-queue a job the operator had cancelled.
- registry.cancel(): a queued record is marked cancelled BEFORE
scheduler.cancel() fires its onCancel (markTerminal, shared with
finalize, which follows and stays idempotent).
- streamCommand's two onCancels persist the sidecar when the record is
terminal.
- Graceful shutdown is the exception: the reaper's cancels are nobody's
decision, and a queued sidecar must stay `queued` for bootQueuedJobs.
shutdownCancel calls registry.beginShutdown() first; a queued job then
reaches onCancel still queued, and nothing is written.
- A job's meta writes are chained (metaWriter), so the enqueue's "queued"
write can never land after the terminal one.
Tests: registry.test.ts +1 (onCancel sees the record cancelled);
streamCommand.test.ts (new, 3): function and command jobs cancelled while
queued end `cancelled` on disk; at shutdown the sidecar stays `queued`.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Diffstat:
5 files changed, 280 insertions(+), 13 deletions(-)
diff --git a/common/jobs/registry.test.ts b/common/jobs/registry.test.ts
@@ -126,3 +126,24 @@ test("cancelling a queued background job fires onCancel and removes it", () => {
"removed from its queue",
);
});
+
+// Release 13 slice W1: a managed job persists its sidecar in onCancel, so a
+// queued job must already be terminal when the scheduler fires it. It used to
+// be marked after, and the sidecar kept "queued".
+test("a queued job is already cancelled when its onCancel runs", () => {
+ const q = `test:cancel-order:${newJobId()}`;
+ enqueue(rec(q)); // the head, running
+ const queued = rec(q);
+ let seen: { status: string; endedAt?: number } | null = null;
+ getRegistry().register(queued);
+ getRegistry().enqueue(queued, {
+ start: () => {},
+ onCancel: () => {
+ seen = { status: queued.status, endedAt: queued.endedAt };
+ },
+ });
+ assert.equal(getRegistry().cancel(queued.id), true);
+ assert.equal(seen!.status, "cancelled");
+ assert.equal(typeof seen!.endedAt, "number");
+ assert.equal(getRegistry().positionInQueue(queued.id), -1);
+});
diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts
@@ -166,6 +166,20 @@ export type QueueSnapshot = {
type StartFn = () => void;
type CancelFn = () => void;
+// The terminal transition itself: status, endedAt, exitCode — once. A record
+// already terminal keeps the state it ended in (a cancel racing a child's exit
+// stays "cancelled").
+function markTerminal(
+ job: JobRecord,
+ status: "done" | "failed" | "cancelled",
+ exitCode?: number,
+): void {
+ if (job.status !== "queued" && job.status !== "running") return;
+ job.status = status;
+ job.endedAt = Date.now();
+ if (typeof exitCode === "number") job.exitCode = exitCode;
+}
+
// 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
@@ -174,6 +188,14 @@ type CancelFn = () => void;
// preserving the long-standing one-running-job-per-queue behavior.
class JobRegistry {
private jobs = new Map<string, JobRecord>();
+ // Set by the graceful-shutdown reaper (shutdownCancel.ts) before it cancels
+ // every live job; see cancel().
+ private shuttingDown = false;
+
+ // The server is going down: the cancels that follow are nobody's decision.
+ beginShutdown(): void {
+ this.shuttingDown = true;
+ }
register(record: JobRecord): void {
this.jobs.set(record.id, record);
@@ -245,11 +267,7 @@ class JobRegistry {
): void {
const job = this.jobs.get(id);
if (!job) return;
- if (job.status === "queued" || job.status === "running") {
- job.status = status;
- job.endedAt = Date.now();
- if (typeof exitCode === "number") job.exitCode = exitCode;
- }
+ markTerminal(job, status, exitCode);
// Release per-task and drain references on terminal jobs.
job.tasks = [];
job.drainController = undefined;
@@ -260,8 +278,22 @@ class JobRegistry {
const job = this.jobs.get(id);
if (!job) return false;
if (job.status === "queued") {
- // Remove from the scheduler (fires the queued job's onCancel), then mark
- // it terminal. finalize's scheduler.complete is a no-op by then.
+ // Terminal FIRST, then out of the scheduler — which fires the queued
+ // job's onCancel, and a managed job persists its sidecar there
+ // (streamCommand.ts), so the record it writes must already say
+ // cancelled. It used to be marked after, and every job cancelled while
+ // queued kept "queued" in its .meta.json (release 13 slice W1). Not
+ // finalize() first: its scheduler.complete would take the entry out of
+ // the queue WITHOUT firing onCancel, and the job's stream and `done`
+ // would never settle. finalize() after is idempotent — it releases the
+ // task references, and scheduler.complete is a no-op by then.
+ //
+ // EXCEPT AT SHUTDOWN. The reaper cancels a queued job only so the exit
+ // cannot promote it into a child; nobody cancelled it. Its sidecar must
+ // stay `queued`, which is what the boot pass (bootQueuedJobs.ts) settles
+ // or re-queues on the next start — so it reaches onCancel still queued,
+ // and onCancel persists only a terminal record.
+ if (!this.shuttingDown) markTerminal(job, "cancelled");
getScheduler().cancel(id);
this.finalize(id, "cancelled");
return true;
diff --git a/common/jobs/shutdownCancel.ts b/common/jobs/shutdownCancel.ts
@@ -37,6 +37,12 @@ export function armShutdownCancel(): void {
// building a registry during shutdown just to cancel nothing would be
// worse than useless. No registry means nothing ever ran.
const registry = globalThis.__yttJobRegistry__;
+ // First: these cancels are the exit's, not an operator's, so a queued
+ // job keeps its `queued` sidecar for the boot pass (registry.cancel).
+ // Optional-called: under `next dev` the registry on globalThis can
+ // predate this method (HMR swaps the class, not the instance), and a
+ // throw here would skip every cancel below.
+ registry?.beginShutdown?.();
for (const job of registry?.list() ?? []) {
if (job.status !== "running" && job.status !== "queued") continue;
try {
diff --git a/common/jobs/streamCommand.test.ts b/common/jobs/streamCommand.test.ts
@@ -0,0 +1,177 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { mkdtemp, readFile, rm } from "node:fs/promises";
+import { tmpdir } from "node:os";
+import path from "node:path";
+import type { Paths } from "../lib/paths";
+import type { JobMeta } from "./jobMeta";
+import { getRegistry, newJobId } from "./registry";
+import {
+ runManagedCommand,
+ runManagedFunction,
+ type StreamActionResult,
+} from "./streamCommand";
+
+// Run with:
+// pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/streamCommand.test.ts
+//
+// A JOB CANCELLED WHILE STILL QUEUED (release 13 slice W1). Its start() never
+// runs, so the terminal sidecar write in start()'s .finally never happens;
+// onCancel is the only place that can write it, and it used to write nothing —
+// `<id>.meta.json` kept the "queued" the enqueue wrote, and a boot pass could
+// re-queue a job the operator had cancelled. Each test holds its own queue key
+// with a job that runs until released, so the job under test only ever queues
+// and nothing is spawned. Temp .jobs dirs only.
+
+type Started = Extract<StreamActionResult, { ok: true }>;
+
+function ok(r: StreamActionResult): Started {
+ if (!r.ok) throw new Error(r.error);
+ return r;
+}
+
+async function jobsDir(): Promise<{ paths: Paths; root: string }> {
+ const root = await mkdtemp(path.join(tmpdir(), "stream-command-"));
+ return { paths: { jobsDir: path.join(root, ".jobs") } as Paths, root };
+}
+
+async function readMeta(paths: Paths, id: string): Promise<JobMeta | null> {
+ try {
+ return JSON.parse(
+ await readFile(path.join(paths.jobsDir, `${id}.meta.json`), "utf8"),
+ ) as JobMeta;
+ } catch {
+ return null;
+ }
+}
+
+// The sidecar is written asynchronously and nothing hands out its promise, so
+// wait (briefly) for it to reach `status`; report the last one seen otherwise.
+async function metaReaches(
+ paths: Paths,
+ id: string,
+ status: string,
+): Promise<JobMeta | null> {
+ let last: JobMeta | null = null;
+ for (let i = 0; i < 100; i++) {
+ last = await readMeta(paths, id);
+ if (last?.status === status) return last;
+ await new Promise((r) => setTimeout(r, 20));
+ }
+ return last;
+}
+
+// A job on `queueKey` that runs until released: everything submitted behind it
+// queues.
+async function hold(paths: Paths, queueKey: string) {
+ let release!: () => void;
+ const held = new Promise<void>((r) => {
+ release = r;
+ });
+ const job = ok(
+ await runManagedFunction({
+ kind: "test-holder",
+ queueKey,
+ paths,
+ fn: () => held,
+ }),
+ );
+ return {
+ async release() {
+ release();
+ await job.done;
+ },
+ };
+}
+
+test("a function job cancelled while queued ends with a cancelled sidecar", async () => {
+ const { paths, root } = await jobsDir();
+ const queueKey = `test:queued-cancel:${newJobId()}`;
+ const holder = await hold(paths, queueKey);
+ try {
+ const job = ok(
+ await runManagedFunction({
+ kind: "test-queued",
+ queueKey,
+ paths,
+ fn: async () => {
+ throw new Error("must never start");
+ },
+ }),
+ );
+ assert.equal(getRegistry().get(job.jobId)?.status, "queued");
+ assert.equal((await metaReaches(paths, job.jobId, "queued"))?.status, "queued");
+
+ assert.equal(getRegistry().cancel(job.jobId), true);
+ assert.deepEqual(await job.done, { status: "cancelled", jobId: job.jobId });
+ const meta = await metaReaches(paths, job.jobId, "cancelled");
+ assert.equal(meta?.status, "cancelled");
+ assert.equal(typeof meta?.endedAt, "number");
+ assert.equal(meta?.startedAt, undefined, "it never started");
+ } finally {
+ await holder.release();
+ await rm(root, { recursive: true, force: true });
+ }
+});
+
+test("a command job cancelled while queued ends with a cancelled sidecar", async () => {
+ const { paths, root } = await jobsDir();
+ const queueKey = `test:queued-cancel-cmd:${newJobId()}`;
+ const holder = await hold(paths, queueKey);
+ try {
+ const job = ok(
+ await runManagedCommand({
+ kind: "test-queued-command",
+ queueKey,
+ paths,
+ cwd: root,
+ // Never run: it only queues behind the holder.
+ command: "false",
+ args: [],
+ }),
+ );
+ assert.equal((await metaReaches(paths, job.jobId, "queued"))?.status, "queued");
+ assert.equal(getRegistry().cancel(job.jobId), true);
+ assert.equal((await job.done).status, "cancelled");
+ assert.equal((await metaReaches(paths, job.jobId, "cancelled"))?.status, "cancelled");
+ } finally {
+ await holder.release();
+ await rm(root, { recursive: true, force: true });
+ }
+});
+
+// THE ONE CANCEL THAT MUST NOT: the graceful-shutdown reaper cancels every
+// queued job only so the exit cannot promote one into a child. Nobody cancelled
+// it, and its `queued` sidecar is what the boot pass (bootQueuedJobs.ts)
+// settles or re-queues on the next start. Last in the file: it puts the
+// process's registry into shutdown, and replaces it afterwards.
+test("at shutdown a queued job's sidecar stays queued, for the boot pass", async () => {
+ const { paths, root } = await jobsDir();
+ const queueKey = `test:shutdown-cancel:${newJobId()}`;
+ const holder = await hold(paths, queueKey);
+ try {
+ const job = ok(
+ await runManagedFunction({
+ kind: "test-queued",
+ queueKey,
+ paths,
+ fn: async () => {
+ throw new Error("must never start");
+ },
+ }),
+ );
+ assert.equal((await metaReaches(paths, job.jobId, "queued"))?.status, "queued");
+ getRegistry().beginShutdown();
+ assert.equal(getRegistry().cancel(job.jobId), true);
+ // The job itself still settles, and the registry records the cancel.
+ assert.equal((await job.done).status, "cancelled");
+ assert.equal(getRegistry().get(job.jobId)?.status, "cancelled");
+ // Nothing rewrote the sidecar (a wrong write would land within ms).
+ await new Promise((r) => setTimeout(r, 200));
+ assert.equal((await readMeta(paths, job.jobId))?.status, "queued");
+ } finally {
+ await holder.release();
+ globalThis.__yttJobRegistry__ = undefined;
+ await rm(root, { recursive: true, force: true });
+ }
+});
diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts
@@ -164,6 +164,19 @@ function makeDoneDeferred(jobId: string): {
return { done, settle };
}
+// ONE WRITER PER JOB, IN ORDER. A job's sidecar is written at least twice —
+// at enqueue, then with its terminal state — and each write is async, so two
+// in flight at once could land in either order: a cancel a moment after the
+// enqueue could leave "queued" on disk, or a shorter JSON over a longer one's
+// tail. Chained, each write starts when the previous one has ended and
+// serializes the record as it is THEN. writeJobMeta never rejects.
+function metaWriter(paths: Paths, record: JobRecord): () => void {
+ let last: Promise<void> = Promise.resolve();
+ return () => {
+ last = last.then(() => writeJobMeta(paths, record));
+ };
+}
+
export async function runManagedCommand(
opts: RunManagedCommandOpts,
): Promise<StreamActionResult> {
@@ -180,6 +193,7 @@ export async function runManagedCommand(
);
const { done, settle } = makeDoneDeferred(id);
+ const persistMeta = metaWriter(opts.paths, record);
const safe = makeSafeController<string>();
let fileStream: WriteStream | null = null;
@@ -244,23 +258,31 @@ export async function runManagedCommand(
safe.safeClose();
requestSnapshotOnFinish(id, opts);
// Persist terminal state (status/endedAt/exitCode now set by finalize).
- void writeJobMeta(opts.paths, record);
+ persistMeta();
// Throttled retention so the .jobs directory stays bounded on its own.
void maybePruneJobLogs(opts.paths);
settle(record.status);
});
};
+ // Cancelled while still queued: start() never runs, so its .finally never
+ // writes the terminal sidecar — this does (release 13 slice W1; it used to
+ // stay "queued", and the next boot could re-queue a job the operator had
+ // cancelled). registry.cancel() marks the record cancelled BEFORE the
+ // scheduler fires this; at a graceful shutdown it does not, and a record
+ // still `queued` is left on disk as it is, for the boot pass.
const onCancel = () => {
cancelledBeforeStart = true;
+ if (record.status === "cancelled") persistMeta();
safe.safeClose();
settle("cancelled");
};
registry.enqueue(record, { start, onCancel });
// Persist queued/running identity up front so a mid-run crash still leaves a
- // sidecar; the .finally above rewrites it with the terminal state.
- void writeJobMeta(opts.paths, record);
+ // sidecar; the .finally above (or onCancel) rewrites it with the terminal
+ // state.
+ persistMeta();
return { ok: true, jobId: id, stream, done };
}
@@ -305,6 +327,7 @@ export async function runManagedFunction(
);
const { done, settle } = makeDoneDeferred(id);
+ const persistMeta = metaWriter(opts.paths, record);
const safe = makeSafeController<string>();
let fileStream: WriteStream | null = null;
@@ -384,22 +407,30 @@ export async function runManagedFunction(
safe.safeClose();
requestSnapshotOnFinish(id, opts);
// Persist terminal state (status/endedAt/exitCode now set by finalize).
- void writeJobMeta(opts.paths, record);
+ persistMeta();
// Throttled retention so the .jobs directory stays bounded on its own.
void maybePruneJobLogs(opts.paths);
settle(record.status);
});
};
+ // Cancelled while still queued: start() never runs, so its .finally never
+ // writes the terminal sidecar — this does (release 13 slice W1; it used to
+ // stay "queued", and the next boot could re-queue a job the operator had
+ // cancelled). registry.cancel() marks the record cancelled BEFORE the
+ // scheduler fires this; at a graceful shutdown it does not, and a record
+ // still `queued` is left on disk as it is, for the boot pass.
const onCancel = () => {
cancelledBeforeStart = true;
+ if (record.status === "cancelled") persistMeta();
safe.safeClose();
settle("cancelled");
};
registry.enqueue(record, { start, onCancel });
// Persist queued/running identity up front so a mid-run crash still leaves a
- // sidecar; the .finally above rewrites it with the terminal state.
- void writeJobMeta(opts.paths, record);
+ // sidecar; the .finally above (or onCancel) rewrites it with the terminal
+ // state.
+ persistMeta();
return { ok: true, jobId: id, stream, done };
}