import path from "node:path"; import { appendFile, readdir, stat } 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"; import { isPublishStageJobKind } from "./jobKinds"; // THE BOOT PASS OVER STALE `queued` METAS (release 9, B4b). // // The registry and the scheduler are in memory; a job's `.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, 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; nor is a publish stage (release 18, `publish-*`): the publish // lane and Publish now re-derive stages from the stamps on disk, and a // stage re-queued with its run's preconditions would wait on a run that // is gone; // 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, not held by the live registry, and whose // writer is gone (`writerIsGone`: another live process — `archilyzer run` — // writes into the same `.jobs/`) are touched: the storage boot pass can // enqueue a job of its own while this runs. `running` metas are not re-queued // — whether a job that was mid-flight should be re-run is not a decision a // boot pass can make — but they are CLOSED, by the second pass below // (`settleRunningJobMetas`, release 17 slice D0). 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. export type RequeueFn = ( spec: JobSpec, ) => Promise<{ ok: true; jobId: string } | { ok: false; error: string }>; export type CancelCategory = "idle" | "sync" | "publish" | "stale" | "no-spec" | "superseded" | "refused"; export type BootQueuedResult = { requeued: { id: string; newId: string }[]; cancelled: { id: string; reason: string; category: CancelCategory }[]; }; export const RESTART_REASON = "server restarted before it ran"; export const PUBLISH_RESTART_REASON = "server restarted; the publish lane re-derives stages from on-disk state"; 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; 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 ?? {}, }); } // WHO WROTE THE META, AND ARE THEY STILL THERE. // // The registry is in memory, so after a restart nothing in this process knows a // job the last one was running or queuing — and the meta alone cannot say // whether its process died or is another process that is very much alive: // `archilyzer run` (bin/run-operation.ts) runs a job offline and writes its // meta into the same `.jobs/`. So a meta names its writer (`pid`, release 17) // and the writer is gone when: // - the meta names none: it predates release 17, so the only writer it can // have had is a process older than this one; // - it names THIS process's pid: a previous process with the same number — // a container's editor comes back as the same pid every restart, and the // caller already dropped every meta this process wrote (`bootedAt`, // `isLive`); // - no process has that pid (`kill(pid, 0)` → ESRCH). // A pid that answers is left alone, even if the number was reused by an // unrelated process: a job left `running` on /jobs is the cost, a live job // closed under its own feet would be the alternative. export type WriterCheck = { // This process's pid; defaults to process.pid. selfPid?: number; // Defaults to processIsAlive. Injected by the tests. isProcessAlive?: (pid: number) => boolean; }; export function processIsAlive(pid: number): boolean { try { process.kill(pid, 0); return true; } catch (err) { // EPERM: it exists, it is just not ours to signal. return (err as NodeJS.ErrnoException).code === "EPERM"; } } export function writerIsGone(meta: JobMeta, check: WriterCheck = {}): boolean { if (typeof meta.pid !== "number") return true; if (meta.pid === (check.selfPid ?? process.pid)) return true; return !(check.isProcessAlive ?? processIsAlive)(meta.pid); } export type SettleQueuedOpts = WriterCheck & { 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; maxAgeMs?: number; }; // THE WAIT ON THE STORAGE BOOT PASS, BOUNDED (release 10, L2). // // The settle waits for the storage pass (release 9 review, LOW-1) so a channel // whose location is mid-autoRepoint is reachable when its job is re-queued. It // waited with no bound, and one probe in that pass is not bounded either: every // findmnt has FINDMNT_TIMEOUT_MS (3 s), but the `stat` of a location's root and // the statfs for its free space are plain syscalls, and on a hung network mount // they never return. The boot pass then never ran, and every `queued` meta // stayed `queued` on /jobs for the life of the process. // // 60 s. A healthy pass takes milliseconds (findmnt answers in under 10 ms here). // Its bounded worst case is two or three 3 s findmnts per location — identity, // fstab, or where the uuid is mounted — plus, for a location being re-pointed, // the preflight's second probe and a stat per channel on it: ~12 s a location. // 60 s covers several locations at that worst case; past it the pass is stuck on // a syscall that will not answer, and waiting longer buys nothing. // // On a timeout the settle runs anyway. What a job re-queued for a channel the // pass has not reached then meets depends on why the pass is slow: // - an UNMOUNTED drive (the root is simply absent): the media guard // (jobs/jobKinds.ts `needsMedia`, lib/channelMedia.ts) gets ENOENT at // once, so the job is refused — at submission, which closes the old meta // `cancelled` with the error, or when it starts, with the reason in its // log. /jobs says so, and Retry is one click. // - a HUNG mount, the case this bound exists for: the guard's own `stat` // (lib/channelMedia.ts) hangs on the same syscall, so that job waits with // it. Re-queues run one at a time (below), so every re-queue after it stays // `queued` until the mount answers. Every cancel has run by then — cancels // come first — and re-queues are few (one at the first live boot), so this // is left, not fixed: see the release 10 record, L2 "found and left". // The storage pass is NOT cancelled (it has no signal, and a re-point it // already enqueued must finish): it runs on, and a line says when it ends. export const STORAGE_PASS_WAIT_MS = 60_000; export async function waitForStoragePass( pass: Promise, opts: { timeoutMs?: number; log?: (line: string) => void; now?: () => number; } = {}, ): Promise<"done" | "timed-out"> { const timeoutMs = opts.timeoutMs ?? STORAGE_PASS_WAIT_MS; const log = opts.log ?? (() => {}); const now = opts.now ?? Date.now; const startedAt = now(); // A pass that throws has still finished; the settle does not care how. This // is a DERIVED promise — racing it never touches the pass itself. const finished = pass.then( () => "done" as const, () => "done" as const, ); let timer: ReturnType | undefined; const timedOut = new Promise<"timed-out">((resolve) => { timer = setTimeout(() => resolve("timed-out"), timeoutMs); }); const outcome = await Promise.race([finished, timedOut]); clearTimeout(timer); if (outcome === "timed-out") { const secs = (ms: number) => `${Math.round(ms / 1000)} s`; log( `[boot] storage pass still running after ${secs(timeoutMs)}; settling queued jobs without it ` + `(a hung mount? a job re-queued for a channel on it will wait on the same mount, ` + `and the re-queues after it with it)`, ); void finished.then(() => log( `[boot] storage pass finished ${secs(now() - startedAt)} after the queued-job pass began waiting ` + `(it stopped waiting at ${secs(timeoutMs)})`, ), ); } return outcome; } // What instrumentation.ts runs: the settle, after the storage pass or after // STORAGE_PASS_WAIT_MS, whichever comes first. export async function settleAfterStoragePass( storagePass: Promise, opts: SettleQueuedOpts & { waitMs?: number }, ): Promise { await waitForStoragePass(storagePass, { timeoutMs: opts.waitMs, log: opts.log, }); return settleQueuedJobMetas(opts); } export async function settleQueuedJobMetas( opts: SettleQueuedOpts, ): Promise { const log = opts.log ?? (() => {}); const maxAgeMs = opts.maxAgeMs ?? REQUEUE_MAX_AGE_MS; 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(); const stale: JobMeta[] = []; 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; if (!writerIsGone(meta, opts)) continue; stale.push(meta); } 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(); // spec key -> kept meta id const toRequeue: { meta: JobMeta; spec: JobSpec }[] = []; for (const meta of stale) { if (!opts.requeue) { 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 (isPublishStageJobKind(meta.kind)) { await cancel(meta, "publish", PUBLISH_RESTART_REASON); } 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) { await cancel( meta, "no-spec", `${RESTART_REASON} (kind ${meta.kind} cannot be re-queued)`, ); } else { const key = specKey(meta.spec); const newer = kept.get(key); if (newer) { await cancel( meta, "superseded", `superseded by a newer queued job (${newer})`, ); } else { kept.set(key, meta.id); toRequeue.push({ meta, spec: meta.spec }); } } } // Oldest first, so the survivors keep their relative queue order. for (const { meta, spec } of toRequeue.reverse()) { let res: Awaited>; 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(); for (const c of result.cancelled) by.set(c.category, (by.get(c.category) ?? 0) + 1); log( `[boot] stale queued jobs: re-queued ${result.requeued.length}, cancelled ${result.cancelled.length}` + (by.size ? ` (${[...by].map(([k, n]) => `${k} ${n}`).join(", ")})` : ""), ); } return result; } async function closeMeta( paths: Paths, meta: JobMeta, reason: string, endedAt: number = Date.now(), ): Promise { const closed: JobMeta = { ...meta, status: "cancelled", endedAt, cancelReason: reason, }; try { // tmp + rename: this pass can close hundreds of metas in one burst. await writeFileAtomic(metaPath(paths, meta.id), JSON.stringify(closed)); } 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; } // THE BOOT PASS OVER STALE `running` METAS (release 17 slice D0). // // A job running when its process died never got its terminal write either — // a SIGKILL, an OOM, a crash; and a SIGTERM too, whenever the job's function // is still unwinding when Next exits (shutdownCancel.ts cancels every live // job but does not wait for one to finish, so the `cancelled` write that // streamCommand makes when the function returns is usually lost with the // process). Its meta says `running` forever. On 2026-10-01 three // `refresh-report` metas from processes that had been gone for hours still // read `running`. // // Each is CLOSED as the queued pass closes one — `cancelled`, with // INTERRUPTED_REASON — never re-run (whether a half-finished job should run // again is the operator's call; Retry is one click). There is no separate // `interrupted` status: `cancelled` + `cancelReason` is the terminal state this // file already writes for "the server went down under it", and every reader of // a meta (listJobs, /jobs, Retry) already handles it. `endedAt` is the job // log's last write — the last moment the job is known to have been alive — // else this boot. // // Same filters as the queued pass: before this boot, not in the live registry, // writer gone. It does not wait for the storage pass: it re-queues nothing. export const INTERRUPTED_REASON = "interrupted: the process running it stopped before it finished"; export type SettleRunningOpts = WriterCheck & { paths: Paths; bootedAt: number; isLive?: (id: string) => boolean; log?: (line: string) => void; }; export type BootRunningResult = { interrupted: { id: string; kind: string; channelSlug?: string }[]; }; export async function settleRunningJobMetas( opts: SettleRunningOpts, ): Promise { const result: BootRunningResult = { interrupted: [] }; 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 !== "running") continue; if (typeof meta.queuedAt === "number" && meta.queuedAt >= opts.bootedAt) { continue; } if (opts.isLive?.(id)) continue; if (!writerIsGone(meta, opts)) continue; let lastAlive = Date.now(); try { lastAlive = (await stat(path.join(opts.paths.jobsDir, `${id}.log`))).mtimeMs; } catch { /* no log: this boot is the best bound there is */ } if (await closeMeta(opts.paths, meta, INTERRUPTED_REASON, Math.round(lastAlive))) { result.interrupted.push({ id, kind: meta.kind, ...(meta.channelSlug ? { channelSlug: meta.channelSlug } : {}), }); } } if (result.interrupted.length > 0) { const by = new Map(); for (const j of result.interrupted) by.set(j.kind, (by.get(j.kind) ?? 0) + 1); opts.log?.( `[boot] jobs a previous process left running, closed as interrupted: ${result.interrupted.length}` + ` (${[...by].map(([k, n]) => `${k} ${n}`).join(", ")})`, ); } return result; }