Archilyzer · Source

archilyzer

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

commit 7f261121c96fdc48a8c5e9415632fd77ec54bd2b
parent 4afbc3b72204d4b4172fdf153da8b5c06b1a2252
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri, 25 Sep 2026 19:29:41 -0400

Merge main (release 9, 4386f369) into one-core/c1-hub-look

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

Diffstat:
Mcommon/controller/autoRunner.ts | 12++++++++----
Mcommon/controller/videoTitles.test.ts | 28++++++++++++++++++++++------
Mcommon/controller/videoTitles.ts | 25++++++++++++++++++++++---
Acommon/jobs/bootQueuedJobs.test.ts | 338+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/bootQueuedJobs.ts | 243+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/downloadBackoff.ts | 14++++++++------
Mcommon/jobs/jobMeta.ts | 4++++
Mcommon/jobs/platformBackoff.test.ts | 19+++++++++++++++++++
Mcommon/jobs/platformBackoff.ts | 13+++++++++++++
Mcommon/jobs/registry.ts | 27+++++++++++++++++++++++++++
Mcommon/jobs/streamCommand.ts | 3++-
Mcommon/lib/homepageSummary.test.ts | 42+++++++++++++++++++++++++++++++++++++++++-
Mcommon/lib/homepageSummary.ts | 34++++++++++++++++++++++++++++------
Mcommon/views/jobRows.test.ts | 54+++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/views/jobRows.ts | 10++++++++--
Mcommon/ytdlp/channelArgs.test.ts | 25++++++++++++++++++++++---
Acommon/ytdlp/managedDownloadsSleep.test.ts | 165+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/ytdlp/platformArgs.mjs | 16++++++++++++++++
Mcommon/ytdlp/runYtdlp.ts | 43+++++++++++++++++++++++++++++++++++++++----
Meditor/CHANGELOG.md | 6++++++
Meditor/app/api/test/stuck-job/route.ts | 28+++++++++++++++++++++++++++-
Meditor/app/channels/[slug]/availabilityActions.ts | 8++++----
Meditor/app/channels/[slug]/backfillActions.ts | 4++--
Meditor/app/channels/[slug]/digestActions.ts | 5+++--
Meditor/app/channels/[slug]/incompleteTranscriptActions.ts | 5+++--
Meditor/app/channels/[slug]/normalizeActions.ts | 4++--
Meditor/app/channels/[slug]/persistActions.ts | 5++---
Meditor/app/channels/[slug]/pipelineActions.ts | 27+++++++++++++++------------
Meditor/app/channels/[slug]/socialActions.ts | 5+++--
Meditor/app/channels/[slug]/videos/[id]/operationActions.ts | 8+++++---
Meditor/app/channels/[slug]/videos/[id]/videoActions.ts | 31++++++++++++++++++++-----------
Meditor/app/channels/[slug]/whisperActions.ts | 22+++++++++++-----------
Meditor/app/channels/groupActions.ts | 4++--
Meditor/app/channels/lib/relocationJob.ts | 6++----
Aeditor/app/lib/safeRevalidate.test.ts | 61+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/lib/safeRevalidate.ts | 82+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/saved-videos/backupActions.ts | 3++-
Meditor/app/sites/lib/buildAction.ts | 4++--
Meditor/app/storage/lib/evictClipsJob.ts | 6+++---
Meditor/app/storage/lib/repointJob.ts | 6++----
Meditor/app/storage/lib/savedVideosJob.ts | 5++---
Meditor/e2e/ops-api.spec.ts | 59+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/instrumentation.ts | 58+++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mhomepage/app/components/ArchiveGrowthChart.tsx | 3++-
Mhomepage/app/components/FamilyStats.tsx | 9++++++++-
Mhomepage/app/page.tsx | 4++--
Mplans/STATE.md | 10++++++++++
Mplans/release-9.md | 232++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Mumtool/report-to-video/platform-args.test.mjs | 17++++++++++++++---
49 files changed, 1717 insertions(+), 125 deletions(-)

diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts @@ -72,6 +72,7 @@ import { import { isCoolingDown, isVideoDeferred, + mergeBackoffEntry, pruneDeferred, pruneExpired, } from "../jobs/platformBackoff"; @@ -1612,14 +1613,17 @@ async function runLoop( if (kind === "download") { // Merge in any cooldown a manual sync/import wrote to the shared state // (read-modify-write from outside the runner) since our last persist. - // Take the later `until` so a sync-set 429 cooldown is honored even though - // our own next persist would otherwise clobber it. + // `until` and `fails` each merge to their max, so a sync-set 429 cooldown + // is honored even though our own next persist would otherwise clobber it, + // and neither side's escalation count is lost. try { const persisted = (await readAutoQueueState(paths)).download .platformBackoff; for (const [pf, e] of Object.entries(persisted)) { - const cur = kindState.platformBackoff[pf]; - if (!cur || e.until > cur.until) kindState.platformBackoff[pf] = e; + kindState.platformBackoff[pf] = mergeBackoffEntry( + kindState.platformBackoff[pf], + e, + ); } } catch { // Best-effort: a transient read failure just skips this iteration's merge. diff --git a/common/controller/videoTitles.test.ts b/common/controller/videoTitles.test.ts @@ -289,22 +289,38 @@ test("metadata titles are memoized per channel until data/ changes", async () => assert.equal(second.get("b")?.title, "Bravo"); assert.equal(videoTitleMetadataReadCount() - r1, 0); - // A new video dir changes data/'s mtime: the channel's memo is dropped and - // every id is read again ("a" now has no file, so no title). + // A new video dir changes data/'s mtime: the memo is REBASED — the dirs + // still there keep their titles and only the new one is read (release 8 + // review, V). "a"'s dir is still there, so its memoized title stands. await writeInfo(paths, "c", { id: "c", title: "Charlie" }); const r2 = videoTitleMetadataReadCount(); const third = await readChannelVideoTitles(paths, SLUG, ["a", "b", "c"]); - assert.equal(third.has("a"), false); + assert.equal(third.get("a")?.title, "Alpha"); + assert.equal(third.get("b")?.title, "Bravo"); assert.equal(third.get("c")?.title, "Charlie"); - assert.equal(videoTitleMetadataReadCount() - r2, 3); + assert.equal(videoTitleMetadataReadCount() - r2, 1); + + // A REMOVED dir takes its title with it: the rebase keeps only the dirs + // still present, so "a" falls through, is read, and has no title. + await rm(path.join(dataDir, "a"), { recursive: true, force: true }); + // Two changes inside one filesystem timestamp tick would share an mtime; + // pin a distinct one so the rebase certainly fires. + const later = new Date("2026-01-02T00:00:00Z"); + await utimes(dataDir, later, later); + const r3 = videoTitleMetadataReadCount(); + const fourth = await readChannelVideoTitles(paths, SLUG, ["a", "b", "c"]); + assert.equal(fourth.has("a"), false); + assert.equal(fourth.get("c")?.title, "Charlie"); + assert.equal(videoTitleMetadataReadCount() - r3, 1); // A miss is not memoized: a title that lands later is found. + await mkdir(path.join(dataDir, "a"), { recursive: true }); await writeFile( path.join(dataDir, "a", "metadata.info.json"), JSON.stringify({ id: "a", title: "Alpha again" }), ); - const fourth = await readChannelVideoTitles(paths, SLUG, ["a"]); - assert.equal(fourth.get("a")?.title, "Alpha again"); + const fifth = await readChannelVideoTitles(paths, SLUG, ["a"]); + assert.equal(fifth.get("a")?.title, "Alpha again"); } finally { resetVideoTitleMemo(); await rm(dir, { recursive: true, force: true }); diff --git a/common/controller/videoTitles.ts b/common/controller/videoTitles.ts @@ -29,7 +29,7 @@ // invariant). import { existsSync } from "node:fs"; -import { open as openFile, stat } from "node:fs/promises"; +import { open as openFile, readdir, stat } from "node:fs/promises"; import path from "node:path"; import { open } from "lmdb"; import type { Paths } from "../lib/paths"; @@ -156,7 +156,10 @@ async function readMetadataTitle(videoDir: string): Promise<string | null> { // 8,049 such dirs: 5,024 ms cold / 224 ms warm per render, measured 2026-09-25. // // Per channel, keyed by the `data/` directory's mtime: adding or removing a -// video dir changes it and drops the channel's entry. Rewriting a +// video dir changes it, and the entry is then REBASED rather than dropped — the +// titles of the dirs still there carry over (one readdir, no file reads), and +// only the new dirs are read (release 8 review, V: a busy download lane adds a +// dir every few minutes, and each one used to re-read all 8,049). Rewriting a // metadata.info.json INSIDE an existing dir does not — accepted, because a // video's title does not change after download. Misses (no file yet, e.g. a dir // mid-download) are NOT memoized, so a title that lands later is picked up. @@ -196,8 +199,24 @@ async function memoForChannel( } const memo = titleMemo(); let entry = memo.get(key); - if (!entry || entry.dataDirMtimeMs !== mtimeMs) { + if (!entry) { entry = { dataDirMtimeMs: mtimeMs, titles: new Map() }; + } else if (entry.dataDirMtimeMs !== mtimeMs) { + // Keep the titles of the dirs that are still there; a removed dir's title + // goes with it, so a re-listed id falls through to the other sources. + let present: Set<string> | null = null; + try { + present = new Set(await readdir(dataDir)); + } catch { + present = null; + } + const kept = new Map<string, string>(); + if (present) { + for (const [id, title] of entry.titles) { + if (present.has(id)) kept.set(id, title); + } + } + entry = { dataDirMtimeMs: mtimeMs, titles: kept }; } // Re-insert so iteration order is least-recently-used first. memo.delete(key); diff --git a/common/jobs/bootQueuedJobs.test.ts b/common/jobs/bootQueuedJobs.test.ts @@ -0,0 +1,338 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mkdtemp, mkdir, readFile, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import type { JobMeta } from "./jobMeta"; +import type { JobSpec } from "./jobSpec"; +import { settleQueuedJobMetas, type RequeueFn } from "./bootQueuedJobs"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/bootQueuedJobs.test.ts +// +// THE BOOT PASS (release 9, B4b), over a temp .jobs dir with an injected +// re-queue. It re-queues LITTLE: never a sync, never a meta older than 24 h, +// only the newest of duplicate specs; everything else is closed `cancelled` +// with the reason. An idle boot only ever cancels. Never run against a real +// corpus. + +const BOOT = 2_000_000_000_000; + +async function fixture(metas: Partial<JobMeta>[]): Promise<{ + paths: Paths; + root: string; + read: (id: string) => Promise<JobMeta>; + log: (id: string) => Promise<string>; +}> { + const root = await mkdtemp(path.join(tmpdir(), "boot-queued-")); + const jobsDir = path.join(root, ".jobs"); + await mkdir(jobsDir, { recursive: true }); + for (const m of metas) { + const meta = { + kind: "whisper-all", + queueKey: "platform:youtube", + status: "queued", + queuedAt: BOOT - 60_000, + ...m, + }; + await writeFile( + path.join(jobsDir, `${meta.id}.meta.json`), + JSON.stringify(meta), + ); + } + return { + paths: { jobsDir } as Paths, + root, + read: async (id) => + JSON.parse( + await readFile(path.join(jobsDir, `${id}.meta.json`), "utf8"), + ) as JobMeta, + log: (id) => + readFile(path.join(jobsDir, `${id}.log`), "utf8").catch(() => ""), + }; +} + +const SPEC: JobSpec = { kind: "whisper-all", slug: "teamrcn", params: {} }; +const SYNC: JobSpec = { + kind: "sync", + slug: "teamrcn", + params: { fullSweep: false }, +}; +const HOUR = 60 * 60 * 1000; + +function recordingRequeue(): { fn: RequeueFn; calls: JobSpec[] } { + const calls: JobSpec[] = []; + let n = 0; + return { + calls, + fn: async (spec) => { + calls.push(spec); + return { ok: true, jobId: `NEW${++n}` }; + }, + }; +} + +test("a fresh non-sync meta with a spec is re-queued, and the old one is closed naming the new id", async () => { + const f = await fixture([{ id: "A1", spec: SPEC, channelSlug: "teamrcn" }]); + try { + const rq = recordingRequeue(); + const lines: string[] = []; + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: rq.fn, + bootedAt: BOOT, + log: (l) => lines.push(l), + }); + assert.deepEqual(rq.calls, [SPEC]); + assert.deepEqual(res.requeued, [{ id: "A1", newId: "NEW1" }]); + assert.deepEqual(res.cancelled, []); + const meta = await f.read("A1"); + assert.equal(meta.status, "cancelled"); + assert.match( + meta.cancelReason ?? "", + /server restarted before it ran; re-queued as NEW1/, + ); + assert.ok(typeof meta.endedAt === "number"); + assert.match(await f.log("A1"), /re-queued as NEW1/); + // One line per re-queued job, then the summary. + assert.equal(lines.length, 2); + assert.match(lines[0], /re-queued whisper-all teamrcn: A1 -> NEW1/); + assert.match(lines[1], /re-queued 1, cancelled 0/); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +// The review's sizing of the live corpus: 346 queued syncs over 51 channels, +// ~7 per channel, plus months-old whisper-all sweeps. Exactly one job may come +// back from this. +test("7 duplicate syncs + 1 stale whisper-all + 1 fresh non-sync: exactly one re-queue", async () => { + const syncs = Array.from({ length: 7 }, (_, i) => ({ + id: `S${i}`, + kind: "sync", + spec: SYNC, + channelSlug: "teamrcn", + queuedAt: BOOT - (i + 1) * HOUR, + })); + const f = await fixture([ + ...syncs, + { + id: "W1", + spec: { kind: "whisper-all", slug: "old", params: {} }, + queuedAt: BOOT - 60 * 24 * HOUR, + }, + { id: "F1", spec: SPEC, queuedAt: BOOT - 2 * HOUR }, + ]); + try { + const rq = recordingRequeue(); + const lines: string[] = []; + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: rq.fn, + bootedAt: BOOT, + log: (l) => lines.push(l), + }); + assert.deepEqual(rq.calls, [SPEC]); + assert.deepEqual(res.requeued, [{ id: "F1", newId: "NEW1" }]); + const byCat = (c: string) => + res.cancelled.filter((x) => x.category === c).length; + assert.equal(byCat("sync"), 7); + assert.equal(byCat("stale"), 1); + assert.match( + (await f.read("S0")).cancelReason ?? "", + /the scheduler re-derives syncs/, + ); + assert.match( + (await f.read("W1")).cancelReason ?? "", + /queued before the last restart, stale/, + ); + assert.match( + lines.at(-1) ?? "", + /re-queued 1, cancelled 8 \(sync 7, stale 1\)/, + ); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("of duplicate specs only the newest is re-queued; params are compared key-order-free", async () => { + const f = await fixture([ + { + id: "D1", + spec: { kind: "whisper-all", slug: "c", params: { a: 1, b: 2 } }, + queuedAt: BOOT - 3 * HOUR, + }, + { + id: "D2", + spec: { kind: "whisper-all", slug: "c", params: { b: 2, a: 1 } }, + queuedAt: BOOT - 1 * HOUR, + }, + { + id: "D3", + spec: { kind: "whisper-all", slug: "c", params: { a: 9 } }, + queuedAt: BOOT - 2 * HOUR, + }, + ]); + try { + const rq = recordingRequeue(); + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: rq.fn, + bootedAt: BOOT, + }); + assert.deepEqual(res.requeued.map((r) => r.id).sort(), ["D2", "D3"]); + assert.deepEqual( + res.cancelled.map((c) => [c.id, c.category]), + [["D1", "superseded"]], + ); + assert.match( + (await f.read("D1")).cancelReason ?? "", + /superseded by a newer queued job \(D2\)/, + ); + // Survivors are re-queued oldest first, keeping their relative order. + assert.deepEqual( + res.requeued.map((r) => r.id), + ["D3", "D2"], + ); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("a queued meta with no spec is cancelled with the reason, not re-queued", async () => { + const f = await fixture([{ id: "B1", kind: "build-export" }]); + try { + const rq = recordingRequeue(); + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: rq.fn, + bootedAt: BOOT, + }); + assert.equal(rq.calls.length, 0); + assert.deepEqual( + res.cancelled.map((c) => c.category), + ["no-spec"], + ); + const meta = await f.read("B1"); + assert.equal(meta.status, "cancelled"); + assert.match(meta.cancelReason ?? "", /^server restarted before it ran/); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("an idle boot cancels every queued meta and re-queues none", async () => { + const f = await fixture([{ id: "C1", spec: SPEC }, { id: "C2" }]); + try { + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: null, + bootedAt: BOOT, + idleReason: "idle boot", + }); + assert.deepEqual(res.requeued, []); + assert.deepEqual(res.cancelled.map((c) => c.id).sort(), ["C1", "C2"]); + assert.match( + (await f.read("C1")).cancelReason ?? "", + /idle boot: not re-queued/, + ); + assert.equal((await f.read("C2")).status, "cancelled"); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("terminal, running, this-boot and live metas are left alone", async () => { + const f = await fixture([ + { id: "D1", status: "done", spec: SPEC }, + { id: "D2", status: "running", spec: SPEC }, + { id: "D3", spec: SPEC, queuedAt: BOOT + 5 }, + { id: "D4", spec: SPEC }, + ]); + try { + const rq = recordingRequeue(); + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: rq.fn, + bootedAt: BOOT, + isLive: (id) => id === "D4", + }); + assert.equal(rq.calls.length, 0); + assert.deepEqual(res, { requeued: [], cancelled: [] }); + assert.equal((await f.read("D1")).status, "done"); + assert.equal((await f.read("D2")).status, "running"); + assert.equal((await f.read("D3")).status, "queued"); + assert.equal((await f.read("D4")).status, "queued"); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("a malformed meta is skipped and the pass continues", async () => { + const f = await fixture([{ id: "M2", spec: SPEC }]); + try { + await writeFile( + path.join(f.paths.jobsDir, "M1.meta.json"), + '{"id":"M1","kind":', + ); + await writeFile( + path.join(f.paths.jobsDir, "M0.meta.json"), + '{"status":"queued"}', + ); + const rq = recordingRequeue(); + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: rq.fn, + bootedAt: BOOT, + }); + assert.deepEqual(res.requeued, [{ id: "M2", newId: "NEW1" }]); + assert.deepEqual(res.cancelled, []); + assert.equal( + await readFile(path.join(f.paths.jobsDir, "M1.meta.json"), "utf8"), + '{"id":"M1","kind":', + ); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("a refused or throwing re-queue closes the meta with the error", async () => { + const f = await fixture([ + { id: "E1", spec: { kind: "whisper-all", slug: "a", params: {} } }, + { id: "E2", spec: { kind: "whisper-all", slug: "b", params: {} } }, + ]); + try { + let calls = 0; + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: async () => { + calls++; + if (calls === 1) { + return { ok: false, error: 'Channel "teamrcn" not found' }; + } + throw new Error("boom"); + }, + bootedAt: BOOT, + }); + assert.equal(res.cancelled.length, 2); + assert.ok(res.cancelled.every((c) => c.category === "refused")); + const reasons = [ + (await f.read("E1")).cancelReason ?? "", + (await f.read("E2")).cancelReason ?? "", + ].join("\n"); + assert.match(reasons, /re-queue refused: Channel "teamrcn" not found/); + assert.match(reasons, /re-queue refused: boom/); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("a missing .jobs dir is nothing to do", async () => { + const res = await settleQueuedJobMetas({ + paths: { jobsDir: path.join(tmpdir(), "no-such-jobs-dir-r9") } as Paths, + requeue: null, + bootedAt: BOOT, + }); + assert.deepEqual(res, { requeued: [], cancelled: [] }); +}); diff --git a/common/jobs/bootQueuedJobs.ts b/common/jobs/bootQueuedJobs.ts @@ -0,0 +1,243 @@ +import path from "node:path"; +import { appendFile, readdir } from "node:fs/promises"; +import { writeFileAtomic } from "../lib/jsonFile-server"; +import type { Paths } from "../lib/paths"; +import { metaPath, readJobMeta, type JobMeta } from "./jobMeta"; +import type { JobSpec } from "./jobSpec"; + +// THE BOOT PASS OVER STALE `queued` METAS (release 9, B4b). +// +// The registry and the scheduler are in memory; a job's `<id>.meta.json` is +// written `queued` when it is submitted and rewritten when it finishes. A job +// still waiting in its queue when the server stopped therefore never gets its +// terminal write, and /jobs reads that sidecar forever as `queued` — a job that +// will never run, in a list the operator trusts to say what is pending. +// +// One pass at boot settles each of them, exactly once, and RE-QUEUES LITTLE. +// The first version re-queued every `queued` meta with a spec. The live corpus +// held 392 (review, 2026-09-25): 346 `sync` over 51 channels, 29 months-old +// `whisper-all` sweeps and a few others. That would have been ~382 jobs +// submitted at boot, ahead of the auto-download runner and around the +// heartbeat's maxConcurrentSyncs stagger: the burst the pacing work forbids. +// So, in order: +// 1. an idle boot (`requeue: null`) cancels everything — it must not resume +// work; +// 2. `sync` is never re-queued: the heartbeat re-derives due syncs itself, +// paced; +// 3. a meta queued more than REQUEUE_MAX_AGE_MS before this boot is stale; +// 4. a meta with no replay spec cannot be re-queued; +// 5. of what is left, only the NEWEST meta per kind + channel + params is +// re-queued (through the same path Retry uses); the rest are superseded. +// Every meta this pass does not re-queue is closed `cancelled` with a +// `cancelReason`; a re-queued one is closed naming its new id. +// +// Only metas from BEFORE this boot, and not held by the live registry, are +// touched: the storage boot pass can enqueue a job of its own while this runs. +// `running` metas are left alone — whether a job that was mid-flight should be +// re-run is not a decision a boot pass can make. A malformed meta (readJobMeta +// → null) is skipped. +// +// Best-effort throughout: a meta that cannot be read or written is skipped, +// and the caller voids the promise so readiness never waits on it. + +export type RequeueFn = ( + spec: JobSpec, +) => Promise<{ ok: true; jobId: string } | { ok: false; error: string }>; + +export type CancelCategory = + "idle" | "sync" | "stale" | "no-spec" | "superseded" | "refused"; + +export type BootQueuedResult = { + requeued: { id: string; newId: string }[]; + cancelled: { id: string; reason: string; category: CancelCategory }[]; +}; + +export const RESTART_REASON = "server restarted before it ran"; +export const REQUEUE_MAX_AGE_MS = 24 * 60 * 60 * 1000; + +// Key a spec by what it would DO: kind + channel + bucket + params, with the +// params serialized key-sorted so two equal specs written in a different key +// order collide. +function stableJson(v: unknown): string { + if (Array.isArray(v)) return `[${v.map(stableJson).join(",")}]`; + if (v && typeof v === "object") { + const o = v as Record<string, unknown>; + return `{${Object.keys(o) + .sort() + .filter((k) => o[k] !== undefined) + .map((k) => `${JSON.stringify(k)}:${stableJson(o[k])}`) + .join(",")}}`; + } + return JSON.stringify(v) ?? "null"; +} + +export function specKey(spec: JobSpec): string { + return stableJson({ + kind: spec.kind, + slug: spec.slug, + bucket: spec.bucket, + params: spec.params ?? {}, + }); +} + +export async function settleQueuedJobMetas(opts: { + paths: Paths; + // null: cancel, never re-queue (an idle boot, or the e2e test server). + requeue: RequeueFn | null; + // Metas queued at or after this instant belong to THIS process. + bootedAt: number; + // True for an id the live registry holds (enqueued during this boot). + isLive?: (id: string) => boolean; + log?: (line: string) => void; + idleReason?: string; + maxAgeMs?: number; +}): Promise<BootQueuedResult> { + const log = opts.log ?? (() => {}); + const maxAgeMs = opts.maxAgeMs ?? REQUEUE_MAX_AGE_MS; + const result: BootQueuedResult = { requeued: [], cancelled: [] }; + let names: string[]; + try { + names = await readdir(opts.paths.jobsDir); + } catch { + return result; + } + const ids = names + .filter((n) => n.endsWith(".meta.json")) + .map((n) => n.slice(0, -".meta.json".length)) + .sort(); + + const stale: JobMeta[] = []; + for (const id of ids) { + const meta = await readJobMeta(opts.paths, id); + if (!meta || meta.status !== "queued") continue; + if (typeof meta.queuedAt === "number" && meta.queuedAt >= opts.bootedAt) { + continue; + } + if (opts.isLive?.(id)) continue; + stale.push(meta); + } + + const cancel = async ( + meta: JobMeta, + category: CancelCategory, + reason: string, + ) => { + if (await closeMeta(opts.paths, meta, reason)) { + result.cancelled.push({ id: meta.id, reason, category }); + } + }; + + // Newest first, so the first meta seen for a spec key is the one kept. + stale.sort( + (a, b) => (b.queuedAt ?? 0) - (a.queuedAt ?? 0) || b.id.localeCompare(a.id), + ); + const kept = new Map<string, string>(); // spec key -> kept meta id + const toRequeue: { meta: JobMeta; spec: JobSpec }[] = []; + for (const meta of stale) { + if (!opts.requeue) { + await cancel( + meta, + "idle", + `${RESTART_REASON}${opts.idleReason ? ` (${opts.idleReason}: not re-queued)` : ""}`, + ); + } else if (meta.kind === "sync") { + await cancel( + meta, + "sync", + "server restarted; the scheduler re-derives syncs", + ); + } else if ( + typeof meta.queuedAt !== "number" || + opts.bootedAt - meta.queuedAt > maxAgeMs + ) { + await cancel(meta, "stale", "queued before the last restart, stale"); + } else if (!meta.spec) { + await cancel( + meta, + "no-spec", + `${RESTART_REASON} (kind ${meta.kind} cannot be re-queued)`, + ); + } else { + const key = specKey(meta.spec); + const newer = kept.get(key); + if (newer) { + await cancel( + meta, + "superseded", + `superseded by a newer queued job (${newer})`, + ); + } else { + kept.set(key, meta.id); + toRequeue.push({ meta, spec: meta.spec }); + } + } + } + + // Oldest first, so the survivors keep their relative queue order. + for (const { meta, spec } of toRequeue.reverse()) { + let res: Awaited<ReturnType<RequeueFn>>; + try { + res = await opts.requeue!(spec); + } catch (err) { + res = { ok: false, error: (err as Error).message }; + } + if (res.ok) { + const reason = `${RESTART_REASON}; re-queued as ${res.jobId}`; + if (await closeMeta(opts.paths, meta, reason)) { + result.requeued.push({ id: meta.id, newId: res.jobId }); + log( + `[boot] re-queued ${meta.kind}${meta.channelSlug ? ` ${meta.channelSlug}` : ""}: ${meta.id} -> ${res.jobId}`, + ); + } + } else { + await cancel( + meta, + "refused", + `${RESTART_REASON}; re-queue refused: ${res.error}`, + ); + } + } + + if (result.requeued.length || result.cancelled.length) { + const by = new Map<string, number>(); + for (const c of result.cancelled) + by.set(c.category, (by.get(c.category) ?? 0) + 1); + log( + `[boot] stale queued jobs: re-queued ${result.requeued.length}, cancelled ${result.cancelled.length}` + + (by.size + ? ` (${[...by].map(([k, n]) => `${k} ${n}`).join(", ")})` + : ""), + ); + } + return result; +} + +async function closeMeta( + paths: Paths, + meta: JobMeta, + reason: string, +): Promise<boolean> { + const closed: JobMeta = { + ...meta, + status: "cancelled", + endedAt: Date.now(), + cancelReason: reason, + }; + try { + // tmp + rename: this pass can close hundreds of metas in one burst. + await writeFileAtomic(metaPath(paths, meta.id), JSON.stringify(closed)); + } catch { + return false; + } + // The job's own log says it too, for whoever opens the row. It may not exist + // (a job that never started never opened one); appendFile creates it. + try { + await appendFile( + path.join(paths.jobsDir, `${meta.id}.log`), + `[boot] ${reason}\n`, + ); + } catch { + /* the meta is the record; the log line is a courtesy */ + } + return true; +} diff --git a/common/jobs/downloadBackoff.ts b/common/jobs/downloadBackoff.ts @@ -25,7 +25,7 @@ import { readAutoQueueState, writeAutoQueueState, } from "./autoQueueState"; -import { nextBackoff, pruneExpired } from "./platformBackoff"; +import { mergeBackoffEntry, nextBackoff, pruneExpired } from "./platformBackoff"; // Milliseconds remaining in the platform's current cooldown window, or 0 if it // is not cooling down. Prefers the live shared object; else reads the file. @@ -50,13 +50,15 @@ export async function recordDownloadBackoff( const live = await liveAutoQueueState(paths); let state: AutoQueueState; if (live) { - // Fold in whatever is on disk first (the later `until` wins, as the - // runner's own merge does), so a cooldown written while no runner held the - // object is escalated from, not forgotten. + // Fold in whatever is on disk first (`until` and `fails` each to their + // max, as the runner's own merge does), so a cooldown written while no + // runner held the object is escalated from, not forgotten. const onDisk = (await readAutoQueueState(paths)).download.platformBackoff; for (const [pf, e] of Object.entries(onDisk)) { - const cur = live.download.platformBackoff[pf]; - if (!cur || e.until > cur.until) live.download.platformBackoff[pf] = e; + live.download.platformBackoff[pf] = mergeBackoffEntry( + live.download.platformBackoff[pf], + e, + ); } state = live; } else { diff --git a/common/jobs/jobMeta.ts b/common/jobs/jobMeta.ts @@ -25,6 +25,10 @@ export type JobMeta = { // Replay descriptor for replayable jobs, preserved so an archived (evicted // or post-restart) job can still be retried from disk. See jobSpec.ts. spec?: JobSpec; + // Why a job ended `cancelled` without anyone pressing Cancel — written only + // by the boot pass (bootQueuedJobs.ts) for a job that was still `queued` + // when the server went down. Absent on every other meta. + cancelReason?: string; }; export function metaPath(paths: Paths, id: string): string { diff --git a/common/jobs/platformBackoff.test.ts b/common/jobs/platformBackoff.test.ts @@ -7,6 +7,7 @@ import { clearBackoff, coercePlatformBackoff, isCoolingDown, + mergeBackoffEntry, nextBackoff, pruneExpired, VIDEO_RATE_LIMIT_DEFER_MS, @@ -145,3 +146,21 @@ test("coerceVideoDeferrals tolerates corrupt/missing shapes", () => { { ok: { until: 5, channelSlug: "a" } }, ); }); + +// Release 8 review, S: the two records merge field by field. +test("mergeBackoffEntry takes the max of until and of fails separately", () => { + assert.deepEqual(mergeBackoffEntry(undefined, { until: 5, fails: 2 }), { + until: 5, + fails: 2, + }); + // A short manual cooldown on disk with a low count must not erase the live + // escalation, and a later disk `until` must not be lost either. + assert.deepEqual( + mergeBackoffEntry({ until: 100, fails: 13 }, { until: 200, fails: 1 }), + { until: 200, fails: 13 }, + ); + assert.deepEqual( + mergeBackoffEntry({ until: 300, fails: 1 }, { until: 200, fails: 7 }), + { until: 300, fails: 7 }, + ); +}); diff --git a/common/jobs/platformBackoff.ts b/common/jobs/platformBackoff.ts @@ -39,6 +39,19 @@ export function nextBackoff( return { until: now + Math.round(base * jitter), fails }; } +// Fold two records of the same platform's backoff — the live one and the one +// on disk — into one. `until` and `fails` merge SEPARATELY, each to its max: +// taking one record whole (the later `until`) could pair a short manual +// cooldown with a lost escalation count, so the next 429 restarted at the +// base delay (release 8 review, S). +export function mergeBackoffEntry( + a: PlatformBackoffEntry | undefined, + b: PlatformBackoffEntry, +): PlatformBackoffEntry { + if (!a) return b; + return { until: Math.max(a.until, b.until), fails: Math.max(a.fails, b.fails) }; +} + // True when the platform is currently in a cooldown window. export function isCoolingDown( state: PlatformBackoffState, diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts @@ -129,8 +129,34 @@ export type JobRecord = { // task-count average cannot express it — the corpus is 77k videos and 77k // audio-hours, and those distribute completely differently. completedTaskAudioSeconds?: number; + // When the job last MOVED (epoch ms): its progress snapshot changed + // (`noteProgress`) or a sub-operation finished (`recordTaskDuration`). + // Undefined until the first. What the /jobs stall heuristic measures quiet + // time from — a slow scan at 10 videos/min that holds a slot for an hour with + // no in-flight tasks is still moving (release 9, B4a). + progressAt?: number; }; +// Store a progress snapshot on the record, stamping `progressAt` only when the +// numbers changed — a runner that re-reports the same count is not moving. +export function noteProgress( + record: JobRecord, + snap: JobProgress, + now: number = Date.now(), +): void { + const prev = record.progress; + if ( + !prev || + prev.current !== snap.current || + prev.target !== snap.target || + prev.initial !== snap.initial || + prev.remainingAudioSeconds !== snap.remainingAudioSeconds + ) { + record.progressAt = now; + } + record.progress = snap; +} + export type QueueSnapshot = { name: string; running?: JobRecord; @@ -360,6 +386,7 @@ class JobRegistry { if (!job || durationMs < 0) return; job.completedTaskCount = (job.completedTaskCount ?? 0) + 1; job.completedTaskMs = (job.completedTaskMs ?? 0) + durationMs; + job.progressAt = Date.now(); if (typeof audioSeconds === "number" && audioSeconds > 0) { job.completedTaskAudioSeconds = (job.completedTaskAudioSeconds ?? 0) + audioSeconds; diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts @@ -5,6 +5,7 @@ import { execa } from "execa"; import { getRegistry, newJobId, + noteProgress, type JobProgress, type JobRecord, type JobStatus, @@ -342,7 +343,7 @@ export async function runManagedFunction( }; const setProgress = (snap: JobProgress) => { - record.progress = snap; + noteProgress(record, snap); }; const ctx: JobRunContext = { diff --git a/common/lib/homepageSummary.test.ts b/common/lib/homepageSummary.test.ts @@ -80,7 +80,7 @@ const STATS: VideoStat[] = [ stat({ channelSlug: "p1", id: "z1", uploadDate: "20200101", status: "deleted" }), ]; -test("monthly: upload months, zero-filled from first upload to the build month, every site keyed", () => { +test("monthly: upload months, zero-filled from first upload to the last complete month, every site keyed", () => { const s = buildHomepageSummary(STATS, CHANNEL_SITES, SITES, NOW); assert.equal(s.version, HOMEPAGE_SUMMARY_VERSION); assert.deepEqual( @@ -144,3 +144,43 @@ test("the build month is never emitted, even in January (the year rolls back)", assert.deepEqual(s.monthly!.map((m) => m.month), ["2025-12"]); assert.equal(s.monthlyUnplaced, 1); }); + +// Release 8 review of slice H, the two LOWs (release 9). +test("a transcript bucketize drops (a future transcribedDate) still keeps placed + unplaced = official", () => { + const stats = [ + ...STATS, + // Transcribed "next year" (clock skew or a bad sidecar): the site card's + // bucketized total drops it, but it is placed by its upload month. + stat({ channelSlug: "a1", id: "x9", uploadDate: "20251201", transcribedDate: "20270101" }), + ]; + const s = buildHomepageSummary(stats, CHANNEL_SITES, SITES, NOW); + const placed = s.monthly!.reduce((a, m) => a + m.bySite.alpha + m.bySite.beta, 0); + assert.equal(s.official!.transcripts, 9); + assert.equal(placed + s.monthlyUnplaced!, s.official!.transcripts); + assert.deepEqual(s.monthly![1].bySite, { alpha: 1, beta: 0 }); +}); + +test("a site the activity filter drops gains no monthly key and no official count", () => { + const sites = [...SITES, site("gamma", ["g1"], "https://gamma.example")]; + const stats = [ + ...STATS, + // gamma's only transcript has a transcribedDate bucketize cannot place, + // so gamma has no activity and is filtered out of `sites`. + stat({ channelSlug: "g1", id: "w1", uploadDate: "20251101", transcribedDate: "20270101", downloadedDate: null }), + ]; + const s = buildHomepageSummary(stats, { ...CHANNEL_SITES, g1: ["gamma"] }, sites, NOW); + assert.ok(!s.sites.some((x) => x.siteId === "gamma")); + assert.ok(s.monthly!.every((m) => !("gamma" in m.bySite))); + const placed = s.monthly!.reduce((a, m) => a + m.bySite.alpha + m.bySite.beta, 0); + assert.equal(placed + s.monthlyUnplaced!, s.official!.transcripts); +}); + +test("gone counts only a deleted record we hold (it has a downloadedDate)", () => { + const stats = [ + ...STATS, + stat({ channelSlug: "b1", id: "y9", status: "deleted", downloadedDate: null, hasTranscript: false, transcribedDate: null }), + ]; + const s = buildHomepageSummary(stats, CHANNEL_SITES, SITES, NOW); + assert.equal(s.sites.find((x) => x.siteId === "beta")!.gone, 0); + assert.equal(s.official!.gone, 1); +}); diff --git a/common/lib/homepageSummary.ts b/common/lib/homepageSummary.ts @@ -347,12 +347,26 @@ export function buildHomepageSummary( let transcribedThisMonth = 0; let downloadedThisMonth = 0; const channelSet = new Set<string>(); - // Per public site: channels with a transcript, recordings, seconds, gone. - type SiteAcc = { channels: Set<string>; recordings: number; seconds: number; gone: number }; + // Per public site: channels with a transcript, recordings, transcripts, + // seconds, gone. `transcripts` is counted HERE, in the same pass that places + // or leaves unplaced each transcript, so `placed + unplaced = + // official.transcripts` holds even for a `transcribedDate` that bucketize + // drops (a future or malformed month). + type SiteAcc = { + channels: Set<string>; + recordings: number; + transcripts: number; + seconds: number; + gone: number; + }; const siteAcc = new Map<string, SiteAcc>(); const accOf = (id: string): SiteAcc => { let a = siteAcc.get(id); - if (!a) siteAcc.set(id, (a = { channels: new Set(), recordings: 0, seconds: 0, gone: 0 })); + if (!a) + siteAcc.set( + id, + (a = { channels: new Set(), recordings: 0, transcripts: 0, seconds: 0, gone: 0 }), + ); return a; }; // Upload-month placement of transcribed public-universe records. @@ -376,8 +390,10 @@ export function buildHomepageSummary( channelName.set(s.channelSlug, s.channel); const acc = accOf(siteId); if (s.downloadedDate) acc.recordings += 1; - if (s.status === "deleted") acc.gone += 1; + // "Gone at the source, still here": only a record we actually hold. + if (s.status === "deleted" && s.downloadedDate) acc.gone += 1; if (hasTx) { + acc.transcripts += 1; acc.channels.add(s.channelSlug); acc.seconds += s.duration > 0 ? s.duration : 0; const um = uploadMonthOf(s.uploadDate); @@ -483,7 +499,10 @@ export function buildHomepageSummary( for (const it of uploadItems) { const i = monthIndex.get(it.month); if (i === undefined) continue; // unreachable: every item's month is in range - monthly[i].bySite[it.siteId] = (monthly[i].bySite[it.siteId] ?? 0) + 1; + // A site the activity filter dropped has no key here and gets none: its + // transcripts are not placed, and not counted in `official` either. + if (!(it.siteId in monthly[i].bySite)) continue; + monthly[i].bySite[it.siteId] += 1; } // The family's own totals — hours summed in seconds, then rounded once. @@ -493,7 +512,10 @@ export function buildHomepageSummary( sites: summarySites.length, channels: summarySites.reduce((a, s) => a + (s.channels ?? 0), 0), recordings: summarySites.reduce((a, s) => a + (s.recordings ?? 0), 0), - transcripts: summarySites.reduce((a, s) => a + s.transcribed.total, 0), + transcripts: keptIds.reduce( + (a, id) => a + (siteAcc.get(id)?.transcripts ?? 0), + 0, + ), hoursArchived: Math.round(officialSeconds / 3600), gone: summarySites.reduce((a, s) => a + (s.gone ?? 0), 0), }; diff --git a/common/views/jobRows.test.ts b/common/views/jobRows.test.ts @@ -1,6 +1,6 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { newJobId, type JobRecord } from "../jobs/registry"; +import { newJobId, noteProgress, type JobRecord } from "../jobs/registry"; import type { JobListEntry } from "../jobs/listJobs"; import { createScheduler } from "../jobs/scheduler"; import { @@ -399,3 +399,55 @@ test("fromInFlight: a note replaces the detail and drops the synthetic video id" assert.equal(row.videoId, undefined); assert.equal(row.id, "download:metadata-scan alpha"); }); + +// THE STALL LABEL MEASURES QUIET, NOT AGE (release 9, B4a). A metadata scan +// never has a task in flight and advances its count every few seconds; on +// 2026-09-25 one at ~10 videos/min read STUCK · POSSIBLY-STALLED after 10 min. +test("reconcileSlots: a long-running idle job whose progress advances is never stuck", () => { + const { scheduler, ids } = scheduledQueue(); + const id = ids[0]; + const start = 10_000_000; + const r = rec({ + id, + queueKey: "qFacts", + kind: "metadata-scan", + status: "running", + startedAt: start, + tasks: [], + }); + // One scanned video every 6 s for an hour, checked each minute. + for (let t = 0; t <= 60 * 60 * 1000; t += 6_000) { + noteProgress( + r, + { metric: "scans", initial: 0, target: 1_382, current: t / 6_000 }, + start + t, + ); + if (t % 60_000 === 0) { + const [fact] = reconcileSlots(scheduler.queues(), () => r, start + t); + assert.equal(fact.stuck, undefined, `flagged at ${t / 60_000} min`); + } + } +}); + +test("reconcileSlots: idle tasks and progress frozen past 10 min is possibly-stalled", () => { + const { scheduler, ids } = scheduledQueue(); + const id = ids[0]; + const start = 10_000_000; + const r = rec({ + id, + queueKey: "qFacts", + status: "running", + startedAt: start, + tasks: [], + }); + noteProgress(r, { metric: "scans", initial: 0, target: 100, current: 40 }, start + 20 * 60_000); + // The same count re-reported is not a move. + noteProgress(r, { metric: "scans", initial: 0, target: 100, current: 40 }, start + 25 * 60_000); + assert.equal(r.progressAt, start + 20 * 60_000); + const at = (min: number) => + reconcileSlots(scheduler.queues(), () => r, start + min * 60_000)[0]; + assert.equal(at(29).stuck, undefined); + const late = at(31); + assert.equal(late.stuck?.reason, "possibly-stalled"); + assert.equal(late.heal, false); +}); diff --git a/common/views/jobRows.ts b/common/views/jobRows.ts @@ -23,7 +23,11 @@ import type { JobRowView, StuckReason } from "./jobRowView"; export const RECENT_MS = 30_000; // A running slot is treated as "possibly stalled" (the soft case) once a -// running/draining record has held its slot this long with no in-flight tasks. +// running/draining record has held its slot this long with no in-flight tasks +// AND no progress: the clock runs from the later of its start and its last move +// (`JobRecord.progressAt`). Age alone flagged a healthy metadata scan at ~10 +// videos/min as STUCK on 2026-09-25 — it never has a task in flight, and it +// advances its count every few seconds (release 9, B4a). // Distinct from the hard cases (terminal-but-holding-slot / record-evicted), // which are always safe to auto-heal; the soft case is only surfaced/force- // releasable, never auto-healed. @@ -78,8 +82,10 @@ export function reconcileSlots( heal = true; } else { const ageMs = now - (rec.startedAt ?? rec.queuedAt); + const quietMs = + now - Math.max(rec.startedAt ?? rec.queuedAt, rec.progressAt ?? 0); const idle = !rec.tasks || rec.tasks.length === 0; - if (idle && ageMs > STUCK_AGE_MS) { + if (idle && quietMs > STUCK_AGE_MS) { // Soft case: a live running/draining record sitting idle far too // long. Force-releasable, but deliberately NOT auto-healed. stuck = { reason: "possibly-stalled", heldMs: ageMs }; diff --git a/common/ytdlp/channelArgs.test.ts b/common/ytdlp/channelArgs.test.ts @@ -19,15 +19,34 @@ test("a rumble URL gets the four platform args", () => { ); }); -test("youtube gets no platform args", () => { +test("youtube gets the request pace and nothing else (release 9, B3)", () => { assert.deepEqual( channelExtraArgs(cfg({ url: "https://www.youtube.com/@x" })), - [], + ["--sleep-requests", "1"], ); - assert.deepEqual(platformArgs("youtube"), []); + assert.deepEqual(platformArgs("youtube"), ["--sleep-requests", "1"]); assert.deepEqual(platformArgs(null), []); }); +test("a platform with no entry gets no platform args", () => { + assert.deepEqual(platformArgs("twitch"), []); + assert.deepEqual( + channelExtraArgs(cfg({ url: "https://example.com/feed" })), + [], + ); +}); + +test("a youtube channel's own --sleep-requests comes after, and wins", () => { + const args = channelExtraArgs( + cfg({ + url: "https://www.youtube.com/@x", + ytdlpExtraArgs: ["--sleep-requests", "3"], + }), + ); + assert.deepEqual(args, ["--sleep-requests", "1", "--sleep-requests", "3"]); + assert.equal(args.lastIndexOf("--sleep-requests"), args.length - 2); +}); + test("an explicit config.platform wins over the URL's host", () => { assert.deepEqual( channelExtraArgs( diff --git a/common/ytdlp/managedDownloadsSleep.test.ts b/common/ytdlp/managedDownloadsSleep.test.ts @@ -0,0 +1,165 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mkdtempSync, writeFileSync, mkdirSync } from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import type { + DownloadAttempt, + DownloadOutcomeRecord, + DownloadOutcomeStatus, +} from "../lib/downloadOutcome"; +import type { DownloadFailureClass } from "../lib/availability"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test ytdlp/managedDownloadsSleep.test.ts +// +// THE BETWEEN-DOWNLOADS SLEEP PACES MEDIA FETCHES (release 9, B2). On +// 2026-09-25 a filtered channel's download-missing slept 30 s after every video +// the download filter declined — 193 sleeps for 14 archives. The loop is driven +// here with an injected per-video download (canned outcomes, no yt-dlp) and an +// injected sleeper that only counts. getPaths() memoizes, so the env is set +// before anything imports it. +const ROOT = mkdtempSync(path.join(os.tmpdir(), "managed-sleep-")); +process.env.TRANSCRIPTS_DIR = ROOT; +process.env.SETTINGS_FILE = path.join(ROOT, "settings.json"); +writeFileSync( + process.env.SETTINGS_FILE, + JSON.stringify({ minFreeDiskGB: 0, sleepBetweenDownloadsSeconds: 30 }) + "\n", +); +mkdirSync(path.join(ROOT, "channels", "c", "data"), { recursive: true }); + +const { runManagedDownloads, declinedWithoutMediaFetch } = + await import("./runYtdlp"); +const { getPaths } = await import("../lib/paths"); + +const PREFETCH: DownloadAttempt = { + n: 0, + kind: "metadata-prefetch", + handling: "youtube", + usedCookies: false, + ytdlpExitCode: 0, +}; +const PRIMARY: DownloadAttempt = { + n: 1, + kind: "primary", + handling: "youtube", + usedCookies: false, + ytdlpExitCode: 0, +}; + +function outcome( + status: DownloadOutcomeStatus, + attempts: DownloadAttempt[], + failureClass?: DownloadFailureClass, +): DownloadOutcomeRecord { + return { + videoId: "x", + status, + startedAt: "2026-09-25T00:00:00.000Z", + finishedAt: "2026-09-25T00:00:01.000Z", + attempts, + ...(failureClass ? { failureClass } : {}), + }; +} + +const FETCHED = outcome("ok", [PREFETCH, PRIMARY]); +const FILTERED = outcome("skipped-filtered", [PREFETCH]); +const MEMBERS_ONLY = outcome( + "failed", + [{ ...PREFETCH, ytdlpExitCode: 1, availabilityClass: "members_only" }], + "per_video", +); + +// Runs the loop over one URL per outcome and returns how many times it slept. +async function sleepsFor(outcomes: DownloadOutcomeRecord[]): Promise<number> { + let sleeps = 0; + let i = 0; + const urls = outcomes.map( + (_, n) => + `https://www.youtube.com/watch?v=vid${String(n).padStart(8, "0")}`, + ); + const channelConfig = { + handling: "youtube", + url: "https://www.youtube.com/@c/videos", + } as never; + await runManagedDownloads( + { + channelSlug: "c", + mode: "download-missing" as never, + channelConfig, + paths: getPaths(), + onLog: () => {}, + signal: new AbortController().signal, + abortOnError: false, + }, + urls, + channelConfig, + undefined, + { + downloadOne: async () => outcomes[i++], + sleep: async () => { + sleeps++; + }, + }, + ); + return sleeps; +} + +test("a fetched video sleeps before the next one", async () => { + assert.equal(await sleepsFor([FETCHED, FETCHED]), 1); +}); + +test("a video the download filter declined does not sleep", async () => { + assert.equal(await sleepsFor([FILTERED, FETCHED]), 0); + assert.equal(await sleepsFor([FILTERED, FILTERED, FILTERED]), 0); +}); + +// Release 9 review: a per-video failure KEEPS the pace, even one that never +// got past the prefetch. YouTube's soft block ("This content isn't available, +// try again later") classifies as deleted → per_video; skipping the sleep there +// would fire prefetches back to back into the block. +test("a per-video failure at the prefetch still sleeps", async () => { + assert.equal(await sleepsFor([MEMBERS_ONLY, FETCHED]), 1); +}); + +test("the last video never sleeps, whatever it was", async () => { + assert.equal(await sleepsFor([FETCHED]), 0); + assert.equal(await sleepsFor([FILTERED, FETCHED]), 0); +}); + +test("a real fetch that failed still sleeps, per-video or not", async () => { + const perVideoAfterFetch = outcome( + "failed", + [PREFETCH, { ...PRIMARY, ytdlpExitCode: 1 }], + "per_video", + ); + const network = outcome( + "failed", + [{ ...PREFETCH, ytdlpExitCode: 1 }], + "network", + ); + assert.equal(await sleepsFor([perVideoAfterFetch, FETCHED]), 1); + // A network failure at the prefetch is not the video's fault: keep pacing. + assert.equal(await sleepsFor([network, FETCHED]), 1); +}); + +test("declinedWithoutMediaFetch: a chat-only pass is a fetch", () => { + const chatPass: DownloadAttempt = { ...PRIMARY, kind: "live-chat-only" }; + // A failed chat pass leaves the status skipped-filtered, but it asked. + assert.equal( + declinedWithoutMediaFetch( + outcome("skipped-filtered", [PREFETCH, chatPass]), + ), + false, + ); + assert.equal( + declinedWithoutMediaFetch(outcome("chat-only", [PREFETCH, chatPass])), + false, + ); + assert.equal(declinedWithoutMediaFetch(FILTERED), true); + assert.equal(declinedWithoutMediaFetch(MEMBERS_ONLY), false); + assert.equal( + declinedWithoutMediaFetch(outcome("failed", [PREFETCH], "rate_limit")), + false, + ); +}); diff --git a/common/ytdlp/platformArgs.mjs b/common/ytdlp/platformArgs.mjs @@ -15,13 +15,29 @@ import { detectPlatform } from "../lib/detectPlatform.mjs"; // AFTER these, so a channel override wins (yt-dlp is last-flag-wins for // `--sleep-requests` and `--impersonate`). // +// Every entry paces at the REQUEST level with `--sleep-requests 1`: one +// second between the HTTP requests a single yt-dlp process makes (listing +// pages, the prefetch, subtitle fetches, retries). That is a different layer +// from `sleepBetweenDownloadsSeconds`, which only spaces whole per-video +// downloads apart. A spawn that already passes its own `--sleep-requests 1` +// (metadataScan.ts, fetchWindowManaged.ts) now carries it twice — harmless, +// yt-dlp keeps the last. +// // rumble: every request 403s at Cloudflare without a browser TLS fingerprint // (yt-dlp #17496); a probe on 2026-09-24 got 200 with `--impersonate chrome` // and 403 without. `--sleep-requests 1` paces the listing walk — the // the-quartering-rumble full sweep 429'd at page 155 unpaced the same day. +// +// youtube: the 2026-09-25 429 investigation +// (~/reports/release-7/data/q-429-report.md, finding 1) found YouTube had NO +// request-level pacing at all — each video attempt fires 2–4 requests +// (prefetch, subtitles, retries) back to back, and the two channels that 429'd +// were the two most-downloaded. Mirrors rumble's pace; a channel's own +// `ytdlpExtraArgs` still wins because it comes after. /** @type {Readonly<Partial<Record<Platform, readonly string[]>>>} */ export const PLATFORM_ARGS = Object.freeze({ rumble: Object.freeze(["--impersonate", "chrome", "--sleep-requests", "1"]), + youtube: Object.freeze(["--sleep-requests", "1"]), }); /** diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts @@ -23,6 +23,7 @@ import { classifyDownloadFailure, type Availability, } from "../lib/availability"; +import type { DownloadOutcomeRecord } from "../lib/downloadOutcome"; import { alwaysCookies, resolveCookiePolicy, @@ -888,6 +889,28 @@ async function downloadPlaylistManaged( await safeBackfillAvailability(opts); } +export type ManagedDownloadsDeps = { + downloadOne?: typeof downloadOneManaged; + sleep?: (ms: number, signal: AbortSignal) => Promise<void>; +}; + +// True when the download filter declined the video before any media request: +// every attempt it made was the n: 0 metadata prefetch (or its cookie retry) +// and it ended `skipped-filtered`. Attempts n >= 1 are the real fetches — the +// download attempts and the chat-only pass (a failed chat pass leaves the +// status skipped-filtered, and still asked). +// +// A per_video FAILURE is deliberately NOT here, even one that never got past +// the prefetch (release 9 review): YouTube's soft block ("This content isn't +// available, try again later") classifies as `deleted` → per_video, so +// skipping the sleep there would fire prefetches back to back into the block. +export function declinedWithoutMediaFetch( + outcome: DownloadOutcomeRecord, +): boolean { + if (outcome.attempts.some((a) => a.n >= 1)) return false; + return outcome.status === "skipped-filtered"; +} + type ManagedRunResult = { okCount: number; failedCount: number; @@ -906,7 +929,7 @@ type ManagedRunResult = { // cost is the price of being able to retry each independently. Honors the // hard `signal`, the soft `drainSignal`, the per-channel sleep, and the // per-operation tracker. -async function runManagedDownloads( +export async function runManagedDownloads( opts: RunYtdlpOpts, urls: ReadonlyArray<string>, effectiveChannelConfig: ChannelConfig, @@ -914,7 +937,12 @@ async function runManagedDownloads( // already resolved it (for their own prefilter) pass it through so the // forceCookies-without-value warning isn't logged twice. cookiePolicyOverride?: ResolvedCookiePolicy, + // Test seams: the per-video download and the between-downloads sleep. Every + // production caller passes neither. + deps: ManagedDownloadsDeps = {}, ): Promise<ManagedRunResult> { + const downloadOne = deps.downloadOne ?? downloadOneManaged; + const sleep = deps.sleep ?? abortableSleep; const settings = getSettings(); const cookiePolicy = cookiePolicyOverride ?? @@ -987,7 +1015,7 @@ async function runManagedDownloads( }); let outcome; try { - outcome = await downloadOneManaged({ + outcome = await downloadOne({ channelSlug: opts.channelSlug, channelConfig: effectiveChannelConfig, paths: opts.paths, @@ -1057,14 +1085,21 @@ async function runManagedDownloads( processedCount++; const isLast = processedCount >= urls.length; const willAbortLoop = firstFailure !== null && abortOnError; + // A VIDEO THE FILTER DECLINED DOES NOT PAY FOR A FETCH IT NEVER MADE + // (release 9, B2). On 2026-09-25 paramount-tactical's download-missing + // slept 193 times for 14 archives: every video the download filter + // declined slept 30 s after nothing but its metadata prefetch. Every + // other outcome still sleeps — a real fetch, success or failure, and + // every failure, per-video ones included (see declinedWithoutMediaFetch). if ( sleepSeconds > 0 && !isLast && !willAbortLoop && - !opts.signal.aborted + !opts.signal.aborted && + !declinedWithoutMediaFetch(outcome) ) { opts.onLog(`Sleeping ${sleepSeconds}s before next download...\n`); - await abortableSleep(sleepSeconds * 1000, opts.signal); + await sleep(sleepSeconds * 1000, opts.signal); } }), ), diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,12 @@ # Changelog ## [Unreleased] +- **A job that waited in a queue no longer ends `failed` after doing its work.** A sync, download or other job queued behind another on the same platform ran its final page refresh outside any request, where Next refuses it, so the job read `failed` and `pnpm ops … --wait` exited 1 even though the work was done (the teamrcn sync on 2026-09-25). The refresh is now skipped there with one warning in the server log; the pages re-read disk on their next load anyway. +- **Downloads no longer sleep after a video the download filter declined.** The "sleep between downloads" (30 s by default) ran after every video, including each one the channel's download filter declined before fetching anything. A filtered channel's download-missing slept 193 times for 14 downloads on 2026-09-25. It still sleeps after every real fetch and after every failure, per-video ones included. +- **YouTube requests are paced at one per second.** Every yt-dlp run against YouTube now carries `--sleep-requests 1`, as Rumble's already did: the 429 investigation found YouTube had no request-level pacing at all. A channel's own `ytdlpExtraArgs` still wins, because it comes after. +- **`/jobs` no longer calls a slow but moving job stuck.** "STUCK · POSSIBLY-STALLED" now needs the job's progress to have stood still for 10 minutes, not just the job to be 10 minutes old with nothing in flight. A metadata scan at ~10 videos a minute read stuck on 2026-09-25. +- **Jobs left `queued` by a restart are settled at boot.** A job still waiting when the server stopped used to sit on `/jobs` as queued forever. On boot each one ends `cancelled` with its reason, and only a few are queued again, through the same path as Retry. A sync is never re-queued, because the scheduler re-derives syncs at its own pace. Nothing queued more than a day before the restart is re-queued. Of several identical jobs, only the newest is re-queued. With `ARCHILYZER_IDLE_BOOT` set, none are. The server log gets one summary line, plus one line per job queued again. +- Smaller fixes from the release 8 reviews: a manual 429 cooldown merged with the runner's keeps the higher failure count as well as the later end; a channel's video-title memo keeps its titles when a new video directory appears, instead of re-reading every one; and the homepage counts "Transcripts" from the same pass that places them on the chart, counts "gone at the source, still here" only for recordings it holds, and drops the two empty columns its stats strip had when nothing is gone. No number on today's homepage changes. - **The hub looks and reads like the homepage, and its list of archives is "Official instances" with the homepage's numbers.** The hub now opens in the project's own dark look (the *Archilyzer* theme), as the homepage does; published sites keep theirs. Its front page opens on one sentence, "<hours> hours of speech from every official archive, searchable at once." (the official instances' hours, as on the homepage), and one paragraph. **Official instances** has a card for each built-in archive: its name linking out, its description, and its channels, recordings, transcripts and hours, the same figures the homepage shows. A card whose site sets no accent colour takes the homepage card's colour for that site. Archives a visitor adds are listed under **Archives you added**, with their channel count and a Remove button, above the add form; they are kept only in that visitor's browser. The line under the cards reads "Searching 5 archives · … channels · … transcripts right now." and counts what the browser actually loaded, added archives included. The figures come from a new file the hub build writes, `hub-summary.json`, made by the same code that writes the homepage's summary. A hub built without it (an older build, or a machine with no index) shows the cards without figures. "The shelf", "Your archives" and the Member/Added badges are gone. Needs a rebuild and deploy of the hub. - **The Archilyzer theme's third chart colour is no longer the "gone" red.** It is now a green (`#5a9e3a` light, `#86c86a` dark). The homepage chart and cards draw the third official instance in it, and a red layer read as "recordings gone". Needs a rebuild and deploy of the homepage. - **The yt-dlp clip command is back on sites with transcript downloads turned off.** Turning off `transcriptDownloads` (site.json, or the hub's homepage.json) hid three buttons in the transcript viewer. One of them, the yt-dlp button, only copies a `yt-dlp --download-sections` command for a marked clip to the clipboard and serves no file, so it is not a download. It now shows on every site. The switch still hides the Download menu (txt / srt / json) and Copy MD. The site and hub form labels in the editor say so. No setting changed; a site picks this up at its next build and deploy. diff --git a/editor/app/api/test/stuck-job/route.ts b/editor/app/api/test/stuck-job/route.ts @@ -57,5 +57,31 @@ export async function GET(request: Request) { /* best-effort */ } - return NextResponse.json({ id, queueKey }); + // `releaseAfterMs`: finish the fake job after that long FROM OUTSIDE ANY + // REQUEST, the way a job the auto-runner or the sync heartbeat started + // finishes (release 9, B1). The scheduler starts the next queued job + // synchronously inside `complete()`, so that job inherits the async context + // the release runs in — and a job started with no Next work store is exactly + // the one whose revalidatePath threw "static generation store missing". A + // plain setTimeout here would inherit THIS request's store, so the timer is + // armed inside `workAsyncStorage.exit`. `detached` reports whether that + // worked, so a spec can refuse to pass for the wrong reason. + const releaseAfterMs = Number(url.searchParams.get("releaseAfterMs") ?? ""); + let detached = false; + if (Number.isFinite(releaseAfterMs) && releaseAfterMs > 0) { + const release = () => registry.finalize(id, "done"); + try { + // A Next internal (typed, but not public API) — test harness only. + const { workAsyncStorage } = + await import("next/dist/server/app-render/work-async-storage.external"); + workAsyncStorage.exit(() => { + detached = workAsyncStorage.getStore() === undefined; + setTimeout(release, releaseAfterMs); + }); + } catch { + setTimeout(release, releaseAfterMs); + } + } + + return NextResponse.json({ id, queueKey, detached }); } diff --git a/editor/app/channels/[slug]/availabilityActions.ts b/editor/app/channels/[slug]/availabilityActions.ts @@ -1,6 +1,6 @@ "use server"; -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { platformQueueKey, @@ -63,7 +63,7 @@ export async function checkAvailabilityAction( onLog, signal, }); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -91,7 +91,7 @@ export async function quickAvailabilityCheckAction( onLog, signal, }); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -135,7 +135,7 @@ export async function checkMaybeMissingAction( onLog, signal, }); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } diff --git a/editor/app/channels/[slug]/backfillActions.ts b/editor/app/channels/[slug]/backfillActions.ts @@ -1,6 +1,6 @@ "use server"; -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; import { runBackfillChannelJob } from "yt-dlp-transcript-common/controller/operationJobs"; @@ -35,6 +35,6 @@ export async function backfillChannelAction( kindIds, ids, queueKey, - onDone: () => revalidatePath(`/channels/${slug}`), + onDone: () => safeRevalidate([`/channels/${slug}`]), }); } diff --git a/editor/app/channels/[slug]/digestActions.ts b/editor/app/channels/[slug]/digestActions.ts @@ -1,6 +1,7 @@ "use server"; import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { @@ -74,7 +75,7 @@ export async function digestChannelAction( order, limitCount, force, - onDone: () => revalidatePath(`/channels/${slug}`), + onDone: () => safeRevalidate([`/channels/${slug}`]), }); } @@ -133,7 +134,7 @@ export async function digestBucketAction( `Digest bucket: ${result.succeeded} generated, ${result.fresh} already current, ${result.failed} failed.`, ); requestChannelSnapshot(paths, slug); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } diff --git a/editor/app/channels/[slug]/incompleteTranscriptActions.ts b/editor/app/channels/[slug]/incompleteTranscriptActions.ts @@ -1,6 +1,7 @@ "use server"; import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { TRANSCRIPTION_QUEUE, @@ -111,7 +112,7 @@ export async function redownloadIncompleteBucketAction( onLog( `Re-download incomplete: ${succeeded} fixed, ${failed} failed of ${cleaned.length}.`, ); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -175,7 +176,7 @@ export async function redownloadShortAudioBucketAction( onLog( `Re-download short-audio: ${succeeded} fixed, ${failed} failed of ${cleaned.length}.`, ); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } diff --git a/editor/app/channels/[slug]/normalizeActions.ts b/editor/app/channels/[slug]/normalizeActions.ts @@ -15,7 +15,7 @@ // putting it anywhere else would make the fix for a stalled digest lane queue // up BEHIND the digest lane it is meant to unblock. -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { channelQueueKey, @@ -56,7 +56,7 @@ export async function normalizeChannelAction( // The digest lane's classification is derived from disk, so the count this // run just moved is only visible once the snapshot is rebuilt. requestChannelSnapshot(paths, slug); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } diff --git a/editor/app/channels/[slug]/persistActions.ts b/editor/app/channels/[slug]/persistActions.ts @@ -1,6 +1,6 @@ "use server"; -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { diskGate } from "yt-dlp-transcript-common/lib/diskSpace"; @@ -53,8 +53,7 @@ export async function persistKeptAction( onLog, signal, }); - revalidatePath(`/channels/${slug}`); - revalidatePath("/saved-videos"); + safeRevalidate([`/channels/${slug}`, "/saved-videos"]); }, }); } diff --git a/editor/app/channels/[slug]/pipelineActions.ts b/editor/app/channels/[slug]/pipelineActions.ts @@ -1,7 +1,7 @@ "use server"; import path from "node:path"; -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { HANDLING_VALUES, type AudioFormat, @@ -217,11 +217,13 @@ async function runPipelineAction( // The channel report (snapshot) is regenerated automatically after this // job finishes, via the global debounced scheduler hooked into // runManagedFunction's completion. See common/jobs/snapshotScheduler.ts. - revalidatePath(`/channels/${slug}`); - revalidatePath("/channels"); - revalidatePath("/operations/[id]", "page"); - revalidatePath("/cleanup"); - revalidatePath("/"); + safeRevalidate([ + `/channels/${slug}`, + "/channels", + ["/operations/[id]", "page"], + "/cleanup", + "/", + ]); }, }); } @@ -362,9 +364,11 @@ export async function runMetadataScanAction( channelConfig, queueKey, afterRun: () => { - revalidatePath(`/channels/${slug}`); - revalidatePath("/channels"); - revalidatePath("/operations/[id]", "page"); + safeRevalidate([ + `/channels/${slug}`, + "/channels", + ["/operations/[id]", "page"], + ]); }, }); } @@ -519,10 +523,9 @@ export async function importVideoAction( ).catch(() => { /* the download succeeded; a roster write failure must not fail it */ }); - revalidatePath(`/channels/${slug}/videos/${videoId}`); + safeRevalidate([`/channels/${slug}/videos/${videoId}`]); } - revalidatePath(`/channels/${slug}`); - revalidatePath("/channels"); + safeRevalidate([`/channels/${slug}`, "/channels"]); } finally { task.end(); } diff --git a/editor/app/channels/[slug]/socialActions.ts b/editor/app/channels/[slug]/socialActions.ts @@ -7,6 +7,7 @@ // (Fetch → Index), and this file owns the first. import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { @@ -114,7 +115,7 @@ export async function checkPostAvailabilityAction( onLog, signal, }); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); if (!result.ok) throw new Error(result.error ?? "Availability check failed"); }, }); @@ -154,7 +155,7 @@ export async function fetchPostsAction( onLog, signal, }); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); // Surface a failed fetch as a failed JOB (the managed wrapper turns a // throw into status "failed"), so it shows up in the jobs list the same // way a failed download does rather than silently logging. diff --git a/editor/app/channels/[slug]/videos/[id]/operationActions.ts b/editor/app/channels/[slug]/videos/[id]/operationActions.ts @@ -13,7 +13,7 @@ // the export 500s at runtime while tsc stays quiet. digestActions.ts has the // long version of this note. -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { @@ -48,8 +48,10 @@ export async function runOperationForVideoAction( // No queueKey — the lane's own key. No `background` either: a hand-clicked // run goes in front of the sweep's, exactly as backfillChannelAction's does. onDone: () => { - revalidatePath(`/channels/${slug}/videos/${videoId}`); - revalidatePath(`/channels/${slug}`); + safeRevalidate([ + `/channels/${slug}/videos/${videoId}`, + `/channels/${slug}`, + ]); }, }); } diff --git a/editor/app/channels/[slug]/videos/[id]/videoActions.ts b/editor/app/channels/[slug]/videos/[id]/videoActions.ts @@ -3,6 +3,7 @@ import path from "node:path"; import { readdir, readFile, rm, stat } from "node:fs/promises"; import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../../../lib/safeRevalidate"; import { redirect } from "next/navigation"; import type { AudioFormat, @@ -128,7 +129,7 @@ export async function transcodeAudioAction( onLog, signal, }); - revalidatePath(`/channels/${slug}/videos/${videoId}`); + safeRevalidate([`/channels/${slug}/videos/${videoId}`]); }, }); } @@ -159,7 +160,7 @@ export async function transcribeOneAction( onLog, signal, }); - revalidatePath(`/channels/${slug}/videos/${videoId}`); + safeRevalidate([`/channels/${slug}/videos/${videoId}`]); }, }); } @@ -220,8 +221,10 @@ export async function downloadVideoPipelineAction( }), appendArchive: true, }); - revalidatePath(`/channels/${slug}/videos/${videoId}`); - revalidatePath(`/channels/${slug}`); + safeRevalidate([ + `/channels/${slug}/videos/${videoId}`, + `/channels/${slug}`, + ]); } finally { task.end(); } @@ -299,8 +302,10 @@ async function archiveSourceVideo( keepSourceVideoOverride: true, persistOrigin, }); - revalidatePath(`/channels/${slug}/videos/${videoId}`); - revalidatePath(`/channels/${slug}`); + safeRevalidate([ + `/channels/${slug}/videos/${videoId}`, + `/channels/${slug}`, + ]); } finally { task.end(); } @@ -394,8 +399,10 @@ export async function whisperVideoAction( onLog, signal, }); - revalidatePath(`/channels/${slug}/videos/${videoId}`); - revalidatePath(`/channels/${slug}`); + safeRevalidate([ + `/channels/${slug}/videos/${videoId}`, + `/channels/${slug}`, + ]); }, }); } @@ -429,8 +436,10 @@ export async function redownloadIncompleteTranscriptAction( signal, tracker: makeTaskTracker(ctx, onLog), }); - revalidatePath(`/channels/${slug}/videos/${videoId}`); - revalidatePath(`/channels/${slug}`); + safeRevalidate([ + `/channels/${slug}/videos/${videoId}`, + `/channels/${slug}`, + ]); }, }); } @@ -970,7 +979,7 @@ export async function fetchWindowAction(req: { signal, onPlatformBackoff: () => recordDownloadBackoff(platform, paths), }); - revalidatePath(`/channels/${slug}/videos/${videoId}`); + safeRevalidate([`/channels/${slug}/videos/${videoId}`]); }, }); if (!res.ok) { diff --git a/editor/app/channels/[slug]/whisperActions.ts b/editor/app/channels/[slug]/whisperActions.ts @@ -1,6 +1,6 @@ "use server"; -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { AUDIO_FORMAT_VALUES, @@ -91,7 +91,7 @@ export async function transcribeMissingAction( onLog( `Whisper batch: ${result.succeeded} succeeded, ${result.failed} failed, ${result.skipped} skipped, ${result.attempted} attempted.`, ); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -150,7 +150,7 @@ export async function transcribeBucketAction( onLog( `Whisper bucket: ${result.succeeded} succeeded, ${result.failed} failed, ${result.skipped} skipped, ${result.attempted} attempted.`, ); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -203,7 +203,7 @@ export async function transcribeAutoSubsBucketAction( onLog( `Replace auto-captions: ${result.succeeded} succeeded, ${result.failed} failed, ${result.skipped} skipped, ${result.attempted} attempted.`, ); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -224,7 +224,7 @@ export async function clearFailedTranscriptionsAction( onLog( `Cleared ${cleared} entr${cleared === 1 ? "y" : "ies"} from failed-transcriptions; they'll be retried next "Transcribe missing".`, ); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -250,7 +250,7 @@ export async function cleanExtraAudioFormatsAction( onLog( `Clean extra audio formats: ${result.removedFiles} file(s) removed across ${result.cleanedDirs} dir(s); inspected ${result.inspected}.`, ); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -279,7 +279,7 @@ export async function removeWrongFormatAudioAction( onLog( `Remove wrong-format audio: ${result.removedFiles} file(s) removed across ${result.cleanedDirs} dir(s); inspected ${result.inspected}.`, ); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -302,7 +302,7 @@ export async function cleanAudioAction( onLog, signal, }); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -334,7 +334,7 @@ export async function diarizeChannelAction( `Diarize ${slug}: ${result.diarized} diarized, ${result.alreadyExists} already done, ` + `${result.noAudio} without audio, ${result.failed} failed.`, ); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -367,7 +367,7 @@ export async function purgeSupersededAutoSubsAction( onLog, signal, }); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } @@ -397,7 +397,7 @@ export async function checkKeptDeletedAction( onLog( `Kept-deletion check: inspected ${result.kept}, ${result.deleted} gone from source, pinned ${result.pinned}.`, ); - revalidatePath(`/channels/${slug}`); + safeRevalidate([`/channels/${slug}`]); }, }); } diff --git a/editor/app/channels/groupActions.ts b/editor/app/channels/groupActions.ts @@ -1,6 +1,6 @@ "use server"; -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../lib/safeRevalidate"; import { listChannelBriefs, type ChannelBrief, @@ -147,7 +147,7 @@ const RUN_FOR: Record< channelSlug: slug, operation: DIGEST_OPERATION_ID, digest: { lane: "local" }, - onDone: () => revalidatePath(`/channels/${slug}`), + onDone: () => safeRevalidate([`/channels/${slug}`]), }), // EVERY ENABLED LANE KIND, and so NOT runOperationChannelJob: this station is // the LANE, not one operation on it — KIND_FOR above says the same thing by diff --git a/editor/app/channels/lib/relocationJob.ts b/editor/app/channels/lib/relocationJob.ts @@ -1,4 +1,4 @@ -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { relocationQueueKey } from "yt-dlp-transcript-common/lib/queueKeys"; import { @@ -109,9 +109,7 @@ export async function enqueueRelocation(opts: { // // What DOES have to change is what the pages read per render — the badge, // the location line, the free-space figure — so those are revalidated. - revalidatePath(`/channels/${slug}`); - revalidatePath("/channels"); - revalidatePath("/"); + safeRevalidate([`/channels/${slug}`, "/channels", "/"]); }, }); } diff --git a/editor/app/lib/safeRevalidate.test.ts b/editor/app/lib/safeRevalidate.test.ts @@ -0,0 +1,61 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + resetSafeRevalidateWarning, + runGuarded, + safeRevalidate, +} from "./safeRevalidate"; + +// Run with: pnpm -C editor exec tsx --test "app/**/*.test.ts" +// +// Outside a Next request (which is where a QUEUED job's body runs, and where +// this test runs) revalidatePath throws the missing-store invariant. The helper +// must swallow exactly that, warn once per process, and rethrow anything else. + +function captureWarn<T>(fn: () => T): { result: T; warnings: string[] } { + const warnings: string[] = []; + const orig = console.warn; + console.warn = (...args: unknown[]) => { + warnings.push(args.map(String).join(" ")); + }; + try { + return { result: fn(), warnings }; + } finally { + console.warn = orig; + } +} + +test("outside a request, safeRevalidate does not throw and warns once", () => { + resetSafeRevalidateWarning(); + const { warnings } = captureWarn(() => { + safeRevalidate(["/channels/x", "/channels", ["/operations/[id]", "page"]]); + safeRevalidate(["/"]); + }); + const ours = warnings.filter((w) => w.includes("[safeRevalidate]")); + assert.equal(ours.length, 1); + assert.match(ours[0], /\/channels\/x/); + assert.match(ours[0], /\/operations\/\[id\] \(page\)/); +}); + +test("the real revalidatePath does throw here (the premise of the helper)", async () => { + const { revalidatePath } = await import("next/cache"); + assert.throws( + () => revalidatePath("/channels"), + /static generation store missing/, + ); +}); + +test("any other error is rethrown", () => { + resetSafeRevalidateWarning(); + assert.throws( + () => + runGuarded( + () => { + throw new Error("disk on fire"); + }, + "/x", + ["/x"], + ), + /disk on fire/, + ); +}); diff --git a/editor/app/lib/safeRevalidate.ts b/editor/app/lib/safeRevalidate.ts @@ -0,0 +1,82 @@ +import { revalidatePath, revalidateTag } from "next/cache"; + +// REVALIDATION FROM INSIDE A JOB (release 9, B1). A job body (`fn` of +// runManagedFunction) or a job hook (`onDone` / `afterRun` / `afterDone`) runs +// in the request that submitted it only when the job STARTS at once. When the +// job was QUEUED behind another on the same platform queue, the queue's drain +// loop starts it later, outside any request, and Next has no work store: +// `revalidatePath` throws `Invariant: static generation store missing in +// revalidatePath …`. That throw escaped the body AFTER the work was done, so +// the job was finalized `failed` and `pnpm ops … --wait` exited 1 (the teamrcn +// sync, 2026-09-25 15:21). +// +// Nothing is lost by skipping the revalidation there: every page these paths +// name is dynamic and re-reads disk on the next request, and the snapshot +// scheduler revalidates the channel pages itself when it regenerates. So the +// missing-store invariant is swallowed (logged once per process, naming the +// paths), and ANY OTHER error is rethrown — a real failure stays a failure. +// +// Use it ONLY inside job bodies and job hooks. A plain server action (a form +// handler that revalidates and returns) runs inside its request and keeps +// calling `revalidatePath` directly. + +export type RevalidateTarget = string | [string, "page" | "layout"]; + +const MISSING_STORE = /static generation store missing/; + +let warned = false; + +function isMissingStore(err: unknown): boolean { + return err instanceof Error && MISSING_STORE.test(err.message); +} + +function describe(target: RevalidateTarget): string { + return typeof target === "string" ? target : `${target[0]} (${target[1]})`; +} + +// The guard itself, exported for the unit test so the rethrow path can be +// exercised without a Next runtime. `call` is one revalidation. +export function runGuarded( + call: () => void, + label: string, + allLabels: readonly string[], +): void { + try { + call(); + } catch (err) { + if (!isMissingStore(err)) throw err; + if (!warned) { + warned = true; + console.warn( + `[safeRevalidate] no request store (a queued job ran outside a request); ` + + `skipped revalidating ${allLabels.join(", ")} — first skip was ${label}. ` + + `Logged once per process; the pages re-read disk on their next request.`, + ); + } + } +} + +export function safeRevalidate( + paths: RevalidateTarget[], + tags: string[] = [], +): void { + const labels = [...paths.map(describe), ...tags.map((t) => `tag:${t}`)]; + for (const target of paths) { + runGuarded( + () => + typeof target === "string" + ? revalidatePath(target) + : revalidatePath(target[0], target[1]), + describe(target), + labels, + ); + } + for (const tag of tags) { + runGuarded(() => revalidateTag(tag, "max"), `tag:${tag}`, labels); + } +} + +// Test seam: the once-per-process latch. +export function resetSafeRevalidateWarning(): void { + warned = false; +} diff --git a/editor/app/saved-videos/backupActions.ts b/editor/app/saved-videos/backupActions.ts @@ -1,6 +1,7 @@ "use server"; import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { saveSettings } from "../settings/saveSettings"; @@ -45,7 +46,7 @@ export async function backupSavedVideosAction( onLog( `Saved-video backup: ${result.backedUp}/${result.entries} container(s), ${result.bytes} bytes.`, ); - revalidatePath("/saved-videos"); + safeRevalidate(["/saved-videos"]); }, }); } diff --git a/editor/app/sites/lib/buildAction.ts b/editor/app/sites/lib/buildAction.ts @@ -1,6 +1,6 @@ "use server"; -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { buildIndex } from "yt-dlp-transcript-common/controller/buildIndex"; import { buildStats } from "yt-dlp-transcript-common/controller/buildStats"; import { normalizeAllTranscripts } from "yt-dlp-transcript-common/controller/normalizeAll"; @@ -66,7 +66,7 @@ export async function buildStatsAction( fn: async (onLog, signal) => { await buildStats({ paths, onLog, signal }); // every site's Charts tab previews the dataset this just rebuilt. - revalidatePath("/sites/[siteId]/charts", "page"); + safeRevalidate([["/sites/[siteId]/charts", "page"]]); }, }); } diff --git a/editor/app/storage/lib/evictClipsJob.ts b/editor/app/storage/lib/evictClipsJob.ts @@ -1,4 +1,4 @@ -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { runManagedFunction, @@ -63,8 +63,8 @@ export async function enqueueEvictClipWindows(opts: { signal, }); onLog(evictClipWindowsSummary(result)); - revalidatePath("/storage"); - if (opts.slug) revalidatePath(`/channels/${opts.slug}`); + safeRevalidate(["/storage"]); + if (opts.slug) safeRevalidate([`/channels/${opts.slug}`]); }, }); } diff --git a/editor/app/storage/lib/repointJob.ts b/editor/app/storage/lib/repointJob.ts @@ -1,4 +1,4 @@ -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; import { enqueueRepoint } from "yt-dlp-transcript-common/controller/storageLocations"; @@ -35,9 +35,7 @@ export async function enqueueRepointJob(opts: { // Every channel on the location now reads from a different path, so the // pages that draw a media location are stale. No snapshot regen — the // kind is in NO_REGEN_KINDS, and a re-point changes no count. - revalidatePath("/storage"); - revalidatePath("/channels"); - revalidatePath("/"); + safeRevalidate(["/storage", "/channels", "/"]); }, }); } diff --git a/editor/app/storage/lib/savedVideosJob.ts b/editor/app/storage/lib/savedVideosJob.ts @@ -1,4 +1,4 @@ -import { revalidatePath } from "next/cache"; +import { safeRevalidate } from "../../lib/safeRevalidate"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { relocationQueueKey } from "yt-dlp-transcript-common/lib/queueKeys"; import { @@ -73,8 +73,7 @@ export async function enqueueSavedVideosRelocation(opts: { // where the store's bytes are, not what any channel has: every count in // every report is identical afterwards. What DOES change is what these // pages read per render. - revalidatePath("/storage"); - revalidatePath("/saved-videos"); + safeRevalidate(["/storage", "/saved-videos"]); }, }); } diff --git a/editor/e2e/ops-api.spec.ts b/editor/e2e/ops-api.spec.ts @@ -430,6 +430,65 @@ test("metadata-scan starts a job, and the job says it is a metadata-scan", async expect(unknownChannel.body.error).toContain('Channel "nope" not found'); }); +// A QUEUED JOB MUST NOT END `failed` FOR REVALIDATING (release 9, B1). +// +// The scheduler starts the next queued job synchronously when the one ahead +// finishes, so the next job inherits the async context of whatever finished +// it. When that is a job the auto-runner or the sync heartbeat started — no +// request anywhere up the chain — Next has no work store, and revalidatePath +// throws "static generation store missing". The sync body calls it last, so the +// throw landed after the work was done: the job read `failed` and `pnpm ops … +// --wait` exited 1 (the teamrcn sync, 2026-09-25). +// +// The holder is /api/test/stuck-job on platform:youtube, released after 10 s +// from OUTSIDE any request (`releaseAfterMs`; `detached` says the harness +// managed it). The sync is submitted behind it, so the release is what starts +// it — with no store, as live. +test("a sync queued behind another platform job ends done, not failed", async ({ + request, +}) => { + test.setTimeout(90_000); + await resetData("two-slow-channels"); + await settings(); + // slow-b without its --test-slow flag: a sync that finishes in seconds. + await writeChannelConfig("slow-b", { + name: "Slow B", + url: "https://www.youtube.com/@slow-b/videos", + }); + + type Meta = { status: string; queueKey?: string }; + const meta = (id: string) => + readJson<Meta>(`test-transcripts/.jobs/${id}.meta.json`).catch(() => null); + + const hold = await request.get( + `${baseUrl}/api/test/stuck-job?queue=${encodeURIComponent("platform:youtube")}&releaseAfterMs=10000`, + ); + expect(hold.ok()).toBe(true); + const holder = (await hold.json()) as { id: string; detached: boolean }; + expect(holder.detached).toBe(true); + + const sync = await ops(request, "sync", { slug: "slow-b" }); + expect(sync.body.ok).toBe(true); + const jobId = sync.body.jobId!; + // Behind the holder on the same platform queue: the release starts it. The + // 10 s window leaves room for a cold-compiled /api/ops POST and this poll to + // land before the release fires. + await expect.poll(async () => (await meta(jobId))?.status ?? null).toBe("queued"); + expect((await meta(jobId))?.queueKey).toBe("platform:youtube"); + + await expect + .poll(async () => (await meta(jobId))?.status ?? null, { + timeout: 60_000, + intervals: [1_000], + }) + .toMatch(/^(done|failed|cancelled)$/); + const log = await request.get(`${baseUrl}/api/jobs/${jobId}/log?from=0`); + expect(JSON.stringify(await log.json())).not.toContain( + "static generation store missing", + ); + expect((await meta(jobId))?.status).toBe("done"); +}); + test("refresh-report regenerates snapshot.json", async ({ request }) => { test.setTimeout(120_000); await resetData("title-filter-channel"); diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts @@ -17,6 +17,8 @@ export async function register() { // its transitive imports (runTick -> server actions, lmdb, fs) are Node-only, // so guard the import — and run nothing on Edge. if (process.env.NEXT_RUNTIME !== "nodejs") return; + // Anything a previous process left `queued` was queued before this instant. + const bootedAt = Date.now(); // Cancel in-flight children on a graceful shutdown. Armed FIRST, before any of // the runners below: a failure while starting those must not leave the server @@ -96,15 +98,69 @@ export async function register() { // Lazy-imported and `void`ed like everything else here: the controller pulls // in execa transitively (Node-only), and a disk that cannot be probed must // never block server readiness. + // Kept so the queued-meta pass below can wait for it: a re-queue for a + // channel whose location is mid-autoRepoint would be refused as unreachable. + let storagePass: Promise<unknown> = Promise.resolve(); try { const { runStorageBootPass } = await import( "yt-dlp-transcript-common/controller/storageLocations" ); - void runStorageBootPass({ enqueue: !idle }).catch(() => {}); + storagePass = runStorageBootPass({ enqueue: !idle }).catch(() => {}); } catch { /* a storage probe that fails to start must not block server readiness */ } + // STALE `queued` METAS FROM THE LAST PROCESS (release 9, B4b). A job still + // waiting in its queue when the server stopped never got its terminal meta + // write, so /jobs would read it as queued forever. Each is closed + // `cancelled` with a reason; only a few are re-queued (through the path Retry + // uses): never a sync (the heartbeat re-derives those, paced), nothing older + // than 24 h, and only the newest of duplicate specs. Runs on an idle boot + // too, but there it ONLY cancels — an idle boot must not resume work — and so + // does the e2e test server, whose leftover metas belong to a previous run's + // fixture. It waits for the storage pass above, so a channel being + // re-pointed is reachable when its job is re-queued. See + // common/jobs/bootQueuedJobs.ts. Lazy, voided, best-effort: never blocks + // readiness. + try { + const { settleQueuedJobMetas } = await import( + "yt-dlp-transcript-common/jobs/bootQueuedJobs" + ); + const { getPaths } = await import("yt-dlp-transcript-common/lib/paths"); + const { getRegistry } = await import( + "yt-dlp-transcript-common/jobs/registry" + ); + const testServer = process.env.EDITOR_TEST_ROUTES === "1"; + const cancelOnly = idle || testServer; + void storagePass + .then(() => + settleQueuedJobMetas({ + paths: getPaths(), + bootedAt, + isLive: (id) => getRegistry().get(id) !== undefined, + log: (line) => console.log(line), + idleReason: idle + ? "idle boot" + : testServer + ? "test server" + : undefined, + requeue: cancelOnly + ? null + : async (spec) => { + const { runJobSpec } = await import("./app/jobs/runJobSpec"); + const res = await runJobSpec(spec); + if (!res.ok) return { ok: false, error: res.error }; + // Nobody reads this stream; release it as the ops routes do. + void res.stream.cancel(); + return { ok: true, jobId: res.jobId }; + }, + }), + ) + .catch(() => {}); + } catch { + /* a boot pass that fails to start must not block server readiness */ + } + // Everything past here STARTS work. On an idle boot, nothing does. if (idle) return; diff --git a/homepage/app/components/ArchiveGrowthChart.tsx b/homepage/app/components/ArchiveGrowthChart.tsx @@ -8,7 +8,8 @@ import { } from "yt-dlp-transcript-common/lib/homepageChart"; // The front page's showpiece: every official instance's back catalogue as -// stacked strata, one month per step, from the oldest upload to this month. +// stacked strata, one month per step, from the oldest upload to the last +// complete month (the build month is partial, so the summary never emits it). // // BUILD-TIME ONLY. A server component that renders inline SVG from the summary; // no chart library and no client JS ships on `/`. The hover layer is native: diff --git a/homepage/app/components/FamilyStats.tsx b/homepage/app/components/FamilyStats.tsx @@ -25,7 +25,14 @@ export function FamilyStats({ const builtOk = built && !Number.isNaN(built.getTime()); return ( <div className="flex flex-col gap-3"> - <dl className="grid grid-cols-2 sm:grid-cols-3 lg:grid-cols-7 border-t border-l border-[var(--border)]"> + {/* Seven lg columns when the gone cell (span 2) is drawn, five when it + is not — otherwise the row ends in two empty columns. Both class + strings are complete literals (Tailwind v4 scans source). */} + <dl + className={`grid grid-cols-2 sm:grid-cols-3 ${ + official.gone > 0 ? "lg:grid-cols-7" : "lg:grid-cols-5" + } border-t border-l border-[var(--border)]`} + > {cells.map((c, i) => ( <div key={c.label} diff --git a/homepage/app/page.tsx b/homepage/app/page.tsx @@ -122,8 +122,8 @@ export default function Home() { )} {/* ── What it does ─────────────────────────────────────────────────── */} - {/* One paragraph, not three titled panels: without titles, three - one-sentence items are just a paragraph with bullets. */} + {/* Three short paragraphs, not three titled panels: without titles, + three one-sentence items read as prose, not as a feature grid. */} <section className="border-t border-[var(--border)]"> <div className={`${CONTAINER} py-14 sm:py-20`}> <h2 className="font-display text-2xl font-semibold text-[var(--foreground)]"> diff --git a/plans/STATE.md b/plans/STATE.md @@ -15,6 +15,16 @@ with it, `BUILD_ID` `a9YAe_dwhuM6DiCPzIwMD`. Still live from release 7: See the [release-8 rollout record](release-8.md#rollout-2026-09-25-late-afternoon--0e72ef73-live-on-3001-second-restart-of-the-day). +**Evening 2026-09-25 (in progress).** Homepage refresh is production (`a684a0ce`). Paramount Tactical +finished on the operator's "queue step 1": metadata-scan job `01M3DBFV07NC0Z6KWZ1BG3K9BX` scanned the +last 201 (0 errors; store 1,370 entries + 12 members-only errors of 1,473); download-missing +`01M3DD3WBVYPPXB1NCQ1HDX8BM` found 1,381 missing, 92 complete, excluded 12 members-only, and +**settled 1,369 by the include filter with no request and no sleep**; `keep-videos` matched 1 +(`NV1QqS9NOiU`, already kept). Both ops jobs read `failed` on the revalidate invariant (release 9 B1). +The include filter is still `Quartering`: clearing it (Configure) re-evaluates the 1,369 for the +auto-download lane — the operator's call. Release 9 (`one-core/r9-fixes`) and hub slice C1 +(`one-core/c1-hub-look`) are reviewed and in fixes; see `plans/release-9.md` once merged. + **Last updated:** 2026-09-25 (late afternoon). **Release 8 is live, the day's second restart.** - Release 8 (`bb3dbb4c` → `0e72ef73`) has two slices, each Opus-reviewed: SHIP AFTER FIXES, then SHIP. diff --git a/plans/release-9.md b/plans/release-9.md @@ -1,13 +1,233 @@ -# Release 9 — the hub, polished (+ fixes) +# Release 9 — four fixes behind one restart -`main` at `9247211e`, release 8 live 2026-09-25. The operator withheld the homepage's link to the hub -(`HUB_LINK_ENABLED = false`) "until it is more polished". Release 9 is that polish — C1 brings the hub -into the homepage's look and copy with the homepage's own numbers, C2 the federated-search UX — plus -the fixes slice (`one-core/r9-fixes`). Flipping the homepage's hub link is a later step, not a slice. -Rules: `plans/tools/implementer-rules.md`. Record file: this file. +`main` at `9247211e`, release 8 live on :3001 since 2026-09-25 14:43. Release 9 is the four small +bugs the release-8 rollout found on the live editor, batched so they cost one editor restart: a +queued job that ends `failed` after doing its work, a 30 s sleep after every video the download +filter declined, YouTube with no request-level pacing, and a `/jobs` page that calls a slow scan +stuck and keeps restart-orphaned jobs `queued` forever. The release-8 reviews' optional lows ride +along. The same release also carries the hub's polish, which the operator asked for before the homepage +links the hub again (`HUB_LINK_ENABLED = false`): C1 brings the hub into the homepage's look and copy +with the homepage's own numbers, and C2 the federated-search UX. Flipping the homepage's hub link is a +later step, not a slice. Rules: `plans/tools/implementer-rules.md`. Record file: this file. ## Record +### Slice F, as shipped — four fixes (2026-09-25) + +Branch `one-core/r9-fixes` off `main` `9247211e`, one Opus implementer, no sibling slices. + +**B1 — a queued job no longer ends `failed` after succeeding.** The scheduler starts the next +queued job synchronously inside `complete()`, so that job inherits the async context of whatever +finished the one ahead. When the chain began at a job with no request behind it (an auto-runner +unit, a heartbeat sync), Next has no work store and `revalidatePath` throws `Invariant: static +generation store missing`. Job bodies call it last, so the throw landed after the work was done: +the teamrcn sync at 15:21 read `failed` and `pnpm ops … --wait` exited 1. +`editor/app/lib/safeRevalidate.ts` — `safeRevalidate(paths, tags = [])` — swallows exactly that +invariant (one `console.warn` per process, naming the paths) and rethrows anything else. It is +swapped in at every job body (`fn`) and job hook (`onDone` / `afterRun` / `afterDone`), found by +grep and by indentation, not from the prompt's list alone: 18 files, 63 call lines, including +`repointJob.ts`'s `afterDone`, which the list did not name. Server actions that revalidate in +their own request are untouched, including the ones that revalidate after `drainStream` +(`review/actions.ts`, `refreshAllChannelSnapshotsAction`), because they are still inside it. + +The first e2e version queued the sync behind a `--test-slow` sync. That held the queue for minutes, +and it could not fail: a job an ops request started carries that request's store. The final version +holds `platform:youtube` with `/api/test/stuck-job?releaseAfterMs=4000`. The new parameter finishes +the fake holder from a timer armed inside `workAsyncStorage.exit`, i.e. outside any request, and +the route reports `detached` so the spec cannot pass for the wrong reason. The spec then queues an +ops sync behind it. **Against the pre-fix sync body it fails with the live line** — `[error] +Invariant: static generation store missing in revalidatePath /channels/slow-b`, status `failed` +(`r9-e2e-prefix.log`) — and it passes with the fix. + +**B2 — no 30 s sleep after a video the download filter declined.** `runManagedDownloads` slept +`sleepBetweenDownloadsSeconds` after every video. `declinedWithoutMediaFetch(outcome)` is true when +every attempt was the `n: 0` metadata prefetch (or its cookie retry) and the video ended +`skipped-filtered`. Any attempt `n >= 1` is a real fetch: the download attempts and the chat-only +pass. A failed chat pass leaves the status `skipped-filtered` and still sleeps. As first shipped, a +`per_video` failure that never got past the prefetch also skipped the sleep. The review narrowed +that (see "Review fixes"): **every failure keeps the pace.** `runManagedDownloads` is exported with a +`deps` seam (`downloadOne`, `sleep`). `managedDownloadsSleep.test.ts` drives the loop with canned +outcomes: +- a fetched video sleeps; +- a filtered one does not; +- a per-video failure at the prefetch sleeps; +- the last video never sleeps; +- a failed real fetch and a network failure at the prefetch both sleep; +- a unit table for the predicate. + +**B3 — YouTube request pacing.** `PLATFORM_ARGS.youtube = ["--sleep-requests", "1"]`, with the +comment citing `~/reports/release-7/data/q-429-report.md` finding 1. The prompt named +`release-8/data`; the report is in `release-7/data`. The comment above the table now says what the +pace is, and that the metadata scan and the clip window, which already pass `--sleep-requests 1`, +carry it twice, harmlessly (yt-dlp keeps the last). The argv tests were updated deliberately: +- `channelArgs.test.ts`: youtube gets the pace; a platform with no entry gets nothing; a channel's + own `--sleep-requests 3` comes after and wins. +- `platform-args.test.mjs`: +1 case, youtube carries the pace and a URL on no known platform + carries nothing. The YouTube clip fetch and simulate now expect the pace. + +No e2e asserts a YouTube argv exactly. `rumble-sweep.spec.ts` only asserts Rumble lines, and +`fake-ytdlp.mjs`'s `lastNonFlag` never sees the trailing `1`, because every spawn ends in a URL, +`--load-info-json <path>` or `-a <file>`. + +**B4a — `/jobs` flags a stall by quiet time, not age.** `JobRecord.progressAt` is stamped at the +source. `noteProgress` (the `setProgress` every `runManagedFunction` job gets) stamps it only when +the snapshot's numbers changed, and `recordTaskDuration` stamps it when a sub-operation finishes. +`reconcileSlots` measures quiet time from the later of the start and that stamp. This deviates from +the prompt's suggested view-side map, on purpose. A map is only fed while someone has `/jobs` open, +so the first render after an hour away would either flag a healthy scan (if it seeded from the +start) or hide a real stall for 10 more minutes (if it seeded from now). The stamp has neither +problem, and the registry is in memory, so it is still evicted with the record. +`jobRows.test.ts` +2: an hour of 6 s steps is never stuck, and progress frozen past 10 min is +possibly-stalled. The same count re-reported is not a move. **Both fail on the old rule.** The +`/api/test/stuck-job` fixture backdates `startedAt` with no `progressAt`, so `queue.spec.ts`'s +force-release case still sees its stall. + +**B4b — boot settles the metas a restart left `queued`, and re-queues little.** +`common/jobs/bootQueuedJobs.ts` (`settleQueuedJobMetas`) runs from `instrumentation.ts`: +lazy-imported, voided and best-effort. It runs after the storage boot pass, and since the review +fix it WAITS for that pass. Its scope is each `*.meta.json` with status `queued` that was queued +before this boot and is not held by the live registry. A malformed meta is skipped. In order: +1. An idle boot (`ARCHILYZER_IDLE_BOOT`) and the e2e test server (`EDITOR_TEST_ROUTES=1`) cancel + everything and re-queue nothing. +2. A `sync` is never re-queued: "server restarted; the scheduler re-derives syncs". +3. A meta queued more than 24 h (`REQUEUE_MAX_AGE_MS`) before the boot is cancelled: "queued + before the last restart, stale". +4. A meta with no spec is cancelled with the reason. +5. Of what is left, only the NEWEST meta per kind + channel + bucket + params (params key-sorted) + is re-queued, through `runJobSpec` (the path Retry uses). The old meta is closed `cancelled`, + naming the new id. The others are cancelled as "superseded by a newer queued job (<id>)". A + refused or throwing re-queue is cancelled with its error. + +The log gets one line per re-queued job and one summary line: re-queued N, cancelled M, by reason. +The reason is also appended to each job's `.log`, which makes the pair prunable, since pruning walks +`.log` ids. Metas are written tmp + rename. `running` metas are left alone. `JobMeta` gains the +optional `cancelReason`. + +**Live size, from the review's read-only count.** 1,396 metas, 392 `queued`: +- 346 `sync` over 51 channels, in batches from 08-03 to 09-24; +- 29 `whisper-all` from June and July; +- about 7 others with a spec; +- 10 with none. + +The first version would have called `runJobSpec` about 382 times at boot. Under the current rules, +the syncs and everything older than 24 h are cancelled, so at most a handful of fresh non-sync jobs +come back. + +`bootQueuedJobs.test.ts` has 9 cases over a temp `.jobs` dir with an injected re-queue. One is 7 +duplicate syncs + 1 stale `whisper-all` + 1 fresh non-sync, which gives exactly 1 re-queue. Other +cases cover duplicate specs written in a different key order, and a malformed meta next to a good +one. It was not run against the real corpus. + +**Lows, all three done.** +- **S:** `mergeBackoffEntry` merges `until` and `fails` separately, each to its max, in both the + runner's merge and `recordDownloadBackoff`'s. `platformBackoff.test.ts` +1. +- **V:** a changed `data/` mtime REBASES the channel's title memo instead of dropping it. One + `readdir` keeps the titles of the dirs still present, only new dirs are read, and a removed + dir's title goes with it. The memo case in `videoTitles.test.ts` now asserts 1 read for a new + dir and a fall-through for a removed one. It pins a distinct mtime, because two changes inside + one filesystem tick share one. +- **H:** + - `official.transcripts` is a per-site accumulator from the same pass that places each + transcript. + - A site outside `keptIds` gains no monthly key. + - `gone` counts only with a `downloadedDate`. + - `FamilyStats` is `lg:grid-cols-5` without a gone cell. + - The stale comments in `ArchiveGrowthChart.tsx` :10-11 and `page.tsx` :125, and the test + title, are fixed. + + `homepageSummary.test.ts` +3. No number on today's data moves. The review measured placed + + unplaced = `official.transcripts` (49,767) and all 480 deleted records with a `downloadedDate`. + +| sha | what | +|---|---| +| `d927dad4` | B1: `app/lib/safeRevalidate.ts` + its unit test (3), swapped in at every job body and hook (18 files); first ops-api e2e case | +| `37d263cb` | B2: `declinedWithoutMediaFetch`, the `deps` seam on `runManagedDownloads`, `managedDownloadsSleep.test.ts` (6) | +| `ac55a23c` | B3: `PLATFORM_ARGS.youtube`, table comment, `channelArgs.test.ts` (+2), `platform-args.test.mjs` (+1) | +| `efd75cb0` | B4a: `JobRecord.progressAt`, `noteProgress`, quiet-time stall rule, `jobRows.test.ts` (+2) | +| `cc884951` | B4b: `bootQueuedJobs.ts` + test (6), `JobMeta.cancelReason`, the instrumentation boot pass | +| `db18b73d` | low S: `mergeBackoffEntry` in both merges, test (+1) | +| `49993934` | low V: title memo rebased on a new dir, test updated | +| `c714bc14` | low H: `homepageSummary` counts + guard, `FamilyStats` grid, comments, test (+3) | +| `78d165f6` | B1 e2e reworked: `stuck-job?releaseAfterMs` (detached release), the spec now reproduces the live failure pre-fix | +| `0c3e0b9a` | `plans:` this record, the `[Unreleased]` bullets | +| `be48f610` | (review HIGH-1, LOW-1, LOW-2, NIT-1) boot pass: no syncs, 24 h age cap, newest per spec, summary log; storage pass awaited; atomic meta writes; test 6 → 9 | +| `11c947ab` | (review MED) B2 skips the sleep for `skipped-filtered` only; the per-video test now expects a sleep | +| `23134785` | (review LOW-3) the B1 e2e releases its holder after 10 s, not 4 | +| _this_ | `plans:` the review fixes in this record, the CHANGELOG and the gates | + +**Gates**, all from the worktree root. Heavy steps started at ≥ 3 GB available memory. +- **tsc** (`pnpm -r --no-bail --workspace-concurrency=1 exec tsc --noEmit`) was clean before every + code commit. It was checked once to catch a planted error, so the silent `exit=0` is real. +- **common: 1,836/1,836**, from 1,816 + 20: B2 +6, B3 +2, B4a +2, B4b +6, S +1, H +3; V's case was + rewritten in place. +- **editor unit: 78/78**, from 75 + 3 (`safeRevalidate.test.ts`). +- **`test:scripts`: 162 pass + 1 skip of 163**, from 161 + 1 skip, +1 in `platform-args.test.mjs`. + A first run during e2e run 1 failed `E2E_QUEUE=0 bypasses the queue entirely`, because this + worktree's own e2e held the lock; the rerun with no e2e running is clean. +- **mcp: 219/219.** +- **Builds:** `pnpm --filter editor exec next build` ok, `pnpm --filter export exec next build` ok, + and, because H touches `homepage/`, `pnpm --filter homepage exec next build` ok. No dangling + `export/public` links. +- **EDITOR e2e.** Spec list `r9-specs.txt`: `ops-api queues auto-queue lane-runner jobs-filters + queue` (the grep for `possibly-stalled|STUCK|data-job-id` adds `queue.spec.ts`), plus + `jobs-retry jobs-active-order jobs-channel` (the retry path B4b reuses), `rumble-sweep` (argv), + `title-filter` (B2's filtered path), `pacing metadata-scan-botcheck sync-deep`, all `.spec.ts`. + - Run 1 (`r9-e2e1.log`), on `c714bc14`: **84 passed, 2 failed, 13.4 min**. One failure was the + first B1 spec, which timed out behind the slow sync (the rework above). The other was + `auto-queue.spec.ts:218`, the run's first test, where the `/api/auto-queue/status` GET outlived + the 30 s test timeout on a cold dev compile. It passed in run 2 with no change. + - Pre-fix proof (`r9-e2e-prefix.log`): the reworked B1 spec with the sync body's + `revalidatePath` restored: **1 failed**, with the live invariant in the job log. The fix was + restored before any commit. + - Run 2 (`r9-e2e2.log`), on `78d165f6`: **86 passed, 0 failed, 5.7 min**, with no queue wait. + +**Review fixes (review verdict SHIP AFTER FIXES, `r9-review.md`).** +- **HIGH-1** (the boot-flood size), **LOW-1** (await the storage pass), **LOW-2** (malformed-meta + test) and **NIT-1** (atomic meta writes) are all in `be48f610`. +- **MED** (B2 keeps the pace after a per-video failure) is `11c947ab`. +- **LOW-3** (10 s release window) is `23134785`. +- **LOW-4** is left; see "Found and left". + +Re-gate on `23134785`: +- tsc clean. +- common **1,839/1,839** (1,836 + 3 boot-pass cases). +- editor unit **78/78**. +- EDITOR e2e `jobs-filters.spec.ts ops-api.spec.ts queues.spec.ts queue.spec.ts` + (`r9-e2e3.log`): **57 passed, 0 failed, 2.7 min**. + +No e2e covers the boot pass. The test server is cancel-only by design, so the pass is pinned by its +9 unit cases. + +**Numbers: none** (per the prompt). No `settings.json`, `site.json` or `config.json` key changed. +`JobMeta.cancelReason` and `JobRecord.progressAt` are additive and optional. + +**Found and left.** +- **B2 lets filtered prefetches run back to back.** A filtered video's metadata prefetch is still + one yt-dlp process against YouTube. Without the 30 s sleep, a run of declined videos is a run of + prefetches, spaced only by process start-up and B3's one-second request pace inside each. The + metadata scan does the same work at ~10/min with no 429s. If a filtered channel's + download-missing ever 429s, the lever is a smaller sleep for prefetch-only videos, not the full + one. +- **Classification quirk, left as is.** YouTube's soft block ("This content isn't available, try + again later") classifies as `deleted` → `per_video` in `classifyDownloadFailure`. So a soft block + neither triggers the platform backoff nor aborts a batch, and it reads as a per-video property. It + is why B2 keeps the sleep after every per-video failure. `classifyDownloadFailure` was not changed + (review instruction). +- **B4b leaves `running` metas.** A job running at a hard crash (no graceful shutdown) still reads + `running` → `archived` on `/jobs`. Whether a half-done job should re-run is not a boot pass's + call. +- **`cancelReason` is not drawn on `/jobs`.** It is in the meta and appended to the job's log, + which the row's log view shows. +- **The rollout's first boot settles the 392 live `queued` metas.** Nearly all are cancelled with + a reason, as above; only fresh (< 24 h) non-sync jobs, newest per spec, come back. Booting once with + `ARCHILYZER_IDLE_BOOT=1` re-queues none. +- **LOW-4, left:** `safeRevalidate`'s "once per process" warning is once per module instance. Next + can load a module more than once, so it may warn a second time. +- **The stuck-job harness reaches into a Next internal** + (`next/dist/server/app-render/work-async-storage.external`). It is typed, test-route-only, and + reported through `detached`, so a Next upgrade that moves it fails the spec loudly rather than + letting it pass vacuously. + ### Slice C1, as shipped — the hub in the family's register (2026-09-25) Branch `one-core/c1-hub-look` off `main` `9247211e`. **The operator's decision (2026-09-25 evening):** diff --git a/umtool/report-to-video/platform-args.test.mjs b/umtool/report-to-video/platform-args.test.mjs @@ -30,16 +30,27 @@ test("a Rumble clip fetch carries the rumble platform args, before the retry's e assert.equal(argv[argv.indexOf("--download-sections") + 1], "*12.00-20.50"); }); -test("a YouTube clip fetch gets no platform args", () => { +test("a YouTube clip fetch gets youtube's request pace and no impersonation", () => { const argv = clipFetchArgs({ url: YOUTUBE, from: 1, to: 2, fmt: "b", dest: "/x/out.mp4" }); - assert.deepEqual(slice(argv, "--force-keyframes-at-cuts", "-f"), []); + assert.deepEqual(slice(argv, "--force-keyframes-at-cuts", "-f"), ["--sleep-requests", "1"]); assert.ok(!argv.includes("--impersonate")); }); -test("the availability simulate carries the platform args for Rumble only", () => { +test("youtube carries --sleep-requests 1; a URL on no known platform carries nothing", () => { + assert.deepEqual([...PLATFORM_ARGS.youtube], ["--sleep-requests", "1"]); + const other = "https://example.com/some/video"; + assert.deepEqual(slice(simulateArgs(other), "--no-warnings", "--"), []); + assert.deepEqual( + slice(clipFetchArgs({ url: other, from: 1, to: 2, fmt: "b", dest: "/x/o.mp4" }), "--force-keyframes-at-cuts", "-f"), + [], + ); +}); + +test("the availability simulate carries each platform's args", () => { const rumble = simulateArgs(RUMBLE); assert.deepEqual(slice(rumble, "--no-warnings", "--"), [...PLATFORM_ARGS.rumble]); assert.deepEqual(rumble.slice(-2), ["--", RUMBLE]); assert.ok(rumble.includes("--simulate")); assert.ok(!simulateArgs(YOUTUBE).includes("--impersonate")); + assert.deepEqual(slice(simulateArgs(YOUTUBE), "--no-warnings", "--"), ["--sleep-requests", "1"]); });