Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit 3879dac23914fa1b5628effe60ec921e891da555
parent 11f127474fc7c5edf9a18a153bd5b5a3ca85c692
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri, 25 Sep 2026 19:24:12 -0400

editor: the boot pass re-queues little — no syncs, nothing stale, newest per spec (review fix)

Review HIGH-1: the live .jobs holds 392 queued metas (346 sync over 51
channels, 29 whisper-all from June/July), so the first boot would have
submitted ~382 jobs ahead of the auto-download runner and around the
heartbeat's maxConcurrentSyncs stagger. settleQueuedJobMetas now:
- never re-queues a sync ("server restarted; the scheduler re-derives
  syncs");
- cancels anything queued more than REQUEUE_MAX_AGE_MS (24 h) before the
  boot ("queued before the last restart, stale");
- re-queues only the newest meta per kind + channel + bucket + params
  (key-sorted), cancelling the rest as "superseded by a newer queued job";
- logs one summary line (re-queued N, cancelled M by reason) plus one line
  per re-queued job.
Also: closeMeta writes tmp + rename (NIT-1); instrumentation waits for the
storage boot pass before settling, so a channel mid-autoRepoint is not
refused as unreachable (LOW-1). bootQueuedJobs.test.ts: 9 cases, incl. 7
duplicate syncs + 1 stale whisper-all + 1 fresh non-sync -> exactly 1
re-queue, and a malformed meta skipped while the pass continues (LOW-2).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

Diffstat:
Mcommon/jobs/bootQueuedJobs.test.ts | 179++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Mcommon/jobs/bootQueuedJobs.ts | 172++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------
Meditor/instrumentation.ts | 60+++++++++++++++++++++++++++++++++++++-----------------------
3 files changed, 332 insertions(+), 79 deletions(-)

diff --git a/common/jobs/bootQueuedJobs.test.ts b/common/jobs/bootQueuedJobs.test.ts @@ -12,9 +12,10 @@ import { settleQueuedJobMetas, type RequeueFn } from "./bootQueuedJobs"; // 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. +// re-queue. It re-queues LITTLE: never a sync, never a meta older than 24 h, +// only the newest of duplicate specs; everything else is closed `cancelled` +// with the reason. An idle boot only ever cancels. Never run against a real +// corpus. const BOOT = 2_000_000_000_000; @@ -29,7 +30,7 @@ async function fixture(metas: Partial<JobMeta>[]): Promise<{ await mkdir(jobsDir, { recursive: true }); for (const m of metas) { const meta = { - kind: "sync", + kind: "whisper-all", queueKey: "platform:youtube", status: "queued", queuedAt: BOOT - 60_000, @@ -52,7 +53,13 @@ async function fixture(metas: Partial<JobMeta>[]): Promise<{ }; } -const SPEC: JobSpec = { kind: "sync", slug: "teamrcn", params: {} }; +const SPEC: JobSpec = { kind: "whisper-all", slug: "teamrcn", params: {} }; +const SYNC: JobSpec = { + kind: "sync", + slug: "teamrcn", + params: { fullSweep: false }, +}; +const HOUR = 60 * 60 * 1000; function recordingRequeue(): { fn: RequeueFn; calls: JobSpec[] } { const calls: JobSpec[] = []; @@ -66,7 +73,7 @@ function recordingRequeue(): { fn: RequeueFn; calls: JobSpec[] } { }; } -test("a queued meta with a spec is re-queued, and the old one is closed naming the new id", async () => { +test("a fresh non-sync 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(); @@ -88,8 +95,106 @@ test("a queued meta with a spec is re-queued, and the old one is closed naming t ); 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\)/); + // One line per re-queued job, then the summary. + assert.equal(lines.length, 2); + assert.match(lines[0], /re-queued whisper-all teamrcn: A1 -> NEW1/); + assert.match(lines[1], /re-queued 1, cancelled 0/); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +// The review's sizing of the live corpus: 346 queued syncs over 51 channels, +// ~7 per channel, plus months-old whisper-all sweeps. Exactly one job may come +// back from this. +test("7 duplicate syncs + 1 stale whisper-all + 1 fresh non-sync: exactly one re-queue", async () => { + const syncs = Array.from({ length: 7 }, (_, i) => ({ + id: `S${i}`, + kind: "sync", + spec: SYNC, + channelSlug: "teamrcn", + queuedAt: BOOT - (i + 1) * HOUR, + })); + const f = await fixture([ + ...syncs, + { + id: "W1", + spec: { kind: "whisper-all", slug: "old", params: {} }, + queuedAt: BOOT - 60 * 24 * HOUR, + }, + { id: "F1", spec: SPEC, queuedAt: BOOT - 2 * HOUR }, + ]); + 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: "F1", newId: "NEW1" }]); + const byCat = (c: string) => + res.cancelled.filter((x) => x.category === c).length; + assert.equal(byCat("sync"), 7); + assert.equal(byCat("stale"), 1); + assert.match( + (await f.read("S0")).cancelReason ?? "", + /the scheduler re-derives syncs/, + ); + assert.match( + (await f.read("W1")).cancelReason ?? "", + /queued before the last restart, stale/, + ); + assert.match( + lines.at(-1) ?? "", + /re-queued 1, cancelled 8 \(sync 7, stale 1\)/, + ); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("of duplicate specs only the newest is re-queued; params are compared key-order-free", async () => { + const f = await fixture([ + { + id: "D1", + spec: { kind: "whisper-all", slug: "c", params: { a: 1, b: 2 } }, + queuedAt: BOOT - 3 * HOUR, + }, + { + id: "D2", + spec: { kind: "whisper-all", slug: "c", params: { b: 2, a: 1 } }, + queuedAt: BOOT - 1 * HOUR, + }, + { + id: "D3", + spec: { kind: "whisper-all", slug: "c", params: { a: 9 } }, + queuedAt: BOOT - 2 * HOUR, + }, + ]); + try { + const rq = recordingRequeue(); + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: rq.fn, + bootedAt: BOOT, + }); + assert.deepEqual(res.requeued.map((r) => r.id).sort(), ["D2", "D3"]); + assert.deepEqual( + res.cancelled.map((c) => [c.id, c.category]), + [["D1", "superseded"]], + ); + assert.match( + (await f.read("D1")).cancelReason ?? "", + /superseded by a newer queued job \(D2\)/, + ); + // Survivors are re-queued oldest first, keeping their relative order. + assert.deepEqual( + res.requeued.map((r) => r.id), + ["D3", "D2"], + ); } finally { await rm(f.root, { recursive: true, force: true }); } @@ -105,7 +210,10 @@ test("a queued meta with no spec is cancelled with the reason, not re-queued", a bootedAt: BOOT, }); assert.equal(rq.calls.length, 0); - assert.equal(res.cancelled.length, 1); + assert.deepEqual( + res.cancelled.map((c) => c.category), + ["no-spec"], + ); const meta = await f.read("B1"); assert.equal(meta.status, "cancelled"); assert.match(meta.cancelReason ?? "", /^server restarted before it ran/); @@ -124,10 +232,7 @@ test("an idle boot cancels every queued meta and re-queues none", async () => { idleReason: "idle boot", }); assert.deepEqual(res.requeued, []); - assert.deepEqual( - res.cancelled.map((c) => c.id), - ["C1", "C2"], - ); + assert.deepEqual(res.cancelled.map((c) => c.id).sort(), ["C1", "C2"]); assert.match( (await f.read("C1")).cancelReason ?? "", /idle boot: not re-queued/, @@ -164,10 +269,38 @@ test("terminal, running, this-boot and live metas are left alone", async () => { } }); +test("a malformed meta is skipped and the pass continues", async () => { + const f = await fixture([{ id: "M2", spec: SPEC }]); + try { + await writeFile( + path.join(f.paths.jobsDir, "M1.meta.json"), + '{"id":"M1","kind":', + ); + await writeFile( + path.join(f.paths.jobsDir, "M0.meta.json"), + '{"status":"queued"}', + ); + const rq = recordingRequeue(); + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: rq.fn, + bootedAt: BOOT, + }); + assert.deepEqual(res.requeued, [{ id: "M2", newId: "NEW1" }]); + assert.deepEqual(res.cancelled, []); + assert.equal( + await readFile(path.join(f.paths.jobsDir, "M1.meta.json"), "utf8"), + '{"id":"M1","kind":', + ); + } 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 }, + { id: "E1", spec: { kind: "whisper-all", slug: "a", params: {} } }, + { id: "E2", spec: { kind: "whisper-all", slug: "b", params: {} } }, ]); try { let calls = 0; @@ -182,18 +315,14 @@ test("a refused or throwing re-queue closes the meta with the error", async () = }, bootedAt: BOOT, }); - assert.deepEqual( - res.cancelled.map((c) => c.id), - ["E1", "E2"], - ); - assert.match( + assert.equal(res.cancelled.length, 2); + assert.ok(res.cancelled.every((c) => c.category === "refused")); + const reasons = [ (await f.read("E1")).cancelReason ?? "", - /re-queue refused: Channel "teamrcn" not found/, - ); - assert.match( (await f.read("E2")).cancelReason ?? "", - /re-queue refused: boom/, - ); + ].join("\n"); + assert.match(reasons, /re-queue refused: Channel "teamrcn" not found/); + assert.match(reasons, /re-queue refused: boom/); } finally { await rm(f.root, { recursive: true, force: true }); } diff --git a/common/jobs/bootQueuedJobs.ts b/common/jobs/bootQueuedJobs.ts @@ -1,5 +1,6 @@ import path from "node:path"; -import { appendFile, readdir, writeFile } from "node:fs/promises"; +import { appendFile, readdir } from "node:fs/promises"; +import { writeFileAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; import { metaPath, readJobMeta, type JobMeta } from "./jobMeta"; import type { JobSpec } from "./jobSpec"; @@ -12,20 +13,29 @@ import type { JobSpec } from "./jobSpec"; // 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. +// One pass at boot settles each of them, exactly once, and RE-QUEUES LITTLE. +// The first version re-queued every `queued` meta with a spec. The live corpus +// held 392 (review, 2026-09-25): 346 `sync` over 51 channels, 29 months-old +// `whisper-all` sweeps and a few others. That would have been ~382 jobs +// submitted at boot, ahead of the auto-download runner and around the +// heartbeat's maxConcurrentSyncs stagger: the burst the pacing work forbids. +// So, in order: +// 1. an idle boot (`requeue: null`) cancels everything — it must not resume +// work; +// 2. `sync` is never re-queued: the heartbeat re-derives due syncs itself, +// paced; +// 3. a meta queued more than REQUEUE_MAX_AGE_MS before this boot is stale; +// 4. a meta with no replay spec cannot be re-queued; +// 5. of what is left, only the NEWEST meta per kind + channel + params is +// re-queued (through the same path Retry uses); the rest are superseded. +// Every meta this pass does not re-queue is closed `cancelled` with a +// `cancelReason`; a re-queued one is closed naming its new id. // // 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. +// re-run is not a decision a boot pass can make. A malformed meta (readJobMeta +// → null) is skipped. // // 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. @@ -34,12 +44,41 @@ export type RequeueFn = ( spec: JobSpec, ) => Promise<{ ok: true; jobId: string } | { ok: false; error: string }>; +export type CancelCategory = + "idle" | "sync" | "stale" | "no-spec" | "superseded" | "refused"; + export type BootQueuedResult = { requeued: { id: string; newId: string }[]; - cancelled: { id: string; reason: string }[]; + cancelled: { id: string; reason: string; category: CancelCategory }[]; }; export const RESTART_REASON = "server restarted before it ran"; +export const REQUEUE_MAX_AGE_MS = 24 * 60 * 60 * 1000; + +// Key a spec by what it would DO: kind + channel + bucket + params, with the +// params serialized key-sorted so two equal specs written in a different key +// order collide. +function stableJson(v: unknown): string { + if (Array.isArray(v)) return `[${v.map(stableJson).join(",")}]`; + if (v && typeof v === "object") { + const o = v as Record<string, unknown>; + return `{${Object.keys(o) + .sort() + .filter((k) => o[k] !== undefined) + .map((k) => `${JSON.stringify(k)}:${stableJson(o[k])}`) + .join(",")}}`; + } + return JSON.stringify(v) ?? "null"; +} + +export function specKey(spec: JobSpec): string { + return stableJson({ + kind: spec.kind, + slug: spec.slug, + bucket: spec.bucket, + params: spec.params ?? {}, + }); +} export async function settleQueuedJobMetas(opts: { paths: Paths; @@ -51,8 +90,10 @@ export async function settleQueuedJobMetas(opts: { isLive?: (id: string) => boolean; log?: (line: string) => void; idleReason?: string; + maxAgeMs?: number; }): Promise<BootQueuedResult> { const log = opts.log ?? (() => {}); + const maxAgeMs = opts.maxAgeMs ?? REQUEUE_MAX_AGE_MS; const result: BootQueuedResult = { requeued: [], cancelled: [] }; let names: string[]; try { @@ -64,6 +105,8 @@ export async function settleQueuedJobMetas(opts: { .filter((n) => n.endsWith(".meta.json")) .map((n) => n.slice(0, -".meta.json".length)) .sort(); + + const stale: JobMeta[] = []; for (const id of ids) { const meta = await readJobMeta(opts.paths, id); if (!meta || meta.status !== "queued") continue; @@ -71,33 +114,99 @@ export async function settleQueuedJobMetas(opts: { continue; } if (opts.isLive?.(id)) continue; + stale.push(meta); + } - let reason: string; + const cancel = async ( + meta: JobMeta, + category: CancelCategory, + reason: string, + ) => { + if (await closeMeta(opts.paths, meta, reason)) { + result.cancelled.push({ id: meta.id, reason, category }); + } + }; + + // Newest first, so the first meta seen for a spec key is the one kept. + stale.sort( + (a, b) => (b.queuedAt ?? 0) - (a.queuedAt ?? 0) || b.id.localeCompare(a.id), + ); + const kept = new Map<string, string>(); // spec key -> kept meta id + const toRequeue: { meta: JobMeta; spec: JobSpec }[] = []; + for (const meta of stale) { if (!opts.requeue) { - reason = `${RESTART_REASON}${opts.idleReason ? ` (${opts.idleReason}: not re-queued)` : ""}`; + await cancel( + meta, + "idle", + `${RESTART_REASON}${opts.idleReason ? ` (${opts.idleReason}: not re-queued)` : ""}`, + ); + } else if (meta.kind === "sync") { + await cancel( + meta, + "sync", + "server restarted; the scheduler re-derives syncs", + ); + } else if ( + typeof meta.queuedAt !== "number" || + opts.bootedAt - meta.queuedAt > maxAgeMs + ) { + await cancel(meta, "stale", "queued before the last restart, stale"); } else if (!meta.spec) { - reason = `${RESTART_REASON} (kind ${meta.kind} cannot be re-queued)`; + await cancel( + meta, + "no-spec", + `${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 }); + const key = specKey(meta.spec); + const newer = kept.get(key); + if (newer) { + await cancel( + meta, + "superseded", + `superseded by a newer queued job (${newer})`, + ); } else { - reason = `${RESTART_REASON}; re-queue refused: ${res.error}`; + kept.set(key, meta.id); + toRequeue.push({ meta, spec: meta.spec }); } } - const closed = await closeMeta(opts.paths, meta, reason); - if (!closed) continue; - if (!reason.includes("; re-queued as ")) { - result.cancelled.push({ id, reason }); + } + + // Oldest first, so the survivors keep their relative queue order. + for (const { meta, spec } of toRequeue.reverse()) { + let res: Awaited<ReturnType<RequeueFn>>; + try { + res = await opts.requeue!(spec); + } catch (err) { + res = { ok: false, error: (err as Error).message }; } + if (res.ok) { + const reason = `${RESTART_REASON}; re-queued as ${res.jobId}`; + if (await closeMeta(opts.paths, meta, reason)) { + result.requeued.push({ id: meta.id, newId: res.jobId }); + log( + `[boot] re-queued ${meta.kind}${meta.channelSlug ? ` ${meta.channelSlug}` : ""}: ${meta.id} -> ${res.jobId}`, + ); + } + } else { + await cancel( + meta, + "refused", + `${RESTART_REASON}; re-queue refused: ${res.error}`, + ); + } + } + + if (result.requeued.length || result.cancelled.length) { + const by = new Map<string, number>(); + for (const c of result.cancelled) + by.set(c.category, (by.get(c.category) ?? 0) + 1); log( - `[boot] job ${id} (${meta.kind}${meta.channelSlug ? ` ${meta.channelSlug}` : ""}) was queued: ${reason}`, + `[boot] stale queued jobs: re-queued ${result.requeued.length}, cancelled ${result.cancelled.length}` + + (by.size + ? ` (${[...by].map(([k, n]) => `${k} ${n}`).join(", ")})` + : ""), ); } return result; @@ -115,7 +224,8 @@ async function closeMeta( cancelReason: reason, }; try { - await writeFile(metaPath(paths, meta.id), JSON.stringify(closed), "utf8"); + // tmp + rename: this pass can close hundreds of metas in one burst. + await writeFileAtomic(metaPath(paths, meta.id), JSON.stringify(closed)); } catch { return false; } diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts @@ -98,22 +98,28 @@ export async function register() { // Lazy-imported and `void`ed like everything else here: the controller pulls // in execa transitively (Node-only), and a disk that cannot be probed must // never block server readiness. + // Kept so the queued-meta pass below can wait for it: a re-queue for a + // channel whose location is mid-autoRepoint would be refused as unreachable. + let storagePass: Promise<unknown> = Promise.resolve(); try { const { runStorageBootPass } = await import( "yt-dlp-transcript-common/controller/storageLocations" ); - void runStorageBootPass({ enqueue: !idle }).catch(() => {}); + storagePass = runStorageBootPass({ enqueue: !idle }).catch(() => {}); } catch { /* 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 + // write, so /jobs would read it as queued forever. Each is closed + // `cancelled` with a reason; only a few are re-queued (through the path Retry + // uses): never a sync (the heartbeat re-derives those, paced), nothing older + // than 24 h, and only the newest of duplicate specs. 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. It waits for the storage pass above, so a channel being + // re-pointed is reachable when its job is re-queued. See // common/jobs/bootQueuedJobs.ts. Lazy, voided, best-effort: never blocks // readiness. try { @@ -126,23 +132,31 @@ export async function register() { ); 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(() => {}); + void storagePass + .then(() => + 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 */ }