Archilyzer · Source

archilyzer

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

commit 1cd09317405758b6257d1b732a68390df71502f0
parent 7447eb6b69b0c1a1f959e09ce8986d7ec3dd2d34
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri,  2 Oct 2026 15:20:54 -0400

common: an X post fetch keeps what it read when cancelled, timed out or failed, and resumes from gallery-dl's cursor

gallery-dl's --dump-json printed nothing until a clean exit, so a fetch that
the rate limit held past its 30-minute budget, or that was cancelled, ended
with zero posts saved. The fetcher now runs gallery-dl with output.jsonl and
--verbose and reads both pipes as lines:

- each record is normalized as it arrives (GalleryDlStream, pure and tested);
- gallery-dl's per-page "Cursor:" debug line is the resume point; checkpoints
  every 200 posts or 60 s write the posts and the previous page's cursor
  through a new PostFetchInput.onCheckpoint;
- a timeout, cancel or non-zero exit returns what was read plus the resume
  point (an error is reported as PostFetchResult.error, after the controller
  has saved the posts); the next run passes it back as
  extractor.twitter.cursor;
- an incremental or resumed run stops after 100 archived-or-older posts in a
  row instead of walking the whole timeline (stopAtKnown; off for --full);
- a backfill run's budget is 3 h, an incremental one's stays 30 min;
- non-debug stderr (rate-limit waits) goes to the job log as it happens.

Tests: stream rules (cursor lag, early stop, watermark, jsonl tuples, pretty
array fallback) and fetchPosts end to end against a fake gallery-dl: cancel
keeps posts + cursor, the next run resumes and finishes, early stop kills the
walk, a crash part-way fails the job but keeps its posts, an auth failure
neither drops nor invents a resume point.

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

Diffstat:
Acommon/controller/fetchPosts.test.ts | 199+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/fetchPosts.ts | 53+++++++++++++++++++++++++++++++++++++++++++++++++----
Mcommon/social/fetchers.ts | 15+++++++++++++--
Mcommon/social/xGalleryDlFetcher.test.ts | 125+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/social/xGalleryDlFetcher.ts | 411++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------
Meditor/CHANGELOG.md | 1+
6 files changed, 722 insertions(+), 82 deletions(-)

