Archilyzer · Source

archilyzer

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

commit 35015500e43d1b37691e6edfe950efd9cb4d6418
parent 90ea0a8ca0a60dbfc2af889269df14de36b3d857
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Thu,  1 Oct 2026 16:51:47 -0400

common, editor: boot closes the running metas a dead process left behind

A job meta records its writer's pid. At boot, settleRunningJobMetas closes every
running meta from before this boot, not in the registry, whose writer is gone
(no pid, this pid, or no such process) as cancelled with an "interrupted"
reason, ended at its log's last write; a meta another live process owns (an
archilyzer run) is left alone, and the queued pass now applies the same check.
channelWriters reads the registry only, so such a ghost never held a move; a
test pins it.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

Diffstat:
Mcommon/controller/channelWriters.test.ts | 45+++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/channelWriters.ts | 8++++++++
Mcommon/jobs/bootQueuedJobs.test.ts | 106++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/jobs/bootQueuedJobs.ts | 145++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Mcommon/jobs/jobMeta.ts | 9++++++++-
Meditor/instrumentation.ts | 13++++++++++++-
6 files changed, 315 insertions(+), 11 deletions(-)

diff --git a/common/controller/channelWriters.test.ts b/common/controller/channelWriters.test.ts @@ -1,5 +1,10 @@ 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 { settleRunningJobMetas } from "../jobs/bootQueuedJobs"; import { getRegistry, newJobId, type JobRecord } from "../jobs/registry"; import type { AutoQueueKind } from "../lib/autoQueueTypes"; import type { AutoRunnerInFlight } from "./autoRunner"; @@ -214,3 +219,43 @@ test("through the live registry: cancel leaves it stopping, the job's end stamps assert.equal(typeof record.endedAt, "number"); assert.deepEqual(channelWriters(slug), []); }); + +// A GHOST NEVER HOLDS A MOVE (release 17 slice D0): a `running` meta a dead +// process left on disk is no writer — before the boot pass closes it, and +// after. +test("a running meta a dead process left on disk is never a writer", async () => { + const root = await mkdtemp(path.join(tmpdir(), "channel-writers-ghost-")); + try { + const jobsDir = path.join(root, ".jobs"); + await mkdir(jobsDir, { recursive: true }); + const id = newJobId(); + const slug = `ghost-${id}`; + await writeFile( + path.join(jobsDir, `${id}.meta.json`), + JSON.stringify({ + id, + kind: "refresh-report", + queueKey: "", + channelSlug: slug, + status: "running", + queuedAt: Date.now() - 60_000, + startedAt: Date.now() - 60_000, + pid: 2 ** 22 + 1, // past pid_max: no such process + }), + ); + assert.deepEqual(channelWriters(slug, { includeQueued: true }), []); + const res = await settleRunningJobMetas({ + paths: { jobsDir } as Paths, + bootedAt: Date.now(), + isLive: (j) => getRegistry().get(j) !== undefined, + }); + assert.deepEqual(res.interrupted.map((j) => j.id), [id]); + const meta = JSON.parse( + await readFile(path.join(jobsDir, `${id}.meta.json`), "utf8"), + ) as { status: string }; + assert.equal(meta.status, "cancelled"); + assert.deepEqual(channelWriters(slug, { includeQueued: true }), []); + } finally { + await rm(root, { recursive: true, force: true }); + } +}); diff --git a/common/controller/channelWriters.ts b/common/controller/channelWriters.ts @@ -21,6 +21,14 @@ import { getAutoRunnerStatus, type AutoRunnerInFlight } from "./autoRunner"; // enqueued it, and the refusal names the writer so the operator knows what to // wait for or cancel. // +// A GHOST NEVER HOLDS A MOVE (release 17 slice D0). Only THIS process's +// registry is read, never a `<id>.meta.json`: a meta a dead process left +// `running` (three `refresh-report`s on 2026-10-01) is not a record here and +// cannot name a writer. The boot pass closes such metas as interrupted +// (jobs/bootQueuedJobs.ts `settleRunningJobMetas`) so /jobs stops showing them; +// nothing here needs to know. Keep it that way: a reader of metas here would +// have to ask whether the writer is alive (`writerIsGone`) first. +// // TWO HALVES, for the reason editor/app/channels/lib/mediaBusy.ts gives: the // job registry is half the truth. The auto-queue lanes run their per-video // units in-process and make no job record (the omnimirror incident, diff --git a/common/jobs/bootQueuedJobs.test.ts b/common/jobs/bootQueuedJobs.test.ts @@ -1,16 +1,20 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { mkdtemp, mkdir, readFile, rm, writeFile } from "node:fs/promises"; +import { mkdtemp, mkdir, readFile, rm, utimes, 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 { + INTERRUPTED_REASON, STORAGE_PASS_WAIT_MS, + processIsAlive, settleAfterStoragePass, settleQueuedJobMetas, + settleRunningJobMetas, waitForStoragePass, + writerIsGone, type RequeueFn, } from "./bootQueuedJobs"; @@ -451,3 +455,103 @@ test("a storage pass that throws counts as finished: no timeout, the settle runs test("the production bound is 60 s", () => { assert.equal(STORAGE_PASS_WAIT_MS, 60_000); }); + +// THE RUNNING PASS (release 17 slice D0). A job running when its process died +// keeps `running` on disk forever; the next boot closes it `cancelled` with +// INTERRUPTED_REASON — unless the process that wrote it is still alive (an +// `archilyzer run` beside the editor), in which case it is left alone. + +const ALIVE_OTHER = 4242; // a live process that is not this one +const DEAD = 4343; +const SELF = 4444; +const check = { + selfPid: SELF, + isProcessAlive: (pid: number) => pid === ALIVE_OTHER || pid === SELF, +}; + +test("writerIsGone: no pid, this pid, a dead pid are gone; another live pid is not", () => { + const m = (pid?: number) => + ({ id: "X", kind: "k", queueKey: "", status: "running", queuedAt: 0, ...(pid === undefined ? {} : { pid }) }) as JobMeta; + assert.equal(writerIsGone(m(), check), true, "a pre-release-17 meta"); + assert.equal(writerIsGone(m(SELF), check), true, "a previous process with this pid"); + assert.equal(writerIsGone(m(DEAD), check), true); + assert.equal(writerIsGone(m(ALIVE_OTHER), check), false); +}); + +test("processIsAlive: this process is; a pid past pid_max is not", () => { + assert.equal(processIsAlive(process.pid), true); + assert.equal(processIsAlive(2 ** 22 + 1), false); +}); + +test("a running meta from a dead process is closed as interrupted, ended at its log's last write", async () => { + const f = await fixture([ + { id: "R1", kind: "refresh-report", status: "running", channelSlug: "the-quartering-rumble", pid: DEAD, startedAt: BOOT - 3 * HOUR }, + { id: "R2", kind: "refresh-report", status: "running", channelSlug: "the-quartering" }, // no pid: pre-release-17 + { id: "R3", kind: "whisper-all", status: "running", channelSlug: "x", pid: ALIVE_OTHER }, // an `archilyzer run` + { id: "R4", kind: "whisper-all", status: "running", queuedAt: BOOT + 1, pid: SELF }, // this process's own + { id: "R5", kind: "whisper-all", status: "running", pid: SELF }, // live in the registry + { id: "D1", kind: "whisper-all", status: "done", pid: DEAD }, + ]); + try { + const logMtime = BOOT - 2 * HOUR; + const logFile = path.join(f.paths.jobsDir, "R1.log"); + await writeFile(logFile, "Regenerating report for the-quartering-rumble…\n"); + await utimes(logFile, logMtime / 1000, logMtime / 1000); + const lines: string[] = []; + const res = await settleRunningJobMetas({ + paths: f.paths, + bootedAt: BOOT, + isLive: (id) => id === "R5", + log: (l) => lines.push(l), + ...check, + }); + assert.deepEqual( + res.interrupted.map((j) => j.id), + ["R1", "R2"], + ); + const r1 = await f.read("R1"); + assert.equal(r1.status, "cancelled"); + assert.equal(r1.cancelReason, INTERRUPTED_REASON); + assert.equal(r1.endedAt, logMtime); + assert.equal(r1.pid, DEAD, "the rest of the meta is kept"); + assert.match(await f.log("R1"), /\[boot\] interrupted/); + const r2 = await f.read("R2"); + assert.equal(r2.status, "cancelled"); + assert.ok(typeof r2.endedAt === "number", "no log: ended at this boot"); + for (const id of ["R3", "R4", "R5"]) assert.equal((await f.read(id)).status, "running", id); + assert.equal((await f.read("D1")).status, "done"); + assert.equal(lines.length, 1); + assert.match(lines[0], /closed as interrupted: 2 \(refresh-report 2\)/); + // Idempotent: the next boot finds nothing. + const again = await settleRunningJobMetas({ + paths: f.paths, + bootedAt: BOOT, + isLive: (id) => id === "R5", + ...check, + }); + assert.deepEqual(again.interrupted, []); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("the queued pass leaves a queued meta whose writer is still alive", async () => { + const f = await fixture([ + { id: "Q1", spec: SPEC, pid: ALIVE_OTHER }, + { id: "Q2", spec: SPEC, pid: DEAD, kind: "sync" }, + ]); + try { + const rq = recordingRequeue(); + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: rq.fn, + bootedAt: BOOT, + ...check, + }); + assert.deepEqual(rq.calls, []); + assert.deepEqual(res.cancelled.map((c) => c.id), ["Q2"]); + assert.equal((await f.read("Q1")).status, "queued"); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); diff --git a/common/jobs/bootQueuedJobs.ts b/common/jobs/bootQueuedJobs.ts @@ -1,5 +1,5 @@ import path from "node:path"; -import { appendFile, readdir } from "node:fs/promises"; +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"; @@ -31,11 +31,14 @@ import type { JobSpec } from "./jobSpec"; // 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. A malformed meta (readJobMeta -// → null) is skipped. +// 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. @@ -80,7 +83,48 @@ export function specKey(spec: JobSpec): string { }); } -export type SettleQueuedOpts = { +// 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; @@ -206,6 +250,7 @@ export async function settleQueuedJobMetas( continue; } if (opts.isLive?.(id)) continue; + if (!writerIsGone(meta, opts)) continue; stale.push(meta); } @@ -308,11 +353,12 @@ async function closeMeta( paths: Paths, meta: JobMeta, reason: string, + endedAt: number = Date.now(), ): Promise<boolean> { const closed: JobMeta = { ...meta, status: "cancelled", - endedAt: Date.now(), + endedAt, cancelReason: reason, }; try { @@ -333,3 +379,86 @@ async function closeMeta( } 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<BootRunningResult> { + 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<string, number>(); + 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; +} diff --git a/common/jobs/jobMeta.ts b/common/jobs/jobMeta.ts @@ -27,8 +27,14 @@ export type JobMeta = { 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. + // when the server went down, or still `running` when the process that ran + // it died. Absent on every other meta. cancelReason?: string; + // The process that wrote this meta (`process.pid`). The boot pass reads it + // to tell a job a dead process left `running` or `queued` from one another + // LIVE process owns right now — `archilyzer run` writes into the same + // `.jobs/`. Absent on metas written before release 17. + pid?: number; }; export function metaPath(paths: Paths, id: string): string { @@ -55,6 +61,7 @@ export async function writeJobMeta( endedAt: record.endedAt, exitCode: record.exitCode, spec: record.spec, + pid: process.pid, }; await writeFile(metaPath(paths, record.id), JSON.stringify(meta), "utf8"); } catch { diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts @@ -146,7 +146,7 @@ export async function register() { // common/jobs/bootQueuedJobs.ts. Lazy, voided, best-effort: never blocks // readiness. try { - const { settleAfterStoragePass } = await import( + const { settleAfterStoragePass, settleRunningJobMetas } = await import( "yt-dlp-transcript-common/jobs/bootQueuedJobs" ); const { getPaths } = await import("yt-dlp-transcript-common/lib/paths"); @@ -154,6 +154,17 @@ export async function register() { "yt-dlp-transcript-common/jobs/registry" ); const testServer = process.env.E2E_TEST_ROUTES === "1"; + // STALE `running` METAS FROM A DEAD PROCESS (release 17 slice D0): closed + // `cancelled` as interrupted, never re-run — on every boot, idle and test + // server included, because closing one starts nothing. Does not wait for + // the storage pass: it touches no channel. A meta another live process + // still owns (`archilyzer run`) is left alone — see writerIsGone. + void settleRunningJobMetas({ + paths: getPaths(), + bootedAt, + isLive: (id) => getRegistry().get(id) !== undefined, + log: (line) => console.log(line), + }).catch(() => {}); const cancelOnly = idle || testServer; void settleAfterStoragePass(storagePass, { paths: getPaths(),