// 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, CLIP_WINDOW_PLATFORM_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 { 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 { 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 = {}) { const rec: Recorded = { sleeps: [], backoffs: [], cleans: [], progress: [] }; const deps: Partial = { 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; }, gapRemainingMs: () => 0, noteGap: () => {}, ...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, 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("E2E_CLIP_WINDOW_GAP_MS replaces the default gap; an explicit gapMs still wins", async () => { const saved = process.env.E2E_CLIP_WINDOW_GAP_MS; try { process.env.E2E_CLIP_WINDOW_GAP_MS = "2000"; const h = harness(); await fetchWindows({ paths: await corpus(), items: [item("v1", "/ok"), item("v2", "/ok")], provenance: { requestedBy: "test" }, deps: h.deps, }); assert.deepEqual(h.rec.sleeps, [2000]); const explicit = harness(); await go(await corpus(), [item("v1", "/ok"), item("v2", "/ok")], explicit, { gapMs: 1234 }); assert.deepEqual(explicit.rec.sleeps, [1234]); } finally { if (saved === undefined) delete process.env.E2E_CLIP_WINDOW_GAP_MS; else process.env.E2E_CLIP_WINDOW_GAP_MS = saved; } }); 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); }); test("rumble windows are further apart, whatever platform their channel names", async () => { const h = harness(); const rumble = (id: string): FetchWindowsItem => ({ slug: "c", id, from: 10, to: 20, webpageUrl: `https://rumble.com/${id}-x.html` }); const r = await fetchWindows({ paths: await corpus(), items: [rumble("v1"), rumble("v2")], provenance: { requestedBy: "test" }, deps: h.deps, }); assert.equal(r.fetched.length, 2); assert.equal(h.rec.sleeps.length, 1); const floor = CLIP_WINDOW_PLATFORM_MIN_GAP_SECONDS.rumble * 1000; assert.ok(h.rec.sleeps[0] >= floor, `${h.rec.sleeps[0]} is at least rumble's ${floor}`); assert.ok(h.rec.sleeps[0] <= floor * 1.5 + 60_000); }); test("a second run on the same platform waits out the gap the first one set", async () => { const next = new Map(); const shared = { gapRemainingMs: (key: string) => next.get(key) ?? 0, noteGap: (key: string, ms: number) => void next.set(key, ms), }; const first = harness(shared); await go(await corpus(), [item("v1", "/ok"), item("v2", "/403")], first); assert.deepEqual(first.rec.sleeps, [1000], "a run's own first fetch owes nothing"); assert.equal(next.get("clip-window:youtube"), 1000, "a refused fetch sets the gap too"); const second = harness(shared); await go(await corpus(), [item("v3", "/ok"), item("v4", "/ok")], second); assert.deepEqual(second.rec.sleeps, [1000, 1000], "the first fetch waits for the earlier run's gap"); }); // THE SAVED CONTAINER (release 21 D2): a window it holds is cut, not fetched; // a pointer it cannot be read for fails the item, never a fetch in its place; // a container that ends too early is fetched as before. test("saved containers: a covered window is cut with no request and no pause; unreadable fails; not-covering is fetched", async () => { const paths = await corpus(); const cuts: { id: string; from: number; to: number; clipId?: string }[] = []; const h = harness({ savedWindowSource: async ({ videoDir }) => { const id = path.basename(videoDir); if (id === "held") { return { kind: "covers", pointer: { storedAt: "", dir: "/store", file: "source-media.mp4", bytes: 1 }, container: "/store/source-media.mp4", durationSec: 600, probe: { durationSec: 600, width: 160, height: 90, hasVideo: true, hasAudio: true }, }; } if (id === "unplugged") { return { kind: "unreadable", error: `Channel "c": media is not reachable — its saved video is not there` }; } if (id === "short") return { kind: "not-covering", container: "/store/x.mp4", durationSec: 5 }; return { kind: "none" }; }, cutSavedWindow: async ({ videoDir, from, to, provenance }) => { cuts.push({ id: path.basename(videoDir), from, to, clipId: provenance.clipId }); return { file: "x.mp4", path: "/x.mp4", from, to, bytes: 1, provenance: { requestedBy: "test", source: "saved-video" } }; }, }); const before = await spawns(); const r = await go( paths, [ { ...item("held", "/a"), clipId: "r#1" }, item("unplugged", "/b"), item("short", "/c"), { ...item("held", "/d", 30, 40), clipId: "r#2" }, ], h, ); assert.deepEqual(cuts, [ { id: "held", from: 10, to: 20, clipId: "r#1" }, { id: "held", from: 30, to: 40, clipId: "r#2" }, ]); assert.deepEqual(r.cut.map((i) => i.id), ["held", "held"]); assert.equal(r.failed.length, 1); assert.equal(r.failed[0].class, "unreachable"); assert.equal(r.failed[0].item.id, "unplugged"); assert.match(r.failed[0].message, /media is not reachable/); assert.deepEqual(r.fetched.map((i) => i.id), ["short"]); assert.equal((await spawns()) - before, 1, "only the not-covering window touched the network"); assert.deepEqual(h.rec.sleeps, [], "a cut owes no pause, and one fetch owes none either"); assert.deepEqual(h.rec.progress.at(-1), { metric: "clips", initial: 0, target: 4, current: 4 }); }); test("saved containers: a cut that fails is a cut-failed item and the run carries on", async () => { const h = harness({ savedWindowSource: async () => ({ kind: "covers", pointer: { storedAt: "", dir: "/store", file: "source-media.mp4", bytes: 1 }, container: "/store/source-media.mp4", durationSec: 600, probe: { durationSec: 600, width: 160, height: 90, hasVideo: true, hasAudio: true }, }), cutSavedWindow: async ({ from }) => { if (from === 10) throw new Error("ffmpeg failed cutting 10.00–20.00"); return { file: "x.mp4", path: "/x.mp4", from, to: from + 10, bytes: 1, provenance: null }; }, }); const r = await go(await corpus(), [item("v1", "/a"), item("v2", "/b", 30, 40)], h); assert.equal(r.failed.length, 1); assert.equal(r.failed[0].class, "cut-failed"); assert.deepEqual(r.cut.map((i) => i.id), ["v2"]); assert.equal(r.stopped, undefined); });