diff --git a/common/controller/fetchPosts.test.ts b/common/controller/fetchPosts.test.ts @@ -0,0 +1,199 @@ +// 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. +// +// 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. +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 (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<string[]> { + 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); +}); diff --git a/common/controller/fetchPosts.ts b/common/controller/fetchPosts.ts @@ -12,6 +12,7 @@ import { readChannelConfig, } from "./channels"; import type { Paths } from "../lib/paths"; +import type { Post } from "../lib/posts"; import { isSocialChannel } from "../lib/channelConfig"; import { alwaysCookies, @@ -156,6 +157,27 @@ export async function fetchPosts( lastFetchedAt: new Date().toISOString(), }; + // Progress a long run saves as it goes (a fetcher that streams calls this + // between pages): the posts so far, then the resume point they are safe + // under, so an interrupted run — editor restart, crash — continues from it. + let checkpointWritten = 0; + let checkpointSkipped = 0; + const shards = new Set<string>(); + const onCheckpoint = async (cp: { posts: Post[]; cursor: string }) => { + const w = await writePosts(channelRoot, cp.posts); + checkpointWritten += w.written; + checkpointSkipped += w.skipped; + for (const shard of w.shards) shards.add(shard); + await writePostFetchState(channelRoot, { + ...state, + cursor: cp.cursor, + lastFetchedCount: checkpointWritten, + }); + if (w.written > 0) { + log(`Saved ${w.written} new post(s) (${checkpointWritten} so far this run); resume point recorded.`); + } + }; + let result; try { result = await fetcher.fetch({ @@ -169,6 +191,8 @@ export async function fetchPosts( cookieSource: xLogin?.source, browserCookies: xLogin?.browserSpec, limit: opts.limit, + stopAtKnown: !opts.full, + onCheckpoint, signal: effectiveSignal, onLog: log, }); @@ -176,20 +200,40 @@ export async function fetchPosts( const message = (err as Error).message; log(`[error] ${message}`); state.lastError = message; + // Whatever was checkpointed is on disk; keep the resume point that covers + // it rather than dropping back to the watermark. + const prior = await readPostFetchState(channelRoot); + if (prior?.cursor) state.cursor = prior.cursor; + if (checkpointWritten > 0) state.lastFetchedCount = checkpointWritten; await writePostFetchState(channelRoot, state); - return { ok: false, written: 0, skipped: 0, complete: false, error: message }; + return { + ok: false, + written: checkpointWritten, + skipped: checkpointSkipped, + complete: false, + error: message, + }; } - const written = await writePosts(channelRoot, result.posts); + const final = await writePosts(channelRoot, result.posts); + for (const shard of final.shards) shards.add(shard); + const written = { + written: checkpointWritten + final.written, + skipped: checkpointSkipped + final.skipped, + }; log( `Wrote ${written.written} new post(s), skipped ${written.skipped} already archived` + - (written.shards.length > 0 ? ` (shards: ${written.shards.join(", ")})` : "") + + (shards.size > 0 ? ` (shards: ${[...shards].join(", ")})` : "") + (result.complete ? "." : " — more history remains."), ); state.lastFetchedCount = written.written; if (result.cursor && !result.complete) state.cursor = result.cursor; if (result.needsCookies) state.needsCookies = true; + if (result.error) { + log(`[error] ${result.error}`); + state.lastError = result.error; + } await writePostFetchState(channelRoot, state); // Reuse the existing lastSyncedAt field so the scheduler @@ -206,10 +250,11 @@ export async function fetchPosts( }); return { - ok: true, + ok: !result.error, written: written.written, skipped: written.skipped, complete: result.complete, needsCookies: result.needsCookies, + error: result.error, }; } diff --git a/common/social/fetchers.ts b/common/social/fetchers.ts @@ -52,6 +52,14 @@ export type PostFetchInput = { // Soft cap on how many posts to return in one run. Undefined = no cap // beyond the watermark/seen-id stop conditions. limit?: number; + // Stop once the walk reaches posts that are already archived. False for a + // --full repair, which re-walks everything on purpose. Undefined: the + // fetcher's own default. + stopAtKnown?: boolean; + // Save progress during a long run: called with posts not yet handed back and + // the resume point they are safe under. Posts passed here are NOT returned + // again in the result. A fetcher that pages in one shot never calls it. + onCheckpoint?: (checkpoint: { posts: Post[]; cursor: string }) => Promise<void>; signal: AbortSignal; // Progress/diagnostic sink, wired to the job log. onLog?: (line: string) => void; @@ -59,14 +67,17 @@ export type PostFetchInput = { export type PostFetchResult = { posts: Post[]; - // False when the run stopped early (limit hit, or aborted) and more history - // remains behind `cursor`. + // False when the run stopped early (limit hit, budget or cancel, or an error + // part-way) and more history remains behind `cursor`. complete: boolean; cursor?: string; // Set when the fetcher stopped because credentials are missing or expired. // Maps onto the existing "needs cookies" snapshot bucket rather than // inventing a new failure surface. needsCookies?: boolean; + // The run failed part-way. `posts` and `cursor` still hold what it got, so + // the controller saves them before reporting the failure. + error?: string; }; export type SocialFetcherProbe = { diff --git a/common/social/xGalleryDlFetcher.test.ts b/common/social/xGalleryDlFetcher.test.ts @@ -84,3 +84,128 @@ test("no source resolved (a caller from before the choice): the jar, else the al "firefox", ); }); + +// --------------------------------------------------------------------------- +// Streaming, resume and early stop (gallery-dl output read as it runs) +// --------------------------------------------------------------------------- + +import { + GalleryDlStream, + RESTART_FROM_TOP_CURSOR, + STOP_AFTER_KNOWN, +} from "./xGalleryDlFetcher"; + +function tweet(n: number, date = "2026-06-01 12:00:00") { + // Ids above 2^53, as X's are, so a numeric parse would show. + const id = String(1899000000000000000n + BigInt(n)); + return { + tweet_id: id, + conversation_id: id, + date, + content: `fake tweet ${n}`, + author: { name: "faketester", nick: "Fake" }, + user: { name: "faketester", nick: "Fake" }, + }; +} +// gallery-dl's jsonl: a [type, kwdict] (or [type, url, kwdict]) tuple per line, +// with the id as a bare (unquoted) number. +const jsonl = (n: number, date?: string) => + JSON.stringify([2, tweet(n, date)]).replace(/"tweet_id":"(\d+)"/, '"tweet_id":$1'); + +function newStream(opts: Partial<ConstructorParameters<typeof GalleryDlStream>[0]> = {}) { + return new GalleryDlStream({ + channelSlug: "fake-x", + seenIds: new Set(), + stopAtKnown: true, + ...opts, + }); +} + +test("streaming argv: jsonl always; --verbose and the cursor only when asked", () => { + const probe = buildGalleryDlArgs({ accountUrl: ACCOUNT, limit: 1 }); + assert.equal(probe.includes("output.jsonl=true"), true); + assert.equal(probe.includes("--verbose"), false); + assert.equal(probe.some((a) => a.startsWith("extractor.twitter.cursor=")), false); + + const run = buildGalleryDlArgs({ accountUrl: ACCOUNT, cursor: "2_123/DAAB", stream: true }); + assert.equal(run.includes("--verbose"), true); + assert.equal(flagValue(run, "-o") !== undefined, true); + assert.equal(run.includes("extractor.twitter.cursor=2_123/DAAB"), true); + assert.equal(run[run.length - 1], "https://x.com/someaccount/timeline"); +}); + +test("stream: records parse from jsonl tuples, keep big ids exact, count each tweet once", () => { + const s = newStream(); + s.pushStdout(jsonl(1)); + s.pushStdout(JSON.stringify([3, "https://pbs.twimg.com/x.jpg", tweet(1)])); // same tweet, media record + s.pushStdout(jsonl(2)); + const posts = s.finish(); + assert.equal(s.distinct, 2); + assert.equal(posts.length, 2); + // The id survives exactly (a JSON.parse of the bare number would round it). + assert.ok(posts.some((p) => p.id.includes("1899000000000000001")), posts.map((p) => p.id).join(",")); +}); + +test("stream: a checkpoint carries the PREVIOUS cursor; the final cursor is the latest", () => { + const s = newStream(); + s.pushStdout(jsonl(1)); + s.pushStderr("[twitter][debug] Cursor: 1/AAA"); + // Only one cursor so far: nothing is provably covered yet. + assert.equal(s.takeCheckpoint(), null); + s.pushStdout(jsonl(2)); + s.pushStderr("[twitter][debug] Cursor: 1/BBB"); + const cp = s.takeCheckpoint(); + assert.equal(cp?.cursor, "1/AAA"); + assert.equal(cp?.posts.length, 2); + // Nothing new since: no second checkpoint. + assert.equal(s.takeCheckpoint(), null); + assert.equal(s.cursor, "1/BBB"); +}); + +test("stream: debug noise is dropped; info lines are logged and rate limits noticed", () => { + const logged: string[] = []; + const s = newStream({ onLog: (l) => logged.push(l) }); + s.pushStderr("[urllib3.connectionpool][debug] https://x.com:443 \"GET /i/api HTTP/1.1\" 200"); + s.pushStderr("[twitter][info] Waiting for 14 minutes until 15:11:22 (rate limit)"); + assert.equal(s.rateLimited, true); + assert.deepEqual(logged, ["gallery-dl: [twitter][info] Waiting for 14 minutes until 15:11:22 (rate limit)"]); + assert.match(s.stderrTail(), /rate limit/); +}); + +test("stream: stops after STOP_AFTER_KNOWN archived-or-older posts IN A ROW, not before", () => { + // Archive ids 0..2*STOP_AFTER_KNOWN+10 (as normalized post ids). + const archived = newStream(); + for (let n = 0; n <= 2 * STOP_AFTER_KNOWN + 10; n++) archived.pushStdout(jsonl(n)); + const seen = new Set(archived.finish().map((p) => p.id)); + + const s = newStream({ seenIds: seen }); + s.pushStdout(jsonl(9000)); // new + for (let n = 0; n < STOP_AFTER_KNOWN - 1; n++) s.pushStdout(jsonl(n)); // 99 known + assert.equal(s.stopReached, false); + s.pushStdout(jsonl(9001)); // a new post resets the run + for (let n = STOP_AFTER_KNOWN; n < 2 * STOP_AFTER_KNOWN - 1; n++) s.pushStdout(jsonl(n)); // 99 known + assert.equal(s.stopReached, false); + s.pushStdout(jsonl(2 * STOP_AFTER_KNOWN)); // the 100th in a row + assert.equal(s.stopReached, true); + assert.equal(s.finish().length, 2); +}); + +test("stream: posts older than the watermark are not new; a --full run never stops early", () => { + const s = newStream({ since: "2026-05-01T00:00:00.000Z", stopAtKnown: false }); + s.pushStdout(jsonl(1, "2026-06-01 12:00:00")); + for (let n = 2; n < STOP_AFTER_KNOWN + 10; n++) s.pushStdout(jsonl(n, "2026-04-01 12:00:00")); + assert.equal(s.stopReached, false); + assert.equal(s.older, STOP_AFTER_KNOWN + 8); + assert.equal(s.finish().length, 1); +}); + +test("stream: a gallery-dl that ignores output.jsonl (one pretty array at exit) still parses", () => { + const s = newStream(); + const pretty = JSON.stringify([[2, tweet(1)], [2, tweet(2)]], null, 2); + for (const line of pretty.split("\n")) s.pushStdout(line); + assert.equal(s.finish().length, 2); +}); + +test("restart-from-top cursor is the timeline extractor's first state", () => { + assert.equal(RESTART_FROM_TOP_CURSOR, "1/"); +}); diff --git a/common/social/xGalleryDlFetcher.ts b/common/social/xGalleryDlFetcher.ts @@ -24,10 +24,12 @@ // `--dump-json` contract and are covered by a fixture binary in the editor e2e // suite; treat the first real run as the spike. +import { createInterface } from "node:readline"; +import type { Readable } from "node:stream"; import { execa } from "execa"; import { getPaths } from "../lib/paths"; import type { Post } from "../lib/posts"; -import { normalizeXTweets, type XTweetRaw } from "./xNormalize"; +import { normalizeXTweet, type XTweetRaw } from "./xNormalize"; import { readXSessionStatus, xCookieFile } from "./xSessionBroker"; import type { XCookieSource } from "./xCookieSource"; import { @@ -87,10 +89,23 @@ export function buildGalleryDlArgs(opts: { // the session fresh, which is the fix for X cookies expiring in days. cookieFile?: string; limit?: number; + // gallery-dl's own resume point (`extractor.twitter.cursor`), as logged by a + // previous streamed run — see GalleryDlStream. + cursor?: string; + // A long fetch, read as it runs: log at debug level so gallery-dl prints its + // "Cursor: …" line after every page. Off for the probe, whose error message + // is the tail of stderr and must not be buried in debug noise. + stream?: boolean; }): string[] { const args = [ // One JSON object per item on stdout — we want metadata, never files. "--dump-json", + // …and print each one AS IT ARRIVES. Without this gallery-dl collects the + // whole walk in memory and prints it only when it exits cleanly, so a run + // that is timed out or cancelled (a long, rate-limited backfill) hands back + // nothing at all. + "-o", + "output.jsonl=true", // Never write anything to disk: v1 archives text + links + thread // structure only, no media. "--no-download", @@ -126,10 +141,177 @@ export function buildGalleryDlArgs(opts: { if (opts.limit && opts.limit > 0) { args.push("--range", `1-${Math.floor(opts.limit)}`); } + if (opts.cursor) { + args.push("-o", `extractor.twitter.cursor=${opts.cursor}`); + } + if (opts.stream) args.push("--verbose"); args.push(timelineUrlFor(opts.accountUrl)); return args; } +// The timeline extractor's cursor is `<state>[_<tweet id>]/<API cursor>`, and +// "1/" is its first state with no API cursor: walk from the top. A run that +// stopped before gallery-dl logged any cursor resumes from here rather than +// falling back to the watermark, which would skip whatever lies between the +// posts it did get and the archive. +export const RESTART_FROM_TOP_CURSOR = "1/"; + +// A watermarked or resumed run stops once this many posts IN A ROW are already +// archived (or older than the watermark). Consecutive, so a pinned old post at +// the top of the timeline does not end the run; well above one page (~20 +// tweets, plus their quotes), so re-reading the page a resumed cursor points at +// does not end it either. +export const STOP_AFTER_KNOWN = 100; + +// How long one run may take. gallery-dl waits out X's rate limit in-process +// ("Waiting for 14 minutes … (rate limit)"), so a backfill spends most of its +// time asleep; with progress saved as it goes, a run that hits its budget +// simply resumes on the next one. +export const INCREMENTAL_BUDGET_MS = 30 * 60_000; +export const BACKFILL_BUDGET_MS = 3 * 60 * 60_000; + +// Write what has arrived every so often, so an editor restart mid-backfill +// loses at most this much. +const CHECKPOINT_EVERY_POSTS = 200; +const CHECKPOINT_EVERY_MS = 60_000; + +const CURSOR_LINE = /\]\[debug\]\s+Cursor:\s+(\S+)\s*$/; +const DEBUG_LINE = /\]\[debug\]/; +const RATE_LIMIT_LINE = /Waiting for .*rate limit|429 Too Many Requests/i; +const STDERR_TAIL_LINES = 20; + +export type GalleryDlCheckpoint = { posts: Post[]; cursor: string }; + +// One streamed gallery-dl run, fed line by line: stdout carries a JSON record +// per line (output.jsonl), stderr carries the log — including, with --verbose, +// "[twitter][debug] Cursor: <value>" after each page. Pure (no I/O), so the +// resume rules are testable without spawning anything. +// +// Why a checkpoint saves the PREVIOUS cursor: gallery-dl flushes a page's +// records to stdout before it logs that page's cursor on stderr, but they are +// two pipes, and this side may read the cursor first. The cursor one page back +// is certainly covered by what has been read, so a crash re-reads at most a +// page (deduped). Once the process has exited and both pipes are drained, the +// latest cursor is safe — that is `cursor`. +export class GalleryDlStream { + // Every distinct post seen this run, new or not. + distinct = 0; + newPosts = 0; + known = 0; + older = 0; + rateLimited = false; + stopReached = false; + cursor?: string; + private previousCursor?: string; + private savedCursor?: string; + private pending: Post[] = []; + private ids = new Set<string>(); + private knownRun = 0; + private tail: string[] = []; + private unparsed: string[] = []; + + constructor( + private readonly opts: { + channelSlug: string; + seenIds: ReadonlySet<string>; + since?: string; + stopAtKnown: boolean; + onLog?: (line: string) => void; + }, + ) {} + + pushStdout(line: string): void { + if (!line.trim()) return; + const records = parseGalleryDlOutput(line); + if (records.length === 0) { + // Not a JSON line: a gallery-dl that ignored output.jsonl and prints one + // pretty array at exit. Parsed whole in finish(). + this.unparsed.push(line); + return; + } + for (const rec of records) this.addRecord(rec); + } + + pushStderr(line: string): void { + const cursor = CURSOR_LINE.exec(line); + if (cursor) { + this.previousCursor = this.cursor; + this.cursor = cursor[1]; + return; + } + if (DEBUG_LINE.test(line) || !line.trim()) return; + if (RATE_LIMIT_LINE.test(line)) this.rateLimited = true; + this.tail.push(line); + if (this.tail.length > STDERR_TAIL_LINES) this.tail.shift(); + this.opts.onLog?.(`gallery-dl: ${line.trim()}`); + } + + // The last few non-debug stderr lines, for an error message. + stderrTail(lines = 5): string { + return this.tail.slice(-lines).join("\n"); + } + + pendingCount(): number { + return this.pending.length; + } + + // Posts to write now, with the resume point they are safe under — or null + // when there is nothing new to record. + takeCheckpoint(): GalleryDlCheckpoint | null { + const cursor = this.previousCursor; + if (!cursor) return null; + if (this.pending.length === 0 && cursor === this.savedCursor) return null; + this.savedCursor = cursor; + const posts = this.pending; + this.pending = []; + return { posts, cursor }; + } + + // A checkpoint that could not be written goes back into the queue, so the + // final write retries it. + requeue(posts: Post[]): void { + this.pending = [...posts, ...this.pending]; + } + + // Everything not yet checkpointed. Call after the process has exited and both + // pipes have closed. + finish(): Post[] { + if (this.distinct === 0 && this.unparsed.length > 0) { + for (const rec of parseGalleryDlOutput(this.unparsed.join("\n"))) { + this.addRecord(rec); + } + } + this.unparsed = []; + const posts = this.pending; + this.pending = []; + return posts; + } + + private addRecord(rec: XTweetRaw): void { + const post = normalizeXTweet(rec, this.opts.channelSlug); + // gallery-dl can emit a tweet more than once (a directory record, then one + // per media URL); count it once. + if (!post || this.ids.has(post.id)) return; + this.ids.add(post.id); + this.distinct++; + const isKnown = this.opts.seenIds.has(post.id); + const isOlder = + !isKnown && Boolean(this.opts.since && post.createdAt <= this.opts.since); + if (isKnown || isOlder) { + if (isKnown) this.known++; + else this.older++; + this.knownRun++; + if (this.opts.stopAtKnown && this.knownRun >= STOP_AFTER_KNOWN) { + this.stopReached = true; + } + return; + } + this.knownRun = 0; + this.newPosts++; + this.pending.push(post); + } +} + // WHICH LOGIN gallery-dl is handed this run (release 16 slice XL), pure so the // argv per source is testable: // "browser" — `--cookies-from-browser <spec>`, always (gallery-dl reads the @@ -234,6 +416,19 @@ function flattenRecords(items: unknown[]): XTweetRaw[] { return out; } +// Feed a pipe to `onLine` a line at a time; resolves once it has closed. +function readLines( + input: Readable | null | undefined, + onLine: (line: string) => void, +): Promise<void> { + if (!input) return Promise.resolve(); + return new Promise((resolve) => { + const rl = createInterface({ input, crlfDelay: Infinity }); + rl.on("line", onLine); + rl.on("close", resolve); + }); +} + export const xGalleryDlFetcher: SocialFetcher = { id: "x-gallery-dl", label: "X / Twitter (gallery-dl)", @@ -320,56 +515,105 @@ export const xGalleryDlFetcher: SocialFetcher = { }); onLog?.(choice.note); const { cookies, cookieFile } = choice; - const args = buildGalleryDlArgs({ accountUrl, cookies, cookieFile, limit }); - onLog?.(`Running ${bin} for ${accountUrl}`); + const args = buildGalleryDlArgs({ + accountUrl, + cookies, + cookieFile, + limit, + cursor: input.cursor, + stream: true, + }); + // An incremental run reads only what is new; a backfill (or its resume) + // walks history and waits out rate limits as it goes. + const budgetMs = since ? INCREMENTAL_BUDGET_MS : BACKFILL_BUDGET_MS; + onLog?.( + `Running ${bin} for ${accountUrl}` + + (input.cursor ? ` from the saved resume point` : "") + + ` (up to ${Math.round(budgetMs / 60_000)} min this run; progress is saved as it goes)`, + ); - let stdout = ""; - let stderr = ""; - let exitCode: number | null = null; - try { - const res = await execa(bin, args, { - reject: false, - // gallery-dl walks the timeline newest-first; a long backfill is - // legitimately slow, so the cap is generous. - timeout: 30 * 60_000, - cancelSignal: signal, + const stream = new GalleryDlStream({ + channelSlug, + seenIds, + since, + // A --full repair re-walks everything on purpose. + stopAtKnown: input.stopAtKnown ?? Boolean(since), + onLog, + }); + + const subprocess = execa(bin, args, { + reject: false, + // Records are consumed as lines below; buffering a whole backfill's + // stdout would hold all of it in memory for nothing. + buffer: false, + stdin: "ignore", + timeout: budgetMs, + cancelSignal: signal, + // gallery-dl flushes each record itself, but keep Python from block- + // buffering anything else on the pipe. + env: { PYTHONUNBUFFERED: "1" }, + }); + + // Checkpoints run one at a time, in order, while the walk continues. + let checkpoints: Promise<void> = Promise.resolve(); + let checkpointsOn = Boolean(input.onCheckpoint); + let lastCheckpointAt = Date.now(); + const maybeCheckpoint = () => { + if (!checkpointsOn || !input.onCheckpoint) return; + const due = + stream.pendingCount() >= CHECKPOINT_EVERY_POSTS || + Date.now() - lastCheckpointAt >= CHECKPOINT_EVERY_MS; + if (!due) return; + const cp = stream.takeCheckpoint(); + if (!cp) return; + lastCheckpointAt = Date.now(); + const save = input.onCheckpoint; + checkpoints = checkpoints.then(async () => { + if (!checkpointsOn) { + stream.requeue(cp.posts); + return; + } + try { + await save(cp); + } catch (err) { + // Keep the posts for the final write and stop checkpointing, so a + // saved resume point never runs ahead of what is on disk. + checkpointsOn = false; + stream.requeue(cp.posts); + onLog?.(`[warn] could not save progress mid-run: ${(err as Error).message}`); + } }); - stdout = res.stdout ?? ""; - stderr = res.stderr ?? ""; - exitCode = res.exitCode ?? null; - } catch (err) { - const message = (err as Error).message; - if (/ENOENT/.test(message)) { - throw new Error( - `gallery-dl not found (looked for "${bin}"; set GALLERY_DL_BIN)`, + }; + + let stoppedAtKnown = false; + const stdoutDone = readLines(subprocess.stdout, (line) => { + stream.pushStdout(line); + if (stream.stopReached && !stoppedAtKnown) { + stoppedAtKnown = true; + onLog?.( + `Reached ${STOP_AFTER_KNOWN} already-archived posts in a row — caught up; stopping gallery-dl.`, ); + subprocess.kill(); } - throw err; - } + }); + const stderrDone = readLines(subprocess.stderr, (line) => { + stream.pushStderr(line); + maybeCheckpoint(); + }); - if (exitCode !== 0) { - const tail = stderr.trim().split("\n").slice(-5).join("\n"); - if (looksLikeAuthFailure(tail)) { - onLog?.(`[auth] gallery-dl could not authenticate to X:\n${tail}`); - // Not a crash — a known, recoverable state the Needs-cookies bucket - // already models. Keep whatever it did manage to emit. - return { - posts: normalizeXTweets(parseGalleryDlOutput(stdout), channelSlug).filter( - (p) => !seenIds.has(p.id), - ), - complete: false, - needsCookies: true, - }; - } - throw new Error(`gallery-dl exited ${exitCode}: ${tail || "(no output)"}`); + const res = await subprocess; + await Promise.all([stdoutDone, stderrDone]); + await checkpoints; + + const spawnError = `${res.code ?? ""} ${res.message ?? ""}`; + if (res.failed && /ENOENT/.test(spawnError) && stream.distinct === 0) { + throw new Error( + `gallery-dl not found (looked for "${bin}"; set GALLERY_DL_BIN)`, + ); } - // gallery-dl handles X rate limits by BLOCKING ("Waiting for 7 minutes - // until ... (rate limit)"), which otherwise looks like a hung job. Surface - // it so the log explains the stall. Guest (cookie-less) runs hit this far - // sooner than authenticated ones. - const rateLimited = /Waiting for .*rate limit|429 Too Many Requests/i.test(stderr); - if (rateLimited) { + const posts = stream.finish(); + if (stream.rateLimited) { onLog?.( "[rate-limit] X throttled this run — gallery-dl paused between pages. " + (cookieFile || cookies @@ -377,42 +621,57 @@ export const xGalleryDlFetcher: SocialFetcher = { : "This is a guest (no-credential) run; guest quota is much smaller."), ); } + onLog?.( + `gallery-dl read ${stream.distinct} post(s): ${stream.newPosts} new` + + (stream.known ? `, ${stream.known} already archived` : "") + + (stream.older ? `, ${stream.older} older than the watermark` : ""), + ); - const records = parseGalleryDlOutput(stdout); - onLog?.(`gallery-dl returned ${records.length} record(s)`); - const all = normalizeXTweets(records, channelSlug); - - // Apply the same stop conditions the Bluesky fetcher uses. gallery-dl walks - // the whole timeline itself (there is no cursor to hand back), so we filter - // its output rather than stopping mid-stream. - const posts: Post[] = []; - let known = 0; - let older = 0; - for (const post of all) { - if (seenIds.has(post.id)) { - known++; - continue; - } - if (since && post.createdAt <= since) { - older++; - continue; + // Where the next run should pick up when this one did not finish: the last + // page gallery-dl completed, else the resume point it was given. A run that + // saved new posts before any page completed restarts from the top (never + // the watermark — see RESTART_FROM_TOP_CURSOR); one that read nothing new + // leaves the watermark as it was. + const timeline = /\/timeline$/.test(timelineUrlFor(accountUrl)); + const resumeAt = + stream.cursor ?? + input.cursor ?? + (timeline && stream.newPosts > 0 ? RESTART_FROM_TOP_CURSOR : undefined); + + if (stoppedAtKnown) return { posts, complete: true }; + + if (res.isCanceled || signal.aborted || res.timedOut) { + onLog?.( + res.timedOut + ? `Stopped at this run's ${Math.round(budgetMs / 60_000)}-minute budget; the next run resumes where this one left off.` + : "Cancelled; the next run resumes where this one left off.", + ); + return { posts, complete: false, cursor: resumeAt }; + } + + if (res.exitCode !== 0) { + const tail = stream.stderrTail(); + if (looksLikeAuthFailure(tail)) { + onLog?.(`[auth] gallery-dl could not authenticate to X:\n${tail}`); + // Not a crash — a known, recoverable state the Needs-cookies bucket + // already models. Keep whatever it did manage to emit. + return { posts, complete: false, cursor: resumeAt, needsCookies: true }; } - posts.push(post); + return { + posts, + complete: false, + cursor: resumeAt, + error: `gallery-dl exited ${res.exitCode ?? res.signal ?? "abnormally"}: ${tail || "(no output)"}`, + }; } - onLog?.( - `${posts.length} new post(s)` + - (known ? `, ${known} already archived` : "") + - (older ? `, ${older} older than the watermark` : ""), - ); - // gallery-dl has no resumable cursor: it either walked the timeline it was - // given or it didn't. `complete` is therefore true whenever it exited - // cleanly — except when a --range cap may have cut it short. - // A capped OR rate-limited run stopped early, so it must not report - // `complete` — that would let the controller advance the watermark past - // posts it never saw. - const capped = Boolean(limit && all.length >= limit); - return { posts, complete: !capped && !rateLimited }; + // A clean exit walked everything it was given: gallery-dl waits out rate + // limits rather than skipping pages. Only a --range cap leaves history + // behind — resume from the last page it finished. + const capped = Boolean(limit && stream.distinct >= limit); + return capped + ? { posts, complete: false, cursor: resumeAt } + : { posts, complete: true }; }, }; diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **An X post fetch keeps what it has read when it is cancelled, times out or fails part-way, and the next fetch picks up where it stopped.** The gallery-dl fetcher used to receive an account's posts all at once when gallery-dl finished, so a long fetch that X's rate limit held past the 30-minute limit — or one you cancelled — ended with nothing saved. gallery-dl now hands over each post as it reads it: the fetch saves posts every 200 posts or every minute along with gallery-dl's own resume point, and the next fetch of the channel continues from that point instead of starting again. A fetch of new posts stops once it reaches 100 already-archived posts in a row instead of reading the whole timeline, and a fetch of an account's history may run for up to 3 hours (new-post fetches keep the 30-minute limit). gallery-dl's rate-limit waits now appear in the job's log as they happen. - **umtool can keep each report's render folder on a media drive.** With `UMTOOL_MEDIA_DIR` set, in umtool's environment (restart umtool after setting it), to a directory inside that drive, a report project's `out/` (its fetched windows, segments and finished video) is a link to the same path under that directory: a project's first build makes it there, and `umtool storage move-out <project>` (or `--all`) moves an existing one, copying it, checking the copy and only then leaving the link; `--dry-run` says how much would move, and `umtool storage move-back` brings one home. The manifest, its revisions, notes and sources stay where they are, and nothing in umtool reads a project differently. When the drive is not mounted, a build or source check refuses and says so instead of starting a new folder on the main disk; umtool never creates the media directory itself. `umtool storage` lists where each project's `out/` is. With `UMTOOL_MEDIA_DIR` unset nothing changes. - **umtool can move a report's cut clips and share batches to the media drive, one project at a time.** The **Move deliverables** button in a report's deliver panel, or `umtool storage deliverables <project> --to media` (`--to local` to bring them back, `--dry-run` to see what would move), moves the project's `clips/` folder and every `share-…` batch folder to the same path under `UMTOOL_MEDIA_DIR`, checks each copy, leaves a link in its place, and then records the choice in the report's `video.manifest.json` as `"storage": { "deliverables": "media" }`. From then on a cut or a new batch is made on the media drive, and everything that opens `clips/<id>.mp4` keeps working through the link. The panel shows where the deliverables are, and the button is disabled, with the reason beside it, while a cut or a batch is running, while the drive is not mounted, or (to the media drive) when `UMTOOL_MEDIA_DIR` is not set. When the drive is not mounted, a cut or a batch refuses and says so instead of starting again on the main disk, and `umtool check` reports it as blocking, for the render folder too. Nothing moves until you ask: without the switch, a report's deliverables stay in the project. - **umtool's cache moves to `~/.cache/archilyzer/umtool`** (`$XDG_CACHE_HOME/archilyzer/umtool` when that is set, or `UMTOOL_CACHE_DIR`). It was inside the song project's data folder, so it followed that folder onto whatever drive it was on. Run `umtool index` once after updating to rebuild the project index in its new place; umtool works without it, only slower, and the rest of the cache is remade as it is needed. `umtool doctor` now also shows the reports, media and cache folders, and the old cache folder while it is still there; it can be deleted.