Archilyzer · Source

archilyzer

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

commit 1aac49840c13bbb31e960ac20198be30af7067de
parent 3ead8acf930c2262f53336ee10172b43c0387507
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Thu,  1 Oct 2026 23:54:55 -0400

Merge r17/dashboard-answers (release 17 slice D0) — the dashboard answers while a report regenerates: the snapshot walk yields between chunks and regenerations run one at a time with one queued successor, the auto-queue status poll shares one memo keyed on the settings it reads, boot closes the running job metas a dead process left behind, Refresh report waits at most 15 s then names its job and queue position, the ops forms return job ids at once; reviewed SHIP

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

Diffstat:
Mcommon/controller/channelSnapshot.ts | 42+++++++++++++++++++++++++++++++++++++++---
Mcommon/controller/channelWriters.test.ts | 45+++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/channelWriters.ts | 8++++++++
Acommon/controller/snapshotYield.test.ts | 39+++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/bootQueuedJobs.test.ts | 106++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/jobs/bootQueuedJobs.ts | 145++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Mcommon/jobs/jobMeta.ts | 9++++++++-
Acommon/jobs/snapshotScheduler.test.ts | 185+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/snapshotScheduler.ts | 263+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mcommon/lib/queueKeys.ts | 4++++
Mcommon/views/autoQueueStatus.test.ts | 100+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/views/autoQueueStatus.ts | 94+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/CHANGELOG.md | 3+++
Meditor/app/api/ops/refresh-report/route.ts | 40++++++++++++++++++++++++++++------------
Meditor/app/api/test/invalidate-cache/route.ts | 8++++++++
Aeditor/app/api/test/settle-running-metas/route.ts | 26++++++++++++++++++++++++++
Meditor/app/channels/[slug]/components/RefreshSnapshotButton.tsx | 61++++++++++++++++++++++++++++++++++++++++++++++---------------
Meditor/app/channels/actions.ts | 131+++++++++++++++++++++++++++++++------------------------------------------------
Meditor/app/components/actions/InlineActionButton.tsx | 7++++++-
Meditor/app/lib/requestCache.ts | 3+++
Meditor/app/operations/[id]/page.tsx | 7++++---
Meditor/app/operations/status.ts | 44+++++++++++++++++++++++++++++++++++++++-----
Meditor/e2e/bulk-actions.spec.ts | 8+++++++-
Aeditor/e2e/dashboard-answers.spec.ts | 356+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/e2e/ops-api.spec.ts | 35+++++++++++++++++++++++++++++------
Meditor/instrumentation.ts | 13++++++++++++-
Mplans/release-17.md | 220+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
27 files changed, 1819 insertions(+), 183 deletions(-)

diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts @@ -674,6 +674,41 @@ export { SNAPSHOT_FILENAME, snapshotPath, readChannelSnapshot }; const SNAPSHOT_VIDEO_CONCURRENCY = 16; +// THE WALK YIELDS TO THE EVENT LOOP every SNAPSHOT_YIELD_EVERY videos. +// +// It runs in the editor's own process, on the main thread, and a video's unit +// parses its cues.json and (in the reconcile pass and for a non-YouTube +// archive) its metadata.info.json — ~0.6 MB each on a long VOD, a few ms of +// JSON.parse apiece. Sixteen units in flight keep every turn of the loop busy +// with that work; on 2026-10-01 two regenerations of 2,000–3,000-video +// channels ran for over an hour while `/` and `/jobs` did not answer. So the +// ids are walked in chunks, each fanned out under the same limit, and between +// chunks the walk waits one `setImmediate`: the loop gets a turn with no +// snapshot work queued in it, and every request whose I/O completed meanwhile +// runs before the next chunk starts. +// +// Two full waves of the limit (32 at 16 wide): a chunk is ~0.1–0.2 s of +// parsing on such a channel, and a multiple of the width keeps every wave full +// — 25 left the second wave nine wide, a fifth of the walk's throughput. +// Measured in the slice D0 record (plans/release-17.md). +const SNAPSHOT_YIELD_EVERY = 2 * SNAPSHOT_VIDEO_CONCURRENCY; + +// `fn` over `items` in chunks of `size`, in order, one setImmediate between +// chunks. The concurrency inside a chunk is whatever `fn` imposes. Exported for +// snapshotYield.test.ts, which pins the yield. +export async function mapInYieldingChunks<T, R>( + items: readonly T[], + size: number, + fn: (item: T) => Promise<R>, +): Promise<R[]> { + const out: R[] = []; + for (let i = 0; i < items.length; i += size) { + if (i > 0) await new Promise<void>((resolve) => setImmediate(resolve)); + out.push(...(await Promise.all(items.slice(i, i + size).map(fn)))); + } + return out; +} + // ONE LEVEL of a video dir's `clips/` — the files in it, and nothing deeper. // Deliberately not recursive: `clipWindow-server.ts` writes `<from>-<to>.<ext>` // and `<from>-<to>.json` flat into it and nothing else does, so a recursion @@ -839,8 +874,10 @@ export async function generateChannelSnapshot( } const limit = pLimit(SNAPSHOT_VIDEO_CONCURRENCY); - const perVideo = await Promise.all( - videoDirNames.map((id) => + const perVideo = await mapInYieldingChunks( + videoDirNames, + SNAPSHOT_YIELD_EVERY, + (id) => // One video directory's reads are one unit through the watchdog. limit(() => through(async () => { const dir = path.join(dataDir, id); @@ -974,7 +1011,6 @@ export async function generateChannelSnapshot( digest, }; })), - ), ); const filesById = new Map<string, VideoFiles>(); 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/controller/snapshotYield.test.ts b/common/controller/snapshotYield.test.ts @@ -0,0 +1,39 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mapInYieldingChunks } from "./channelSnapshot"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test controller/snapshotYield.test.ts +// +// THE WALK YIELDS BETWEEN CHUNKS (release 17 slice D0). A macrotask queued +// while the first chunk runs — a request's I/O callback, here a setImmediate +// probe — runs before the second chunk's first unit starts. Without the yield +// every unit below would finish in microtasks before the probe ever ran. + +test("a macrotask queued during chunk 1 runs before chunk 2 starts", async () => { + const events: string[] = []; + const items = Array.from({ length: 64 }, (_, i) => i); + const out = await mapInYieldingChunks(items, 32, async (i) => { + if (i === 0) setImmediate(() => events.push("probe")); + events.push(`unit ${i}`); + await Promise.resolve(); + return i * 2; + }); + assert.deepEqual(out, items.map((i) => i * 2), "results in order"); + const probe = events.indexOf("probe"); + assert.ok(probe > events.indexOf("unit 31"), "after the first chunk"); + assert.ok(probe < events.indexOf("unit 32"), "before the second chunk"); +}); + +test("one chunk, no yield; an empty list, no call", async () => { + let calls = 0; + assert.deepEqual( + await mapInYieldingChunks([1, 2, 3], 32, async (i) => { + calls++; + return i; + }), + [1, 2, 3], + ); + assert.equal(calls, 3); + assert.deepEqual(await mapInYieldingChunks([], 32, async (i: number) => i), []); +}); 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/common/jobs/snapshotScheduler.test.ts b/common/jobs/snapshotScheduler.test.ts @@ -0,0 +1,185 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/snapshotScheduler.test.ts +// +// ONE REGENERATION AT A TIME, ONCE PER CHANNEL (release 17 slice D0). Every +// refresh-report used to run on the empty queue key — all at once — and two +// passes a moment apart could both enqueue the same channel. Temp corpus only. + +const root = await mkdtemp(path.join(tmpdir(), "snapshot-scheduler-")); +process.env.TRANSCRIPTS_DIR = path.join(root, "transcripts"); +process.env.SETTINGS_FILE = path.join(root, "settings.json"); +await writeFile(process.env.SETTINGS_FILE, "{}"); +for (const slug of ["alpha", "beta", "gamma"]) { + const dir = path.join(process.env.TRANSCRIPTS_DIR, "channels", slug); + await mkdir(path.join(dir, "data"), { recursive: true }); + await writeFile(path.join(dir, "config.json"), JSON.stringify({ handling: "youtube" })); +} +// A channel relocated to a drive that is not there: its walk refuses (guard 3) +// with a sentence naming the unreachable media. +{ + const dir = path.join(process.env.TRANSCRIPTS_DIR, "channels", "unmounted"); + await mkdir(dir, { recursive: true }); + await writeFile( + path.join(dir, "config.json"), + JSON.stringify({ handling: "youtube", dataDir: path.join(root, "no-such-drive", "unmounted", "data") }), + ); +} + +const { getPaths } = await import("../lib/paths"); +const { getRegistry, newJobId } = await import("./registry"); +const { + REFRESH_REPORT_ACTIVE, + REFRESH_REPORT_QUEUE, + isRefreshReportPending, + refreshReportWaitNotice, + requestRefreshReport, + startRefreshReport, + waitForRefreshReport, +} = await import("./snapshotScheduler"); + +// A refresh-report that is RUNNING until released: a record on the queue whose +// start does nothing. +function holdQueue(slug: string): () => void { + const registry = getRegistry(); + const id = newJobId(); + const holder = { + id, + kind: "refresh-report", + queueKey: REFRESH_REPORT_QUEUE, + channelSlug: slug, + status: "queued" as const, + queuedAt: Date.now(), + logPath: "/dev/null", + }; + registry.register(holder); + registry.enqueue(holder, { start: () => {}, onCancel: () => {} }); + return () => registry.finalize(id, "done"); +} + +test.after(async () => { + await rm(root, { recursive: true, force: true }); +}); + +test("two requests for one channel at once enqueue one regeneration", async () => { + const paths = getPaths(); + const [a, b] = await Promise.all([ + startRefreshReport(paths, "alpha"), + startRefreshReport(paths, "alpha"), + ]); + const started = [a, b].filter((r) => r.ok); + const refused = [a, b].filter((r) => !r.ok); + assert.equal(started.length, 1, "exactly one starts"); + assert.equal(refused.length, 1); + const r = refused[0]; + assert.ok(!r.ok && r.info === true && r.error === REFRESH_REPORT_ACTIVE); + const s = started[0]; + assert.ok(s.ok); + assert.equal((await s.done).status, "done"); + assert.equal(isRefreshReportPending("alpha"), false); + // Once it has finished, a new request starts a new one. + const again = await startRefreshReport(paths, "alpha"); + assert.ok(again.ok); + await again.done; +}); + +test("every regeneration runs on the one serial refresh-report queue", async () => { + const paths = getPaths(); + const a = await startRefreshReport(paths, "alpha"); + const b = await startRefreshReport(paths, "beta"); + assert.ok(a.ok && b.ok); + for (const id of [a.jobId, b.jobId]) { + assert.equal(getRegistry().get(id)?.queueKey, REFRESH_REPORT_QUEUE); + } + // Serial: the second is behind the first unless the first already finished. + const first = getRegistry().get(a.jobId); + const second = getRegistry().get(b.jobId); + if (first?.status === "running") assert.equal(second?.status, "queued"); + assert.equal((await a.done).status, "done"); + assert.equal((await b.done).status, "done"); +}); + +test("a running regeneration gets one queued successor, and no second", async () => { + const paths = getPaths(); + // A regeneration of beta that is RUNNING and stays so until released: a + // record on the refresh-report queue whose start does nothing. + const registry = getRegistry(); + const holderId = newJobId(); + const holder = { + id: holderId, + kind: "refresh-report", + queueKey: REFRESH_REPORT_QUEUE, + channelSlug: "beta", + status: "queued" as const, + queuedAt: Date.now(), + logPath: "/dev/null", + }; + registry.register(holder); + registry.enqueue(holder, { start: () => {}, onCancel: () => {} }); + assert.equal(registry.get(holderId)?.status, "running"); + + // A change during the walk: it may already have passed it, so one successor + // is queued behind it… + const successor = await startRefreshReport(paths, "beta"); + assert.ok(successor.ok, "queued behind the running one, not dropped"); + assert.equal(registry.get(successor.jobId)?.status, "queued"); + assert.equal(isRefreshReportPending("beta"), true); + // …and a further change finds it queued: it has not started reading yet. + const third = await startRefreshReport(paths, "beta"); + assert.ok(!third.ok && third.info === true && third.error === REFRESH_REPORT_ACTIVE); + + registry.finalize(holderId, "done"); + assert.equal((await successor.done).status, "done"); +}); + +// A PERSON'S REFRESH WAITS A BOUNDED TIME (re-review R1, R2). + +test("a refresh behind the queue answers with where it is, once the wait runs out", async () => { + const paths = getPaths(); + const release = holdQueue("beta"); + try { + const asked = await requestRefreshReport(paths, "gamma"); + assert.ok(asked.ok && asked.started); + const wait = await waitForRefreshReport(paths, asked.jobId, { timeoutMs: 300, pollMs: 50 }); + assert.equal(wait.state, "waiting"); + assert.ok(wait.state === "waiting" && wait.ahead === 1, "one regeneration ahead of it"); + assert.equal( + wait.state === "waiting" && refreshReportWaitNotice(wait), + `Queued behind 1 report regeneration — the report updates when it finishes (job ${asked.jobId}).`, + ); + // A second person's click finds the same job, queued. + const again = await requestRefreshReport(paths, "gamma"); + assert.deepEqual(again, { ok: true, jobId: asked.jobId, started: false }); + release(); + const finished = await waitForRefreshReport(paths, asked.jobId, { pollMs: 20 }); + assert.deepEqual(finished, { state: "done", jobId: asked.jobId }); + } finally { + release(); + } +}); + +test("a failed walk is a failure with its sentence, for the click that found it queued too", async () => { + const paths = getPaths(); + const release = holdQueue("beta"); + let queuedId = ""; + try { + const first = await requestRefreshReport(paths, "unmounted"); + assert.ok(first.ok && first.started); + queuedId = first.jobId; + // Another click while it waits: it did not start this job, and still + // hears how it ended. + const second = await requestRefreshReport(paths, "unmounted"); + assert.deepEqual(second, { ok: true, jobId: queuedId, started: false }); + release(); + const wait = await waitForRefreshReport(paths, queuedId, { pollMs: 20 }); + assert.equal(wait.state, "failed"); + assert.ok(wait.state === "failed" && /unmounted/.test(wait.error), "the walk's own sentence"); + } finally { + release(); + } +}); diff --git a/common/jobs/snapshotScheduler.ts b/common/jobs/snapshotScheduler.ts @@ -1,3 +1,5 @@ +import path from "node:path"; +import { readFile } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import { DEFAULT_REPORT_DEBOUNCE_PRESET, @@ -10,6 +12,9 @@ import { } from "../controller/channelSnapshot"; import { getRegistry } from "./registry"; import { drainStream } from "./drainStream"; +import type { StreamActionResult } from "./streamCommand"; +import { REFRESH_REPORT_QUEUE } from "../lib/queueKeys"; +import { readJobMeta } from "./jobMeta"; // Global, debounced "channel report" (snapshot) regeneration scheduler. // @@ -93,6 +98,13 @@ type SchedulerState = { // counts (/channels, the dashboard) would sit stale until the next unrelated // change. One integer closes that gap at zero cost. generation: number; + // Slugs whose refresh-report is being enqueued right now. Between the + // registry check and the record's registration runManagedFunction awaits + // (the media guard, the jobs dir), so two passes a moment apart both found + // nothing queued and both enqueued — two regenerations of the same + // 3,000-video channel, seen live on 2026-10-01. Held until the record is in + // the registry, which answers from then on. + starting: Set<string>; }; declare global { @@ -107,6 +119,7 @@ function getState(): SchedulerState { timer: null, firstDirtyAt: null, generation: 0, + starting: new Set(), }; } return globalThis.__yttSnapshotScheduler__; @@ -146,11 +159,210 @@ export function requestChannelSnapshot(paths: Paths, slug: string): void { state.timer.unref?.(); } +// THE ONE QUEUE EVERY REFRESH-REPORT RUNS ON (lib/queueKeys.ts): one +// regeneration at a time, corpus-wide. +// +// It was the empty queueKey — "a local filesystem scan, so it runs in +// parallel" — and every dirty channel's walk ran at once in the editor's own +// process. On 2026-10-01 two of them (2,000 and 3,260 videos) ran side by side +// for over an hour and `/` and `/jobs` did not answer. A non-empty key is how +// the registry serializes (registry.ts `enqueue`: concurrency 1 per key), so +// nothing new was needed; a report that waits its turn is a report a minute +// late, and the pages keep answering meanwhile. +export { REFRESH_REPORT_QUEUE }; + +// Why startRefreshReport started nothing for a slug that already has one. +export const REFRESH_REPORT_ACTIVE = "already queued"; + +// The QUEUED regeneration of `slug` (not yet started), if there is one. +export function queuedRefreshReportId(slug: string): string | null { + const queued = getRegistry() + .list() + .find( + (j) => + j.kind === "refresh-report" && + j.channelSlug === slug && + j.status === "queued", + ); + return queued?.id ?? null; +} + +// Is a regeneration of `slug` queued or being enqueued — one that has not +// started reading yet, and so will see any change made before it does? +// +// A RUNNING one does not count. It may already have walked past the change +// that asks for this regeneration, so one successor is queued behind it (the +// queue is serial, so it starts when the running one ends). At most one: a +// second request finds the successor queued. +export function isRefreshReportPending(slug: string): boolean { + if (getState().starting.has(slug)) return true; + return queuedRefreshReportId(slug) !== null; +} + +// Enqueue one channel's report regeneration on REFRESH_REPORT_QUEUE — unless +// one is already queued or being enqueued for it, which is answered `info` +// with REFRESH_REPORT_ACTIVE (it has not started reading, so it will see the +// change). Every walk goes through here — the debounced pass, "Update all +// reports" and a channel's own Refresh report — so none runs outside the +// queue, and they dedup against each other. `onError` hears the walk's error +// sentence (an unmounted drive's) when the job fails. +export async function startRefreshReport( + paths: Paths, + slug: string, +): Promise<StreamActionResult> { + if (isRefreshReportPending(slug)) { + return { ok: false, error: REFRESH_REPORT_ACTIVE, info: true }; + } + const state = getState(); + state.starting.add(slug); + try { + const { runManagedFunction } = await import("./streamCommand"); + return await runManagedFunction({ + kind: "refresh-report", + queueKey: REFRESH_REPORT_QUEUE, + paths, + channelSlug: slug, + fn: async (onLog) => { + onLog(`Regenerating report for ${slug}…`); + const snap = await generateChannelSnapshot(paths, slug); + getState().generation++; + const excluded = excludedDownloadIdSet(snap); + const awaitingTranscription = excluded.size + ? snap.buckets.downloadedNoTranscript.filter( + (id) => !excluded.has(id), + ).length + : snap.buckets.downloadedNoTranscript.length; + onLog( + `Done. ${snap.totals.videos} videos · ` + + `${snap.undownloadedIds.length} undownloaded · ` + + `${awaitingTranscription} awaiting transcription.`, + ); + // No revalidatePath here — the caller revalidates once, after its + // streams drain. + }, + }); + } finally { + // The record is registered by now (or nothing was made): the registry + // answers for this slug from here on. A reset (resetSnapshotScheduler) + // may have swapped the state meanwhile; deleting from the old set is + // harmless. + state.starting.delete(slug); + } +} + +// ONE CHANNEL'S REPORT, ASKED FOR BY A PERSON OR A SCRIPT: the job that will +// regenerate it — started now, or the one already queued for the channel +// (which has not started reading, so it is as fresh). A caller mid-enqueue for +// the same slug has no id yet; it is waited for, briefly. The stream is +// cancelled: nothing reads it (the log is on disk). +export async function requestRefreshReport( + paths: Paths, + slug: string, +): Promise<{ ok: true; jobId: string; started: boolean } | { ok: false; error: string }> { + const started = await startRefreshReport(paths, slug); + if (started.ok) { + void started.stream.cancel(); + return { ok: true, jobId: started.jobId, started: true }; + } + if (!started.info) return { ok: false, error: started.error }; + for (let i = 0; i < 40; i++) { + const queued = queuedRefreshReportId(slug); + if (queued) return { ok: true, jobId: queued, started: false }; + if (!getState().starting.has(slug)) break; + await new Promise((resolve) => setTimeout(resolve, 50)); + } + // The one we would have waited for started in the meantime (or another + // caller's enqueue failed): ask once more — now nothing is queued, so this + // queues one, or says why it cannot. + const again = await startRefreshReport(paths, slug); + if (again.ok) { + void again.stream.cancel(); + return { ok: true, jobId: again.jobId, started: true }; + } + const queued = queuedRefreshReportId(slug); + return queued + ? { ok: true, jobId: queued, started: false } + : { ok: false, error: again.error }; +} + +// HOW LONG A PERSON'S "Refresh report" WAITS before it answers with where its +// job is instead. The queue is serial: behind a 3,000-video walk, or during +// Update all reports, the report may be minutes away, and a button that says +// "Refreshing…" for minutes says nothing. +export const REFRESH_REPORT_WAIT_MS = 15_000; + +export type RefreshReportWait = + | { state: "done"; jobId: string } + | { state: "failed"; jobId: string; error: string } + // Still queued or running when the wait ran out. `ahead`: the jobs before it + // on the refresh-report queue — 0 means it is the one regenerating now. + | { state: "waiting"; jobId: string; ahead: number }; + +// The `[error]` line a failed job's log ends with (streamCommand writes the +// thrown sentence there — an unmounted drive's, for a walk), else null. +async function lastErrorLine(paths: Paths, jobId: string): Promise<string | null> { + try { + const raw = await readFile(path.join(paths.jobsDir, `${jobId}.log`), "utf8"); + const lines = raw.split("\n").filter((l) => l.startsWith("[error] ")); + const last = lines[lines.length - 1]; + return last ? last.slice("[error] ".length) : null; + } catch { + return null; + } +} + +// Wait up to `timeoutMs` for a refresh-report job to end, and say how it ended +// — or where it is. Polls the registry (it hands out no completion promise for +// a job someone else started); a job evicted from it is read from its meta. +export async function waitForRefreshReport( + paths: Paths, + jobId: string, + opts: { timeoutMs?: number; pollMs?: number } = {}, +): Promise<RefreshReportWait> { + const timeoutMs = opts.timeoutMs ?? REFRESH_REPORT_WAIT_MS; + const pollMs = opts.pollMs ?? 250; + const until = Date.now() + timeoutMs; + for (;;) { + const job = getRegistry().get(jobId); + let status = job?.status; + if (!job) { + status = (await readJobMeta(paths, jobId))?.status ?? "failed"; + } + if (status === "done") return { state: "done", jobId }; + if (status !== "queued" && status !== "running") { + const error = + (await lastErrorLine(paths, jobId)) ?? + `Refresh report ${status} (job ${jobId})`; + return { state: "failed", jobId, error }; + } + if (Date.now() >= until) { + return { + state: "waiting", + jobId, + ahead: Math.max(0, getRegistry().positionInQueue(jobId)), + }; + } + await new Promise((resolve) => setTimeout(resolve, pollMs)); + } +} + +// The operator's sentence for a wait that ran out. +export function refreshReportWaitNotice(w: Extract<RefreshReportWait, { state: "waiting" }>): string { + const where = + w.ahead === 0 + ? "Regenerating now" + : `Queued behind ${w.ahead} report regeneration${w.ahead === 1 ? "" : "s"}`; + return `${where} — the report updates when it finishes (job ${w.jobId}).`; +} + async function fire(): Promise<void> { const state = getState(); // Snapshot and clear before the async work: any action that fires during // regeneration re-marks the slug dirty and schedules a fresh pass (trailing - // edge), instead of being swallowed by this in-flight batch. + // edge), instead of being swallowed by this in-flight batch. That pass is + // not dropped when the slug's regeneration is RUNNING: startRefreshReport + // queues one successor behind it (it is skipped only when one is already + // queued, which has not started reading and so will see the change). const batch = [...state.dirty.entries()]; state.dirty.clear(); state.timer = null; @@ -158,61 +370,18 @@ async function fire(): Promise<void> { if (batch.length === 0) return; try { - const { runManagedFunction } = await import("./streamCommand"); - - // Skip channels that already have a refresh-report job queued/running — - // a fresh snapshot is already on its way (mirrors - // refreshAllChannelSnapshotsAction's dedup). - const active = new Set( - getRegistry() - .list() - .filter( - (j) => - j.kind === "refresh-report" && - (j.status === "queued" || j.status === "running") && - j.channelSlug, - ) - .map((j) => j.channelSlug as string), - ); - const regenerated: string[] = []; const streams: ReadableStream<string>[] = []; for (const [slug, paths] of batch) { - if (active.has(slug)) continue; - const result = await runManagedFunction({ - kind: "refresh-report", - // Empty queueKey: bypass queue serialization. Snapshot regen is a local - // filesystem scan, so it runs in parallel rather than waiting behind - // sync/download work (see registry.ts enqueue). - queueKey: "", - paths, - channelSlug: slug, - fn: async (onLog) => { - onLog(`Regenerating report for ${slug}…`); - const snap = await generateChannelSnapshot(paths, slug); - getState().generation++; - const excluded = excludedDownloadIdSet(snap); - const awaitingTranscription = excluded.size - ? snap.buckets.downloadedNoTranscript.filter( - (id) => !excluded.has(id), - ).length - : snap.buckets.downloadedNoTranscript.length; - onLog( - `Done. ${snap.totals.videos} videos · ` + - `${snap.undownloadedIds.length} undownloaded · ` + - `${awaitingTranscription} awaiting transcription.`, - ); - // No revalidatePath here — it's done once after all drains below. - }, - }); + const result = await startRefreshReport(paths, slug); if (!result.ok) continue; regenerated.push(slug); streams.push(result.stream); } // Wait for every snapshot to finish writing before revalidating so the - // re-rendered pages read fresh counts. queueKey "" runs them in parallel, - // so this is ~the slowest snapshot, not the sum. + // re-rendered pages read fresh counts. They run one at a time on + // REFRESH_REPORT_QUEUE, so this is the sum of them. await Promise.all(streams.map(drainStream)); if (regenerated.length > 0) { diff --git a/common/lib/queueKeys.ts b/common/lib/queueKeys.ts @@ -30,6 +30,10 @@ export const DIGEST_REMOTE_QUEUE = "digest:remote"; // is a separate question, answered by backfillLimit() rather than by the queue. export const BACKFILL_QUEUE = "backfill"; +// Every channel report regeneration (`refresh-report`), one at a time +// corpus-wide — jobs/snapshotScheduler.ts says why (release 17 slice D0). +export const REFRESH_REPORT_QUEUE = "refresh-report"; + // Per-channel queue for channel-local bookkeeping jobs (clean/clear/verify). export function channelQueueKey(slug: string): string { return `channel:${slug}`; diff --git a/common/views/autoQueueStatus.test.ts b/common/views/autoQueueStatus.test.ts @@ -8,7 +8,9 @@ import type { AutoRunnerStatus, LeafPending } from "../controller/autoRunner"; import { buildAutoQueueLanes } from "./autoQueueLanes"; import type { PriorityView } from "./channelPriority"; import { + AUTO_QUEUE_STATUS_MEMO_MS, buildAutoQueueStatusPayload, + singleFlightMemo, type AutoQueueStatusInputs, } from "./autoQueueStatus"; @@ -203,3 +205,101 @@ test("the pending fold is per lane and reaches the focus banner", () => { // A lane with nothing pending is still an entry, with the same shape. assert.deepEqual(payload.download.pendingByLeaf, {}); }); + +// THE POLL'S MEMO (slice D0, release 17): N pollers of the 3 s poll cost one +// fold of the snapshots, concurrent callers share the one in flight, and time +// is the only thing that expires it. + +function deferred<T>() { + let resolve!: (v: T) => void; + let reject!: (e: unknown) => void; + const promise = new Promise<T>((res, rej) => { + resolve = res; + reject = rej; + }); + return { promise, resolve, reject }; +} + +test("memo: concurrent callers share the computation in flight", async () => { + const memo = singleFlightMemo<number>({ now: () => 0 }); + const d = deferred<number>(); + let calls = 0; + const compute = () => { + calls++; + return d.promise; + }; + const all = Promise.all([memo.get("k", compute), memo.get("k", compute), memo.get("k", compute)]); + d.resolve(42); + assert.deepEqual(await all, [42, 42, 42]); + assert.equal(calls, 1); +}); + +test("memo: a landed value is reused for the window, then recomputed", async () => { + let t = 1_000; + const memo = singleFlightMemo<number>({ now: () => t }); + let calls = 0; + const compute = async () => ++calls; + assert.equal(await memo.get("k", compute), 1); + t += AUTO_QUEUE_STATUS_MEMO_MS - 1; + assert.equal(await memo.get("k", compute), 1, "inside the window: the memo"); + t += 1; + assert.equal(await memo.get("k", compute), 2, "at the window's end: computed again"); + assert.equal(calls, 2); +}); + +test("memo: the window runs from when the value LANDED, not when it started", async () => { + let t = 0; + const memo = singleFlightMemo<number>({ ttlMs: 3_000, now: () => t }); + const d = deferred<number>(); + const first = memo.get("k", () => d.promise); + t = 10_000; // a slow fold: ten seconds + d.resolve(7); + assert.equal(await first, 7); + t = 12_000; + assert.equal(await memo.get("k", async () => 8), 7); +}); + +test("memo: a rejection is shared by its waiters and never memoized", async () => { + const memo = singleFlightMemo<number>({ now: () => 0 }); + const d = deferred<number>(); + const a = memo.get("k", () => d.promise); + const b = memo.get("k", async () => 99); + d.reject(new Error("drive not answering")); + await assert.rejects(a, /drive not answering/); + await assert.rejects(b, /drive not answering/); + assert.equal(await memo.get("k", async () => 5), 5); +}); + +test("memo: clear() drops the value and detaches the computation in flight", async () => { + const memo = singleFlightMemo<string>({ now: () => 0 }); + assert.equal(await memo.get("k", async () => "old fixture"), "old fixture"); + memo.clear(); + const slow = deferred<string>(); + const before = memo.get("k", () => slow.promise); + memo.clear(); + // A caller after the clear does not join the detached computation… + assert.equal(await memo.get("k", async () => "new fixture"), "new fixture"); + // …and when it lands, it neither answers the memo nor evicts the new value. + slow.resolve("stale"); + assert.equal(await before, "stale"); + assert.equal(await memo.get("k", async () => "unused"), "new fixture"); +}); + +test("memo: a different key misses the memo and the computation in flight", async () => { + const memo = singleFlightMemo<string>({ now: () => 0 }); + assert.equal(await memo.get("tree A", async () => "counts under A"), "counts under A"); + // A rule added, a focus set: the settings the fold reads changed. + assert.equal(await memo.get("tree B", async () => "counts under B"), "counts under B"); + assert.equal(await memo.get("tree B", async () => "unused"), "counts under B"); + // Asking under A again is a miss too — the stored value is B's. + assert.equal(await memo.get("tree A", async () => "A again"), "A again"); + + // In flight: a caller with another key does not join it, and the older + // computation landing late does not overwrite the newer key's value. + const slowA = deferred<string>(); + const a = memo.get("tree A2", () => slowA.promise); + assert.equal(await memo.get("tree C", async () => "C"), "C"); + slowA.resolve("late A2"); + assert.equal(await a, "late A2"); + assert.equal(await memo.get("tree C", async () => "unused"), "C"); +}); diff --git a/common/views/autoQueueStatus.ts b/common/views/autoQueueStatus.ts @@ -208,3 +208,97 @@ export function buildAutoQueueStatusPayload( lanes: inputs.lanes, }; } + +// THE POLL'S EXPENSIVE HALF, SHARED: single-flight plus a short memo. +// +// The payload above is cheap; what it is built FROM is not. Each lane's +// pending work is a fold over every channel's snapshot (the shell's +// `computeLeafPending`, four lanes), and every surface that shows the lanes +// polls it every 3 s (`useOperationsStatus`). Each poll used to compute its +// own: on 2026-10-01, with the main thread busy regenerating a report, one +// computation took 96 s, and every tab's poll started another behind it — a +// queue that only grew. Now concurrent callers share the computation in flight, +// and a result is reused for AUTO_QUEUE_STATUS_MEMO_MS after it lands, so N +// pollers cost one fold per window. +// +// KEYED BY THE SETTINGS THE FOLD READS, AND OTHERWISE BY TIME. The fold reads +// settings itself (`computeLeafPending`: each lane's policy and rule tree, +// and `channelPriority` — paused channels, the focus, the compiled leaf ids), +// so the shell passes a key built from those (`autoQueue` + `channelPriority`): +// a rule added, a focus set or a channel paused changes the key and misses the +// memo, so a count is never keyed by a tree that is no longer the one drawn. +// What nothing tells the memo about — a snapshot rewritten, the runner picking +// a video (its in-flight set is subtracted inside the fold) — is up to +// AUTO_QUEUE_STATUS_MEMO_MS behind; the next poll is the correction. A lane's +// hold, the runner's status and its picks are not behind it at all: the shell +// reads them on every call (editor/app/operations/status.ts). +export const AUTO_QUEUE_STATUS_MEMO_MS = 3_000; + +export type SingleFlightMemo<T> = { + // The memoized value while it is fresh AND was computed under `key`; else the + // computation in flight under `key`; else `compute()`, started now and shared + // with every caller asking with `key` until it settles. A rejection is not + // memoized: the next caller computes again. + get(key: string, compute: () => Promise<T>): Promise<T>; + // Forget the value AND detach the computation in flight, whose result is + // then dropped rather than stored (a test reset must not be answered with + // the previous fixture's numbers). + clear(): void; +}; + +export function singleFlightMemo<T>( + opts: { ttlMs?: number; now?: () => number } = {}, +): SingleFlightMemo<T> { + const ttlMs = opts.ttlMs ?? AUTO_QUEUE_STATUS_MEMO_MS; + const now = opts.now ?? Date.now; + let value: { key: string; v: T; at: number } | null = null; + let inFlight: { key: string; p: Promise<T>; seq: number } | null = null; + // Every computation started gets the next number; only the latest started + // may store its value or clear the in-flight slot, so a computation under an + // old key that lands late never overwrites a newer one. clear() bumps it too. + let seq = 0; + return { + get(key, compute) { + if (value && value.key === key && now() - value.at < ttlMs) { + return Promise.resolve(value.v); + } + if (inFlight && inFlight.key === key) return inFlight.p; + const mine = ++seq; + const p = compute().then( + (v) => { + if (seq === mine) { + value = { key, v, at: now() }; + inFlight = null; + } + return v; + }, + (err: unknown) => { + if (seq === mine) inFlight = null; + throw err; + }, + ); + inFlight = { key, p, seq: mine }; + return p; + }, + clear() { + seq++; + value = null; + inFlight = null; + }, + }; +} + +// The editor's one instance, on globalThis like the registry: the polled route +// and the operations pages are separate bundles, and the e2e reset route +// (editor/app/api/test/invalidate-cache) drops it through the same global. +declare global { + // eslint-disable-next-line no-var + var __yttAutoQueueStatusMemo__: SingleFlightMemo<unknown> | undefined; +} + +export function autoQueueStatusMemo<T>(): SingleFlightMemo<T> { + if (!globalThis.__yttAutoQueueStatusMemo__) { + globalThis.__yttAutoQueueStatusMemo__ = singleFlightMemo<unknown>(); + } + return globalThis.__yttAutoQueueStatusMemo__ as SingleFlightMemo<T>; +} diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -26,6 +26,9 @@ - **A site can be private: built for reading on this machine, never deployed and never listed.** A site's settings have a new **Audience** choice (`audience` in `site.json`; only `"private"` is written). A private site is refused by every deploy — **Build & deploy**, **Deploy**, `archilyzer deploy site`, `pnpm ops build-deploy` and `deploy-site`, **Build & deploy all** (which builds it and skips its deploy) and `docker/publish-site.sh` — before anything is uploaded, in a sentence naming the audience; **Build** still builds it. It is left off the homepage, the hub and every other site's footer whatever **List on the Archilyzer homepage and hub** says, it publishes no hub URL, and its `corpus.json` says `"audience": "private"`, so a build of it is refused too if the site is switched back to public before it is rebuilt. Point the MCP at a private site's build to ask about what only it holds. Needs a rebuild and restart of the editor. - **A site built on the host right after another site no longer ships that site's channel list.** A site's **Build** (and Build all without containers) composes every site into the same folder, and the step that copies a site's summaries, stats and duplicates was skipped when that site's own data had not changed, even though another site had composed there since. The site then went out with the other site's channel list and stats. Those steps now run again whenever another site composed last. Needs a rebuild and restart of the editor. - **The hub no longer ships the data of the last site built before it.** **Build hub** built from the folder a site's build had just filled, so the hub carried that site's summaries, transcripts, posts and other data, and served them. The hub's build now clears every site's data first and uses the global search aliases, and **Deploy hub** refuses a hub build that still carries a site's data. Needs a rebuild and restart of the editor, then a rebuild and deploy of the hub. +- **The dashboard and `/jobs` keep answering while a channel's report is regenerated.** Regenerating a report walks every video of the channel inside the editor, and two regenerations of channels with a few thousand videos, running side by side, kept `/`, `/channels` and `/jobs` from loading for over an hour. Regenerations now run one at a time, on their own `refresh-report` queue on `/jobs` — a channel's own **Refresh report** included, which now waits its turn there too: it waits up to 15 seconds, then says where its job is instead ("Queued behind 3 report regenerations — the report updates when it finishes (job …).") and the page catches up when it runs, and a regeneration that fails now shows its reason under the button; a channel whose report is already waiting is not queued a second time, whether the request came from a finished job, **Refresh report** or **Update all reports**, and a change made while a channel's report is being regenerated queues one more regeneration after it rather than being missed; and the walk pauses between batches of videos so pages are served in between. **Update all reports** answers as soon as the regenerations are queued, and the reports land one after another; `pnpm ops refresh-report` answers with the job ids for both a single channel and `{"all":true}`, which `--wait` follows. Needs a rebuild and restart of the editor. +- **The operations pages share one count of the lanes' pending work.** Every open operations page asks for the lanes' status every 3 seconds, and each request used to count every lane's pending videos afresh from every channel's report. That count is now made once and handed to every request in the next 3 seconds. Changing a lane's rules, a focus or a channel's priority counts again at once; otherwise a pending count can be up to 3 seconds behind a report that was just rewritten or a video a lane just picked. A lane's hold, its runner and its picks are still read fresh on every request. +- **Jobs a stopped editor left "running" are closed when it starts again.** A job that was still running when the editor's process ended (killed, crashed, or shut down before the job had finished unwinding) kept "running" in its record for good, and `/jobs` listed it as archived. On start the editor now marks each one **cancelled**, with "interrupted: the process running it stopped before it finished" as the reason on the job's page, and its end time is the last time its log was written. Nothing is run again; **Retry** works as for any cancelled job. A job that another live process is running, such as `archilyzer run`, is left alone, and the same check now keeps the start-up pass from closing that process's queued jobs. Such leftover jobs never blocked a media move. ## [0.11.0] - 2026-09-30 - **Transcripts that arrived after a video was first seen are counted.** The stats behind the homepage, the hub and every site's charts were cached per video and refreshed only when the video's metadata changed, so a transcript that came later — a Whisper run days after the download, or a video downloaded after the last index build — never reached them, and a video with YouTube captions alone had no transcription date. Counts and charts were low; the homepage could show a site with 0 transcripts, 0 channels and 0 hours while it served its videos. A stat is now also redone whenever the index re-reads the video, every transcript has a date, and a captioned video is dated by when its captions arrived rather than by a later Normalize run, so its place on "Transcribed over time" can move. **After updating, rebuild and restart the editor before anything else:** until then, **Build stats dataset** runs the old code and would undo the new stats, while a site, hub or homepage build already runs the new code — and the first stats build of any kind re-reads every video once (about 10–30 minutes on a large archive; it can be stopped and picks up where it stopped). Then build the index, the stats, the homepage, the hub, and the sites. diff --git a/editor/app/api/ops/refresh-report/route.ts b/editor/app/api/ops/refresh-report/route.ts @@ -1,11 +1,11 @@ import { NextResponse } from "next/server"; +import { refreshAllChannelSnapshotsAction } from "../../../channels/actions"; +import { channelExists } from "yt-dlp-transcript-common/controller/channels"; +import { requestRefreshReport } from "yt-dlp-transcript-common/jobs/snapshotScheduler"; +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { - refreshAllChannelSnapshotsAction, - refreshChannelSnapshotAction, -} from "../../../channels/actions"; -import { - actionResponse, OpsInputError, + opsFail, ops, optBool, optString, @@ -16,10 +16,17 @@ export const dynamic = "force-dynamic"; // POST { slug: string } | { all: true } // -// The single-channel form REGENERATES SYNCHRONOUSLY (it is a filesystem scan, -// not a job) and returns `{ ok: true }` once snapshot.json is on disk. The -// `all` form queues one refresh-report job per channel and returns the bulk -// { queued, skipped } the /channels header button shows. +// BOTH FORMS ANSWER ONCE QUEUED, with the job ids `--wait` follows — the ops +// rule (_lib.ts: a job-starting route returns a jobId and never streams). +// Every regeneration runs on the one serial refresh-report queue (release 17), +// so holding the request until a report is on disk meant waiting for every +// walk queued ahead of it: past five minutes, and fetch's own header timeout +// ends the CLI with an error while the job goes on to succeed. +// +// `{ slug }` → `{ ok, jobId, started }`: the job this call queued, or the one +// already queued for the channel (`started: false`; it has not started +// reading, so it is as fresh). `{ all: true }` → the bulk `{ queued, jobIds, +// skipped }` the /channels header button shows. export async function POST(request: Request) { return ops(request, ["slug", "all"], async (body) => { const all = optBool(body, "all"); @@ -36,8 +43,17 @@ export async function POST(request: Request) { } // Re-read through reqSlug now that we know it is the single-channel form: // the shape check belongs on the value that reaches a path.join. - return actionResponse( - await refreshChannelSnapshotAction(reqSlug(body, "slug")), - ); + const checked = reqSlug(body, "slug"); + const paths = getPaths(); + if (!(await channelExists(paths, checked))) { + return opsFail(`Channel "${checked}" not found`, 404); + } + const requested = await requestRefreshReport(paths, checked); + if (!requested.ok) return opsFail(requested.error); + return NextResponse.json({ + ok: true, + jobId: requested.jobId, + started: requested.started, + }); }); } diff --git a/editor/app/api/test/invalidate-cache/route.ts b/editor/app/api/test/invalidate-cache/route.ts @@ -125,6 +125,14 @@ function invalidate() { // a spec that rewrites a title in place inside one mtime tick would otherwise // see the previous spec's title. resetVideoTitleMemo(); + // And the auto-queue status poll's three-second memo of the snapshot-derived + // pending counts (common/views/autoQueueStatus.ts): the next spec's first + // poll must count ITS fixture, not the previous spec's. Dropped through the + // global like the singletons above, so this route imports nothing that + // reaches the runner; a computation still in flight lands in the dropped + // object. + // eslint-disable-next-line @typescript-eslint/no-explicit-any + (globalThis as any).__yttAutoQueueStatusMemo__ = undefined; revalidatePath("/", "layout"); return NextResponse.json({ ok: true }); } diff --git a/editor/app/api/test/settle-running-metas/route.ts b/editor/app/api/test/settle-running-metas/route.ts @@ -0,0 +1,26 @@ +import { NextResponse } from "next/server"; +import { settleRunningJobMetas } from "yt-dlp-transcript-common/jobs/bootQueuedJobs"; +import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { testRouteDenied } from "../_guard"; + +export const dynamic = "force-dynamic"; + +// E2E test harness only, GUARDED BY `E2E_TEST_ROUTES` like every other +// /api/test route (see _guard.ts). It runs the boot pass over `running` metas +// (common/jobs/bootQueuedJobs.ts `settleRunningJobMetas`) NOW, the way +// editor/instrumentation.ts runs it once at boot: the e2e server boots once +// for the whole suite, so a spec that plants a ghost meta (a dead process's +// `running` job) has no other way to watch the pass close it. "This boot" is +// this instant: every meta from before it, not in the registry, whose writer +// is gone. +export async function POST() { + const denied = testRouteDenied(); + if (denied) return denied; + const result = await settleRunningJobMetas({ + paths: getPaths(), + bootedAt: Date.now(), + isLive: (id) => getRegistry().get(id) !== undefined, + }); + return NextResponse.json(result); +} diff --git a/editor/app/channels/[slug]/components/RefreshSnapshotButton.tsx b/editor/app/channels/[slug]/components/RefreshSnapshotButton.tsx @@ -1,23 +1,54 @@ "use client"; -import { useTransition } from "react"; -import { refreshChannelSnapshotAction } from "../../actions"; +import { useState, useTransition } from "react"; +import { + refreshChannelSnapshotAction, + type RefreshSnapshotResult, +} from "../../actions"; +// THE CHANNEL'S "Refresh report". The regeneration runs on the serial +// refresh-report queue (release 17), so the action answers within ~15 s either +// way: nothing (the report is on disk; the page re-renders), `{ error }` (the +// walk could not run — an unmounted drive's sentence), or `{ notice }` (still +// queued: "Queued behind N report regenerations — …, job …"). Both are drawn +// under the button; they used to be dropped, so a refusal showed nothing. export function RefreshSnapshotButton({ slug }: { slug: string }) { const [pending, startTransition] = useTransition(); + const [result, setResult] = useState<RefreshSnapshotResult>(undefined); return ( - <button - type="button" - onClick={() => - startTransition(async () => { - await refreshChannelSnapshotAction(slug); - }) - } - disabled={pending} - aria-label="refresh channel report" - className="self-start px-3 py-1 rounded-md bg-primary text-primary-foreground text-sm font-medium hover:opacity-90 disabled:opacity-50" - > - {pending ? "Refreshing…" : "Refresh report"} - </button> + <div className="flex flex-col items-start gap-1"> + <button + type="button" + onClick={() => + startTransition(async () => { + setResult(undefined); + setResult(await refreshChannelSnapshotAction(slug)); + }) + } + disabled={pending} + aria-label="refresh channel report" + className="self-start px-3 py-1 rounded-md bg-primary text-primary-foreground text-sm font-medium hover:opacity-90 disabled:opacity-50" + > + {pending ? "Refreshing…" : "Refresh report"} + </button> + {result && "error" in result && ( + <span + role="alert" + aria-label="refresh report error" + className="text-xs text-destructive" + > + {result.error} + </span> + )} + {result && "notice" in result && ( + <span + role="status" + aria-label="refresh report notice" + className="text-xs text-muted-foreground" + > + {result.notice} + </span> + )} + </div> ); } diff --git a/editor/app/channels/actions.ts b/editor/app/channels/actions.ts @@ -24,13 +24,13 @@ import { renameChannel } from "yt-dlp-transcript-common/controller/renameChannel import { inspectChannelMedia } from "yt-dlp-transcript-common/lib/channelMedia"; import { channelMediaBusyReason } from "./lib/mediaBusy"; import { - excludedDownloadIdSet, - generateChannelSnapshot, -} from "yt-dlp-transcript-common/controller/channelSnapshot"; -import { requestChannelSnapshot } from "yt-dlp-transcript-common/jobs/snapshotScheduler"; -import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; -import { runManagedFunction } from "yt-dlp-transcript-common/jobs/streamCommand"; -import { drainStream } from "yt-dlp-transcript-common/jobs/drainStream"; + REFRESH_REPORT_ACTIVE, + refreshReportWaitNotice, + requestChannelSnapshot, + requestRefreshReport, + startRefreshReport, + waitForRefreshReport, +} from "yt-dlp-transcript-common/jobs/snapshotScheduler"; import { siteChannelIndex, type Site, @@ -349,25 +349,41 @@ export async function gotoVideoAction( redirect(`/channels/${slug}/videos/${encodeURIComponent(id)}`); } +// What a channel's Refresh report answers: nothing when the report is on disk, +// `{ error }` when it could not be made, `{ notice }` when it is still queued. +export type RefreshSnapshotResult = ActionResult | { notice: string }; + export async function refreshChannelSnapshotAction( slug: string, -): Promise<ActionResult> { +): Promise<RefreshSnapshotResult> { const paths = getPaths(); if (!(await channelExists(paths, slug))) { return { error: `Channel "${slug}" not found` }; } - // generateChannelSnapshot now THROWS on a channel whose media is not - // reachable (guard 3) rather than writing a snapshot that says every video is - // undownloaded. That is the right behaviour and the wrong exception to let - // out of a server action: an uncaught throw here reaches the client as a - // digest-only "an error occurred", and the one thing the operator needs is - // the sentence naming the unmounted drive. Every sibling in this file returns + // THROUGH THE REFRESH-REPORT QUEUE, like every other walk (release 17 slice + // D0): a walk run here, in the request, beside a queued one was two walks + // side by side — half of the 2026-10-01 outage — and could land an older + // read over a newer one. The job is the one this click starts, or the one + // already queued for the channel (it has not started reading, so it is as + // fresh). + // + // A BOUNDED WAIT (REFRESH_REPORT_WAIT_MS, 15 s). The queue is serial: behind + // a 3,000-video walk, or during Update all reports, this channel's report may + // be minutes away. Past the bound the action answers with where the job is + // ("Queued behind 3 report regenerations — …, job …") and the page catches + // up when it runs (the scheduler's generation moves /api/pulse). + // + // A FAILED WALK IS AN ERROR, whoever started it: generateChannelSnapshot + // throws on a channel whose media is not reachable (guard 3) rather than + // writing a snapshot that says every video is undownloaded, and that + // sentence — the job log's `[error]` line — is what comes back, not a + // digest-only "an error occurred". Every sibling in this file returns // { error }; so does this. - try { - await generateChannelSnapshot(paths, slug); - } catch (e) { - return { error: (e as Error).message }; - } + const requested = await requestRefreshReport(paths, slug); + if (!requested.ok) return { error: requested.error }; + const wait = await waitForRefreshReport(paths, requested.jobId); + if (wait.state === "failed") return { error: wait.error }; + if (wait.state === "waiting") return { notice: refreshReportWaitNotice(wait) }; revalidatePath(`/channels/${slug}`); // The Report column on /channels is read off this snapshot, and the row // action sits next to the marker it flips — so revalidate the list too, not @@ -387,82 +403,39 @@ export type RefreshAllResult = { // `QueueOutcome` grew one: without it an HTTP caller could not tell a // fan-out that started work from an action that started none, so // `pnpm ops refresh-report --json '{"all":true}' --wait` returned the moment - // the response arrived. (This action also AWAITS its streams, so by the time - // it answers the work is done — but the ids are what make the response - // honest about what it started, and identical in shape to every other - // fan-out's.) + // the response arrived. The action answers once the jobs are QUEUED (release + // 17: they run one at a time), so the ids are what `--wait` follows. jobIds: string[]; }; export async function refreshAllChannelSnapshotsAction(): Promise<RefreshAllResult> { const paths = getPaths(); const channels = await listChannelConfigs(paths); - const active = new Set( - getRegistry() - .list() - .filter( - (j) => - j.kind === "refresh-report" && - (j.status === "queued" || j.status === "running") && - j.channelSlug, - ) - .map((j) => j.channelSlug as string), - ); const queued: string[] = []; const jobIds: string[] = []; const skipped: { slug: string; reason: string }[] = []; - const streams: ReadableStream<string>[] = []; for (const c of channels) { - if (active.has(c.slug)) { - skipped.push({ slug: c.slug, reason: "already running" }); - continue; - } - const result = await runManagedFunction({ - kind: "refresh-report", - // Empty queueKey: bypass queue serialization. Snapshot regen is a - // local filesystem scan that never touches the platform, so there's - // no reason for it to wait behind sync/download work. See - // registry.ts:69-72 for the documented escape hatch. - queueKey: "", - paths, - channelSlug: c.slug, - fn: async (onLog) => { - onLog(`Regenerating report for ${c.slug}…`); - const snap = await generateChannelSnapshot(paths, c.slug); - const excluded = excludedDownloadIdSet(snap); - const awaitingTranscription = excluded.size - ? snap.buckets.downloadedNoTranscript.filter( - (id) => !excluded.has(id), - ).length - : snap.buckets.downloadedNoTranscript.length; - onLog( - `Done. ${snap.totals.videos} videos · ` + - `${snap.undownloadedIds.length} undownloaded · ` + - `${awaitingTranscription} awaiting transcription.`, - ); - // Deliberately no revalidatePath here — calling it from a - // background fn races with the in-flight re-render that the action's - // own revalidatePath triggers. The action's single revalidate at the - // end picks up every fresh snapshot. - }, - }); + // The scheduler's one entry point: the refresh-report queue (one + // regeneration at a time — they used to run all at once, in parallel) and + // its per-slug dedup, shared with the debounced regeneration. + const result = await startRefreshReport(paths, c.slug); if (!result.ok) { - skipped.push({ slug: c.slug, reason: result.error }); + skipped.push({ + slug: c.slug, + reason: result.error === REFRESH_REPORT_ACTIVE ? "already queued" : result.error, + }); continue; } + // Nothing reads the stream (the log is on disk). + void result.stream.cancel(); queued.push(c.slug); jobIds.push(result.jobId); - streams.push(result.stream); } - // Wait for all snapshots to finish writing before revalidating so the - // pages that read the snapshots read fresh counts. With queueKey === "" the - // jobs all run in parallel, so this waits roughly the time of the - // slowest snapshot, not the sum. - await Promise.all(streams.map(drainStream)); - revalidatePath("/channels"); - revalidatePath("/operations/[id]", "page"); - revalidatePath("/cleanup"); - revalidatePath("/"); + // ANSWERS ONCE QUEUED, NOT ONCE DONE. The regenerations run one at a time, + // so waiting for them was waiting for every channel's walk added up — + // minutes on a real corpus — behind a button and an ops call that time out. + // The pages catch up as each lands: every regeneration moves the + // scheduler's generation, which /api/pulse carries. return { queued, jobIds, skipped }; } diff --git a/editor/app/components/actions/InlineActionButton.tsx b/editor/app/components/actions/InlineActionButton.tsx @@ -110,7 +110,12 @@ async function runAction(variant: Variant): Promise<StreamActionResult> { if (result && "error" in result) { return { ok: false, error: result.error }; } - // refreshReport runs no managed job, so synthesize an already-complete result. + // Still queued after the action's bounded wait: a neutral line, not a red one. + if (result && "notice" in result) { + return { ok: false, error: result.notice, info: true }; + } + // The report is on disk (the action waited for its job), so synthesize an + // already-complete result. return { ok: true, jobId: "", diff --git a/editor/app/lib/requestCache.ts b/editor/app/lib/requestCache.ts @@ -21,6 +21,9 @@ import { getSettings } from "yt-dlp-transcript-common/lib/settings"; // scheduler rewrites these files on a ~1 s debounce from job runners inside // common/, which cannot import next/cache to invalidate anything. A cache // nothing can invalidate is just a stale number with extra steps. +// ONE EXCEPTION, bounded: the operations status poll's 3-second memo of the +// lanes' pending counts (operations/status.ts, `autoQueueStatusMemo`), keyed by +// the settings it is folded from and otherwise expired by time. // Keyed on the `paths` argument by identity, which works because getPaths() // memoizes its result at module scope and hands back the same object every // call. Pass it straight through; don't spread or rebuild it. diff --git a/editor/app/operations/[id]/page.tsx b/editor/app/operations/[id]/page.tsx @@ -280,9 +280,10 @@ export default async function OperationPage({ const sections = sectionsFor(op.id as SectionConfig["operation"]); const channelWork = sections.length > 0 ? ( - // No extra disk walk: the census is built from the request-cached - // `getChannelBriefs` that `buildAutoQueueStatusPayload` already read - // above (operations/status.ts:38). Keyed for the same reason + // The census reads the request-cached `getChannelBriefs`. It is the + // same listing `buildAutoQueueStatusPayload` read above only when that + // call missed its memo (operations/status.ts); on a hit the payload's + // counts may be up to 3 s older than this table. Keyed for the same reason // settingsFormFor's elements are — a server element handed to a client // component lands in its children array with React's dev-only key check // still to run over it. diff --git a/editor/app/operations/status.ts b/editor/app/operations/status.ts @@ -14,9 +14,11 @@ import { LANES } from "yt-dlp-transcript-common/lib/autoQueueTypes"; import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool"; import { buildAutoQueueLanes } from "yt-dlp-transcript-common/views/autoQueueLanes"; import { + autoQueueStatusMemo, buildAutoQueueStatusPayload as build, type AutoQueueStatusPayload, } from "yt-dlp-transcript-common/views/autoQueueStatus"; +import type { ChannelBrief } from "yt-dlp-transcript-common/controller/channels"; import { getChannelBriefs } from "../lib/requestCache"; import { readPriorityView } from "./channelPriorityView"; @@ -29,12 +31,25 @@ import { readPriorityView } from "./channelPriorityView"; // times per poll, because `buildKind` did its own reading; and the channel // briefs are shared with the lanes builder through the per-request cache, and // with the four computeLeafPending calls through their `shared` argument. -export async function buildAutoQueueStatusPayload(): Promise<AutoQueueStatusPayload> { + +// THE SNAPSHOT-DERIVED HALF, behind the shared single-flight memo +// (`autoQueueStatusMemo`, common/views/autoQueueStatus.ts): the channel briefs +// and the four lanes' pending work. It is the expensive half, a fold over every +// channel's snapshot, and the one every poller used to pay for separately. The +// memo is KEYED by the settings the fold reads (each lane's policy and tree in +// `autoQueue`, and `channelPriority`), so an operator's edit to either misses it +// and is never paired with counts from the tree before; otherwise it holds for +// AUTO_QUEUE_STATUS_MEMO_MS (3 s) — a snapshot rewritten or a video the runner +// just picked shows at most one poll late. That window is the one exception to +// requestCache.ts's "no cache longer than a request" rule, bounded by time. +type SnapshotHalf = { + briefs: ChannelBrief[]; + pendingByKind: LeafPending[]; +}; + +async function computeSnapshotHalf(): Promise<SnapshotHalf> { const paths = getPaths(); - const settings = getSettings(); - const pool = getWorkerPool(); - const [priority, state, briefs] = await Promise.all([ - readPriorityView(), + const [state, briefs] = await Promise.all([ readAutoQueueState(paths), getChannelBriefs(paths), ]); @@ -47,6 +62,25 @@ export async function buildAutoQueueStatusPayload(): Promise<AutoQueueStatusPayl computeLeafPending(lane, paths, { configs: briefs, state }), ), ); + return { briefs, pendingByKind }; +} + +export async function buildAutoQueueStatusPayload(): Promise<AutoQueueStatusPayload> { + const paths = getPaths(); + const settings = getSettings(); + const pool = getWorkerPool(); + // FRESH on every call: the priority view, the state document (picks, + // cooldowns, deferrals), the settings (holds, policies), the pool and the + // runners. The pending counts are memoized under a key of the settings they + // are folded from, so an edit to a policy, a tree or a priority recomputes + // them at once; only what changes without a settings write (a snapshot, the + // runner's in-flight set) can be up to 3 s behind. + const memoKey = JSON.stringify([settings.autoQueue, settings.channelPriority]); + const [priority, state, { briefs, pendingByKind }] = await Promise.all([ + readPriorityView(), + readAutoQueueState(paths), + autoQueueStatusMemo<SnapshotHalf>().get(memoKey, computeSnapshotHalf), + ]); const byLane = <T>(values: readonly T[]): Record<AutoQueueKind, T> => Object.fromEntries(LANES.map((lane, i) => [lane, values[i]])) as Record< AutoQueueKind, diff --git a/editor/e2e/bulk-actions.spec.ts b/editor/e2e/bulk-actions.spec.ts @@ -165,9 +165,15 @@ test("Select failed + Delete directories removes the dirs and queues no job", as await expect(page.getByLabel("select vidB")).toBeVisible(); // Crucially: NO managed job was queued (no transcode/download triggered). + // The report generateReport asked for is a refresh-report job of its own + // (release 17: every report walk runs on the refresh-report queue), and + // /jobs may draw it before its default filter hides it — not this action's. await page.goto("/jobs"); await expect( - page.getByRole("row").filter({ hasText: "test-transcribe" }), + page + .getByRole("row") + .filter({ hasText: "test-transcribe" }) + .filter({ hasNotText: "refresh-report" }), ).toHaveCount(0); }); diff --git a/editor/e2e/dashboard-answers.spec.ts b/editor/e2e/dashboard-answers.spec.ts @@ -0,0 +1,356 @@ +import { link, mkdir, readFile, writeFile } from "node:fs/promises"; +import { join } from "node:path"; +import { test, expect, type APIRequestContext, type Page } from "@playwright/test"; +import { ulid } from "yt-dlp-transcript-common/jobs/ulid"; +import { baseUrl } from "./baseUrl"; +import { + channelStage, + generateReport, + resetData, + resolvePath, +} from "./helpers"; + +// THE DASHBOARD ANSWERS WHILE A REPORT REGENERATES (release 17 slice D0). +// +// On 2026-10-01 the live editor's `/`, `/channels` and `/jobs` gave no +// response for over an hour while two `refresh-report` jobs walked 2,000- and +// 3,260-video channels in its own process, side by side, and three more +// `refresh-report` metas still read `running` hours after the process that ran +// them was gone. The first case below is the first half: two channels whose +// reports take several seconds each, regenerated by "Update all reports" — one +// after the other, never side by side — with `/` and `/jobs` polled every 2 s +// throughout. The second is the ghost: a `running` meta a dead process left +// behind is closed as interrupted by the boot pass, one a live process owns is +// not, and neither stops a media move. +// +// THE BUDGET IS 5 s, OR THREE TIMES WHAT THE SAME PAGE TOOK JUST BEFORE THE +// REGENERATION, WHICHEVER IS LONGER. The suite runs under `next dev` on a +// machine other suites and builds share; at a load average of 30 a page that +// renders in 0.3 s on a quiet machine takes 2–4 s with nothing regenerating at +// all. What this pins is that the regeneration does not starve the pages, not +// how fast the machine is — so the idle measurement taken a moment before sets +// the floor, and on a quiet machine the budget is simply 5 s. CAPPED AT 15 s: +// on a machine so loaded that idle pages take over 5 s, the floor stops +// growing, and the case fails rather than stretching to hide starvation. +// +// What the case pins, honestly: the serial queue (read from the jobs' own +// records) and starvation at the scale of seconds. The yield between chunks is +// pinned by a unit test (common/controller/snapshotYield.test.ts): on a dev +// server under load `main`'s walk answered inside such a budget too. + +const OPS_AUTH = { authorization: "Bearer test-worker-token" }; + +// THE BIG CHANNELS. Every video dir hardlinks ONE ~400 KB metadata.info.json +// — the size of a long VOD's, the file the walk parses — so 600 of them cost +// one file's bytes and a second to make. No `webpage_url`: the reconcile pass +// at the walk's start parses each file for it and, finding none, renames +// nothing (one shared id would merge every dir into one). The archive names a +// non-YouTube extractor, so the walk parses each file a second time for its +// native id — the Rumble-channel path. 600 is several seconds of walk under +// `next start` on a quiet machine and well over a minute under `next dev` at a +// load average of 30; 2,000 did not finish inside five minutes there. +const BIG = ["big-channel-a", "big-channel-b"]; +const BIG_VIDEOS = 600; + +async function seedBigChannel(slug: string): Promise<void> { + const channelDir = resolvePath(`test-transcripts/channels/${slug}`); + const dataDir = join(channelDir, "data"); + await mkdir(dataDir, { recursive: true }); + await writeFile( + join(channelDir, "config.json"), + JSON.stringify({ handling: "youtube", name: "Big Channel" }), + ); + const formats = Array.from({ length: 400 }, (_, i) => ({ + format_id: `hls-${i}`, + url: `https://example.invalid/${"x".repeat(600)}${i}`, + ext: "mp4", + protocol: "m3u8_native", + width: 1280, + height: 720, + tbr: 1234.5 + i, + http_headers: { "User-Agent": `Mozilla/5.0 ${"y".repeat(80)}`, Accept: "*/*" }, + fragments: [ + { url: "seg0", duration: 6 }, + { url: "seg1", duration: 6 }, + ], + })); + const template = resolvePath(`test-transcripts/.${slug}-meta.json`); + await writeFile( + template, + JSON.stringify({ id: "native", title: "A long VOD", duration: 3600, formats }), + ); + const ids = Array.from( + { length: BIG_VIDEOS }, + (_, i) => `${slug.slice(-1)}big${String(i).padStart(8, "0")}`, + ); + for (let i = 0; i < ids.length; i += 100) { + await Promise.all( + ids.slice(i, i + 100).map(async (id) => { + const dir = join(dataDir, id); + await mkdir(dir); + await link(template, join(dir, "metadata.info.json")); + }), + ); + } + await writeFile( + join(channelDir, "archive"), + ids.map((id) => `rumble ${id}\n`).join(""), + ); + await writeFile(join(channelDir, "playlist"), ""); +} + +type Sample = { path: string; ms: number; status: number }; + +async function timedGet( + request: APIRequestContext, + path: string, +): Promise<Sample> { + const t = Date.now(); + const res = await request.get(`${baseUrl}${path}`, { timeout: 60_000 }); + // The body too: a page that sends its head and stalls is not an answer. + await res.body(); + return { path, ms: Date.now() - t, status: res.status() }; +} + +type JobMetaOnDisk = { + id: string; + kind: string; + channelSlug?: string; + status: string; + startedAt?: number; + endedAt?: number; +}; + +test("/ and /jobs answer while two large reports regenerate, one after the other", async ({ + request, +}) => { + test.setTimeout(420_000); + await resetData("empty"); + for (const slug of BIG) await seedBigChannel(slug); + + // Warm both routes first: under `next dev` the first request compiles the + // page, which is the dev server's cost and not what this measures. Then the + // idle floor: the slowest of three more rounds, with nothing regenerating. + for (const path of ["/", "/jobs"]) { + expect((await timedGet(request, path)).status).toBe(200); + } + // And the ops route that starts the regeneration: its first request compiles + // it, which under load held `/` for 15 s in one full-suite run — a dev + // server's compile, not a walk. A body naming neither form is refused (400) + // after the route has loaded, and starts nothing. + const warm = await request.post(`${baseUrl}/api/ops/refresh-report`, { + headers: OPS_AUTH, + data: {}, + timeout: 120_000, + }); + expect(warm.status()).toBe(400); + const idle: Sample[] = []; + for (let i = 0; i < 3; i++) { + for (const path of ["/", "/jobs"]) idle.push(await timedGet(request, path)); + } + const budgetMs = Math.min( + Math.max(5_000, 3 * Math.max(...idle.map((s) => s.ms))), + 15_000, + ); + + // "Update all reports", through the ops API: a refresh-report job per + // channel on the serial queue. The call answers once they are QUEUED, with + // their ids; "while they regenerate" lasts until both jobs' records say + // they have ended. + const startedAt = Date.now(); + const res = await request.post(`${baseUrl}/api/ops/refresh-report`, { + headers: OPS_AUTH, + data: { all: true }, + timeout: 120_000, + }); + expect(res.status()).toBe(200); + const done = { body: (await res.json()) as { queued?: string[]; jobIds?: string[] } }; + expect([...(done.body.queued ?? [])].sort()).toEqual(BIG); + expect(done.body.jobIds).toHaveLength(2); + const ended = async (): Promise<boolean> => { + for (const id of done.body.jobIds ?? []) { + const m = await readFile(resolvePath(`test-transcripts/.jobs/${id}.meta.json`), "utf8") + .then((raw) => JSON.parse(raw) as JobMetaOnDisk) + .catch(() => null); + if (!m || m.status === "queued" || m.status === "running") return false; + } + return true; + }; + + let finishedAt: number | null = null; + const during: Sample[] = []; + while (finishedAt === null) { + if (Date.now() - startedAt > 360_000) throw new Error("the regenerations did not end in 6 min"); + const tick = Date.now(); + for (const path of ["/", "/jobs"]) { + const s = await timedGet(request, path); + during.push(s); + console.log(`[dashboard-answers] +${tick - startedAt} ms ${path} ${s.ms} ms`); + } + if (await ended()) { + finishedAt = Date.now(); + break; + } + const wait = 2_000 - (Date.now() - tick); + if (wait > 0) await new Promise((r) => setTimeout(r, wait)); + } + for (const slug of BIG) { + const snapshot = JSON.parse( + await readFile( + resolvePath(`test-transcripts/channels/${slug}/snapshot.json`), + "utf8", + ), + ) as { totals: { videos: number } }; + expect(snapshot.totals.videos).toBe(BIG_VIDEOS); + } + + // ONE AFTER THE OTHER: the second started when the first had ended. Read + // off the jobs' records on disk; the terminal write lands just after the + // job's stream closes, so it is waited for. + const readMetas = () => + Promise.all( + (done.body.jobIds ?? []).map( + async (id) => + JSON.parse( + await readFile(resolvePath(`test-transcripts/.jobs/${id}.meta.json`), "utf8"), + ) as JobMetaOnDisk, + ), + ); + await expect + .poll(async () => (await readMetas()).map((m) => `${m.kind} ${m.status}`), { + timeout: 10_000, + }) + .toEqual(["refresh-report done", "refresh-report done"]); + const metas = await readMetas(); + const [first, second] = [...metas].sort( + (a, b) => (a.startedAt ?? 0) - (b.startedAt ?? 0), + ); + expect(second.startedAt ?? 0).toBeGreaterThanOrEqual(first.endedAt ?? Infinity); + + // The fixture has to make the walk long enough to be polled through, or the + // budget below proves nothing. + const regenMs = (finishedAt ?? Date.now()) - startedAt; + test.info().annotations.push({ + type: "timings", + description: + `idle max ${Math.max(...idle.map((s) => s.ms))} ms, budget ${budgetMs} ms; ` + + `regeneration ${regenMs} ms; ${during.map((s) => `${s.path} ${s.ms}`).join(", ")}`, + }); + expect(regenMs, "the regenerations took several seconds").toBeGreaterThan(4_000); + for (const path of ["/", "/jobs"]) { + expect( + during.filter((s) => s.path === path).length, + `${path} was polled during the regeneration`, + ).toBeGreaterThanOrEqual(2); + } + for (const s of during) { + expect(s.status, s.path).toBe(200); + expect(s.ms, `${s.path} answered in ${s.ms} ms (budget ${budgetMs} ms)`).toBeLessThan(budgetMs); + } +}); + +// --------------------------------------------------------------------------- +// The ghost. + +const SLUG = "test-youtube"; + +async function quiet(page: Page): Promise<void> { + await expect + .poll( + async () => { + const res = await page.request.get(`${baseUrl}/api/jobs/active`); + const body = await res.json(); + const jobs: { channelSlug?: string; status: string }[] = Array.isArray( + body, + ) + ? body + : (body.jobs ?? []); + return jobs.filter( + (j) => + j.channelSlug === SLUG && + (j.status === "running" || j.status === "queued"), + ).length; + }, + { timeout: 30_000 }, + ) + .toBe(0); +} + +// A `running` refresh-report meta for SLUG, as a process that died mid-walk +// leaves it: queued and started an hour ago, its log last written then. +async function plantGhost(pid: number): Promise<string> { + const id = ulid(Date.now() - 60 * 60 * 1000); + const jobsDir = resolvePath("test-transcripts/.jobs"); + await mkdir(jobsDir, { recursive: true }); + const at = Date.now() - 60 * 60 * 1000; + await writeFile( + join(jobsDir, `${id}.meta.json`), + JSON.stringify({ + id, + kind: "refresh-report", + queueKey: "", + channelSlug: SLUG, + status: "running", + queuedAt: at, + startedAt: at, + pid, + }), + ); + await writeFile( + join(jobsDir, `${id}.log`), + `Regenerating report for ${SLUG}…\n`, + ); + return id; +} + +async function metaOf(id: string): Promise<{ status: string; cancelReason?: string }> { + return JSON.parse( + await readFile(resolvePath(`test-transcripts/.jobs/${id}.meta.json`), "utf8"), + ); +} + +test("a ghost running meta from a dead process is closed as interrupted, and does not block a move", async ({ + page, +}, testInfo) => { + test.setTimeout(180_000); + await resetData("one-youtube-channel-with-data"); + await generateReport(page, SLUG); + await quiet(page); + + // Past pid_max: no process has it. And this test runner's own pid: a live + // process that is not the editor — an `archilyzer run` beside it. + const ghost = await plantGhost(2 ** 22 + 1); + const live = await plantGhost(process.pid); + + const res = await page.request.post(`${baseUrl}/api/test/settle-running-metas`); + expect(res.ok()).toBe(true); + const { interrupted } = (await res.json()) as { interrupted: { id: string }[] }; + expect(interrupted.map((j) => j.id)).toContain(ghost); + expect(interrupted.map((j) => j.id)).not.toContain(live); + + const closed = await metaOf(ghost); + expect(closed.status).toBe("cancelled"); + expect(closed.cancelReason).toMatch(/^interrupted: /); + expect((await metaOf(live)).status).toBe("running"); + + // /jobs says why, on the job's own page. + await page.goto(`/jobs/${ghost}`); + await expect(page.getByTestId("cancel-reason")).toContainText("interrupted"); + + // Neither the closed ghost nor the live one is a writer on the channel: the + // Storage panel offers the move and the move completes. + const root = testInfo.outputPath("ghost-root"); + await mkdir(root, { recursive: true }); + await page.goto(channelStage(SLUG, "storage")); + await page.getByLabel("destination root").fill(root); + await page.getByRole("button", { name: "Preview", exact: true }).click(); + await expect(page.getByLabel("relocation preview")).toBeVisible({ + timeout: 15_000, + }); + const moveButton = page.getByRole("button", { name: "Move media" }); + await expect(moveButton).toBeEnabled(); + await moveButton.click(); + await expect(page.getByLabel("Move media output")).toContainText("Moved", { + timeout: 60_000, + }); +}); diff --git a/editor/e2e/ops-api.spec.ts b/editor/e2e/ops-api.spec.ts @@ -501,23 +501,46 @@ test("refresh-report regenerates snapshot.json", async ({ request }) => { await rm(resolvePath(SNAP), { force: true }); expect(await pathExists(SNAP)).toBe(false); - // No page, no click: the route IS the refresh. It regenerates SYNCHRONOUSLY - // (a filesystem scan, not a job), so { ok: true } means the file is there. + // No page, no click: the route IS the refresh. It answers once the + // regeneration is QUEUED (release 17: one serial refresh-report queue, and + // the ops rule — a job-starting route returns a jobId), so the job is + // followed to `done` before the file is read. + const doneJob = async (jobId: string | undefined) => { + expect(jobId).toBeTruthy(); + await expect + .poll( + async () => + ( + await readJson<{ status: string }>( + `test-transcripts/.jobs/${jobId}.meta.json`, + ).catch(() => null) + )?.status ?? null, + { timeout: 30_000 }, + ) + .toBe("done"); + }; const first = await ops(request, "refresh-report", { slug: SLUG }); - expect(first.body).toEqual({ ok: true }); + expect(first.body.ok).toBe(true); + await doneJob(first.body.jobId); const snapshot = await readJson<{ generatedAt: string; totals: { videos: number } }>(SNAP); expect(snapshot.generatedAt).toBeTruthy(); expect(snapshot.totals.videos).toBeGreaterThanOrEqual(0); - // Re-running REWRITES it. Polled through the action itself because two scans - // of a six-video fixture can land in the same millisecond. + // Re-running REWRITES it. Polled because two scans of a six-video fixture + // can land in the same millisecond. await expect .poll(async () => { - await ops(request, "refresh-report", { slug: SLUG }); + const again = await ops(request, "refresh-report", { slug: SLUG }); + await doneJob(again.body.jobId); return (await readJson<{ generatedAt: string }>(SNAP)).generatedAt; }) .not.toBe(snapshot.generatedAt); + // An unknown channel is a 404 naming it, and starts nothing. + const missing = await ops(request, "refresh-report", { slug: "no-such-channel" }); + expect(missing.status).toBe(404); + expect(missing.body.error).toMatch(/no-such-channel/); + // The bulk form queues a job per channel and reports both lists. const all = await ops(request, "refresh-report", { all: true }); expect(all.status).toBe(200); 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(), diff --git a/plans/release-17.md b/plans/release-17.md @@ -813,4 +813,224 @@ loop (`buildIndex.ts` ~1943 and the fingerprint), a different hunk from T1's `sc The hub suite runs `next dev` in hub mode over `export/public` and never composes, so it shows the hub app is unchanged; the clearing itself is pinned by `compose-hub.test.ts`. +### Slice D0, as shipped — the dashboard answers while a snapshot regenerates (2026-10-01) + +Branch `r17/dashboard-answers` off `main` `7f4901f1`, worktree `~/Projects/r13-lows-editor` (editor 5401, +test 5411, export 5410 — `pnpm wt list`'s block #24), one Opus implementer. Scratch files `D0-*` in the +job's `tmp`. The plan is Step 0b above. + +**What was found before building** (the corpus only read; one video dir's text copied to scratch for the +profile). +- **The walk did yield — on every file read.** Each video's unit awaits its reads, so the loop turned + between them; what the walk does ON the loop is parse. A CPU profile of `generateChannelSnapshot` over + 300 copies of one omnibased video dir's text (metadata 0.62 MB, cues 0.62 MB, two 2.9 MB VTTs): + `readNormalizedTranscript` (the cues parse, `readTranscriptCoverage`) 1,037 ms self, `readWebpageUrl` + 791 ms self, everything else in the walk under 70 ms. `readWebpageUrl` is the reconcile pass at the + walk's start (`reconcileVideoDirs.ts`): it reads and parses EVERY `metadata.info.json` for its + `webpage_url`, one at a time, on every regeneration. +- **One walk alone does not starve the pages.** Built and served by `next start` (loopback, a scratch + corpus of one 2,000-video channel, load average 24), with this slice: the regeneration took 26 s, and + `/` answered in 0.28–1.6 s and `/jobs` in 0.18–0.81 s, polled every 2 s throughout. Under `next dev` + at a load average of 32–35 the same walk took 118 s and the answers ran 0.3–28 s; with `main`'s five + files swapped back, 2.0–12.5 s. On this machine, under the dev server, one regeneration does not + separate the two, and an isolated tsx micro-benchmark (800 such dirs, three interleaved pairs, load + 25–35) put the chunked walk and `main`'s inside each other's noise (loop delay max 300–860 ms either + way). +- **So the hour-long outage needed more than one walk**, and the live metas show the rest: two walks side + by side (the empty queue key), the omnibased one on the platter through `onDrive` (four slots per + location, so every page read of that drive queued behind the walk's units), two regenerations of + omnibased four seconds apart in the restarted process (`01M3WHY8…`, `01M3WHYC…` — two passes racing + between the registry check and the record's registration), and the auto-queue status poll recomputing + behind all of it. D0 closes the three it owns: one queue, a per-slug dedup that covers the enqueue in + flight, and the status memo. The platter half is T1's (text reads leave `onDrive`). +- **The ghosts never held a move.** `channelWriters` reads this process's registry only, never a meta, + and the registry is in memory: a meta a dead process left `running` was never a writer. `/jobs` + listed such a job as `archived` (`listJobs`: a non-terminal meta not in the registry). Harmless to a + move, misleading on `/jobs`; the boot pass now closes them. +- **What a SIGTERM'd process leaves behind today** (`shutdownCancel.ts`, "graceful shutdown only"): the + reaper cancels every live job and re-raises the signal at once, without waiting. A queued job keeps + `queued` on disk on purpose (the boot pass settles it). A running managed job's terminal meta is + written only when its function returns (`streamCommand.ts`, the `.finally`); `generateChannelSnapshot` + takes no signal, so a regeneration never returns before the exit, and its meta stays `running`. A + SIGKILL leaves the same, and orphans children. + +**What it does.** +- **The walk yields between chunks.** `generateChannelSnapshot` maps its video dirs through + `mapInYieldingChunks` — chunks of `SNAPSHOT_YIELD_EVERY` (32: two full waves of the 16-wide limit), one + `setImmediate` between chunks. The unit body is untouched: the diff in that file is the helper, the + constant and the call's two lines, so T1's merge is trivial. +- **One regeneration at a time, once per channel.** `snapshotScheduler.ts` gains `REFRESH_REPORT_QUEUE` + (`"refresh-report"`: a non-empty key is the registry's existing concurrency-1 serialization, so no new + mechanism) and `startRefreshReport(paths, slug)`, the one entry point: a slug with a `refresh-report` + queued, running or being enqueued (a `starting` set on the scheduler's global state, held across + `runManagedFunction`'s awaits) is answered `info` with `REFRESH_REPORT_ACTIVE`. The debounced pass and + **Update all reports** (`refreshAllChannelSnapshotsAction`, `editor/app/channels/actions.ts` — not in + the slice's file list and owned by no other slice) both go through it, so they dedup against each + other; the action still reports a skip as "already running". The job body is the scheduler's (it bumps + `generation`, which the action's copy did not). +- **The auto-queue status poll shares one fold.** `common/views/autoQueueStatus.ts` gains + `singleFlightMemo` (concurrent callers share the computation in flight; a landed value is reused for + `AUTO_QUEUE_STATUS_MEMO_MS` = 3 s from when it LANDED; a rejection is not memoized; `clear()` detaches + an in-flight computation) and `autoQueueStatusMemo()` on `globalThis`. `buildAutoQueueStatusPayload` + is unchanged. The shell (`editor/app/operations/status.ts`) memoizes the snapshot-derived half only — + the channel briefs and the four lanes' `computeLeafPending` — and reads the priority view, the state + document, settings, the pool and the runners fresh. The e2e reset route drops the memo with the other + singletons. +- **Boot closes the running metas a dead process left.** A meta now records `pid` (`jobMeta.ts`, owned by + no slice). `bootQueuedJobs.ts` gains `settleRunningJobMetas`, run from `instrumentation.ts` on every + boot (idle and test server too: it re-queues nothing, and it does not wait for the storage pass): a + `running` meta from before this boot, not in the registry, whose writer is gone is closed `cancelled` + with `cancelReason` "interrupted: the process running it stopped before it finished", `endedAt` = its + log's mtime (the last moment it is known to have run), else now. **The status:** `cancelled` + + `cancelReason` already is this file's terminal state for "the server went down under it", and every + reader (listJobs, `/jobs`, the job page's "Cancelled because", Retry) handles it — so there is no new + `interrupted` status. **The dead-process test** (`writerIsGone`): metas carried no pid or host, and + "predates this boot" alone is not reliable here, because `archilyzer run` (`bin/run-operation.ts`) runs + jobs offline into the same `.jobs/`. Gone = no pid (written before this release), or this process's + own pid (a container's editor restarts as the same pid; this process's own metas are already excluded + by `bootedAt` and `isLive`), or `kill(pid, 0)` → ESRCH (EPERM counts as alive). A live pid is left + alone. The queued pass now applies the same check. `channelWriters.ts` gains a comment saying why a + ghost cannot reach it, and a test pins it. +- **e2e** `dashboard-answers.spec.ts` (two cases) and `/api/test/settle-running-metas` (POST, guarded by + `E2E_TEST_ROUTES`), which runs the running pass on demand: the e2e server boots once per suite. + (a) Two 600-video channels (one hardlinked ~400 KB `metadata.info.json` per dir, a Rumble archive so + each file is parsed twice) regenerated by Update all reports through the ops API while `/` and `/jobs` + are polled every 2 s: both reports land, the second job started when the first had ended (read from + their records), and every answer is inside the budget. (b) A `running` refresh-report meta with a dead + pid is closed as interrupted, one with a live pid (the test runner's) is left `running`, the job page + says "interrupted", and the channel's Move media completes. + +**Commits** + +| Commit | What | +|---|---| +| `8ce3b24d` | `common:` the chunked walk; `REFRESH_REPORT_QUEUE`, `startRefreshReport`, the `starting` set; `snapshotScheduler.test.ts` (2); Update all reports through it | +| `56efcd74` | `common, editor:` `singleFlightMemo` + `autoQueueStatusMemo` + 5 tests; the shell memoizes the snapshot half; the reset route drops it | +| `36babcad` | `common, editor:` meta `pid`; `writerIsGone`, `processIsAlive`, `settleRunningJobMetas`; the queued pass's writer check; boot wiring; 4 boot tests + 1 `channelWriters` test | +| `fef9ebc2` | `editor(e2e):` `dashboard-answers.spec.ts`; `/api/test/settle-running-metas` | +| `fe1676c3` | `editor(e2e):` the spec reworked after run 1 (two 600-video channels, the serial check, the idle-floored budget) | +| `3c90b8e9` | `plans:` this section; the editor changelog | + +#### Gates (logs `$T/D0-*.log`) + +- **tsc** (all workspaces) clean at every commit. +- **common:** **2,496/2,496** — new: `snapshotScheduler.test.ts` 2, `autoQueueStatus.test.ts` +5, + `bootQueuedJobs.test.ts` +4, `channelWriters.test.ts` +1. **Editor unit:** 109/109. + **test:scripts:** run 1: 301 passed, 1 failed, 2 skipped (304); rerun: 300 passed, 2 failed, 2 + skipped. Every failure is in `queue-lock.test.mjs` ("prints a banner naming the holder while waiting", + "serves waiters in arrival order (FIFO)"), timing cases run at load averages of 16–30; the file run + alone passes 11/11, and the slice touches nothing under `scripts/`. +- **Build:** the capped editor build with the corpus linked (`ln -sT`, `systemd-run --scope -p + MemoryMax=6G`, the link removed after): exit 0, 123 s, at `fef9ebc2`. +- **e2e** (`jobs`, `channels`, `channel-storage`, `dashboard-answers`; `next dev`): + - run 1 at `fef9ebc2`: **23 passed, 2 failed, 16.7 min**. `dashboard-answers` (a) timed out at 5 min: + 2,000 videos under `next dev` at a load average of ~30 did not finish (the scratch-server measurement + above: 118 s for the walk alone, every answer slower). `channel-storage` "relocate a channel's media + to another root, and move it back" — the run's first test, on a cold dev server — timed out at 90 s + on its first request to the file route. + - run 2 at `fe1676c3`: **25 passed, 0 failed, 3.7 min** (34 min with the queue wait). (a) took 23.8 s: + nine polls of each page while the two regenerations ran, `/` 0.14–1.9 s and `/jobs` 0.34–1.6 s + (budget 5 s); (b) 7.8 s. The relocate case passed — run 1's failure was the cold first test. +- **Numbers tool:** none. + +**Deviations from the plan, one sentence each.** +- The chunk is 32, not 25, so every wave of the 16-wide limit is full (25 left each chunk's second wave + nine wide). +- The memo covers the shell's snapshot-derived half, not the whole payload: memoizing it all would hold a + lane's hold, runner and picks up to 3 s behind a click, and specs read them straight after one. +- The running-meta pass is in `bootQueuedJobs.ts` beside the queued one, not in `registry.ts`, which is + in memory and holds no metas. +- "Interrupted" is `cancelled` + `cancelReason` (the existing terminal state for a job the server went + down under), not a new status. +- `channelWriters.ts` gets a comment and a test, no code: a ghost cannot reach it. +- The e2e budget is 5 s or three times the slowest idle answer measured just before, whichever is longer: + the dev server on this shared machine at load 30 takes 2–4 s a page with nothing regenerating, and the + spec pins starvation, not the machine's speed. +- The spec regenerates two 600-video channels, not one large one: 2,000 did not finish in five minutes + under the dev server here, and two prove the serial queue from the jobs' records. + +**Found and left.** +- The reconcile pass parses every `metadata.info.json` in full, on every regeneration, for one field — the + largest single cost of a walk. Reading only the head, or skipping a dir whose name is already canonical, + is a follow-up (`reconcileVideoDirs.ts` is not D0's). +- `generateChannelSnapshot` takes no abort signal, so a Cancel or a graceful shutdown cannot stop a walk; + the chunk boundary is now the natural place to check one — a follow-up. +- A meta whose pid has since been reused by an unrelated live process stays `running` on `/jobs` until + that process exits; a start-time check (`/proc/<pid>/stat`) would close it. Not worth it today. +- `editor/app/jobs/[id]/page.tsx` and `JobRow.tsx` describe the cancel reason as the one "a restart left + queued"; it now also covers a dead process's running job. Comments only, left. + +**Changelog** (`editor/CHANGELOG.md` `[Unreleased]`): three bullets — the dashboard and `/jobs` answer +while reports regenerate (one at a time, deduped; Update all reports takes the sum), the operations pages +share one pending count (up to 3 s old; holds, runners and picks fresh), and jobs a stopped editor left +"running" are closed at start. + +**Rollout note.** No `archilyzer run` may be in flight across the first editor restart after this ships: +every meta written before it has no `pid`, so the boot pass reads its writer as gone and would close a +still-running offline job's meta as interrupted (the live job rewrites it at its next meta write, but +`/jobs` offers Retry meanwhile). The same holds across pid namespaces — an editor in a container judging +a pid a host-side `archilyzer run` wrote, or the reverse, sees ESRCH. The release's rollout stops the +editor for the migration anyway. + +#### Review (SHIP AFTER FIXES) and the fixes + +| Finding | Fix | +|---|---| +| H1 — the memo served pending counts folded from settings (lane policy and tree, `channelPriority`) up to 3 s stale against a fresh tree; the payload-reading specs were not run | `5b236f8c`: `singleFlightMemo.get(key, compute)`; the shell keys it on `JSON.stringify([settings.autoQueue, settings.channelPriority])`, time the only other expiry; only the latest-started computation stores (an old key landing late never overwrites); comments say what can be one poll late (a rewritten snapshot, the runner's in-flight set) and what is fresh; +1 test. The status-reading specs ran in the full suite below | +| L2 — a change during a RUNNING regeneration was dropped until the next change (as on main) | `616a9ce2`: dedup only against a queued or starting regeneration (`isRefreshReportPending`), so a running walk gets one queued successor and no second; fire()'s trailing-edge comment is now true; +1 test | +| L3 — the per-channel Refresh report walked in the request, outside the queue | `5509ccf4`: `refreshChannelSnapshotAction` (Refresh report, ops `{slug}`, e2e `generateReport`) starts a job through `startRefreshReport` and awaits it — or polls the one already queued for the channel — and returns the walk's error sentence (`onError`) as its `{error}`; the ops route's comment follows | +| L4 — a queued refresh-report holds the Storage panel's Move (`mediaBusy.ts` counts queued jobs) | Not changed in D0, as ruled: T1's `mediaOnly` (refresh-report does not need media) clears it when T2's courtesy check passes it; to confirm at the T1/T2 merges | +| L3 follow-on — `generateReport` (85 specs) now queues a job, which `/jobs` may draw before its default filter hides it | `7ba8cd41`: `bulk-actions.spec` "queues no job" leaves `refresh-report` rows out; `dashboard-answers` compiles the ops route (a 400 on `{}`) before measuring — its first compile held `/` 15 s once under load | +| L5 — nothing pinned the yield; the e2e floor could grow without bound | `fea5b203`: `snapshotYield.test.ts` — a `setImmediate` probe queued during the first chunk runs before the second chunk's first unit (it fails with the yield removed, checked); 2 tests. The e2e budget is `min(max(5 s, 3 × slowest idle), 15 s)`, and the spec header says it pins the serial queue and gross starvation, not the yield | +| L6 — the first boot after this ships closes a live pre-slice `archilyzer run`'s meta | The rollout note above | +| N7 — `REFRESH_REPORT_QUEUE` beside the other queue keys | `616a9ce2`: defined in `lib/queueKeys.ts`, re-exported by the scheduler | +| N8 — `operations/[id]/page.tsx`'s census comment claimed the payload's listing | `5b236f8c`: says it is the same listing only on a memo miss | +| N9 — `requestCache.ts` said "no exception" | `5b236f8c`: one paragraph naming the memo and its bound | +| N10 — `channelWriters.test.ts` conflicts with T1 (both append) | Keep both, at the merge | +| the record and changelog | this commit: this subsection, the rollout note; the changelog's first two bullets say Refresh report waits its turn, a change during a regeneration queues one more, and a rules/focus/priority edit recounts at once | + +**Gates after the fixes.** tsc clean at every commit. **common 2,500/2,500** (+4: memo key 1, the +successor 1, the yield 2). **Editor unit 109/109.** **e2e** — L3 makes every `generateReport` (85 specs) a +queued job, so the whole editor suite rather than the two lists (it contains both), in three runs on a +machine that ran out of memory (15 GB used, swap 19/19 GB; a parakeet transcription of the live editor +beside several suites): + - run 3, the full suite (702 tests) at `fea5b203`: stopped at 374 passed, 8 failed, ~50 min, when + pages began to crash (`page.goto: Page crashed`). Two failures were D0's and are fixed in `7ba8cd41` + (`bulk-actions` "queues no job"; `dashboard-answers` (a): one `/` of 14.9 s at the moment the ops + route compiled, every other answer ≤ 4.6 s). The other six passed in run 4. + - run 4 at `7ba8cd41` (run 3's failures plus every spec file from `new-channel-onboarding` on, and both + lists — 69 files): **295 passed, 154 failed, 35.1 min** (54 with the queue). Every spec in both lists + passed — `jobs` 2, `channels` 8, `channel-storage` 13, `dashboard-answers` 2, `focus-banner` 3, + `auto-queue` 24, `lane-runner` 5, `channel-priority` 9, `operation-settings` 7, `backfill` 20, + `channel-work` 11, `ops-api` 23 — except `view-route`, which ran after a machine-wide OOM kill at + 22:41 took the test server (the kernel log names it, among browser tabs, the desktop session and the + live editor's transcriber): every test from the 300th on failed in under 2 s against a dead server. + Before it, four failures in `site-scope` (2), `sites-crud` and `social-channel` (the known 22–26 s + fetch-posts case). + - run 5 at `7ba8cd41`, the 25 spec files with a failure in run 4 (`view-route` included): **183 passed, 0 + failed, 12.2 min** (13 with the queue). + So every editor spec has passed with the fixes in, across runs 3–5, and both lists in full. + +**test:scripts** was not re-run: the fixes touch nothing under `scripts/`. + +#### Re-review (SHIP AFTER FIXES) and the fixes + +| Finding | Fix | +|---|---| +| R1 — Refresh report and both ops forms waited for the whole serial queue, with no bound and no feedback (the CLI's fetch gives up at 300 s; the button dropped the action's result) | `5287e8c2`: `requestRefreshReport` (start, or find the queued job) and `waitForRefreshReport` (up to `REFRESH_REPORT_WAIT_MS` = 15 s → done, failed, or waiting with the regenerations ahead), `refreshReportWaitNotice`; `e491ec14`: the action returns `{ notice }` ("Queued behind N report regenerations — the report updates when it finishes (job …).") past the bound, `RefreshSnapshotButton` draws an error and the notice, `InlineActionButton` shows the notice neutrally, and Update all reports answers once queued; `7ae184bd`: ops `{ slug }` returns `{ ok, jobId, started }` (404 for an unknown channel) and `{ all }` returns the ids at once, as `_lib.ts` rules — `--wait` follows them | +| R2 — the "already queued" branch reported success when that job failed, or when the slug was mid-enqueue | `5287e8c2`: the wait reads any job's end — a failure returns the job log's `[error]` sentence, whoever started it; a slug mid-enqueue is waited for (2 s) until its id exists; +2 tests | +| R3 — `started.stream` left open | `5287e8c2`: `requestRefreshReport` cancels every stream it starts | +| R4 — run 4's count | this commit: 154 failed, not 68 (the log's tally) | +| the changelog | this commit: bullet 1 says what a long wait shows | + +**Gates after the re-review fixes.** tsc clean at every commit; **common 2,502/2,502** (+2); **editor unit +109/109**; e2e `dashboard-answers`, `channels`, `ops-api`, `channel-work` (only these, as asked): run 6 at `7ae184bd`: **44 passed, 0 failed, 4.7 min** +(dashboard-answers (a): 30 polls of each page while the two regenerations ran, every answer ≤ 3.1 s). + +**Merge of `main` `1d5c33bf`** (U1, export 0.11.1, XP) at `52446024`: two conflicts, both appends — +this file keeps U1's and XP's record sections before D0's under `## Record`, and `editor/CHANGELOG.md` +keeps every `[Unreleased]` bullet (main's, then D0's). Re-gated on the merged tree: tsc clean; **common +2,519/2,519**; **editor unit 109/109**; e2e `dashboard-answers`, `channels`, `jobs`, `channel-storage` +(run 7): **25 passed, 0 failed, 6.4 min** (23 with the queue wait; every dashboard answer ≤ 3.7 s). + ## Rollout