// fetchPosts over the gallery-dl X fetcher, end to end against a fake binary: // a run that is cancelled (or times out) part-way keeps what it read and the // resume point, and the next run continues from there. A drained run stops at // a page boundary, keeps that resume point, and is not a failure. // // Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test controller/fetchPosts.test.ts import { test } from "node:test"; import assert from "node:assert/strict"; import { chmod, mkdir, mkdtemp, readFile, readdir, writeFile } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; const ROOT = await mkdtemp(path.join(os.tmpdir(), "fetchposts-")); const BIN = path.join(ROOT, "fake-gallery-dl.mjs"); const ARGS_LOG = path.join(ROOT, "argv.jsonl"); // getPaths() is resolved lazily and cached, so these must be set before the // first call — i.e. before any fetch below. process.env.TRANSCRIPTS_DIR = path.join(ROOT, "transcripts"); process.env.GALLERY_DL_BIN = BIN; process.env.FAKE_ARGS_LOG = ARGS_LOG; // Pages of three tweets, newest first. Page N is emitted, then gallery-dl's // verbose "Cursor:" line. With no cursor the fake emits pages 1–2 and then // blocks on a "rate limit" (the real stall); given the cursor after page 2 it // emits page 3 and exits cleanly. FAKE_MODE=known emits page 1 then endless // tweets older than the watermark, to prove the early stop kills it. // FAKE_MODE=pages emits six pages, each followed (after a pause) by an info // line, then its cursor: a walk with a page boundary every few hundred ms. await writeFile( BIN, `#!/usr/bin/env node import { appendFileSync } from "node:fs"; const args = process.argv.slice(2); appendFileSync(process.env.FAKE_ARGS_LOG, JSON.stringify(args) + "\\n"); const cursorArg = args.find((a) => a.startsWith("extractor.twitter.cursor=")); const cursor = cursorArg ? cursorArg.split("=")[1] : undefined; const tw = (n) => '[2, {"tweet_id": ' + (1899000000000000100n - BigInt(n)) + ', "date": "2026-06-' + String(28 - n).padStart(2, "0") + ' 12:00:00", "content": "fake ' + n + '", "author": {"name": "faketester"}, "user": {"name": "faketester"}}]\\n'; const page = (p) => { for (let i = 0; i < 3; i++) process.stdout.write(tw(p * 3 + i)); }; const log = (l) => process.stderr.write(l + "\\n"); log("[gallery-dl][debug] Version 1.32.9"); if (process.env.FAKE_MODE === "known") { page(0); let n = 3; // Distinct tweets from before the archive's newest post: what an incremental // run meets once it is past everything new. const old = (k) => '[2, {"tweet_id": ' + (1899000000000000000n - BigInt(k)) + ', "date": "2026-01-01 12:00:00", "content": "old ' + k + '", "author": {"name": "faketester"}, "user": {"name": "faketester"}}]\\n'; const t = setInterval(() => { for (let i = 0; i < 20; i++) process.stdout.write(old(n++)); }, 5); setTimeout(() => { clearInterval(t); process.exit(0); }, 30_000); } else if (process.env.FAKE_MODE === "authfail") { log("[twitter][error] HTTP redirect to login page (401 Unauthorized)"); process.exit(1); } else if (process.env.FAKE_MODE === "crash") { for (let i = 0; i < 3; i++) { process.stdout.write('[2, {"tweet_id": ' + (1899000000000000500n + BigInt(i)) + ', "date": "2026-07-0' + (i + 1) + ' 12:00:00", "content": "new ' + i + '", "author": {"name": "faketester"}, "user": {"name": "faketester"}}]\\n'); } log("[twitter][debug] Cursor: 1/CRASH-P1"); log("[twitter][error] AbortExtraction: 404 Not Found (GraphQL query id rotated)"); process.exit(4); } else if (process.env.FAKE_MODE === "pages") { let p = 0; const next = () => { page(p); setTimeout(() => { log("[twitter][info] page " + p + " read"); setTimeout(() => { log("[twitter][debug] Cursor: 1/PAGES-P" + p); p++; if (p < 6) setTimeout(next, 300); else process.exit(0); }, 300); }, 100); }; next(); } else if (cursor === "1/AFTER-P2") { page(2); process.exit(0); } else { page(0); log("[twitter][debug] Cursor: 1/AFTER-P1"); page(1); log("[twitter][debug] Cursor: 1/AFTER-P2"); log("[twitter][info] Waiting for 14 minutes until 15:11:22 (rate limit)"); setTimeout(() => process.exit(0), 60_000); } `, ); await chmod(BIN, 0o755); const SLUG = "fake-x"; const channelRoot = path.join(ROOT, "transcripts", "channels", SLUG); await mkdir(channelRoot, { recursive: true }); await writeFile( path.join(channelRoot, "config.json"), JSON.stringify({ handling: "transcribe", sourceKind: "social", postFetcher: "x-gallery-dl", socialHandle: "faketester", platform: "twitter", name: "Fake (X)", url: "https://x.com/faketester", }), ); const { fetchPosts } = await import("./fetchPosts"); const { readPostFetchState, readSeenPostIds } = await import("../lib/posts-server"); async function lastArgv(): Promise { const lines = (await readFile(ARGS_LOG, "utf8")).trim().split("\n"); return JSON.parse(lines[lines.length - 1]) as string[]; } test("a cancelled backfill keeps what it read and where it stopped", async () => { const controller = new AbortController(); const log: string[] = []; const result = await fetchPosts({ paths: (await import("../lib/paths")).getPaths(), slug: SLUG, settings: {}, signal: controller.signal, onLog: (line) => { log.push(line); // The fake is now asleep on its "rate limit" — cancel, as the operator did. if (/rate limit/.test(line)) controller.abort(); }, }); assert.equal(result.ok, true, log.join("\n")); assert.equal(result.written, 6); assert.equal(result.complete, false); assert.equal((await readSeenPostIds(channelRoot)).size, 6); const state = await readPostFetchState(channelRoot); assert.equal(state?.cursor, "1/AFTER-P2"); assert.ok(log.some((l) => /Cancelled; the next run resumes/.test(l)), log.join("\n")); const argv = await lastArgv(); assert.ok(argv.includes("output.jsonl=true")); assert.ok(argv.includes("--verbose")); }); test("the next run resumes from the saved cursor and finishes the history", async () => { const result = await fetchPosts({ paths: (await import("../lib/paths")).getPaths(), slug: SLUG, settings: {}, }); assert.equal(result.ok, true); assert.equal(result.written, 3); assert.equal(result.complete, true); assert.ok((await lastArgv()).includes("extractor.twitter.cursor=1/AFTER-P2")); assert.equal((await readSeenPostIds(channelRoot)).size, 9); // Caught up: no resume point left behind. assert.equal((await readPostFetchState(channelRoot))?.cursor, undefined); }); test("an incremental run stops itself once it is past everything new", async () => { process.env.FAKE_MODE = "known"; const started = Date.now(); const result = await fetchPosts({ paths: (await import("../lib/paths")).getPaths(), slug: SLUG, settings: {}, }); delete process.env.FAKE_MODE; assert.equal(result.ok, true); assert.equal(result.written, 0); assert.equal(result.complete, true); // Killed on the early stop, not left to run its 30 s. assert.ok(Date.now() - started < 15_000, `took ${Date.now() - started} ms`); assert.equal((await readdir(channelRoot)).includes("posts-state.json"), true); }); test("a run that dies part-way fails the job but keeps its posts and resume point", async () => { process.env.FAKE_MODE = "crash"; const result = await fetchPosts({ paths: (await import("../lib/paths")).getPaths(), slug: SLUG, settings: {}, }); delete process.env.FAKE_MODE; assert.equal(result.ok, false); assert.match(result.error ?? "", /exited 4: .*404 Not Found/); assert.equal(result.written, 3); const state = await readPostFetchState(channelRoot); assert.equal(state?.cursor, "1/CRASH-P1"); assert.match(state?.lastError ?? "", /404 Not Found/); }); test("an auth failure that read nothing keeps the resume point it had, and invents none", async () => { const run = async () => { process.env.FAKE_MODE = "authfail"; try { return await fetchPosts({ paths: (await import("../lib/paths")).getPaths(), slug: SLUG, settings: {}, }); } finally { delete process.env.FAKE_MODE; } }; // The crash above left a resume point: an auth failure must not drop it. const kept = await run(); assert.equal(kept.needsCookies, true); assert.equal((await readPostFetchState(channelRoot))?.cursor, "1/CRASH-P1"); // From a caught-up state it must not invent one either: the watermark stands. await writeFile(path.join(channelRoot, "posts-state.json"), "{}"); const fresh = await run(); assert.equal(fresh.ok, true); assert.equal(fresh.written, 0); const state = await readPostFetchState(channelRoot); assert.equal(state?.needsCookies, true); assert.equal(state?.cursor, undefined); }); async function makeChannel(slug: string): Promise { const root = path.join(ROOT, "transcripts", "channels", slug); await mkdir(root, { recursive: true }); await writeFile( path.join(root, "config.json"), JSON.stringify({ handling: "transcribe", sourceKind: "social", postFetcher: "x-gallery-dl", socialHandle: "faketester", platform: "twitter", name: "Fake (X)", url: "https://x.com/faketester", }), ); return root; } test("a drain mid-page waits for the page's cursor, stops there and keeps it; the run is not a failure", async () => { const slug = "fake-x-drain-page"; const root = await makeChannel(slug); const drain = new AbortController(); const log: string[] = []; process.env.FAKE_MODE = "pages"; let result; try { result = await fetchPosts({ paths: (await import("../lib/paths")).getPaths(), slug, settings: {}, drain: drain.signal, onLog: (line) => { log.push(line); // Page 1's records are in; its cursor is not yet. if (/page 1 read/.test(line)) drain.abort(); }, }); } finally { delete process.env.FAKE_MODE; } assert.equal(result.ok, true, log.join("\n")); assert.equal(result.drained, true); assert.equal(result.complete, false); // Pages 0 and 1, and not one more. assert.equal(result.written, 6); assert.equal((await readSeenPostIds(root)).size, 6); const state = await readPostFetchState(root); assert.equal(state?.cursor, "1/PAGES-P1"); assert.equal(state?.lastError, undefined); assert.equal(log.some((l) => /page 2 read/.test(l)), false, log.join("\n")); assert.ok(log.some((l) => /Drained; the next run resumes from the last page read/.test(l)), log.join("\n")); // The next run resumes there (cancelled once it has started: only its argv // is wanted). const cancel = new AbortController(); await fetchPosts({ paths: (await import("../lib/paths")).getPaths(), slug, settings: {}, signal: cancel.signal, onLog: (line) => { if (/^Running /.test(line)) setTimeout(() => cancel.abort(), 200); }, }); assert.ok((await lastArgv()).includes("extractor.twitter.cursor=1/PAGES-P1")); }); test("a drain while gallery-dl waits out a rate limit stops it at once, keeping the last page's cursor", async () => { const slug = "fake-x-drain-wait"; const root = await makeChannel(slug); const drain = new AbortController(); const log: string[] = []; const started = Date.now(); const result = await fetchPosts({ paths: (await import("../lib/paths")).getPaths(), slug, settings: {}, drain: drain.signal, onLog: (line) => { log.push(line); if (/rate limit/.test(line)) drain.abort(); }, }); assert.equal(result.ok, true, log.join("\n")); assert.equal(result.drained, true); assert.equal(result.written, 6); // Not left asleep for its 60 s. assert.ok(Date.now() - started < 15_000, `took ${Date.now() - started} ms`); assert.equal((await readPostFetchState(root))?.cursor, "1/AFTER-P2"); }); test("a drain before the fetch starts reads nothing and spawns nothing", async () => { const slug = "fake-x-drain-early"; const root = await makeChannel(slug); const drain = new AbortController(); drain.abort(); const before = (await readFile(ARGS_LOG, "utf8")).trim().split("\n").length; const result = await fetchPosts({ paths: (await import("../lib/paths")).getPaths(), slug, settings: {}, drain: drain.signal, }); assert.deepEqual(result, { ok: true, written: 0, skipped: 0, complete: false, drained: true }); assert.equal((await readFile(ARGS_LOG, "utf8")).trim().split("\n").length, before); assert.equal(await readPostFetchState(root), null); });