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 { cutSavedVideoWindow, savedWindowSource, } from "../lib/savedVideoWindow-server"; import type { JobProgress } from "../jobs/registry"; import { downloadGapMs } from "../jobs/platformBackoff"; import { notePlatformGap, platformGapRemainingMs } from "../jobs/platformGap"; import { heldPlatformRefusal, platformCooldownRemainingMs, recordDownloadBackoff, recordPlatformClean, } from "../jobs/downloadBackoff"; import { FetchWindowError, fetchWindowManaged, type FetchWindowOpts, type FetchWindowResult, } from "../ytdlp/fetchWindowManaged"; import { channelPaceSeconds, channelPlatform, configForVideoUrl, } 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. // - A WINDOW HELD BY A SAVED CONTAINER (release 21 D2) is CUT from it, with // no request and no pause either (lib/savedVideoWindow-server.ts). A // pointer whose container cannot be read fails the item as `unreachable` // — never a fetch in its place; one that ends too early is fetched. // - 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; // A platform whose windows must be further apart than that. Rumble's // Cloudflare puts the whole IP behind a JS challenge ("Just a moment…", 403 to // every rumble.com request, impersonated or not) after a handful of requests // in a few minutes, and a window costs three (page, embed JSON, HLS manifest): // on 2026-10-07 windows 30–45 s apart drew it within four, and it lifted again // in about five minutes of quiet. export const CLIP_WINDOW_PLATFORM_MIN_GAP_SECONDS: Readonly> = Object.freeze({ rumble: 120 }); // TEST-ONLY: E2E_CLIP_WINDOW_GAP_MS replaces the batch gap above, so the e2e // suite sees the one owed pause without waiting 30–45 s for it // (fetch-window.spec). Unset everywhere but the e2e test server. function testGapMs(): number | undefined { const raw = Number(process.env.E2E_CLIP_WINDOW_GAP_MS); return Number.isFinite(raw) && raw > 0 ? raw : undefined; } export type FetchWindowsItem = { slug: string; id: string; from: number; to: number; // `#`, 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, or the video's // saved container is not. | "unreachable" // ffmpeg could not cut the window out of a covering saved container. | "cut-failed" // 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[]; // Cut from the video's saved container, not fetched (release 21 D2). cut: 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; // The saved-container tier (lib/savedVideoWindow-server.ts). savedWindowSource: typeof savedWindowSource; cutSavedWindow: typeof cutSavedVideoWindow; sleep: (ms: number, signal?: AbortSignal) => Promise; getSettings: () => SiteSettings; readChannelConfig: (paths: Paths, slug: string) => Promise; findVideoSourceUrl: ( paths: Paths, slug: string, id: string, config: ChannelConfig, ) => Promise; assertTextReadable: (paths: Paths, slug: string) => Promise; cooldownRemainingMs: (platform: string, paths: Paths) => Promise; heldRefusal: (platform: string, paths: Paths) => Promise; recordBackoff: ( platform: string, paths: Paths, failureClass: "rate_limit" | "network", ) => Promise; recordClean: (platform: string, paths: Paths) => Promise; // The gap ACROSS runs (jobs/platformGap.ts, keyed `clip-window:`): // two batches back to back on one queue — a umtool manifest each — must not // put the second one's first fetch right after the first one's last. gapRemainingMs: (key: string) => number; noteGap: (key: string, gapMs: number) => void; }; function abortableSleep(ms: number, signal?: AbortSignal): Promise { 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, savedWindowSource, cutSavedWindow: cutSavedVideoWindow, 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), gapRemainingMs: (key) => platformGapRemainingMs(key), noteGap: (key, gapMs) => notePlatformGap(key, gapMs), }; 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(); 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; }): Promise { 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: [], cut: [], 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(); 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(); const textProblem = async (slug: string): Promise => { 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; } // THEN THE SAVED CONTAINER: cut, no request, no pause, not a network // attempt. const saved = await deps.savedWindowSource({ paths, slug: item.slug, videoDir, from: item.from, to: item.to, }); if (saved.kind === "unreadable") { fail(item, "unreachable", saved.error); done += 1; progress(); continue; } if (saved.kind === "covers") { try { const cut = await deps.cutSavedWindow({ paths, videoDir, source: saved, from: item.from, to: item.to, provenance: { requestedBy: provenance.requestedBy, manifest: provenance.manifest, clipId: item.clipId, reason: item.reason, pad: item.pad, requestedAt, }, signal: fetchSignal, }); result.cut.push(item); log(` ✂ ${label}: cut from the saved video into clips/${cut.file}`); } catch (e) { fail(item, "cut-failed", (e as Error).message); } 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"; // Paced as the platform the URL points at, which a channel with no // platform of its own (community-notes holds Rumble videos) does not name. const paced = configForVideoUrl(config, url); const gap = gapMs ?? testGapMs() ?? downloadGapMs( config.sleepBetweenDownloadsSeconds ?? settings.sleepBetweenDownloadsSeconds, channelPaceSeconds(paced), staticSleepRequestsSeconds(channelPlatform(paced)), { minSeconds: Math.max( CLIP_WINDOW_MIN_GAP_SECONDS, platformMinGapSeconds(channelPlatform(paced)), CLIP_WINDOW_PLATFORM_MIN_GAP_SECONDS[platform] ?? 0, ), }, ); const gapKey = `clip-window:${platform}`; // Between this run's own fetches, the gap; before its first, whatever is // left of the gap the last run on this platform set. const wait = networkAttempts > 0 ? gap : deps.gapRemainingMs(gapKey); if (wait > 0) { log( networkAttempts > 0 ? `Sleeping ${Math.round(wait / 1000)}s before the next fetch...` : `${platform} was asked for a window by an earlier run: waiting ${Math.round(wait / 1000)}s.`, ); await deps.sleep(wait, 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; } finally { // Settled, refused or not: the next fetch on this platform, in this run // or the next one, starts no sooner than the gap from now. deps.noteGap(gapKey, gap); } done += 1; progress(); } log( `Fetch windows: ${result.fetched.length} fetched, ` + (result.cut.length ? `${result.cut.length} cut from saved videos, ` : "") + `${result.cached.length} cached, ` + `${result.failed.length} failed` + (result.notAttempted.length ? `, ${result.notAttempted.length} not attempted (stopped: ${result.stopped})` : "") + ".", ); return result; }