Archilyzer · Source

archilyzer

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

commit 0e84914b7fd62b967973d8b315704ffba3d483ee
parent 5896dd4867315fc3b9ae3375a96eb3a30f50041e
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Wed,  7 Oct 2026 17:51:53 -0400

Merge fetch-windows (a paced batch of clip windows, one job per platform)

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

Diffstat:
MPUBLISH.md | 22++++++++++++++++------
Acommon/controller/fetchWindows.test.ts | 241+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/fetchWindows.ts | 469+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/jobDetail.test.ts | 29+++++++++++++++++++++++++++++
Mcommon/jobs/jobDetail.ts | 9++++++++-
Mcommon/jobs/jobKinds.test.ts | 3+++
Mcommon/jobs/jobKinds.ts | 17+++++++++++++++++
Mcommon/jobs/registry.ts | 6+++++-
Acommon/publish/missingEvidenceWindows.test.ts | 131+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/publish/reportMedia.ts | 150+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcommon/views/jobRows.ts | 12+++++++-----
Mcommon/ytdlp/fetchWindowManaged.ts | 30+++++++++++++++++++++++++++++-
Meditor/CHANGELOG.md | 1+
Meditor/app/api/media/fetch-window/[jobId]/route.ts | 6++++--
Aeditor/app/api/ops/fetch-windows/route.test.ts | 128+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/fetch-windows/route.ts | 179+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/channels/[slug]/videos/fetchWindowsAction.ts | 406+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/components/JobProgressBars.tsx | 3+++
Meditor/app/jobs/jobReplayRegistry.ts | 49+++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/sites/[siteId]/reports/page.tsx | 15++++++++++-----
Aeditor/app/sites/components/FetchMissingEvidenceButton.tsx | 91+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/e2e/fetch-window.spec.ts | 84+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Meditor/e2e/fixtures/bin/fake-ytdlp.mjs | 7+++++++
Mscripts/archilyzer-ops.mjs | 17+++++++++++++++++
Mscripts/archilyzer-ops.test.mjs | 9+++++++++
Mumtool/docs/quirks.md | 5+++++
Mumtool/report-to-video/fetch-via-editor.mjs | 137+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
27 files changed, 2207 insertions(+), 49 deletions(-)

