commit 7a0f18c35bad4fe920d9984d66b428cbb534b107
parent cf3c0abab3f35edc741d81f805482e96ea2932b4
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 25 Sep 2026 18:43:14 -0400
editor: at boot, settle the metas a restart left `queued`
A job still waiting in its queue when the server stopped never got its
terminal meta write, so /jobs read it as queued forever. A boot pass
(common/jobs/bootQueuedJobs.ts, lazy-imported and voided from
instrumentation.ts) re-queues each such meta from its replay spec through
runJobSpec — the path Retry uses — and closes the old meta `cancelled`
with `cancelReason` naming the new job id; a meta with no spec, or whose
re-queue is refused, is closed `cancelled` with the reason. An idle boot
(ARCHILYZER_IDLE_BOOT) and the e2e test server only cancel. Metas queued
by this process or held by the live registry, and `running` metas, are
left alone. One console line per job, and the reason appended to the
job's log. JobMeta gains the optional `cancelReason`.
bootQueuedJobs.test.ts: 6 cases over a temp .jobs dir with an injected
re-queue.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Diffstat:
4 files changed, 388 insertions(+), 0 deletions(-)
diff --git a/common/jobs/bootQueuedJobs.test.ts b/common/jobs/bootQueuedJobs.test.ts
@@ -0,0 +1,209 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { mkdtemp, mkdir, readFile, rm, writeFile } 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 type { JobSpec } from "./jobSpec";
+import { settleQueuedJobMetas, type RequeueFn } from "./bootQueuedJobs";
+
+// Run with:
+// pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/bootQueuedJobs.test.ts
+//
+// THE BOOT PASS (release 9, B4b), over a temp .jobs dir with an injected
+// re-queue. A `queued` meta left by a restart is re-queued from its spec, or
+// closed `cancelled` with the reason; an idle boot only ever cancels. Never
+// run against a real corpus.
+
+const BOOT = 2_000_000_000_000;
+
+async function fixture(metas: Partial<JobMeta>[]): Promise<{
+ paths: Paths;
+ root: string;
+ read: (id: string) => Promise<JobMeta>;
+ log: (id: string) => Promise<string>;
+}> {
+ const root = await mkdtemp(path.join(tmpdir(), "boot-queued-"));
+ const jobsDir = path.join(root, ".jobs");
+ await mkdir(jobsDir, { recursive: true });
+ for (const m of metas) {
+ const meta = {
+ kind: "sync",
+ queueKey: "platform:youtube",
+ status: "queued",
+ queuedAt: BOOT - 60_000,
+ ...m,
+ };
+ await writeFile(
+ path.join(jobsDir, `${meta.id}.meta.json`),
+ JSON.stringify(meta),
+ );
+ }
+ return {
+ paths: { jobsDir } as Paths,
+ root,
+ read: async (id) =>
+ JSON.parse(
+ await readFile(path.join(jobsDir, `${id}.meta.json`), "utf8"),
+ ) as JobMeta,
+ log: (id) =>
+ readFile(path.join(jobsDir, `${id}.log`), "utf8").catch(() => ""),
+ };
+}
+
+const SPEC: JobSpec = { kind: "sync", slug: "teamrcn", params: {} };
+
+function recordingRequeue(): { fn: RequeueFn; calls: JobSpec[] } {
+ const calls: JobSpec[] = [];
+ let n = 0;
+ return {
+ calls,
+ fn: async (spec) => {
+ calls.push(spec);
+ return { ok: true, jobId: `NEW${++n}` };
+ },
+ };
+}
+
+test("a queued meta with a spec is re-queued, and the old one is closed naming the new id", async () => {
+ const f = await fixture([{ id: "A1", spec: SPEC, channelSlug: "teamrcn" }]);
+ try {
+ const rq = recordingRequeue();
+ const lines: string[] = [];
+ const res = await settleQueuedJobMetas({
+ paths: f.paths,
+ requeue: rq.fn,
+ bootedAt: BOOT,
+ log: (l) => lines.push(l),
+ });
+ assert.deepEqual(rq.calls, [SPEC]);
+ assert.deepEqual(res.requeued, [{ id: "A1", newId: "NEW1" }]);
+ assert.deepEqual(res.cancelled, []);
+ const meta = await f.read("A1");
+ assert.equal(meta.status, "cancelled");
+ assert.match(
+ meta.cancelReason ?? "",
+ /server restarted before it ran; re-queued as NEW1/,
+ );
+ assert.ok(typeof meta.endedAt === "number");
+ assert.match(await f.log("A1"), /re-queued as NEW1/);
+ assert.equal(lines.length, 1);
+ assert.match(lines[0], /A1 \(sync teamrcn\)/);
+ } finally {
+ await rm(f.root, { recursive: true, force: true });
+ }
+});
+
+test("a queued meta with no spec is cancelled with the reason, not re-queued", async () => {
+ const f = await fixture([{ id: "B1", kind: "build-export" }]);
+ try {
+ const rq = recordingRequeue();
+ const res = await settleQueuedJobMetas({
+ paths: f.paths,
+ requeue: rq.fn,
+ bootedAt: BOOT,
+ });
+ assert.equal(rq.calls.length, 0);
+ assert.equal(res.cancelled.length, 1);
+ const meta = await f.read("B1");
+ assert.equal(meta.status, "cancelled");
+ assert.match(meta.cancelReason ?? "", /^server restarted before it ran/);
+ } finally {
+ await rm(f.root, { recursive: true, force: true });
+ }
+});
+
+test("an idle boot cancels every queued meta and re-queues none", async () => {
+ const f = await fixture([{ id: "C1", spec: SPEC }, { id: "C2" }]);
+ try {
+ const res = await settleQueuedJobMetas({
+ paths: f.paths,
+ requeue: null,
+ bootedAt: BOOT,
+ idleReason: "idle boot",
+ });
+ assert.deepEqual(res.requeued, []);
+ assert.deepEqual(
+ res.cancelled.map((c) => c.id),
+ ["C1", "C2"],
+ );
+ assert.match(
+ (await f.read("C1")).cancelReason ?? "",
+ /idle boot: not re-queued/,
+ );
+ assert.equal((await f.read("C2")).status, "cancelled");
+ } finally {
+ await rm(f.root, { recursive: true, force: true });
+ }
+});
+
+test("terminal, running, this-boot and live metas are left alone", async () => {
+ const f = await fixture([
+ { id: "D1", status: "done", spec: SPEC },
+ { id: "D2", status: "running", spec: SPEC },
+ { id: "D3", spec: SPEC, queuedAt: BOOT + 5 },
+ { id: "D4", spec: SPEC },
+ ]);
+ try {
+ const rq = recordingRequeue();
+ const res = await settleQueuedJobMetas({
+ paths: f.paths,
+ requeue: rq.fn,
+ bootedAt: BOOT,
+ isLive: (id) => id === "D4",
+ });
+ assert.equal(rq.calls.length, 0);
+ assert.deepEqual(res, { requeued: [], cancelled: [] });
+ assert.equal((await f.read("D1")).status, "done");
+ assert.equal((await f.read("D2")).status, "running");
+ assert.equal((await f.read("D3")).status, "queued");
+ assert.equal((await f.read("D4")).status, "queued");
+ } finally {
+ await rm(f.root, { recursive: true, force: true });
+ }
+});
+
+test("a refused or throwing re-queue closes the meta with the error", async () => {
+ const f = await fixture([
+ { id: "E1", spec: SPEC },
+ { id: "E2", spec: SPEC },
+ ]);
+ try {
+ let calls = 0;
+ const res = await settleQueuedJobMetas({
+ paths: f.paths,
+ requeue: async () => {
+ calls++;
+ if (calls === 1) {
+ return { ok: false, error: 'Channel "teamrcn" not found' };
+ }
+ throw new Error("boom");
+ },
+ bootedAt: BOOT,
+ });
+ assert.deepEqual(
+ res.cancelled.map((c) => c.id),
+ ["E1", "E2"],
+ );
+ assert.match(
+ (await f.read("E1")).cancelReason ?? "",
+ /re-queue refused: Channel "teamrcn" not found/,
+ );
+ assert.match(
+ (await f.read("E2")).cancelReason ?? "",
+ /re-queue refused: boom/,
+ );
+ } finally {
+ await rm(f.root, { recursive: true, force: true });
+ }
+});
+
+test("a missing .jobs dir is nothing to do", async () => {
+ const res = await settleQueuedJobMetas({
+ paths: { jobsDir: path.join(tmpdir(), "no-such-jobs-dir-r9") } as Paths,
+ requeue: null,
+ bootedAt: BOOT,
+ });
+ assert.deepEqual(res, { requeued: [], cancelled: [] });
+});
diff --git a/common/jobs/bootQueuedJobs.ts b/common/jobs/bootQueuedJobs.ts
@@ -0,0 +1,133 @@
+import path from "node:path";
+import { appendFile, readdir, writeFile } from "node:fs/promises";
+import type { Paths } from "../lib/paths";
+import { metaPath, readJobMeta, type JobMeta } from "./jobMeta";
+import type { JobSpec } from "./jobSpec";
+
+// THE BOOT PASS OVER STALE `queued` METAS (release 9, B4b).
+//
+// The registry and the scheduler are in memory; a job's `<id>.meta.json` is
+// written `queued` when it is submitted and rewritten when it finishes. A job
+// still waiting in its queue when the server stopped therefore never gets its
+// terminal write, and /jobs reads that sidecar forever as `queued` — a job that
+// will never run, in a list the operator trusts to say what is pending.
+//
+// One pass at boot settles each of them, exactly once:
+// - a meta with a replay descriptor (`spec`) is RE-QUEUED through the same
+// path Retry uses, and the old meta is closed as `cancelled`, naming the
+// new job id — so the work the operator asked for still happens, and the
+// stale row stops claiming to be pending;
+// - a meta with no spec (a kind that cannot be replayed) is closed as
+// `cancelled` with the reason;
+// - on an idle boot (`requeue: null`) EVERY one is closed as `cancelled`:
+// an idle boot must not resume work.
+//
+// Only metas from BEFORE this boot, and not held by the live registry, are
+// touched: the storage boot pass can enqueue a job of its own while this runs.
+// `running` metas are left alone — whether a job that was mid-flight should be
+// re-run is not a decision a boot pass can make.
+//
+// Best-effort throughout: a meta that cannot be read or written is skipped,
+// and the caller voids the promise so readiness never waits on it.
+
+export type RequeueFn = (
+ spec: JobSpec,
+) => Promise<{ ok: true; jobId: string } | { ok: false; error: string }>;
+
+export type BootQueuedResult = {
+ requeued: { id: string; newId: string }[];
+ cancelled: { id: string; reason: string }[];
+};
+
+export const RESTART_REASON = "server restarted before it ran";
+
+export async function settleQueuedJobMetas(opts: {
+ paths: Paths;
+ // null: cancel, never re-queue (an idle boot, or the e2e test server).
+ requeue: RequeueFn | null;
+ // Metas queued at or after this instant belong to THIS process.
+ bootedAt: number;
+ // True for an id the live registry holds (enqueued during this boot).
+ isLive?: (id: string) => boolean;
+ log?: (line: string) => void;
+ idleReason?: string;
+}): Promise<BootQueuedResult> {
+ const log = opts.log ?? (() => {});
+ const result: BootQueuedResult = { requeued: [], cancelled: [] };
+ let names: string[];
+ try {
+ names = await readdir(opts.paths.jobsDir);
+ } catch {
+ return result;
+ }
+ const ids = names
+ .filter((n) => n.endsWith(".meta.json"))
+ .map((n) => n.slice(0, -".meta.json".length))
+ .sort();
+ for (const id of ids) {
+ const meta = await readJobMeta(opts.paths, id);
+ if (!meta || meta.status !== "queued") continue;
+ if (typeof meta.queuedAt === "number" && meta.queuedAt >= opts.bootedAt) {
+ continue;
+ }
+ if (opts.isLive?.(id)) continue;
+
+ let reason: string;
+ if (!opts.requeue) {
+ reason = `${RESTART_REASON}${opts.idleReason ? ` (${opts.idleReason}: not re-queued)` : ""}`;
+ } else if (!meta.spec) {
+ reason = `${RESTART_REASON} (kind ${meta.kind} cannot be re-queued)`;
+ } else {
+ let res: Awaited<ReturnType<RequeueFn>>;
+ try {
+ res = await opts.requeue(meta.spec);
+ } catch (err) {
+ res = { ok: false, error: (err as Error).message };
+ }
+ if (res.ok) {
+ reason = `${RESTART_REASON}; re-queued as ${res.jobId}`;
+ result.requeued.push({ id, newId: res.jobId });
+ } else {
+ reason = `${RESTART_REASON}; re-queue refused: ${res.error}`;
+ }
+ }
+ const closed = await closeMeta(opts.paths, meta, reason);
+ if (!closed) continue;
+ if (!reason.includes("; re-queued as ")) {
+ result.cancelled.push({ id, reason });
+ }
+ log(
+ `[boot] job ${id} (${meta.kind}${meta.channelSlug ? ` ${meta.channelSlug}` : ""}) was queued: ${reason}`,
+ );
+ }
+ return result;
+}
+
+async function closeMeta(
+ paths: Paths,
+ meta: JobMeta,
+ reason: string,
+): Promise<boolean> {
+ const closed: JobMeta = {
+ ...meta,
+ status: "cancelled",
+ endedAt: Date.now(),
+ cancelReason: reason,
+ };
+ try {
+ await writeFile(metaPath(paths, meta.id), JSON.stringify(closed), "utf8");
+ } catch {
+ return false;
+ }
+ // The job's own log says it too, for whoever opens the row. It may not exist
+ // (a job that never started never opened one); appendFile creates it.
+ try {
+ await appendFile(
+ path.join(paths.jobsDir, `${meta.id}.log`),
+ `[boot] ${reason}\n`,
+ );
+ } catch {
+ /* the meta is the record; the log line is a courtesy */
+ }
+ return true;
+}
diff --git a/common/jobs/jobMeta.ts b/common/jobs/jobMeta.ts
@@ -25,6 +25,10 @@ export type JobMeta = {
// Replay descriptor for replayable jobs, preserved so an archived (evicted
// or post-restart) job can still be retried from disk. See jobSpec.ts.
spec?: JobSpec;
+ // Why a job ended `cancelled` without anyone pressing Cancel — written only
+ // by the boot pass (bootQueuedJobs.ts) for a job that was still `queued`
+ // when the server went down. Absent on every other meta.
+ cancelReason?: string;
};
export function metaPath(paths: Paths, id: string): string {
diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts
@@ -17,6 +17,8 @@ export async function register() {
// its transitive imports (runTick -> server actions, lmdb, fs) are Node-only,
// so guard the import — and run nothing on Edge.
if (process.env.NEXT_RUNTIME !== "nodejs") return;
+ // Anything a previous process left `queued` was queued before this instant.
+ const bootedAt = Date.now();
// Cancel in-flight children on a graceful shutdown. Armed FIRST, before any of
// the runners below: a failure while starting those must not leave the server
@@ -105,6 +107,46 @@ export async function register() {
/* a storage probe that fails to start must not block server readiness */
}
+ // STALE `queued` METAS FROM THE LAST PROCESS (release 9, B4b). A job still
+ // waiting in its queue when the server stopped never got its terminal meta
+ // write, so /jobs would read it as queued forever. Each one is re-queued from
+ // its replay spec through the same path Retry uses, or closed `cancelled`
+ // with the reason. Runs on an idle boot too, but there it ONLY cancels — an
+ // idle boot must not resume work — and so does the e2e test server, whose
+ // leftover metas belong to a previous run's fixture. See
+ // common/jobs/bootQueuedJobs.ts. Lazy, voided, best-effort: never blocks
+ // readiness.
+ try {
+ const { settleQueuedJobMetas } = await import(
+ "yt-dlp-transcript-common/jobs/bootQueuedJobs"
+ );
+ const { getPaths } = await import("yt-dlp-transcript-common/lib/paths");
+ const { getRegistry } = await import(
+ "yt-dlp-transcript-common/jobs/registry"
+ );
+ const testServer = process.env.EDITOR_TEST_ROUTES === "1";
+ const cancelOnly = idle || testServer;
+ void settleQueuedJobMetas({
+ paths: getPaths(),
+ bootedAt,
+ isLive: (id) => getRegistry().get(id) !== undefined,
+ log: (line) => console.log(line),
+ idleReason: idle ? "idle boot" : testServer ? "test server" : undefined,
+ requeue: cancelOnly
+ ? null
+ : async (spec) => {
+ const { runJobSpec } = await import("./app/jobs/runJobSpec");
+ const res = await runJobSpec(spec);
+ if (!res.ok) return { ok: false, error: res.error };
+ // Nobody reads this stream; release it as the ops routes do.
+ void res.stream.cancel();
+ return { ok: true, jobId: res.jobId };
+ },
+ }).catch(() => {});
+ } catch {
+ /* a boot pass that fails to start must not block server readiness */
+ }
+
// Everything past here STARTS work. On an idle boot, nothing does.
if (idle) return;