diff --git a/PUBLISH.md b/PUBLISH.md @@ -63,6 +63,7 @@ entry points in `common/publish/build.ts`. | The hub | /sites → Hub → **Build hub** / **Deploy hub** | `build-hub`, `deploy-hub` | `build hub`, `deploy hub [--preview <branch>]` | | The homepage | /sites → Homepage → **Build homepage** (tick *Deploy after build*) / **Deploy homepage**, with an optional preview branch | `build-homepage` (`{"deploy":true}` to deploy after), `deploy-homepage` (`{"preview":"<branch>"}`) | `build homepage [--no-source]`, `deploy homepage [--preview <branch>]` | | The source mirror alone | — (every homepage build runs it) | — | `source publish [--force] [--check] [--keep-scratch]`, `source audit [<git dir>]` | +| Fetch the clip windows a site's reports cite and the disk lacks | a site's **Reports** tab → *Fetch missing evidence* (*Preview missing evidence* lists them) | `fetch-windows` (`{"siteId", "dryRun"?}`) | — | | A site's report evidence media (then its exports) | a site's **Reports** tab → *Prepare evidence media* | `reports-prepare` (`{"siteId"}`) | `reports prepare <id>` | | A site's report exports (HTML, PDF, Markdown, evidence pack) | a site's **Reports** tab → *Export reports* | `reports-export` (`{"siteId", "reportId"?, "formats"?}`) | `reports export <id> [--report <rid>] [--formats html,pdf,md,zip]` | | A report from a /sweep report, an /ask answer or a report-to-video manifest, and a starter manifest from a report | — | — | `reports convert <sweep\|ask\|manifest> <in> --out <report.json> [--channels-dir <dir>]`, `reports to-manifest <report.json> --out <manifest.json>` | @@ -74,16 +75,25 @@ this machine has what a build needs. From `export/`, `pnpm run build` is `archilyzer build site` and `pnpm run deploy` is `archilyzer deploy site`; both take the site from `SITE_ID` when no id is given. -A site with `reports` (site.json) has one more step, before its build and on the -host: **reports prepare** cuts the evidence clip of every span its published reports +A site with `reports` (site.json) has two more steps, before its build and on the +host. First, **fetch the missing evidence**: *Fetch missing evidence* on the Reports +tab, or `pnpm ops fetch-windows --json '{"siteId":"<id>"}' --wait`, fetches the +window of every cited span whose media is not on disk, through the editor's managed +path, as one paced job per platform (YouTube and Rumble side by side; at least 30 s +between two windows on one platform, a 429 or two 403s in a row backing the +platform off and stopping its job). `"dryRun": true` (*Preview missing evidence*) +lists the windows, and the spans no window can fill — a video the source says is +deleted or private, a channel off the site, a drive that is not mounted — without +fetching. Running it again resumes: fetched windows are on disk. Then **reports +prepare** cuts the evidence clip of every span its published reports cite — from the editor's clip windows, the saved-video store or a record's audio, fitted inside 1280×720 (H.264 crf 23, AAC; an audio span is an `.m4a`) — and copies the screenshot and media of every cited post, only those, into `.export-index/sites/<id>/report-media/`, with a manifest (`index.json`) of each -moment's file, size, hash and duration. Nothing is fetched: a citation whose media -is not on disk, a clip over 24 MiB or an invalid report is listed and fails the run -(exit 1, or a failed job) — fetch the window or persist the video, capture the -post, and run it again; what is already cut is reused. +moment's file, size, hash and duration. Prepare fetches nothing: a citation whose +media is not on disk, a clip over 24 MiB or an invalid report is listed and fails +the run (exit 1, or a failed job) — fetch the missing evidence or persist the video, +capture the post, and run it again; what is already cut is reused. When nothing is missing, prepare ends by **exporting** the reports (`reports export` runs the same step alone): each published report is written, as compose would diff --git a/common/controller/fetchWindows.test.ts b/common/controller/fetchWindows.test.ts @@ -0,0 +1,241 @@ +// fetchWindows (the batch) against a fake yt-dlp, through the real +// fetchWindowManaged: what the platform answers is decided by a marker in the +// item's URL, and every shared-state dependency — the cooldown, the hold, the +// backoff, the clean record, the pause — is injected and recorded. +// +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test controller/fetchWindows.test.ts + +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { chmod, mkdir, mkdtemp, readFile, writeFile } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import { + CLIP_WINDOW_MIN_GAP_SECONDS, + fetchWindows, + type FetchWindowsDeps, + type FetchWindowsItem, +} from "./fetchWindows"; +import type { Paths } from "../lib/paths"; +import type { ChannelConfig } from "../lib/channelConfig"; +import type { SiteSettings } from "../lib/settings"; +import type { JobProgress } from "../jobs/registry"; + +const ROOT = await mkdtemp(path.join(os.tmpdir(), "fetchwindows-")); +const BIN = path.join(ROOT, "fake-ytdlp.mjs"); +const ARGS_LOG = path.join(ROOT, "argv.jsonl"); +process.env.FAKE_ARGS_LOG = ARGS_LOG; + +// The URL decides: `/403` a Cloudflare refusal, `/429` a rate limit, `/gone` a +// removed video; anything else writes the window. +await writeFile( + BIN, + `#!/usr/bin/env node +import { appendFileSync, writeFileSync } from "node:fs"; +const args = process.argv.slice(2); +appendFileSync(process.env.FAKE_ARGS_LOG, JSON.stringify(args) + "\\n"); +const url = args[args.length - 1]; +const fail = (msg) => { process.stderr.write(msg + "\\n"); process.exit(1); }; +if (url.includes("/403")) fail("ERROR: [download] Got error: HTTP Error 403: Forbidden"); +if (url.includes("/429")) fail("ERROR: [youtube] x: HTTP Error 429: Too Many Requests"); +if (url.includes("/gone")) fail("ERROR: [youtube] x: Video unavailable. This video has been removed by the uploader"); +writeFileSync(args[args.indexOf("-o") + 1], "mp4"); +`, +); +await chmod(BIN, 0o755); + +let run = 0; +// A fresh corpus per test, so one test's fetched windows are not another's cache. +async function corpus(): Promise<Paths> { + const channelsDir = path.join(ROOT, `corpus-${++run}`, "channels"); + await mkdir(channelsDir, { recursive: true }); + return { ytdlpBin: BIN, channelsDir } as unknown as Paths; +} + +async function spawns(): Promise<number> { + return (await readFile(ARGS_LOG, "utf8").catch(() => "")).trim().split("\n").filter(Boolean).length; +} + +const CONFIG = { url: "https://www.youtube.com/@x", platform: "youtube" } as unknown as ChannelConfig; + +type Recorded = { + sleeps: number[]; + backoffs: [string, string][]; + cleans: string[]; + progress: JobProgress[]; +}; + +function harness(over: Partial<FetchWindowsDeps> = {}) { + const rec: Recorded = { sleeps: [], backoffs: [], cleans: [], progress: [] }; + const deps: Partial<FetchWindowsDeps> = { + sleep: async (ms) => { + rec.sleeps.push(ms); + }, + getSettings: () => ({ sleepBetweenDownloadsSeconds: 0 }) as unknown as SiteSettings, + readChannelConfig: async (_p, slug) => (slug === "c" ? CONFIG : null), + findVideoSourceUrl: async () => null, + assertTextReadable: async () => undefined, + cooldownRemainingMs: async () => 0, + heldRefusal: async () => null, + recordBackoff: async (platform, _p, cls) => { + rec.backoffs.push([platform, cls]); + }, + recordClean: async (platform) => { + rec.cleans.push(platform); + return null; + }, + ...over, + }; + return { rec, deps }; +} + +const item = (id: string, url: string, from = 10, to = 20): FetchWindowsItem => ({ + slug: "c", + id, + from, + to, + webpageUrl: `https://www.youtube.com${url}`, +}); + +async function go(paths: Paths, items: FetchWindowsItem[], h: ReturnType<typeof harness>, extra: { gapMs?: number; signal?: AbortSignal; drainSignal?: AbortSignal } = {}) { + return fetchWindows({ + paths, + items, + provenance: { requestedBy: "test" }, + gapMs: extra.gapMs ?? 1000, + signal: extra.signal, + drainSignal: extra.drainSignal, + setProgress: (p) => h.rec.progress.push(p), + deps: h.deps, + }); +} + +test("a cached window costs no request and no pause", async () => { + const paths = await corpus(); + const clips = path.join(paths.channelsDir, "c", "data", "v1", "clips"); + await mkdir(clips, { recursive: true }); + await writeFile(path.join(clips, "0.00-60.00.mp4"), "mp4"); + const h = harness(); + const before = await spawns(); + const r = await go(paths, [item("v1", "/a", 5, 15), item("v1", "/b", 30, 40), item("v2", "/c")], h); + assert.equal(r.cached.length, 2); + assert.equal(r.fetched.length, 1); + assert.equal((await spawns()) - before, 1, "only the uncached window spawned"); + assert.deepEqual(h.rec.sleeps, [], "one network fetch owes no pause"); + assert.deepEqual(h.rec.cleans, ["youtube"]); + assert.deepEqual(h.rec.progress.at(-1), { metric: "clips", initial: 0, target: 3, current: 3 }); +}); + +test("one 403 is an item failure; the next window is fetched after the gap", async () => { + const h = harness(); + const r = await go(await corpus(), [item("v1", "/403"), item("v2", "/ok")], h); + assert.equal(r.stopped, undefined); + assert.equal(r.failed.length, 1); + assert.equal(r.failed[0].class, "network"); + assert.equal(r.fetched.length, 1); + assert.deepEqual(h.rec.sleeps, [1000]); + assert.deepEqual(h.rec.backoffs, [], "one 403 does not back the platform off"); +}); + +test("two 403s in a row back the platform off and stop the run", async () => { + const h = harness(); + const r = await go(await corpus(), [item("v1", "/403"), item("v2", "/403"), item("v3", "/ok"), item("v4", "/ok")], h); + assert.equal(r.stopped, "network"); + assert.equal(r.failed.length, 2); + assert.deepEqual(r.notAttempted.map((i) => i.id), ["v3", "v4"]); + assert.deepEqual(h.rec.backoffs, [["youtube", "network"]]); +}); + +test("a per-video failure between two 403s breaks the streak", async () => { + const h = harness(); + const r = await go(await corpus(), [item("v1", "/403"), item("v2", "/gone"), item("v3", "/403"), item("v4", "/ok")], h); + assert.equal(r.stopped, undefined); + assert.deepEqual(r.failed.map((f) => f.class), ["network", "per_video", "network"]); + assert.equal(r.fetched.length, 1); + assert.deepEqual(h.rec.backoffs, []); +}); + +test("a 429 records the backoff and stops at once", async () => { + const h = harness(); + const r = await go(await corpus(), [item("v1", "/ok"), item("v2", "/429"), item("v3", "/ok")], h); + assert.equal(r.stopped, "rate-limit"); + assert.equal(r.fetched.length, 1); + assert.deepEqual(r.notAttempted.map((i) => i.id), ["v3"]); + assert.deepEqual(h.rec.backoffs, [["youtube", "rate_limit"]]); +}); + +test("a platform cooling down (or held) ends the run before its next fetch", async () => { + let calls = 0; + const h = harness({ cooldownRemainingMs: async () => (calls++ === 0 ? 0 : 60_000) }); + const before = await spawns(); + const r = await go(await corpus(), [item("v1", "/ok"), item("v2", "/ok"), item("v3", "/ok")], h); + assert.equal(r.stopped, "cooldown"); + assert.equal(r.fetched.length, 1); + assert.deepEqual(r.notAttempted.map((i) => i.id), ["v2", "v3"]); + assert.equal((await spawns()) - before, 1); + + const held = harness({ heldRefusal: async () => "youtube is held." }); + const r2 = await go(await corpus(), [item("v1", "/ok")], held); + assert.equal(r2.stopped, "held"); + assert.deepEqual(r2.notAttempted.map((i) => i.id), ["v1"]); +}); + +test("a drain stops between windows; the one in flight finishes", async () => { + const drain = new AbortController(); + const h = harness({ + recordClean: async () => { + drain.abort(); + return null; + }, + }); + const r = await go(await corpus(), [item("v1", "/ok"), item("v2", "/ok")], h, { drainSignal: drain.signal }); + assert.equal(r.stopped, "drain"); + assert.equal(r.fetched.length, 1); + assert.deepEqual(r.notAttempted.map((i) => i.id), ["v2"]); +}); + +test("the default gap is the platform's, floored at the clip-window minimum and jittered", async () => { + const h = harness(); + const r = await fetchWindows({ + paths: await corpus(), + items: [item("v1", "/ok"), item("v2", "/ok"), item("v3", "/ok")], + provenance: { requestedBy: "test" }, + deps: h.deps, + }); + assert.equal(r.fetched.length, 3); + assert.equal(h.rec.sleeps.length, 2); + for (const ms of h.rec.sleeps) { + assert.ok(ms >= CLIP_WINDOW_MIN_GAP_SECONDS * 1000, `${ms} is at least the floor`); + assert.ok(ms <= CLIP_WINDOW_MIN_GAP_SECONDS * 1500 + 60_000, `${ms} is at most the floor plus half`); + } +}); + +test("an unknown channel, an unreadable one and a missing URL fail their items and the run goes on", async () => { + const h = harness({ + readChannelConfig: async (_p, slug) => (slug === "nope" ? null : CONFIG), + assertTextReadable: async (_p, slug) => { + if (slug === "off") throw new Error("off/data is not a readable directory"); + }, + }); + const r = await go( + await corpus(), + [ + { ...item("v1", "/ok"), slug: "nope" }, + { ...item("v1", "/ok"), slug: "off" }, + { slug: "c", id: "v9", from: 1, to: 2 }, + item("v2", "/ok"), + ], + h, + ); + assert.deepEqual(r.failed.map((f) => f.class), ["unknown-channel", "unreachable", "no-url"]); + assert.equal(r.fetched.length, 1); + assert.deepEqual(h.rec.sleeps, [], "no network fetch preceded the one that ran"); +}); + +test("a duplicated window is fetched once", async () => { + const h = harness(); + const before = await spawns(); + const r = await go(await corpus(), [item("v1", "/ok"), item("v1", "/ok")], h); + assert.equal(r.fetched.length, 1); + assert.equal((await spawns()) - before, 1); +}); diff --git a/common/controller/fetchWindows.ts b/common/controller/fetchWindows.ts @@ -0,0 +1,469 @@ +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import type { ChannelConfig } from "../lib/channelConfig"; +import type { SiteSettings } from "../lib/settings"; +import { getSettings } from "../lib/settings"; +import type { DownloadFailureClass } from "../lib/availability"; +import { resolveCookiePolicy } from "../lib/cookiePolicy"; +import { assertChannelTextReadable } from "../lib/channelMedia"; +import { detectPlatform } from "../lib/platform"; +import { findContainingClipWindow } from "../lib/clipWindow-server"; +import type { JobProgress } from "../jobs/registry"; +import { downloadGapMs } from "../jobs/platformBackoff"; +import { + heldPlatformRefusal, + platformCooldownRemainingMs, + recordDownloadBackoff, + recordPlatformClean, +} from "../jobs/downloadBackoff"; +import { + FetchWindowError, + fetchWindowManaged, + type FetchWindowOpts, + type FetchWindowResult, +} from "../ytdlp/fetchWindowManaged"; +import { + channelPaceSeconds, + channelPlatform, +} from "../ytdlp/channelArgs"; +import { + platformMinGapSeconds, + staticSleepRequestsSeconds, +} from "../ytdlp/platformArgs.mjs"; +import { readChannelConfig } from "./channels"; +import { findVideoSourceUrl } from "./undownloadedVideos"; + +// FETCH A LIST OF CLIP WINDOWS, ONE PLATFORM'S, AS ONE PACED JOB. +// +// The batch form of the single `fetch-window` job. Evidence clips for a report +// or a video used to be fetched one request at a time by hand-written shell +// loops that owned the pacing and the "stop after two failures" rule; this is +// that loop, inside the editor, where the platform's cooldown and pace live. +// +// EACH ITEM GOES THROUGH fetchWindowManaged, unchanged: the channel's cookie +// policy and extra args, the auth retry, the HLS retry, the provenance sidecar. +// What this adds is the space between items: +// +// - A CACHED WINDOW costs nothing: no request, no pause. The cache is asked +// here first (the same containing-window rule fetchWindowManaged applies), +// so a re-run of a half-done list walks straight to what is missing. +// - BETWEEN TWO NETWORK FETCHES the batch downloads' own gap +// (downloadGapMs), with a floor of CLIP_WINDOW_MIN_GAP_SECONDS — jittered +// up to half again, so a platform is never asked on a fixed beat. +// - BEFORE EACH NETWORK FETCH the platform's cooldown and hold. Either one +// ends the run; the rest is left for the next. +// +// THE FAILURE RULE. The class comes from FetchWindowError, classified once in +// fetchWindowManaged: +// rate_limit (429, bot check, soft block) — the backoff is already +// recorded; stop now. +// network (a 403 above all: Rumble's Cloudflare, googlevideo refusing) +// — one is an item failure, because a removed Rumble page +// answers 403 too. A SECOND IN A ROW records a network backoff +// and stops. Any other outcome between them breaks the streak. +// anything else (unavailable, private, a cut that failed) — the item fails +// and the run carries on. +// A success records the platform clean (recordPlatformClean), which is how an +// earlier backoff eases. +// +// RESUMABLE BY RE-RUNNING, as persistVideos: nothing is remembered between +// runs, and a re-run's cache check skips every window the last one fetched. +// +// A JOB OF THIS KIND SPANS CHANNELS, so runManagedFunction's per-channel text +// guard has no channel to ask. It is asked here instead, once per channel: an +// item on a channel whose text is not readable fails as `unreachable` (its +// clips/ directory lives in that text), and the rest carry on. + +// The least a batch waits between two clip-window fetches on one platform. +// A window is a short request, which is exactly why a loop of them looks like +// a scraper: 20 s apart was enough for YouTube to answer 403 (2026-10-07). +export const CLIP_WINDOW_MIN_GAP_SECONDS = 30; + +export type FetchWindowsItem = { + slug: string; + id: string; + from: number; + to: number; + // `<reportId>#<citationId>`, or a manifest's clip id. + clipId?: string; + reason?: string; + pad?: number; + // The page URL when the caller has it; else resolved as the single fetch does. + webpageUrl?: string; +}; + +export type FetchWindowsFailureClass = + | DownloadFailureClass + // No channel config for the slug. + | "unknown-channel" + // The channel's text (and so its clips/) is not readable. + | "unreachable" + // No page URL to fetch from. + | "no-url"; + +export type FetchWindowsFailure = { + item: FetchWindowsItem; + class: FetchWindowsFailureClass; + message: string; +}; + +export type FetchWindowsStop = + | "rate-limit" + | "network" + | "cooldown" + | "held" + | "drain" + | "cancel"; + +export type FetchWindowsResult = { + fetched: FetchWindowsItem[]; + cached: FetchWindowsItem[]; + failed: FetchWindowsFailure[]; + // Due a fetch but not attempted because the run stopped first. + notAttempted: FetchWindowsItem[]; + stopped?: FetchWindowsStop; +}; + +export type FetchWindowsProvenance = { + requestedBy: string; + manifest?: string; + requestedAt?: string; +}; + +// Everything that touches the network, the disk's shared state or the clock. +// Injectable so the tests decide what a platform answers and how long a pause +// is. +export type FetchWindowsDeps = { + fetchWindow: (opts: FetchWindowOpts) => Promise<FetchWindowResult>; + sleep: (ms: number, signal?: AbortSignal) => Promise<void>; + getSettings: () => SiteSettings; + readChannelConfig: (paths: Paths, slug: string) => Promise<ChannelConfig | null>; + findVideoSourceUrl: ( + paths: Paths, + slug: string, + id: string, + config: ChannelConfig, + ) => Promise<string | null>; + assertTextReadable: (paths: Paths, slug: string) => Promise<unknown>; + cooldownRemainingMs: (platform: string, paths: Paths) => Promise<number>; + heldRefusal: (platform: string, paths: Paths) => Promise<string | null>; + recordBackoff: ( + platform: string, + paths: Paths, + failureClass: "rate_limit" | "network", + ) => Promise<void>; + recordClean: (platform: string, paths: Paths) => Promise<string | null>; +}; + +function abortableSleep(ms: number, signal?: AbortSignal): Promise<void> { + if (ms <= 0 || signal?.aborted) return Promise.resolve(); + return new Promise((resolve) => { + const onAbort = () => { + clearTimeout(t); + resolve(); + }; + const t = setTimeout(() => { + signal?.removeEventListener("abort", onAbort); + resolve(); + }, ms); + signal?.addEventListener("abort", onAbort, { once: true }); + }); +} + +const DEFAULT_DEPS: FetchWindowsDeps = { + fetchWindow: fetchWindowManaged, + sleep: abortableSleep, + getSettings, + readChannelConfig, + findVideoSourceUrl: (paths, slug, id, config) => + findVideoSourceUrl(paths, slug, id, config), + assertTextReadable: (paths, slug) => assertChannelTextReadable(paths, slug), + cooldownRemainingMs: (platform, paths) => + platformCooldownRemainingMs(platform, paths), + heldRefusal: (platform, paths) => + heldPlatformRefusal(platform, "This batch", paths), + recordBackoff: (platform, paths, failureClass) => + recordDownloadBackoff(platform, paths, failureClass), + recordClean: (platform, paths) => recordPlatformClean(platform, paths), +}; + +export function fetchWindowsItemLabel(item: FetchWindowsItem): string { + const span = `${item.from.toFixed(2)}–${item.to.toFixed(2)}`; + return `${item.slug}/${item.id} ${span}${item.clipId ? ` (${item.clipId})` : ""}`; +} + +// Duplicates collapse: one window asked for twice is fetched once, and the +// first ask's clipId and reason are the ones recorded. +export function dedupeFetchWindowsItems( + items: readonly FetchWindowsItem[], +): FetchWindowsItem[] { + const seen = new Set<string>(); + const out: FetchWindowsItem[] = []; + for (const item of items) { + const key = `${item.slug}\0${item.id}\0${item.from.toFixed(2)}\0${item.to.toFixed(2)}`; + if (seen.has(key)) continue; + seen.add(key); + out.push(item); + } + return out; +} + +export async function fetchWindows({ + paths, + items, + provenance, + maxHeight, + gapMs, + onLog, + signal, + drainSignal, + setProgress, + deps: depsOverride, +}: { + paths: Paths; + items: FetchWindowsItem[]; + provenance: FetchWindowsProvenance; + maxHeight?: number; + // The pause between two network fetches. Absent = the platform's batch gap, + // floored at CLIP_WINDOW_MIN_GAP_SECONDS. + gapMs?: number; + onLog?: (line: string) => void; + signal?: AbortSignal; + drainSignal?: AbortSignal; + setProgress?: (snap: JobProgress) => void; + deps?: Partial<FetchWindowsDeps>; +}): Promise<FetchWindowsResult> { + const deps: FetchWindowsDeps = { ...DEFAULT_DEPS, ...depsOverride }; + const log = (line: string) => onLog?.(line.endsWith("\n") ? line : `${line}\n`); + const fetchLog = (line: string) => onLog?.(line); + const fetchSignal = signal ?? new AbortController().signal; + const settings = deps.getSettings(); + const requestedAt = provenance.requestedAt ?? new Date().toISOString(); + + const work = dedupeFetchWindowsItems(items); + const result: FetchWindowsResult = { + fetched: [], + cached: [], + failed: [], + notAttempted: [], + }; + log( + `Fetch windows: ${work.length} window(s) for ${provenance.requestedBy}` + + `${provenance.manifest ? ` · ${provenance.manifest}` : ""}.`, + ); + + let done = 0; + const progress = () => + setProgress?.({ metric: "clips", initial: 0, target: work.length, current: done }); + progress(); + + const configs = new Map<string, ChannelConfig | null>(); + const configOf = async (slug: string) => { + if (!configs.has(slug)) configs.set(slug, await deps.readChannelConfig(paths, slug)); + return configs.get(slug) ?? null; + }; + // A channel's text readability, asked once per channel per run. + const readable = new Map<string, string | null>(); + const textProblem = async (slug: string): Promise<string | null> => { + if (!readable.has(slug)) { + readable.set( + slug, + await deps.assertTextReadable(paths, slug).then( + () => null, + (e: unknown) => (e as Error).message, + ), + ); + } + return readable.get(slug) ?? null; + }; + + const fail = ( + item: FetchWindowsItem, + cls: FetchWindowsFailureClass, + message: string, + ) => { + result.failed.push({ item, class: cls, message }); + log(` ✗ ${fetchWindowsItemLabel(item)}: ${cls} — ${message}`); + }; + const stop = (why: FetchWindowsStop, from: number) => { + result.stopped = why; + result.notAttempted.push(...work.slice(from)); + }; + const interrupted = (): FetchWindowsStop | null => + signal?.aborted ? "cancel" : drainSignal?.aborted ? "drain" : null; + + // Network attempts so far (the gap is owed before every one after the + // first), and the run of consecutive network-class failures. + let networkAttempts = 0; + let networkStreak = 0; + + for (let i = 0; i < work.length; i++) { + const item = work[i]; + const why = interrupted(); + if (why) { + stop(why, i); + break; + } + const label = fetchWindowsItemLabel(item); + + const config = await configOf(item.slug); + if (!config) { + fail(item, "unknown-channel", `no channel "${item.slug}"`); + done += 1; + progress(); + continue; + } + const unreadable = await textProblem(item.slug); + if (unreadable) { + fail(item, "unreachable", unreadable); + done += 1; + progress(); + continue; + } + const videoDir = path.join(paths.channelsDir, item.slug, "data", item.id); + + // THE CACHE FIRST, here rather than only inside fetchWindowManaged, so a + // cached window owes no pause and is not counted as a network attempt. + const hit = await findContainingClipWindow(videoDir, item.from, item.to); + if (hit) { + result.cached.push(item); + log(` = ${label}: cached in clips/${hit.file}`); + done += 1; + progress(); + continue; + } + + const url = + item.webpageUrl?.trim() || + (await deps.findVideoSourceUrl(paths, item.slug, item.id, config)); + if (!url) { + fail( + item, + "no-url", + "no metadata.info.json and the playlist does not contain a matching entry", + ); + done += 1; + progress(); + continue; + } + // The cooldown's key, as the single fetch and every download path use it. + const platform = detectPlatform(url) ?? "unknown"; + + if (networkAttempts > 0) { + const gap = + gapMs ?? + downloadGapMs( + config.sleepBetweenDownloadsSeconds ?? settings.sleepBetweenDownloadsSeconds, + channelPaceSeconds(config), + staticSleepRequestsSeconds(channelPlatform(config)), + { + minSeconds: Math.max( + CLIP_WINDOW_MIN_GAP_SECONDS, + platformMinGapSeconds(channelPlatform(config)), + ), + }, + ); + if (gap > 0) { + log(`Sleeping ${Math.round(gap / 1000)}s before the next fetch...`); + await deps.sleep(gap, signal); + const after = interrupted(); + if (after) { + stop(after, i); + break; + } + } + } + + // ASKED BEFORE EVERY NETWORK FETCH: another job (or the lane) may have + // backed this platform off while this one slept. + const held = await deps.heldRefusal(platform, paths); + if (held) { + log(` ${label}: ${held} — stopping; the rest are left for a later run.`); + stop("held", i); + break; + } + const cooldownMs = await deps.cooldownRemainingMs(platform, paths); + if (cooldownMs > 0) { + log( + ` ${label}: ${platform} is in a rate-limit cooldown ` + + `(${Math.ceil(cooldownMs / 1000)}s remaining) — stopping; the rest are ` + + `left for a later run.`, + ); + stop("cooldown", i); + break; + } + + log(` ${label}: fetching…`); + networkAttempts += 1; + try { + const r = await deps.fetchWindow({ + channelSlug: item.slug, + channelConfig: config, + paths, + videoDir, + videoId: item.id, + videoUrl: url, + cwd: path.join(paths.channelsDir, item.slug), + from: item.from, + to: item.to, + provenance: { + requestedBy: provenance.requestedBy, + manifest: provenance.manifest, + clipId: item.clipId, + reason: item.reason, + pad: item.pad, + requestedAt, + }, + cookiePolicy: resolveCookiePolicy(settings, config), + maxHeight, + onLog: fetchLog, + signal: fetchSignal, + onPlatformBackoff: () => deps.recordBackoff(platform, paths, "rate_limit"), + }); + networkStreak = 0; + (r.cached ? result.cached : result.fetched).push(item); + const clean = await deps.recordClean(platform, paths).catch(() => null); + if (clean) log(clean); + } catch (e) { + const message = (e as Error).message; + const cls: FetchWindowsFailureClass = + e instanceof FetchWindowError ? e.failureClass : "unknown"; + fail(item, cls, message); + done += 1; + progress(); + if (cls === "rate_limit") { + // fetchWindowManaged recorded the backoff before it threw. + log(` ${platform} rate-limited this run — stopping; the rest are left for a later run.`); + stop("rate-limit", i + 1); + break; + } + if (cls === "network") { + networkStreak += 1; + if (networkStreak >= 2) { + await deps.recordBackoff(platform, paths, "network").catch(() => {}); + log( + ` ${platform} refused ${networkStreak} fetches in a row — backing it ` + + `off and stopping; the rest are left for a later run.`, + ); + stop("network", i + 1); + break; + } + } else { + networkStreak = 0; + } + continue; + } + done += 1; + progress(); + } + + log( + `Fetch windows: ${result.fetched.length} fetched, ${result.cached.length} cached, ` + + `${result.failed.length} failed` + + (result.notAttempted.length + ? `, ${result.notAttempted.length} not attempted (stopped: ${result.stopped})` + : "") + + ".", + ); + return result; +} diff --git a/common/jobs/jobDetail.test.ts b/common/jobs/jobDetail.test.ts @@ -0,0 +1,29 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { jobSpecDetail } from "./jobDetail"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/jobDetail.test.ts + +const spec = (kind: string, params: Record<string, unknown>) => ({ kind, slug: "c", params }); + +test("a single window names who asked and for which clip", () => { + assert.equal( + jobSpecDetail("fetch-window", spec("fetch-window", { requestedBy: "umtool", manifest: "m", clipId: "c01" })), + "umtool · m/c01", + ); +}); + +test("a batch names who asked, for what, and how many windows", () => { + assert.equal( + jobSpecDetail("fetch-windows", spec("fetch-windows", { requestedBy: "reports", manifest: "demo-site", items: [{}, {}] })), + "reports · demo-site · 2 windows", + ); + assert.equal( + jobSpecDetail("fetch-windows", spec("fetch-windows", { requestedBy: "umtool", items: [{}] })), + "umtool · 1 window", + ); +}); + +test("other kinds add nothing", () => { + assert.equal(jobSpecDetail("sync", spec("sync", { requestedBy: "x" })), undefined); +}); diff --git a/common/jobs/jobDetail.ts b/common/jobs/jobDetail.ts @@ -13,12 +13,19 @@ export function jobSpecDetail( kind: string | undefined, spec: JobSpec | null | undefined, ): string | undefined { - if (kind !== "fetch-window") return undefined; + if (kind !== "fetch-window" && kind !== "fetch-windows") return undefined; const p = spec?.params ?? {}; const s = (v: unknown): string | null => typeof v === "string" && v.trim() !== "" ? v.trim() : null; const by = s(p.requestedBy); if (!by) return undefined; + if (kind === "fetch-windows") { + // A batch: who asked, for what, and how many windows. + const n = Array.isArray(p.items) ? p.items.length : 0; + return [by, s(p.manifest), `${n} window${n === 1 ? "" : "s"}`] + .filter(Boolean) + .join(" · "); + } const what = [s(p.manifest), s(p.clipId)].filter(Boolean).join("/"); return what ? `${by} · ${what}` : by; } diff --git a/common/jobs/jobKinds.test.ts b/common/jobs/jobKinds.test.ts @@ -88,6 +88,8 @@ const ADDED_KINDS: Record<string, { label: string; drainable: boolean }> = { "reports-prepare": { label: "Prepare report media", drainable: false }, // A report site's exports: one pass, cancelled rather than drained. "reports-export": { label: "Export reports", drainable: false }, + // A batch of clip windows: stops between windows, a re-run resumes. + "fetch-windows": { label: "Fetch windows", drainable: true }, // One video's metadata re-read: one spawn, cancelled rather than drained. "refresh-metadata": { label: "Refresh metadata", drainable: false }, }; @@ -146,6 +148,7 @@ const TEXT_KINDS = [ "normalize-transcripts", "purge-superseded-auto-subs", "fetch-window", + "fetch-windows", "evict-clips", "metadata-scan", "refresh-metadata", diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -307,6 +307,23 @@ const JOB_KINDS: Record<string, JobKindMeta> = { needsMedia: false, needsText: true, }, + // A LIST OF CLIP WINDOWS, one platform's, as one paced job + // (controller/fetchWindows.ts): a site's missing evidence, a umtool + // manifest's timeline, an explicit list. One job per platform queue, so + // YouTube and Rumble run side by side, each paced on its own. Drainable: it + // stops between windows and a re-run fetches only what is still missing, + // which is also why it is replayable. It spans channels and so starts with no + // channelSlug — the per-channel text check `needsText` stands for is made by + // the controller, once per channel, before that channel's first window. + "fetch-windows": { + kind: "fetch-windows", + label: "Fetch windows", + drainable: true, + replayable: true, + queueKeyStrategy: "platform", + needsMedia: false, + needsText: true, + }, "redownload-archive": { kind: "redownload-archive", label: "Archive source video", diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts @@ -22,7 +22,11 @@ export type JobProgressMetric = // The metadata scan. Its progress CANNOT be re-counted from disk — the scan // deliberately writes no video directory — so its runner always sets // `current` itself. See ytdlp/metadataScan.ts. - | "scans"; + | "scans" + // Clip windows of a fetch-windows batch (controller/fetchWindows.ts). Spans + // channels, so there is no channel count to re-count; the runner always sets + // `current`. + | "clips"; export type JobProgress = { metric: JobProgressMetric; diff --git a/common/publish/missingEvidenceWindows.test.ts b/common/publish/missingEvidenceWindows.test.ts @@ -0,0 +1,131 @@ +// missingEvidenceWindows over a temp site: which cited spans are on disk, which +// are windows to fetch, and which no window fetch can fill — through the real +// report checker and the real tier lookup (ffprobe over a few frames of lavfi). +// +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test publish/missingEvidenceWindows.test.ts + +import { after, test } from "node:test"; +import assert from "node:assert/strict"; +import { execFileSync } from "node:child_process"; +import { mkdirSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import path from "node:path"; + +const ROOT = mkdtempSync(path.join(tmpdir(), "missing-evidence-")); +Object.assign(process.env, { + TRANSCRIPTS_DIR: path.join(ROOT, "transcripts"), + SAVED_VIDEOS_DIR: path.join(ROOT, "saved-videos"), + SITES_DIR: path.join(ROOT, "transcripts", "sites"), + SETTINGS_FILE: path.join(ROOT, "settings.json"), + EXPORT_PUBLIC_DIR: path.join(ROOT, "public"), + EXPORT_INDEX_DIR: path.join(ROOT, ".export-index"), + EXPORT_BUILDS_DIR: path.join(ROOT, ".export-builds"), + ARCHILYZER_CONFIG_DIR: path.join(ROOT, "config"), +}); +after(() => rmSync(ROOT, { recursive: true, force: true })); + +const { getPaths } = await import("../lib/paths"); +const { missingEvidenceWindows } = await import("./reportMedia"); + +const paths = getPaths(); +const SITE = "demo-site"; +const CH = "demo-channel"; +const OFF = "off-site"; + +const writeJson = (file: string, value: unknown) => { + mkdirSync(path.dirname(file), { recursive: true }); + writeFileSync(file, JSON.stringify(value, null, 2)); +}; +const videoDir = (slug: string, id: string) => path.join(paths.channelsDir, slug, "data", id); + +for (const [slug, extra] of [ + [CH, {}], + [OFF, {}], +] as const) { + writeJson(path.join(paths.channelsDir, slug, "config.json"), { + handling: "transcribe", + name: slug, + url: "https://www.youtube.com/@demo", + ...extra, + }); +} +// A fetched window holding 0–6 s of abc123. +mkdirSync(path.join(videoDir(CH, "abc123"), "clips"), { recursive: true }); +execFileSync("ffmpeg", [ + "-nostdin", "-v", "error", "-y", + "-f", "lavfi", "-i", "testsrc=size=160x90:rate=10:duration=6", + "-f", "lavfi", "-i", "sine=frequency=440:duration=6", + "-c:v", "libx264", "-preset", "ultrafast", "-pix_fmt", "yuv420p", "-c:a", "aac", "-shortest", + path.join(videoDir(CH, "abc123"), "clips", "0.00-6.00.mp4"), +]); +// A video the source says is gone. +writeJson(path.join(videoDir(CH, "gone1"), "availability.json"), { + checkedAt: "2026-10-01T00:00:00.000Z", + availability: "deleted", +}); + +writeJson(path.join(paths.sitesDir, SITE, "site.json"), { + title: "Demo", + channels: [{ slug: CH }], + search: false, + reports: ["r1", "r2"], +}); +const report = (id: string, citations: Record<string, unknown>) => ({ + format: "archilyzer-report", + version: 1, + id, + kind: "sweep", + title: `Report ${id}`, + citations, + sections: [ + { id: "s1", title: "One", body: Object.keys(citations).map((c) => `[${c}](cite:${c})`).join(" ") }, + ], +}); +writeJson( + path.join(paths.sitesDir, SITE, "reports", "r1", "report.json"), + report("r1", { + c01: { kind: "video", channel: CH, id: "abc123", start: 1, end: 2, quote: "on disk" }, + c02: { kind: "video", channel: CH, id: "missing1", start: 10.004, end: 12.001, pad: { before: 1, after: 1 }, quote: "the one to fetch" }, + c03: { kind: "video", channel: CH, id: "gone1", start: 5, end: 6, quote: "removed" }, + c04: { kind: "video", channel: OFF, id: "elsewhere", start: 5, end: 6, quote: "not this site's" }, + }), +); +writeJson( + path.join(paths.sitesDir, SITE, "reports", "r2", "report.json"), + report("r2", { + d01: { kind: "video", channel: CH, id: "missing1", start: 10.004, end: 12.001, pad: { after: 3 }, quote: "again" }, + }), +); + +test("on disk, to fetch, and what no window can fill", async () => { + const r = await missingEvidenceWindows(paths, SITE); + assert.equal(r.siteId, SITE); + assert.equal(r.onDisk, 1); + assert.deepEqual(r.problems, []); + + // One window for the moment two reports cite, with the wider pad on each + // side (from the moment's rounded start and end), and the two decimals a + // window is named by. + assert.deepEqual(r.missing, [ + { + slug: CH, + id: "missing1", + from: 9, + to: 15, + clipId: "r1#c02", + citations: ["r1#c02", "r2#d01"], + reason: "the one to fetch", + }, + ]); + + const why = Object.fromEntries(r.unfetchable.map((u) => [u.id, u.why])); + assert.deepEqual(why, { + gone1: "gone", + elsewhere: "not-in-site", + }); + assert.match(r.unfetchable.find((u) => u.id === "gone1")!.message, /deleted/); +}); + +test("a site that does not exist is an error", async () => { + await assert.rejects(missingEvidenceWindows(paths, "no-such-site"), /no site "no-such-site"/); +}); diff --git a/common/publish/reportMedia.ts b/common/publish/reportMedia.ts @@ -34,8 +34,9 @@ // // A RUN WITH PROBLEMS STILL WRITES THE MANIFEST (the editor lists the problems // from it) and FAILS: a citation without media, a clip over the size limit, an -// invalid or missing report. Nothing here fetches — the editor's fetch-window -// and persist (for spans) and Capture posts (for posts) fill what is missing. +// invalid or missing report. Nothing here fetches — the editor's fetch-windows +// (for spans: `missingEvidenceWindows` below says which), persist, and Capture +// posts (for posts) fill what is missing. // // AN UNMOUNTED DRIVE IS NOT A MISSING CLIP. A span whose media is not found on // a channel whose media tier is not reachable (lib/channelMedia.ts) is @@ -55,10 +56,14 @@ import { momentKey, momentOf, momentProblem, type Moment } from "../lib/citation import type { CitationPad } from "../lib/citations/schema"; import { parseReport } from "../lib/report/validate"; import type { Report } from "../lib/report/schema"; +import { isPermanentlyGone } from "../lib/availability"; +import { loadAvailability } from "../lib/availability-server"; +import { MAX_CLIP_WINDOW_SECONDS } from "../lib/clipWindow"; import { citedEvidenceSpan, isAudioOnlyPlatform, prepareEvidenceClip, + resolveEvidenceSource, widerPad, type EvidenceKind, type EvidenceMedia, @@ -369,3 +374,144 @@ export async function readReportMediaIndex(paths: Paths, siteId: string): Promis return v as ReportMediaIndex; } + +// --------------------------------------------------------------------------- +// WHICH WINDOWS A SITE STILL NEEDS — the list the editor's `fetch-windows` job +// takes, so "fetch what the reports cite" is one request rather than a loop. +// +// The same enumeration prepare runs (loadSiteReports, citedMoments, +// citedEvidenceSpan), and the same question prepare asks first +// (resolveEvidenceSource) — WITHOUT cutting anything. A span prepare would +// call `missing-media` comes back in `missing` as a window to fetch; a span +// that cannot be fetched as a window comes back in `unfetchable` with the +// reason, so a dry run tells the operator what no fetch will fix: +// +// not-in-site the channel is not one of the site's (prepare's not-in-site) +// unreachable the channel's media is on a drive that is not answering — +// not missing, and a fetch would write beside the wrong thing +// gone availability.json says deleted, private or members-only +// audio-only a feed record: its media is the episode's audio, which a +// download fetches; a window selector has no picture to pick +// too-long longer than one window may be — persist the video instead +// +// Spans already on disk count in `onDisk`. Posts are not windows and are left +// to Capture posts. +// --------------------------------------------------------------------------- + +export type EvidenceWindow = { + slug: string; + id: string; + from: number; + to: number; + // `<reportId>#<citationId>` of the first citation of the moment. + clipId: string; + // Every citation of the moment. + citations: string[]; + // The first citation's quote, for the window's provenance. + reason?: string; +}; + +export type UnfetchableReason = "not-in-site" | "unreachable" | "gone" | "audio-only" | "too-long"; + +export type UnfetchableEvidenceWindow = EvidenceWindow & { + why: UnfetchableReason; + message: string; +}; + +export type MissingEvidenceWindows = { + siteId: string; + missing: EvidenceWindow[]; + unfetchable: UnfetchableEvidenceWindow[]; + // Cited spans whose media is already on disk. + onDisk: number; + // The reports that could not be read (prepare's missing/invalid-report). + problems: ReportMediaProblem[]; +}; + +// A window fetched for a span must hold it: its name is two decimals, so the +// start rounds down and the end up. +const floor2 = (n: number) => Math.floor(n * 100 + 1e-6) / 100; +const ceil2 = (n: number) => Math.ceil(n * 100 - 1e-6) / 100; + +export async function missingEvidenceWindows(paths: Paths, siteId: string): Promise<MissingEvidenceWindows> { + if (!listSiteIds(paths).includes(siteId)) { + throw new Error(`no site "${siteId}" (no sites/${siteId}/site.json)`); + } + const site = getSite(siteId, paths); + const { reports, problems } = await loadSiteReports(paths, site); + const quotes = new Map<string, string>(); + for (const report of reports) { + for (const [cid, c] of Object.entries(report.citations ?? {})) { + if (typeof c.quote === "string" && c.quote.trim()) quotes.set(`${report.id}#${cid}`, c.quote.trim()); + } + } + const pool = siteChannelSlugs(site); + const configs = new Map<string, ChannelConfig | null>(); + const configOf = async (slug: string) => { + if (!configs.has(slug)) configs.set(slug, await readChannelConfig(paths, slug)); + return configs.get(slug) ?? null; + }; + const out: MissingEvidenceWindows = { siteId, missing: [], unfetchable: [], onDisk: 0, problems }; + + for (const m of citedMoments(reports)) { + if (m.kind === "post" || m.moment.kind !== "span") continue; + const slug = m.moment.channel; + const id = m.moment.id; + const span = await citedEvidenceSpan(paths.channelsDir, slug, id, { start: m.moment.start, end: m.moment.end, pad: m.pad }); + const clipId = m.citedBy[0]; + const reason = quotes.get(clipId); + const window: EvidenceWindow = { + slug, + id, + from: floor2(span.from), + to: ceil2(span.to), + clipId, + citations: m.citedBy, + ...(reason ? { reason } : {}), + }; + const unfetchable = (why: UnfetchableReason, message: string) => + out.unfetchable.push({ ...window, why, message }); + + if (!pool.has(slug)) { + unfetchable("not-in-site", `channel "${slug}" is not one of this site's channels`); + continue; + } + const config = await configOf(slug); + const audioOnly = isAudioOnlyPlatform(config?.platform); + const source = await resolveEvidenceSource({ + channelsDir: paths.channelsDir, + slug, + id, + span, + audio: m.kind === "audio" || audioOnly, + ffprobeBin: paths.ffprobeBin, + }); + if (source) { + out.onDisk += 1; + continue; + } + const where = await inspectChannelMedia(paths, slug, config, { fresh: true }).catch(() => null); + if (where && where.status !== "ok" && where.status !== "in-place") { + unfetchable("unreachable", `the channel's media is ${where.status}${where.detail ? ` (${where.detail})` : ""}`); + continue; + } + const availability = (await loadAvailability(path.join(paths.channelsDir, slug, "data", id)))?.availability; + if (isPermanentlyGone(availability)) { + unfetchable("gone", `the source says the video is ${availability}`); + continue; + } + if (audioOnly) { + unfetchable("audio-only", "a feed record's media is its audio download, not a window"); + continue; + } + if (window.to - window.from > MAX_CLIP_WINDOW_SECONDS) { + unfetchable( + "too-long", + `${Math.round(window.to - window.from)}s is longer than a window may be (${MAX_CLIP_WINDOW_SECONDS}s) — persist the video`, + ); + continue; + } + out.missing.push(window); + } + return out; +} diff --git a/common/views/jobRows.ts b/common/views/jobRows.ts @@ -167,7 +167,9 @@ function computeJobProgressView( now: number, ): JobRowView["progress"] { const snap = job.progress; - if (!snap || !stat) return undefined; + // A runner-reported `current` needs no channel stat: a batch spanning + // channels (fetch-windows) has none, and still knows its own count. + if (!snap || (!stat && snap.current === undefined)) return undefined; // `current` is RE-COUNTED from disk (readChannelStat), never reported by the // runner — which is why the digest metric needed its own on-disk counter // (digestCount) rather than a number the batch could have just told us. @@ -179,10 +181,10 @@ function computeJobProgressView( const current = snap.current ?? (snap.metric === "downloads" - ? stat.downloadCount + ? stat!.downloadCount : snap.metric === "digests" - ? (stat.digestCount ?? 0) - : stat.transcriptCount); + ? (stat!.digestCount ?? 0) + : stat!.transcriptCount); const range = Math.max(0, snap.target - snap.initial); const advance = Math.max(0, current - snap.initial); const pct = @@ -224,7 +226,7 @@ export function fromRecord(j: JobRecord, ctx: FromRecordContext): JobRowView { channelSlug: j.channelSlug, videoId: j.videoId, progress: - j.status === "running" && j.channelSlug + j.status === "running" && (j.channelSlug || j.progress?.current !== undefined) ? computeJobProgressView(j, ctx.stat, ctx.now) : undefined, tasks: j.tasks?.map((t) => ({ diff --git a/common/ytdlp/fetchWindowManaged.ts b/common/ytdlp/fetchWindowManaged.ts @@ -18,6 +18,8 @@ import { AUTH_RETRY_CLASSES, classifyDownloadFailure, parseUnavailableFromStderr, + type Availability, + type DownloadFailureClass, } from "../lib/availability"; import { DEFAULT_COOKIE_MODE, @@ -63,6 +65,21 @@ export const HLS_PICKY_RETRY_ARGS = ["--downloader-args", "ffmpeg_i:-extension_p // preset is built from it); re-exported so existing importers keep working. export { clipFormatSelector }; +// A fetch that yt-dlp refused, with the refusal already classified. The +// single-window job only shows the message; the batch (controller/fetchWindows) +// reads `failureClass` to decide between carrying on, backing the platform off +// and stopping — without parsing the message a second time. +export class FetchWindowError extends Error { + constructor( + message: string, + readonly failureClass: DownloadFailureClass, + readonly availability: Availability, + ) { + super(message); + this.name = "FetchWindowError"; + } +} + export type FetchWindowProvenance = { requestedBy: string; manifest?: string; @@ -236,6 +253,7 @@ export async function fetchWindowManaged( let cookiesUsed = alwaysCookies(policy); let args = argsWith(cookiesUsed); let outcome = await run(cookiesUsed); + let backedOff = false; if (outcome.exitCode !== 0) { const availability = parseUnavailableFromStderr(outcome.stderrTail); @@ -245,6 +263,7 @@ export async function fetchWindowManaged( // Sync both honour. Recorded before the throw so the NEXT request is // refused at the door rather than re-storming the source. await opts.onPlatformBackoff?.("rate_limit"); + backedOff = true; } else if (AUTH_RETRY_CLASSES.has(availability)) { const retryCookies = authRetryCookies(policy); if (retryCookies) { @@ -270,9 +289,18 @@ export async function fetchWindowManaged( if (outcome.exitCode !== 0) { await rm(part, { force: true }); const tail = outcome.stderrTail.trim().split("\n").slice(-4).join(" / "); - throw new Error( + // Classified from the LAST attempt, which is the one that decided. + const availability = parseUnavailableFromStderr(outcome.stderrTail); + const failureClass = classifyDownloadFailure(outcome.stderrTail, availability); + // A cookie retry that ran into a 429 is a 429 like any other. + if (failureClass === "rate_limit" && !backedOff) { + await opts.onPlatformBackoff?.("rate_limit"); + } + throw new FetchWindowError( `yt-dlp failed fetching ${from.toFixed(2)}–${to.toFixed(2)} of ` + `${opts.videoId} (exit ${outcome.exitCode ?? "null"}): ${tail}`, + failureClass, + availability, ); } diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **Clip windows can be fetched as a list, paced, in one job per platform.** `pnpm ops fetch-windows` (`POST /api/ops/fetch-windows`) takes `{"siteId"}` — every window a site's published reports cite and the disk does not hold — or `{"items": [{"slug", "id", "from", "to", "clipId"?, "reason"?}], "requestedBy", "manifest"?}`, with `"maxHeight"` and `"dryRun"`. A window already on disk is answered at once and joins no job; the rest are grouped by platform queue and fetched as one `fetch-windows` job per platform, so YouTube and Rumble run side by side, each window through the same managed fetch as a single one (cookie policy, auth and HLS retries, provenance). Between two fetches on a platform the job waits that platform's batch gap, at least 30 s and up to half again at random; before each it checks the platform's cooldown and hold and stops when either is set. A 429 backs the platform off and stops the job; one 403 is that window's failure, two in a row back the platform off and stop it. A platform cooling down or held is refused for its group before anything starts. The job is drainable, shows its progress as clip windows, and Retry or running the same body again fetches only what is still missing. A dry run lists the windows per platform, the ones on disk, and the spans no window can fill (a video the source says is deleted, private or members-only, a channel off the site, an unmounted drive). A site's **Reports** tab has **Fetch missing evidence** and **Preview missing evidence** beside **Prepare evidence media**, and umtool's `fetch-via-editor.mjs --all` sends a manifest's whole timeline as one request. A window fetch whose cookie retry runs into a 429 now records the cooldown too. Needs a restart of the editor. - **A social channel can be renamed and deleted from its page.** A posts channel's page (X, Bluesky, a forum thread) had no Danger zone, so it could not be renamed or deleted in the editor at all. It now has the one a video channel has, collapsed under the posts panel and opened by the same `?stage=danger` link: **Rename channel** and **Delete channel**, typing the slug to confirm, refused while the channel is busy, in the same sentences. Needs a restart of the editor. - **A channel can be created, renamed, deleted and put on sites over `pnpm ops`.** `create-channel` is the New channel form (`{"fields": {"name", "handling", "url", …}}`, the same field names as `channel-config`'s `patch`; `"slug"` and `"sites"` optional; the form's "Fetch playlist now", "Fetch posts now" and "Add to top of auto-queue" are off unless asked for, and a job they start is returned so `--wait` follows it). `rename-channel` (`{"slug", "newSlug"}`) and `delete-channel` (`{"slug", "confirm"}`, `confirm` repeating the slug) are the Danger zone's two forms. `channel-config` takes `"sites"` — the whole membership set, `[]` for on no site; an unknown site is refused rather than skipped — and `"excludeFromBuild"` / `"excludeFromCleanup"`, set to the value given rather than toggled, with or without a `patch`. `pnpm ops get channels` lists every channel with its kind, platform and the sites that carry it. Each runs the form's own action, so it refuses what the form refuses, in the same words. Needs a restart of the editor. - **One video's metadata can be read again from its source.** A livestream that has just ended offers one fragmented audio format and no captions; hours later the same URL has plain formats and auto-captions, and the video's `metadata.info.json` still said what it said the first time. "Refresh metadata" under the video page's header, and `pnpm ops refresh-metadata --json '{"slug":"…","id":"…"}'`, re-read that one video with no subtitles and no media, on the platform's queue, with the channel's cookie policy and pace; a held platform or one in a rate-limit cooldown is refused, and a rate limit records the cooldown. The rewrite is recorded in the metadata history as `refresh`, and the job's log ends with what the source now says: `live_status`, how many formats, each audio-only format and its protocol, whether any is non-fragmented, the English subtitle and caption tracks, and the keys that changed. Nothing is downloaded or deleted, and no download attempt is recorded. A video not yet fetched into the archive is refused — a refresh never creates its directory — as are an archive.org record, a Wayback copy and a record completed from a podcast feed, whose metadata is not yt-dlp's. Needs a restart of the editor. diff --git a/editor/app/api/media/fetch-window/[jobId]/route.ts b/editor/app/api/media/fetch-window/[jobId]/route.ts @@ -46,7 +46,9 @@ export async function GET( return NextResponse.json({ error: "no such job" }, { status: 404 }); } const kind = live?.kind ?? meta?.kind ?? entry?.kind; - if (kind !== "fetch-window" && kind !== "redownload-archive") { + // A fetch-windows batch answers with its status and, on failure, its log's + // tail: it names no one file, so `file` stays absent. + if (kind !== "fetch-window" && kind !== "redownload-archive" && kind !== "fetch-windows") { return NextResponse.json( { error: `job ${jobId} is a ${kind ?? "?"} job, not a media fetch` }, { status: 400 }, @@ -59,7 +61,7 @@ export async function GET( const out: Record<string, unknown> = { status, jobId }; if (exitCode !== undefined) out.exitCode = exitCode; - if (status === "done") { + if (status === "done" && kind !== "fetch-windows") { const slug = live?.channelSlug ?? meta?.channelSlug ?? spec?.slug; const videoId = live?.videoId ?? meta?.videoId; const p = spec?.params ?? {}; diff --git a/editor/app/api/ops/fetch-windows/route.test.ts b/editor/app/api/ops/fetch-windows/route.test.ts @@ -0,0 +1,128 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; + +// Run with: +// pnpm -C editor exec tsx --test "app/api/ops/fetch-windows/route.test.ts" +// +// The body's shape, the refusals that come before any job, and a dry run's +// plan — answered from a temp corpus. Every request here is a dry run or a +// refusal: no job is queued and nothing is fetched. + +const ROOT = await mkdtemp(path.join(os.tmpdir(), "fetch-windows-route-")); +// Set before the route (and getPaths, which caches) is first imported. +process.env.WORKER_TOKEN = "test-token"; +process.env.TRANSCRIPTS_DIR = ROOT; +process.env.SETTINGS_FILE = path.join(ROOT, "settings.json"); +const { POST } = await import("./route"); +test.after(() => rm(ROOT, { recursive: true, force: true })); + +// Two channels on two platforms, one video each, the YouTube one with a window +// already on disk. +const CH = path.join(ROOT, "channels"); +async function channel(slug: string, url: string, id: string, webpage: string) { + await mkdir(path.join(CH, slug, "data", id), { recursive: true }); + await writeFile( + path.join(CH, slug, "config.json"), + JSON.stringify({ handling: "youtube", name: slug, url }), + ); + await writeFile( + path.join(CH, slug, "data", id, "metadata.info.json"), + JSON.stringify({ id, webpage_url: webpage }), + ); +} +await channel("yt-chan", "https://www.youtube.com/@yt", "vid1", "https://www.youtube.com/watch?v=vid1"); +await channel("rb-chan", "https://rumble.com/c/rb", "rb1", "https://rumble.com/rb1-x.html"); +await mkdir(path.join(CH, "yt-chan", "data", "vid1", "clips"), { recursive: true }); +await writeFile(path.join(CH, "yt-chan", "data", "vid1", "clips", "0.00-60.00.mp4"), "mp4"); + +async function post( + body: Record<string, unknown>, +): Promise<{ status: number; json: Record<string, unknown> & { error?: string } }> { + const res = await POST( + new Request("http://localhost/api/ops/fetch-windows", { + method: "POST", + headers: { + authorization: "Bearer test-token", + "content-type": "application/json", + }, + body: JSON.stringify(body), + }), + ); + return { status: res.status, json: (await res.json()) as Record<string, unknown> }; +} + +const item = (o: Record<string, unknown> = {}) => ({ slug: "yt-chan", id: "vid1", from: 100, to: 110, ...o }); + +test("siteId or items, never both, never neither; unknown keys refused", async () => { + assert.match((await post({})).json.error!, /"siteId" .* or "items" .* is required/); + const both = await post({ siteId: "s", items: [item()], requestedBy: "t" }); + assert.equal(both.status, 400); + assert.match(both.json.error!, /either "siteId" or "items"/); + const unknown = await post({ items: [item()], requestedBy: "t", gapMs: 1 }); + assert.equal(unknown.status, 400); + assert.match(unknown.json.error!, /unknown key\(s\): gapMs/); +}); + +test("an item list is validated at the door", async () => { + const cases: [Record<string, unknown>, RegExp][] = [ + [{ items: [], requestedBy: "t" }, /non-empty array/], + [{ items: [item({ slug: "../x" })], requestedBy: "t" }, /not a valid channel slug/], + [{ items: [item({ id: ".." })], requestedBy: "t" }, /not a video id/], + [{ items: [item({ from: 20, to: 10 })], requestedBy: "t" }, /must be less than to/], + [{ items: [item({ from: 0, to: 2000 })], requestedBy: "t" }, /at most 900s/], + [{ items: [item({ from: -1 })], requestedBy: "t" }, /non-negative seconds/], + [{ items: [item({ extra: 1 })], requestedBy: "t" }, /items\[0\]" has unknown key\(s\): extra/], + [{ items: [item()] }, /"requestedBy" is required/], + [{ items: [item()], requestedBy: "t", maxHeight: 50 }, /"maxHeight" must be a whole number of pixels/], + ]; + for (const [body, re] of cases) { + const r = await post(body); + assert.equal(r.status, 400, JSON.stringify(body)); + assert.match(r.json.error!, re); + } +}); + +test("a site's windows name the site; a missing site is refused", async () => { + const named = await post({ siteId: "demo-site", requestedBy: "me" }); + assert.equal(named.status, 400); + assert.match(named.json.error!, /"requestedBy" goes with "items"/); + const missing = await post({ siteId: "demo-site" }); + assert.equal(missing.status, 400); + assert.match(missing.json.error!, /No site "demo-site"/); +}); + +test("a dry run answers the cache, groups by platform queue, and names what it cannot resolve", async () => { + const r = await post({ + requestedBy: "test", + dryRun: true, + items: [ + item({ from: 10, to: 20 }), // inside the cached 0–60 window + item({ from: 100, to: 110 }), + item({ from: 100, to: 110 }), // a duplicate collapses + item({ slug: "rb-chan", id: "rb1", from: 5, to: 15 }), + item({ slug: "no-such", id: "x", from: 5, to: 15 }), + ], + }); + assert.equal(r.status, 200); + const j = r.json as { + dryRun: boolean; + groups: { platform: string; queueKey: string; items: { id: string; webpageUrl?: string }[] }[]; + cached: { from: number }[]; + unresolved: { error: string }[]; + jobs?: unknown; + }; + assert.equal(j.dryRun, true); + assert.equal(j.jobs, undefined, "a dry run starts nothing"); + assert.deepEqual(j.cached.map((c) => c.from), [10]); + assert.deepEqual( + j.groups.map((g) => [g.platform, g.items.length]).sort(), + [["rumble", 1], ["youtube", 1]], + ); + const yt = j.groups.find((g) => g.platform === "youtube")!; + assert.equal(yt.items[0].webpageUrl, "https://www.youtube.com/watch?v=vid1"); + assert.equal(j.unresolved.length, 1); + assert.match(j.unresolved[0].error, /Channel "no-such" not found/); +}); diff --git a/editor/app/api/ops/fetch-windows/route.ts b/editor/app/api/ops/fetch-windows/route.ts @@ -0,0 +1,179 @@ +import { NextResponse } from "next/server"; +import { + MAX_CLIP_WINDOW_SECONDS, + MAX_FETCH_MAX_HEIGHT, + MIN_FETCH_MAX_HEIGHT, + isFetchMaxHeight, +} from "yt-dlp-transcript-common/lib/clipWindow"; +import { isValidSiteId } from "yt-dlp-transcript-common/lib/site"; +import type { FetchWindowsItem } from "yt-dlp-transcript-common/controller/fetchWindows"; +import { + fetchMissingEvidenceAction, + fetchWindowsAction, +} from "../../../channels/[slug]/videos/fetchWindowsAction"; +import { + OpsInputError, + ops, + opsFail, + optBool, + optString, + reqSlug, + reqString, + reqVideoId, + type OpsBody, +} from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { siteId, maxHeight?, dryRun? } +// | { items: [{ slug, id, from, to, clipId?, reason?, pad?, webpageUrl? }, …], +// requestedBy, manifest?, maxHeight?, dryRun? } +// dryRun -> { ok, dryRun: true, groups: [{ platform, queueKey, items }], +// cached, unresolved[, onDisk, unfetchable] } +// else -> { ok, dryRun: false, jobs: [{ platform, queueKey, jobId, items }], +// jobIds, jobId (one job only), refused, cached, unresolved +// [, onDisk, unfetchable] } +// +// Fetch clip windows through the managed path, ONE PACED JOB PER PLATFORM +// QUEUE (controller/fetchWindows.ts): the batch form of +// /api/media/fetch-window. `siteId` fetches every window the site's published +// reports cite and the disk does not hold (`unfetchable` lists the spans no +// window can fill, with why; `onDisk` counts the ones already there). `items` +// fetches an explicit list, and must say who asked (`requestedBy`). A window +// already on disk is answered in `cached` and joins no job; a platform cooling +// down or held is `refused` for its group, and the other groups still start. +// Re-running the same body is the resume: fetched windows are cached. +// +// A big list goes in a file: `pnpm ops fetch-windows --file list.json`. + +const ITEM_KEYS = ["slug", "id", "from", "to", "clipId", "reason", "pad", "webpageUrl"]; + +function optItemString(e: OpsBody, key: string, where: string): string | undefined { + const v = e[key]; + if (v === undefined) return undefined; + if (typeof v !== "string") throw new OpsInputError(`"${where}.${key}" must be a string`); + return v.trim() || undefined; +} + +function reqSeconds(e: OpsBody, key: string, where: string): number { + const v = e[key]; + if (typeof v !== "number" || !Number.isFinite(v) || v < 0) { + throw new OpsInputError(`"${where}.${key}" must be finite, non-negative seconds`); + } + return v; +} + +function reqWindowItems(body: OpsBody): FetchWindowsItem[] { + const v = body.items; + if (!Array.isArray(v) || v.length === 0) { + throw new OpsInputError( + `"items" must be a non-empty array of { "slug", "id", "from", "to" }`, + ); + } + return v.map((entry, i) => { + const where = `items[${i}]`; + if (typeof entry !== "object" || entry === null || Array.isArray(entry)) { + throw new OpsInputError(`"${where}" must be an object { "slug", "id", "from", "to" }`); + } + const e = entry as OpsBody; + const extra = Object.keys(e).filter((k) => !ITEM_KEYS.includes(k)); + if (extra.length) { + throw new OpsInputError( + `"${where}" has unknown key(s): ${extra.join(", ")} — accepted: ${ITEM_KEYS.join(", ")}`, + ); + } + const from = reqSeconds(e, "from", where); + const to = reqSeconds(e, "to", where); + if (from >= to) { + throw new OpsInputError(`"${where}": from (${from}) must be less than to (${to})`); + } + if (to - from > MAX_CLIP_WINDOW_SECONDS) { + throw new OpsInputError( + `"${where}": a window may be at most ${MAX_CLIP_WINDOW_SECONDS}s ` + + `(asked for ${Math.round(to - from)}s)`, + ); + } + const pad = e.pad; + if (pad !== undefined && (typeof pad !== "number" || !Number.isFinite(pad))) { + throw new OpsInputError(`"${where}.pad" must be a number`); + } + return { + slug: reqSlug(e, "slug"), + id: reqVideoId(e, "id"), + from, + to, + clipId: optItemString(e, "clipId", where), + reason: optItemString(e, "reason", where)?.slice(0, 400), + ...(typeof pad === "number" ? { pad } : {}), + webpageUrl: optItemString(e, "webpageUrl", where), + }; + }); +} + +function optMaxHeight(body: OpsBody): number | undefined { + const v = body.maxHeight; + if (v === undefined || v === null) return undefined; + if (!isFetchMaxHeight(v)) { + throw new OpsInputError( + `"maxHeight" must be a whole number of pixels from ` + + `${MIN_FETCH_MAX_HEIGHT} to ${MAX_FETCH_MAX_HEIGHT}`, + ); + } + return v; +} + +export async function POST(request: Request) { + return ops( + request, + ["siteId", "items", "requestedBy", "manifest", "maxHeight", "dryRun"], + async (body) => { + const maxHeight = optMaxHeight(body); + const dryRun = optBool(body, "dryRun"); + if (body.siteId !== undefined && body.items !== undefined) { + throw new OpsInputError('send either "siteId" or "items", not both'); + } + let result: Awaited<ReturnType<typeof fetchWindowsAction>>; + if (body.siteId !== undefined) { + // A site's windows say who asked themselves (the reports, by site). + for (const k of ["requestedBy", "manifest"]) { + if (body[k] !== undefined) { + throw new OpsInputError(`"${k}" goes with "items"; a site's windows name the site`); + } + } + const siteId = reqString(body, "siteId"); + if (!isValidSiteId(siteId)) { + throw new OpsInputError( + `"${siteId}" is not a valid site id (lowercase letters, digits and "-"; must start with a letter or digit)`, + ); + } + result = await fetchMissingEvidenceAction(siteId, { dryRun, maxHeight }); + } else if (body.items !== undefined) { + const items = reqWindowItems(body); + // WHO ASKED IS NOT OPTIONAL, as on the single route: a window nobody + // can explain in six months is one nobody can clean up. + const requestedBy = reqString(body, "requestedBy"); + result = await fetchWindowsAction({ + items, + requestedBy, + manifest: optString(body, "manifest")?.trim() || undefined, + maxHeight, + dryRun, + }); + } else { + throw new OpsInputError('"siteId" (a site\'s missing evidence) or "items" (a list of windows) is required'); + } + if (!result.ok) return opsFail(result.error); + if (result.dryRun) { + const { jobs: _jobs, refused: _refused, ...plan } = result; + return NextResponse.json(plan); + } + const { groups: _groups, ...run } = result; + const jobIds = run.jobs.map((j) => j.jobId); + return NextResponse.json({ + ...run, + jobIds, + ...(jobIds.length === 1 ? { jobId: jobIds[0] } : {}), + }); + }, + ); +} diff --git a/editor/app/channels/[slug]/videos/fetchWindowsAction.ts b/editor/app/channels/[slug]/videos/fetchWindowsAction.ts @@ -0,0 +1,406 @@ +"use server"; + +import path from "node:path"; +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"; +import { formatBytes } from "yt-dlp-transcript-common/lib/format"; +import { detectPlatform } from "yt-dlp-transcript-common/lib/platform"; +import { + downloadQueueKey, + resolveQueueKey, +} from "yt-dlp-transcript-common/lib/queueKeys"; +import { + MAX_CLIP_WINDOW_SECONDS, + isFetchMaxHeight, +} from "yt-dlp-transcript-common/lib/clipWindow"; +import { findContainingClipWindow } from "yt-dlp-transcript-common/lib/clipWindow-server"; +import type { ChannelConfig } from "yt-dlp-transcript-common/lib/channelConfig"; +import { + isValidChannelSlug, + readChannelConfig, +} from "yt-dlp-transcript-common/controller/channels"; +import { findVideoSourceUrl } from "yt-dlp-transcript-common/controller/undownloadedVideos"; +import { + fetchWindows, + fetchWindowsItemLabel, + dedupeFetchWindowsItems, + type FetchWindowsItem, +} from "yt-dlp-transcript-common/controller/fetchWindows"; +import { + missingEvidenceWindows, + type UnfetchableEvidenceWindow, +} from "yt-dlp-transcript-common/publish/reportMedia"; +import { isValidSiteId, listSiteIds } from "yt-dlp-transcript-common/lib/site"; +import { + heldPlatformRefusal, + platformCooldownRemainingMs, +} from "yt-dlp-transcript-common/jobs/downloadBackoff"; +import { + runManagedFunction, + type StreamActionResult, +} from "yt-dlp-transcript-common/jobs/streamCommand"; +import { safeRevalidate } from "../../../lib/safeRevalidate"; + +// FETCH A LIST OF CLIP WINDOWS through the managed path, as ONE JOB PER +// PLATFORM QUEUE (controller/fetchWindows.ts walks each one, paced). +// +// The batch form of fetchWindowAction (videos/[id]/videoActions.ts), with the +// same rules at the door: a window of at most MAX_CLIP_WINDOW_SECONDS, a height +// cap in range, a channel that exists, a URL that resolves. What it adds is the +// fan-out: the list is grouped by `downloadQueueKey` and each group starts its +// own job on that queue, so YouTube and Rumble run side by side, each behind +// its own platform's other downloads — persistVideosAction's shape. +// +// A PLATFORM COOLING DOWN OR HELD is refused at the door for ITS group (the +// sentence comes back in `refused`); the other groups still start. A window +// already on disk is answered here (`cached`) and joins no job. +// +// Like the single fetch, NOT GATED BY THE DOWNLOAD PAUSE: an operator (or a +// tool they are driving) asked for these seconds by hand. + +const ID_RE = /^[\w.-]+$/; +const isVideoId = (v: string): boolean => ID_RE.test(v) && v !== "." && v !== ".."; + +export type FetchWindowsRequest = { + items: FetchWindowsItem[]; + // Who asked: the tool or surface. Recorded beside every window. + requestedBy: string; + manifest?: string; + maxHeight?: number; + dryRun?: boolean; + // A queue override for every group (resolveQueueKey's rules). + queueKey?: string; +}; + +export type FetchWindowsUnresolved = { item: FetchWindowsItem; error: string }; + +export type FetchWindowsGroup = { + // The cooldown key (`detectPlatform(url)`), and the queue the job runs on. + platform: string; + queueKey: string; + items: FetchWindowsItem[]; +}; + +export type FetchWindowsActionResult = + | { ok: false; error: string } + | { + ok: true; + dryRun: boolean; + // A real run: one per group that started. + jobs: { platform: string; queueKey: string; jobId: string; items: number }[]; + // A dry run: what each job would be given. + groups: FetchWindowsGroup[]; + // Groups refused at the door, with the refusal's own sentence. + refused: { platform: string; error: string; items: number }[]; + cached: FetchWindowsItem[]; + unresolved: FetchWindowsUnresolved[]; + }; + +// The window rules the HTTP door enforces, asked again here: a replayed spec is +// a file on disk. +function itemProblem(item: FetchWindowsItem): string | null { + if (!isValidChannelSlug(item.slug)) return `"${item.slug}" is not a channel slug`; + if (!isVideoId(item.id)) return `"${item.id}" is not a video id`; + const { from, to } = item; + if ( + !Number.isFinite(from) || + !Number.isFinite(to) || + from < 0 || + from >= to || + to - from > MAX_CLIP_WINDOW_SECONDS + ) { + return ( + `${from}–${to} is not a fetchable window ` + + `(at most ${MAX_CLIP_WINDOW_SECONDS}s, from < to, from >= 0)` + ); + } + return null; +} + +// Validate, answer the cache, resolve every URL and group by queue. Nothing is +// started; nothing touches the network. +async function planFetchWindows( + items: FetchWindowsItem[], + queueOverride: string | undefined, +): Promise<{ + groups: FetchWindowsGroup[]; + cached: FetchWindowsItem[]; + unresolved: FetchWindowsUnresolved[]; +}> { + const paths = getPaths(); + const configs = new Map<string, ChannelConfig | null>(); + const groups = new Map<string, FetchWindowsGroup>(); + const cached: FetchWindowsItem[] = []; + const unresolved: FetchWindowsUnresolved[] = []; + for (const item of dedupeFetchWindowsItems(items)) { + const problem = itemProblem(item); + if (problem) { + unresolved.push({ item, error: problem }); + continue; + } + if (!configs.has(item.slug)) { + configs.set(item.slug, await readChannelConfig(paths, item.slug)); + } + const config = configs.get(item.slug); + if (!config) { + unresolved.push({ item, error: `Channel "${item.slug}" not found` }); + continue; + } + const videoDir = path.join(paths.channelsDir, item.slug, "data", item.id); + if (await findContainingClipWindow(videoDir, item.from, item.to)) { + cached.push(item); + continue; + } + const url = + item.webpageUrl?.trim() || + (await findVideoSourceUrl(paths, item.slug, item.id, config)); + if (!url) { + unresolved.push({ + item, + error: + "Could not determine the video URL: no metadata.info.json and the " + + "playlist does not contain a matching entry.", + }); + continue; + } + const queueKey = resolveQueueKey(downloadQueueKey(config), queueOverride); + const platform = detectPlatform(url) ?? "unknown"; + // One job per queue. Two platforms sharing a queue (an override) share a + // job too; the controller keys the cooldown per item, so each is honoured. + const group = groups.get(queueKey) ?? { platform, queueKey, items: [] }; + group.items.push({ ...item, webpageUrl: url }); + groups.set(queueKey, group); + } + return { groups: [...groups.values()], cached, unresolved }; +} + +// The door each group passes before its job is queued: the platform's hold, +// then its cooldown — fetchWindowAction's checks and sentences. +async function groupRefusal(platform: string): Promise<string | null> { + const paths = getPaths(); + const held = await heldPlatformRefusal(platform, "This batch", paths); + if (held) return held; + const cooldownMs = await platformCooldownRemainingMs(platform, paths); + if (cooldownMs > 0) { + return ( + `${platform} is in a rate-limit cooldown ` + + `(${Math.ceil(cooldownMs / 1000)}s remaining).` + ); + } + return null; +} + +// One platform's windows, as one job. Exported for Retry: the replay hands the +// spec's items back here, and the controller's cache check skips whatever an +// earlier run fetched. +export async function fetchWindowsJobAction(req: { + items: FetchWindowsItem[]; + requestedBy: string; + manifest?: string; + siteId?: string; + maxHeight?: number; + queueKey: string; +}): Promise<StreamActionResult> { + if (req.items.length === 0) return { ok: false, error: "No windows to fetch." }; + if (req.maxHeight !== undefined && !isFetchMaxHeight(req.maxHeight)) { + return { + ok: false, + error: `maxHeight ${req.maxHeight} is not a source height to cap a fetch at.`, + }; + } + for (const item of req.items) { + const problem = itemProblem(item); + if (problem) return { ok: false, error: `${fetchWindowsItemLabel(item)}: ${problem}` }; + } + const paths = getPaths(); + const requestedAt = new Date().toISOString(); + const slugs = [...new Set(req.items.map((i) => i.slug))]; + return runManagedFunction({ + kind: "fetch-windows", + queueKey: req.queueKey, + paths, + // A batch of ONE channel carries it, so /jobs and the media guard can name + // it; one spanning channels carries none, and the controller asks each + // channel's text guard itself. + ...(slugs.length === 1 ? { channelSlug: slugs[0] } : {}), + spec: { + kind: "fetch-windows", + // A spec needs a slug; the replay reads `params.items`. + slug: slugs[0], + params: { + items: req.items, + requestedBy: req.requestedBy, + queueKey: req.queueKey, + ...(req.manifest ? { manifest: req.manifest } : {}), + ...(req.siteId ? { siteId: req.siteId } : {}), + ...(req.maxHeight !== undefined ? { maxHeight: req.maxHeight } : {}), + }, + }, + fn: async (onLog, signal, setProgress, ctx) => { + const result = await fetchWindows({ + paths, + items: req.items, + provenance: { + requestedBy: req.requestedBy, + manifest: req.manifest, + requestedAt, + }, + maxHeight: req.maxHeight, + onLog, + signal, + drainSignal: ctx.drainSignal, + setProgress, + }); + safeRevalidate([ + ...new Set(req.items.map((i) => `/channels/${i.slug}/videos/${i.id}`)), + ]); + // A run that stopped short or lost a window did not do what it was + // asked: the job says so, and a re-run picks up the rest. + if (result.stopped && result.stopped !== "drain" && result.stopped !== "cancel") { + throw new Error( + `Stopped (${result.stopped}) with ${result.notAttempted.length} window(s) ` + + `not attempted — run it again later.`, + ); + } + if (result.failed.length > 0) { + throw new Error(`${result.failed.length} window(s) failed to fetch.`); + } + }, + }); +} + +// Fetch a list of windows ACROSS channels and platforms. A dry run answers with +// the plan and starts nothing. +export async function fetchWindowsAction( + req: FetchWindowsRequest & { siteId?: string }, +): Promise<FetchWindowsActionResult> { + const requestedBy = req.requestedBy.trim(); + if (!requestedBy) { + return { ok: false, error: "requestedBy is required (who is asking for these bytes)." }; + } + if (req.maxHeight !== undefined && !isFetchMaxHeight(req.maxHeight)) { + return { + ok: false, + error: `maxHeight ${req.maxHeight} is not a source height to cap a fetch at.`, + }; + } + const { groups, cached, unresolved } = await planFetchWindows(req.items, req.queueKey); + const base = { groups, cached, unresolved, jobs: [], refused: [] }; + if (req.dryRun) return { ok: true, dryRun: true, ...base }; + if (groups.length === 0) return { ok: true, dryRun: false, ...base }; + + // A window is small, but it lands on the corpus disk: the floor an operator's + // click asks (manual mode), once, against the channels' text. + const paths = getPaths(); + const disk = await diskGate(paths, getSettings(), { + mode: "manual", + dir: paths.channelsDir, + }); + if (!disk.ok) { + return { + ok: false, + error: + `Low disk space: ${formatBytes(disk.freeBytes)} free, ` + + `${formatBytes(disk.thresholdBytes)} required. Free up space or ` + + `lower the floor in Settings.`, + }; + } + + const jobs: { platform: string; queueKey: string; jobId: string; items: number }[] = []; + const refused: { platform: string; error: string; items: number }[] = []; + for (const g of groups) { + const refusal = await groupRefusal(g.platform); + if (refusal) { + refused.push({ platform: g.platform, error: refusal, items: g.items.length }); + continue; + } + let res: StreamActionResult; + try { + res = await fetchWindowsJobAction({ + items: g.items, + requestedBy, + manifest: req.manifest, + siteId: req.siteId, + maxHeight: req.maxHeight, + queueKey: g.queueKey, + }); + } catch (e) { + refused.push({ platform: g.platform, error: (e as Error).message, items: g.items.length }); + continue; + } + if (!res.ok) { + refused.push({ platform: g.platform, error: res.error, items: g.items.length }); + continue; + } + // Nobody reads the stream: the job's log is on disk. + void res.stream.cancel().catch(() => {}); + jobs.push({ + platform: g.platform, + queueKey: g.queueKey, + jobId: res.jobId, + items: g.items.length, + }); + } + if (jobs.length === 0 && refused.length > 0) { + return { + ok: false, + error: refused.map((r) => `${r.platform}: ${r.error}`).join("; "), + }; + } + return { ok: true, dryRun: false, groups, cached, unresolved, jobs, refused }; +} + +export type FetchMissingEvidenceResult = + | { ok: false; error: string } + | (Extract<FetchWindowsActionResult, { ok: true }> & { + siteId: string; + // Cited spans already on disk, and the ones no window fetch can fill. + onDisk: number; + unfetchable: UnfetchableEvidenceWindow[]; + }); + +// Every window a site's published reports cite and the disk does not hold +// (publish/reportMedia.ts, missingEvidenceWindows), fetched as above. What the +// Reports tab's "Fetch missing evidence" and `pnpm ops fetch-windows +// {"siteId": …}` both run. +export async function fetchMissingEvidenceAction( + siteId: string, + opts: { dryRun?: boolean; maxHeight?: number } = {}, +): Promise<FetchMissingEvidenceResult> { + const id = siteId.trim(); + if (!isValidSiteId(id)) return { ok: false, error: `"${id}" is not a valid site id` }; + const paths = getPaths(); + if (!listSiteIds(paths).includes(id)) return { ok: false, error: `No site "${id}"` }; + const need = await missingEvidenceWindows(paths, id); + const extra = { siteId: id, onDisk: need.onDisk, unfetchable: need.unfetchable }; + if (need.missing.length === 0) { + return { + ok: true, + dryRun: opts.dryRun === true, + groups: [], + jobs: [], + refused: [], + cached: [], + unresolved: [], + ...extra, + }; + } + const r = await fetchWindowsAction({ + items: need.missing.map((w) => ({ + slug: w.slug, + id: w.id, + from: w.from, + to: w.to, + clipId: w.clipId, + ...(w.reason ? { reason: w.reason.slice(0, 400) } : {}), + })), + requestedBy: "reports", + manifest: id, + siteId: id, + maxHeight: opts.maxHeight, + dryRun: opts.dryRun, + }); + if (!r.ok) return r; + return { ...r, ...extra }; +} diff --git a/editor/app/jobs/components/JobProgressBars.tsx b/editor/app/jobs/components/JobProgressBars.tsx @@ -158,6 +158,7 @@ const METRIC_LABELS: Record<JobProgressMetric, string> = { digests: "Digests", backfills: "Backfill", scans: "Metadata scan", + clips: "Clip windows", }; // Per-metric glyph for the compact line ("↓ 5/10"). @@ -167,6 +168,7 @@ const METRIC_PREFIX: Record<JobProgressMetric, string> = { digests: "\u00b6 ", backfills: "\u21ba ", scans: "\u2315 ", + clips: "\u2702 ", }; const METRIC_FILL: Record<JobProgressMetric, string> = { @@ -175,6 +177,7 @@ const METRIC_FILL: Record<JobProgressMetric, string> = { digests: "bg-info", backfills: "bg-warning", scans: "bg-info/60", + clips: "bg-success/60", }; // One-line textual summary of a job's batch progress, e.g. "↓ 5/10 · ~2m left" diff --git a/editor/app/jobs/jobReplayRegistry.ts b/editor/app/jobs/jobReplayRegistry.ts @@ -57,6 +57,8 @@ import { refreshVideoMetadataAction, replayFetchWindowAction, } from "../channels/[slug]/videos/[id]/videoActions"; +import { fetchWindowsJobAction } from "../channels/[slug]/videos/fetchWindowsAction"; +import type { FetchWindowsItem } from "yt-dlp-transcript-common/controller/fetchWindows"; import { reportsPrepareAction } from "../sites/lib/reportsPrepareAction"; import { reportsExportAction } from "../sites/lib/reportsExportAction"; @@ -84,6 +86,33 @@ const num = (v: unknown): number | undefined => const strings = (v: unknown): string[] | undefined => Array.isArray(v) ? v.filter((k): k is string => typeof k === "string") : undefined; +// A fetch-windows spec's items, keeping only the well-formed: the action +// re-checks every window, so this only drops what is not even shaped like one. +function windowItems(v: unknown): FetchWindowsItem[] { + if (!Array.isArray(v)) return []; + const out: FetchWindowsItem[] = []; + for (const raw of v) { + if (typeof raw !== "object" || raw === null) continue; + const r = raw as Record<string, unknown>; + const slug = str(r.slug); + const id = str(r.id); + const from = num(r.from); + const to = num(r.to); + if (!slug || !id || from === undefined || to === undefined) continue; + out.push({ + slug, + id, + from, + to, + clipId: str(r.clipId), + reason: str(r.reason), + pad: num(r.pad), + webpageUrl: str(r.webpageUrl), + }); + } + return out; +} + // Flag-style params + the captured queueKey for a spec. function params(spec: JobSpec): { p: Record<string, unknown>; @@ -150,6 +179,26 @@ export const JOB_REPLAY_HANDLERS: Record<string, ReplayHandler> = { maxHeight: num(p.maxHeight), }); }, + // A batch of windows. Replay re-runs the same list on the same queue; the + // controller's cache check skips every window an earlier run fetched. + "fetch-windows": (spec) => { + const { p, queueKey } = params(spec); + const items = windowItems(p.items); + if (items.length === 0) { + return Promise.resolve({ ok: false, error: "Job spec has no windows." }); + } + if (queueKey === undefined) { + return Promise.resolve({ ok: false, error: "Job spec is missing its queue." }); + } + return fetchWindowsJobAction({ + items, + requestedBy: str(p.requestedBy) ?? "unknown", + manifest: str(p.manifest), + siteId: str(p.siteId), + maxHeight: num(p.maxHeight), + queueKey, + }); + }, // Both digest lanes replay through one action; the lane comes from params so a // replayed metered run stays metered (and is refused if the lane has since // been turned off, rather than quietly falling back to local). diff --git a/editor/app/sites/[siteId]/reports/page.tsx b/editor/app/sites/[siteId]/reports/page.tsx @@ -6,6 +6,7 @@ import { formatBytes } from "yt-dlp-transcript-common/lib/format"; import { isCitedSite } from "yt-dlp-transcript-common/lib/site"; import { readReportMediaIndex } from "yt-dlp-transcript-common/publish/reportMedia"; import { PrepareReportMediaButton } from "../../components/PrepareReportMediaButton"; +import { FetchMissingEvidenceButton } from "../../components/FetchMissingEvidenceButton"; import { ExportReportsButton } from "../../components/ExportReportsButton"; import { REPORT_EXPORT_FILENAMES, @@ -125,13 +126,17 @@ export default async function SiteReportsPage({ <div> <h2 className="text-lg font-semibold">Evidence media</h2> <p className="text-sm text-muted-foreground"> - Cuts a clip of every span the published reports cite and copies every - cited post capture, from the media already on disk, ready for the - site&apos;s next build. Nothing is downloaded: a citation whose media - is missing is listed as a problem — fetch its window or persist its - video from the video page, then prepare again. + Fetch missing evidence downloads the window of every cited span + whose media is not on disk, one paced job per platform; Preview + lists them, and the spans no window can fill (a deleted video, a + channel off this site), without fetching. Prepare then cuts a clip + of every cited span and copies every cited post capture, from the + media on disk, ready for the site&apos;s next build — it downloads + nothing, and lists a citation whose media is still missing as a + problem. </p> </div> + <FetchMissingEvidenceButton siteId={siteId} /> <PrepareReportMediaButton siteId={siteId} /> <p className="text-sm text-muted-foreground"> Last prepare job:{" "} diff --git a/editor/app/sites/components/FetchMissingEvidenceButton.tsx b/editor/app/sites/components/FetchMissingEvidenceButton.tsx @@ -0,0 +1,91 @@ +"use client"; + +import Link from "next/link"; +import { useState, useTransition } from "react"; +import { Button } from "yt-dlp-transcript-common/components/ui/button"; +import { + fetchMissingEvidenceAction, + type FetchMissingEvidenceResult, +} from "../../channels/[slug]/videos/fetchWindowsAction"; + +// "Fetch missing evidence": every window the site's published reports cite and +// the disk does not hold, fetched as one paced job per platform +// (fetch-windows). Preview answers the same question without starting +// anything. It is not a streamed log because a run can start two jobs (YouTube +// and Rumble side by side); each is linked, and its log is on /jobs. +export function FetchMissingEvidenceButton({ siteId }: { siteId: string }) { + const [pending, start] = useTransition(); + const [result, setResult] = useState<FetchMissingEvidenceResult | null>(null); + const run = (dryRun: boolean) => { + start(async () => { + setResult(await fetchMissingEvidenceAction(siteId, { dryRun })); + }); + }; + return ( + <div className="flex flex-col gap-2"> + <div className="flex flex-wrap items-center gap-2"> + <Button type="button" variant="outline" disabled={pending} onClick={() => run(true)}> + Preview missing evidence + </Button> + <Button type="button" disabled={pending} onClick={() => run(false)}> + {pending ? "Working…" : "Fetch missing evidence"} + </Button> + </div> + {result && <FetchSummary result={result} />} + </div> + ); +} + +function FetchSummary({ result }: { result: FetchMissingEvidenceResult }) { + if (!result.ok) { + return ( + <p role="alert" className="text-sm text-destructive"> + {result.error} + </p> + ); + } + const toFetch = result.groups.reduce((n, g) => n + g.items.length, 0); + return ( + <div className="flex flex-col gap-1 text-sm" aria-label="Missing evidence"> + <p> + {result.onDisk} cited span(s) on disk · {toFetch} window(s) to fetch + {result.cached.length > 0 && <> · {result.cached.length} already fetched</>} + {result.unfetchable.length > 0 && <> · {result.unfetchable.length} cannot be fetched</>} + </p> + {result.dryRun + ? result.groups.map((g) => ( + <p key={g.queueKey} className="text-muted-foreground"> + {g.platform}: {g.items.length} window(s) on {g.queueKey} + </p> + )) + : result.jobs.map((j) => ( + <p key={j.jobId}> + {j.platform}: {j.items} window(s) —{" "} + <Link href={`/jobs/${j.jobId}`} className="font-mono underline"> + {j.jobId} + </Link> + </p> + ))} + {!result.dryRun && + result.refused.map((r) => ( + <p key={r.platform} role="alert" className="text-destructive"> + {r.platform}: {r.items} window(s) not started — {r.error} + </p> + ))} + {result.unresolved.map((u) => ( + <p key={`${u.item.slug}/${u.item.id}/${u.item.from}`} className="text-destructive"> + {u.item.clipId ?? `${u.item.slug}/${u.item.id}`}: {u.error} + </p> + ))} + {result.unfetchable.length > 0 && ( + <ul className="list-disc pl-5 text-muted-foreground"> + {result.unfetchable.map((u) => ( + <li key={u.clipId}> + {u.clipId} ({u.slug}/{u.id}): {u.message} + </li> + ))} + </ul> + )} + </div> + ); +} diff --git a/editor/e2e/fetch-window.spec.ts b/editor/e2e/fetch-window.spec.ts @@ -9,7 +9,7 @@ // Same token as /api/worker/* (the test server runs with // WORKER_TOKEN=test-worker-token; see package.json dev:test). -import { readdir, readFile, stat, writeFile } from "node:fs/promises"; +import { mkdir, readdir, readFile, stat, writeFile } from "node:fs/promises"; import { test, expect, type APIRequestContext } from "@playwright/test"; import { generateReport, @@ -51,6 +51,8 @@ async function invocations(): Promise<string> { async function pollJob( request: APIRequestContext, jobId: string, + // A batch sleeps the clip-window gap (30–45 s) between two fetches. + timeout = 30_000, ): Promise<Record<string, unknown>> { let last: Record<string, unknown> = {}; await expect @@ -63,7 +65,7 @@ async function pollJob( last = (await r.json()) as Record<string, unknown>; return last.status as string; }, - { timeout: 30_000 }, + { timeout }, ) .not.toMatch(/^(queued|running)$/); return last; @@ -325,6 +327,84 @@ test("a 429 fails the job and puts the platform in cooldown", async ({ expect(body.cooldownMs).toBeGreaterThan(0); }); +// THE BATCH: POST /api/ops/fetch-windows, one paced job per platform queue. +// Three windows: one already on disk (no request, no pause), one the source +// refuses with a 403, one that fetches. A single 403 is an item failure — a +// removed Rumble page answers 403 too — so the run carries on past it and the +// platform is NOT backed off; the job still ends failed, naming the window it +// lost. Exactly one pause is owed (between the two network fetches), at the +// clip-window floor of 30–45 s. +test("a batch skips what is cached, survives one 403, and fetches the rest", async ({ + request, +}) => { + test.setTimeout(150_000); + await mkdir(resolvePath(rel(`data/${VIDEO}/clips`)), { recursive: true }); + await writeFile(resolvePath(clipRel("0.00-30.00.mp4")), "already here"); + + const post = await request.post(`${baseUrl}/api/ops/fetch-windows`, { + headers: AUTH, + data: { + requestedBy: "umtool", + manifest: "demo-batch", + items: [ + { slug: SLUG, id: VIDEO, from: 5, to: 10, clipId: "c01" }, + { + slug: SLUG, + id: "win403vid1", + // The fake reads the sentinel out of the URL. + webpageUrl: "https://www.youtube.com/watch?v=win403vid1", + from: 5, + to: 10, + clipId: "c02", + }, + { slug: SLUG, id: VIDEO, from: 50, to: 60, clipId: "c03" }, + ], + }, + }); + expect(post.status()).toBe(200); + const body = (await post.json()) as { + cached: { clipId: string }[]; + jobs: { platform: string; jobId: string; items: number }[]; + jobId: string; + }; + expect(body.cached.map((c) => c.clipId)).toEqual(["c01"]); + expect(body.jobs).toHaveLength(1); + expect(body.jobs[0]).toMatchObject({ platform: "youtube", items: 2 }); + + const finished = await pollJob(request, body.jobId, 90_000); + expect(finished.status).toBe("failed"); + expect(String(finished.error)).toMatch(/1 window\(s\) failed to fetch/); + expect(String(finished.error)).toMatch(/1 fetched, 0 cached, 1 failed/); + + expect(await exists(clipRel("50.00-60.00.mp4"))).toBe(true); + const sidecar = await readJson<{ requestedBy: string; manifest: string; clipId: string }>( + clipRel("50.00-60.00.json"), + ); + expect(sidecar).toMatchObject({ requestedBy: "umtool", manifest: "demo-batch", clipId: "c03" }); + expect( + await exists(`test-transcripts/channels/${SLUG}/data/win403vid1/clips/5.00-10.00.mp4`), + ).toBe(false); + + // One 403 did not back the platform off: the next ask is not refused. + const next = await request.post(`${baseUrl}/api/ops/fetch-windows`, { + headers: AUTH, + data: { + requestedBy: "umtool", + dryRun: true, + items: [{ slug: SLUG, id: VIDEO, from: 100, to: 110 }], + }, + }); + expect(next.status()).toBe(200); + const plan = (await next.json()) as { groups: { platform: string }[] }; + expect(plan.groups.map((g) => g.platform)).toEqual(["youtube"]); + const single = await request.post(`${baseUrl}/api/media/fetch-window`, { + headers: AUTH, + data: { channelSlug: SLUG, videoId: VIDEO, from: 100, to: 110, requestedBy: "umtool" }, + }); + expect(single.status()).toBe(202); + await pollJob(request, ((await single.json()) as { jobId: string }).jobId); +}); + // THE WHOLE RECORDING, when a window will not do — a tool that needs to re-cut // freely, or a source whose windows would tile the entire runtime. // diff --git a/editor/e2e/fixtures/bin/fake-ytdlp.mjs b/editor/e2e/fixtures/bin/fake-ytdlp.mjs @@ -573,6 +573,13 @@ async function main() { ); process.exit(1); } + // `win403`: Rumble's Cloudflare / googlevideo refusing the media request — + // a `network` failure, which the fetch-windows batch counts toward its + // two-in-a-row backoff. + if (url.toLowerCase().includes("win403")) { + process.stderr.write(`ERROR: [download] Got error: HTTP Error 403: Forbidden\n`); + process.exit(1); + } await ensureDir(path.dirname(dest)); // Deterministic bytes, one chunk's worth, so a size assertion is stable. await writeFile(dest, Buffer.alloc(CHUNK_BYTES, "w")); diff --git a/scripts/archilyzer-ops.mjs b/scripts/archilyzer-ops.mjs @@ -56,6 +56,8 @@ // pnpm ops tags --json '{"op":"define","tag":{"id":"eva-collab","label":"Collab"}}' // pnpm ops tag-videos --file ids.json // pnpm ops persist-videos --file list.json --wait +// pnpm ops fetch-windows --json '{"siteId":"demo-site","dryRun":true}' +// pnpm ops fetch-windows --file windows.json --wait // pnpm ops get tags eva-collab // pnpm ops cut-release --json '{"workspace":"all","version":"next","commit":true}' // @@ -172,6 +174,9 @@ const ACTIONS = [ // Persist specific videos, across channels, to the saved-video store // ({items: [{slug, id}]}), paced and gated; a re-run resumes. "persist-videos", + // Fetch clip windows through the managed path, one paced job per platform + // queue ({siteId} = a site's missing evidence, or {items, requestedBy}). + "fetch-windows", // Cut a changelog's [Unreleased] into a dated release heading (release 10 // slice P). Synchronous. The same writer as `archilyzer release cut`, which // needs no editor at all — this route exists only on an editor built from @@ -434,6 +439,18 @@ export function usage() { " queue; a low disk or a rate limit stops it, and running the same body", " again resumes — saved videos are skipped.", "", + 'fetch-windows fetches clip windows, one paced job per platform queue', + ' (YouTube and Rumble side by side): {"siteId"} fetches every window the', + ' site\'s published reports cite and the disk does not hold; {"items":', + ' [{"slug", "id", "from", "to", "clipId"?, "reason"?, "pad"?,', + ' "webpageUrl"?}, ...], "requestedBy", "manifest"?} fetches a list.', + ' "maxHeight" caps the source height (default 720). "dryRun": true lists', + ' the windows per platform, the ones already on disk ("cached") and the', + ' ones no fetch can fill ("unfetchable": deleted, off the site) and starts', + ' nothing. A platform cooling down or held is refused for its group; a', + ' 429, or two 403s in a row, backs the platform off and stops its job.', + ' Running the same body again resumes — fetched windows are cached.', + "", '"preview": "<branch>" on deploy-site or build-deploy makes it a Cloudflare', " Pages PREVIEW instead of production: the same bundle goes to a branch", " alias, https://<branch>.<project>.pages.dev, and the live site is left", diff --git a/scripts/archilyzer-ops.test.mjs b/scripts/archilyzer-ops.test.mjs @@ -445,6 +445,15 @@ test("persist-videos is a POST to its route, named in the usage", () => { assert.match(usage(), /persist-videos/); }); +test("fetch-windows is a POST to its route, named in the usage", () => { + const p = parseArgs(["fetch-windows", "--json", '{"siteId":"demo-site","dryRun":true}']); + assert.equal(p.method, "POST"); + assert.equal(p.path, "/api/ops/fetch-windows"); + assert.deepEqual(p.body, { siteId: "demo-site", dryRun: true }); + assert.equal(parseArgs(["fetch-windows", "--file", "windows.json"]).bodyFile, "windows.json"); + assert.match(usage(), /fetch-windows fetches clip windows/); +}); + test("feed-metadata posts {slug, dryRun} to /api/ops/feed-metadata", () => { const p = parseArgs(["feed-metadata", "--json", '{"slug":"demo-channel","dryRun":true}']); assert.equal(p.method, "POST"); diff --git a/umtool/docs/quirks.md b/umtool/docs/quirks.md @@ -28,6 +28,11 @@ progressive URL (YouTube's googlevideo mp4) and ffmpeg aborts with "Option extension_picky not found". Adding it unconditionally trades a Rumble failure for a YouTube one. +**A 403 backs the platform off in a batch.** `fetch-via-editor.mjs --all` (the +editor's `fetch-windows` job) treats one 403 as that clip's failure — a removed +Rumble page answers 403 too — but two in a row back the platform off and stop the +job, as a 429 does at once. The windows are 30–45 s apart; 20 s drew YouTube 403s. + **`--force-keyframes-at-cuts` matters because the clip IS the citation.** Without it the cut snaps to the nearest preceding keyframe, which can be seconds early. Fine for scrubbing; not fine when someone is checking your quote. diff --git a/umtool/report-to-video/fetch-via-editor.mjs b/umtool/report-to-video/fetch-via-editor.mjs @@ -56,8 +56,13 @@ function die(message) { const manifestPath = argv.find((a) => !a.startsWith("-") && a.endsWith(".json")); const clipId = flag("--fetch-only") ?? flag("--clip"); -if (!manifestPath || !clipId) { - die("usage: fetch-via-editor.mjs <manifest.json> --fetch-only <clipId> [--full] [--max-height N] [--pad-before N] [--pad-after N] [--progress ndjson]"); +// THE WHOLE TIMELINE AS ONE ASK: every clip entry's window, posted as one +// `fetch-windows` call. The editor answers the windows already on disk at once +// and fetches the rest as one paced job per platform — the pacing, the 403 +// streak and the 429 stop live there, not in a shell loop here. +const wantAll = argv.includes("--all"); +if (!manifestPath || (!clipId && !wantAll) || (clipId && wantAll)) { + die("usage: fetch-via-editor.mjs <manifest.json> (--fetch-only <clipId> [--full] | --all) [--max-height N] [--pad N] [--pad-before N] [--pad-after N] [--progress ndjson]"); } // THE WHOLE RECORDING INSTEAD OF A WINDOW. For a clip whose windows would tile @@ -102,6 +107,111 @@ if (!token) { const whole = JSON.parse(await readFile(manifestPath, "utf8")); const provenance = whole.provenance ?? {}; +const pad = Number(flag("--pad") ?? 3); +const padBefore = Number(flag("--pad-before") ?? pad); +const padAfter = Number(flag("--pad-after") ?? pad); +// TWO DECIMALS, matching the editor's own naming (common/lib/clipWindow.ts) and +// the build's. The name IS the window, so a request that rounds differently +// addresses a different file and the cache misses forever. +const windowOf = (e) => ({ + from: Number(Math.max(0, Number(e.start) - padBefore).toFixed(2)), + to: Number((Number(e.end) + padAfter).toFixed(2)), +}); +const manifestId = provenance.manifestId ?? path.basename(path.dirname(path.resolve(manifestPath))); + +const headers = { + authorization: `Bearer ${token}`, + "content-type": "application/json", +}; + +async function ask(url, init) { + try { + return await fetch(url, init); + } catch (err) { + die(`could not reach the editor at ${editorUrl}: ${err.message}`); + } +} + +if (wantAll) { + if (argv.includes("--full")) die("--all fetches windows; --full is one clip's whole recording"); + const clips = (whole.timeline ?? []).filter((e) => e.type === "clip"); + const items = []; + for (const e of clips) { + const slug = e.channel ?? provenance.channelSlug; + if (!slug || !e.video) die(`${e.id} has no channel or video to fetch`); + items.push({ + slug, + id: e.video, + ...windowOf(e), + clipId: e.id, + pad: Math.max(padBefore, padAfter), + ...(e.webpageUrl ? { webpageUrl: e.webpageUrl } : {}), + reason: String(e.note ?? e.quote ?? `clip window with ${padBefore}s before / ${padAfter}s after`).slice(0, 400), + }); + } + if (items.length === 0) { + EMIT("note", { message: "no clip entries in the timeline — nothing to fetch" }); + EMIT("done", { out: null, nothingToFetch: true }); + process.exit(0); + } + const res = await ask(`${editorUrl}/api/ops/fetch-windows`, { + method: "POST", + headers, + body: JSON.stringify({ + items, + requestedBy: "umtool", + manifest: manifestId, + ...(maxHeight !== undefined ? { maxHeight } : {}), + }), + }); + const body = await res.json().catch(() => ({})); + if (!res.ok || !body.ok) { + die(`the editor refused (HTTP ${res.status}): ${body.error ?? "no reason given"}`); + } + EMIT("note", { + message: + `${items.length} clip(s): ${body.cached.length} already on disk, ` + + `${body.jobs.reduce((n, j) => n + j.items, 0)} queued in ${body.jobs.length} job(s)`, + }); + for (const u of body.unresolved ?? []) { + EMIT("note", { message: ` ${u.item.clipId ?? u.item.id}: ${u.error}` }); + } + for (const r of body.refused ?? []) { + EMIT("note", { message: ` ${r.platform}: ${r.items} window(s) not started — ${r.error}` }); + } + // Every job, polled to its end. A job that stopped short (a rate limit, two + // 403s) fails with the reason; asking again later resumes it, the fetched + // windows answering from the cache. + const pending = new Map(body.jobs.map((j) => [j.jobId, j])); + const failed = []; + const deadline = Date.now() + 6 * 60 * 60_000; + const last = new Map(); + while (pending.size > 0) { + if (Date.now() > deadline) die(`gave up waiting for editor job(s) ${[...pending.keys()].join(", ")}`); + await new Promise((r) => setTimeout(r, POLL_MS * 5)); + for (const [jobId, j] of pending) { + const poll = await ask(`${editorUrl}/api/media/fetch-window/${jobId}`, { headers }); + const p = await poll.json().catch(() => ({})); + if (!poll.ok) die(`polling editor job ${jobId} failed (HTTP ${poll.status}): ${p.error ?? ""}`); + if (p.status !== last.get(jobId)) { + last.set(jobId, p.status); + EMIT("note", { message: ` editor job ${jobId} (${j.platform}, ${j.items} window(s)): ${p.status}` }); + } + if (p.status === "done") pending.delete(jobId); + else if (p.status === "failed" || p.status === "cancelled") { + pending.delete(jobId); + failed.push(`${jobId} ${p.status}: ${String(p.error ?? "").split("\n").slice(-3).join(" / ")}`); + } + } + } + const unfinished = failed.length + (body.refused?.length ?? 0) + (body.unresolved?.length ?? 0); + if (unfinished > 0) { + die(`not every window was fetched:\n ${[...failed, ...(body.refused ?? []).map((r) => `${r.platform}: ${r.error}`)].join("\n ") || "see the notes above"}`); + } + EMIT("done", { out: null, all: true, clips: items.length }); + process.exit(0); +} + // Same resolution build-video.mjs's --fetch-only does, and for the same // reasons: a still has nothing to fetch, a non-clip entry is an error, and a // LEDGER CLAIM is a moment rather than a window (most of a ledger is cited by @@ -136,37 +246,16 @@ if (!channelSlug) { die(`${clipId} has no channel, and the manifest's provenance names none`); } -const pad = Number(flag("--pad") ?? 3); -const padBefore = Number(flag("--pad-before") ?? pad); -const padAfter = Number(flag("--pad-after") ?? pad); -// TWO DECIMALS, matching the editor's own naming (common/lib/clipWindow.ts) and -// the build's. The name IS the window, so a request that rounds differently -// addresses a different file and the cache misses forever. -const from = Number(Math.max(0, Number(entry.start) - padBefore).toFixed(2)); -const to = Number((Number(entry.end) + padAfter).toFixed(2)); +const { from, to } = windowOf(entry); // WHY THESE SECONDS. Stored beside the file so a directory of windows can be // read back months later. The clip's own note is the closest thing the manifest // has to a reason; the manifest id and the clip id say the rest. -const manifestId = provenance.manifestId ?? path.basename(path.dirname(path.resolve(manifestPath))); const reason = entry.note ?? entry.quote ?? `clip window with ${padBefore}s before / ${padAfter}s after`; -const headers = { - authorization: `Bearer ${token}`, - "content-type": "application/json", -}; - -async function ask(url, init) { - try { - return await fetch(url, init); - } catch (err) { - die(`could not reach the editor at ${editorUrl}: ${err.message}`); - } -} - EMIT("fetch", { id: entry.id, video: entry.video,