// The PRIMARY X/Twitter post fetcher: gallery-dl driven as a light headless // subprocess. // // Why gallery-dl leads: it is exactly the shape this codebase already manages // well for yt-dlp and whisper — spawn a binary, stream its stdout, parse // records — so it reuses runManagedCommand's sibling (execa here, since we need // the stdout payload rather than a job log), the job registry, and // cookiePolicy.ts essentially unchanged. No browser process per fetch, no RAM // overhead, and it runs anywhere including containers. // // Accepted weaknesses, both modelled rather than hidden: // - X cookies expire in days. `cookieMode: "defer"` and the existing // needsCookies snapshot bucket already describe "blocked pending // credentials", so an auth failure sets `needsCookies` on the result rather // than inventing a new failure surface. The Playwright session broker // (xSessionBroker.ts) is the fix for the expiry itself. // - gallery-dl hardcodes X's GraphQL query IDs, so it breaks when X rotates // them (historically every 2–4 weeks) until the upstream extractor updates. // That is what the x-playwright fallback exists for. // // NOT VERIFIED AGAINST LIVE X. gallery-dl is not installed in this environment // and a live spike needs a throwaway account plus real cookies. The argv, // parsing and failure mapping below are written to gallery-dl's documented // `--dump-json` contract and are covered by a fixture binary in the editor e2e // suite; treat the first real run as the spike. import { existsSync } from "node:fs"; import path from "node:path"; 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 type { OlderBackfillPosition } from "../lib/posts-server"; import { normalizeXTweet, xCreatedAt, type XTweetRaw } from "./xNormalize"; import { olderFloorDay, stepOlderWindow } from "./olderBackfill"; import { launchXProfile, readXSessionStatus, xCookieFile, } from "./xSessionBroker"; import type { XCookieSource } from "./xCookieSource"; import { listCapturedMediaFiles } from "./postCapture"; import { xArticleLinkFromArchive } from "./xArticle"; import { captureXPosts, xStatusUrl, type MediaDownloadResult, } from "./xPostCapture"; import { registerSocialFetcher, type OlderPostFetchInput, type OlderPostFetchResult, type PostCaptureInput, type PostCaptureResult, type PostFetchInput, type PostFetchResult, type SocialFetcher, type SocialFetcherProbe, } from "./fetchers"; // Substrings in gallery-dl's stderr that mean "we are not authenticated" (or no // longer are). Mapped to needsCookies so the channel lands in the existing // Needs-cookies bucket instead of looking like a generic crash. const AUTH_ERROR_PATTERNS = [ "401", "403", "authorization", "unauthorized", "login required", "requires authentication", "no username", "bad credentials", "could not log in", ]; function looksLikeAuthFailure(text: string): boolean { const lower = text.toLowerCase(); return AUTH_ERROR_PATTERNS.some((p) => lower.includes(p)); } // The argv for a text-only timeline read. Exported so a test can assert the // flags without spawning anything. // gallery-dl's `twitter:user` extractor treats a bare profile URL as a // DISPATCH entry: with --dump-json it emits a single type-6 record pointing at // `//timeline` and yields no tweets at all (verified against // gallery-dl 1.32.9 — a real channel synced with lastFetchedCount: 0 because // of this). Targeting the timeline sub-extractor directly is what actually // returns posts, so normalize here rather than making every caller remember. export function timelineUrlFor(accountUrl: string): string { const trimmed = accountUrl.trim().replace(/\/+$/, ""); // Already a sub-route (/timeline, /with_replies, /media, /status/...) — leave it. try { const u = new URL(trimmed); const segments = u.pathname.split("/").filter(Boolean); if (segments.length !== 1) return trimmed; return `${u.origin}/${segments[0]}/timeline`; } catch { return trimmed; } } type GalleryDlLoginAndCap = { cookies?: string; // A cookies.txt exported by the Playwright session broker // (xSessionBroker.ts). PREFERRED over `cookies`: a live browser profile keeps // the session fresh, which is the fix for X cookies expiring in days. cookieFile?: string; limit?: number; }; export function buildGalleryDlArgs( opts: GalleryDlLoginAndCap & { accountUrl: string; // 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 = galleryDlCommonArgs(opts); 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 random gap gallery-dl waits before each X request, in seconds ("min-max"). export const X_SLEEP_REQUEST = "4.0-10.0"; // The flags every gallery-dl read shares: metadata only, streamed, text tweets // with their retweets and replies, the login, and the --post-range cap. function galleryDlCommonArgs(opts: GalleryDlLoginAndCap): 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", // Text-only tweets are skipped by default (gallery-dl is a media // downloader); this is the flag that makes a timeline text archive work. "-o", "extractor.twitter.text-tweets=true", // Retweets and replies are part of the archived record (they carry the // thread structure), so keep both. "-o", "extractor.twitter.retweets=true", "-o", "extractor.twitter.replies=true", // Media is counted, not archived. "-o", "extractor.twitter.videos=false", "-o", "extractor.twitter.cards=false", ...galleryDlPacingArgs(), ...galleryDlLoginArgs(opts), ]; // A cap on POSTS. `--range` counts files, and with downloads off most tweets // yield none, so it never stopped the walk: a 400-post fetch of a new // channel read 3,803. `--post-range` counts the tweets themselves. if (opts.limit && opts.limit > 0) { args.push("--post-range", `1-${Math.floor(opts.limit)}`); } return args; } // PACING. gallery-dl's default gap between X API requests is 0: it pages as // fast as X answers and only slows down when X's rate-limit headers say to. An // account-history walk is hundreds of requests, so every read waits a RANDOM // 4–10 s before each one (gallery-dl's own "a-b" range, a uniform draw per // request): no fixed rhythm, and well under the rate a person scrolling would // make. Running into the limit anyway is waited out, never pushed through // ("wait" is gallery-dl's default; stated so it stays). Every gallery-dl run // against X carries these — the reads and the post capture alike. function galleryDlPacingArgs(): string[] { return [ "-o", `extractor.twitter.sleep-request=${X_SLEEP_REQUEST}`, "-o", "extractor.twitter.ratelimit=wait", ]; } // Cookies are OPTIONAL. gallery-dl reads X timelines on a guest token with no // account at all (verified: 500 tweets over ~7.5 months for a public account). // Guest access is rate-limited far more aggressively than an authenticated // session, though — gallery-dl will block for minutes on "Waiting for N // minutes (rate limit)" — so credentials remain the path for a deep backfill. // When none are configured we simply run as a guest rather than failing. function galleryDlLoginArgs(opts: GalleryDlLoginAndCap): string[] { if (opts.cookieFile) return ["--cookies", opts.cookieFile]; // Same browser-cookie spec yt-dlp uses (cookiePolicy.ts), same syntax. if (opts.cookies) return ["--cookies-from-browser", opts.cookies]; return []; } // The random gap a post capture waits before each contact with X after its // first (a page load, a gallery-dl run): the same 4–10 s draw as // X_SLEEP_REQUEST, which paces gallery-dl's requests inside one run. export function xRequestPauseMs(rand: () => number = Math.random): number { const [min, max] = X_SLEEP_REQUEST.split("-").map((s) => Number(s) * 1000); return Math.round(min + rand() * (max - min)); } // --------------------------------------------------------------------------- // Post capture: one post's attached media // --------------------------------------------------------------------------- // // The media half of `captureByIds` (the screenshot half is xPostCapture.ts). // One gallery-dl run per post, against the post's own URL, DOWNLOADING — the // one gallery-dl run here without --no-download — into the post's capture // directory, videos included, paced and logged in exactly as the reads are. // gallery-dl prints a tagged line per file (downloaded, or already there), so // each file's source URL is known without a second pass. export const CAPTURE_PRINT_TAG = "ARCHILYZER-CAPTURED"; export function buildGalleryDlCaptureArgs( opts: Omit & { id: string; dir: string }, ): string[] { if (!/^\d+$/.test(opts.id)) throw new Error(`"${opts.id}" is not a post id`); return [ // Videos too: a capture keeps what the post carried. "-o", "extractor.twitter.videos=true", "-o", "extractor.twitter.cards=false", ...galleryDlPacingArgs(), // Exactly this directory, and names that cannot meet shot.png or // capture.json. "-D", opts.dir, "-f", "{tweet_id}_{num}.{extension}", // A line per file once it is on disk, and per file already there (a // forced re-capture): "\t\t". "--Print", `after:${CAPTURE_PRINT_TAG}\t{_url}\t{_path}`, "--Print", `skip:${CAPTURE_PRINT_TAG}\t{_url}\t{_path}`, ...galleryDlLoginArgs(opts), xStatusUrl(opts.id), ]; } // The files a capture run printed, by name (each once). export function parseCapturePrints( stdout: string, ): { name: string; url?: string }[] { const byName = new Map(); for (const line of stdout.split("\n")) { const parts = line.replace(/\r$/, "").split("\t"); if (parts[0] !== CAPTURE_PRINT_TAG || parts.length < 3) continue; const file = parts.slice(2).join("\t").trim(); if (!file) continue; const name = path.basename(file); const url = parts[1].trim(); byName.set(name, { name, ...(url && url !== "None" ? { url } : {}) }); } return [...byName.values()]; } // How long one post's download may take: a long video at a polite pace. const CAPTURE_MEDIA_TIMEOUT_MS = 20 * 60_000; export async function downloadXPostMedia(args: { id: string; dir: string; login: { cookies?: string; cookieFile?: string }; signal: AbortSignal; onLog?: (line: string) => void; }): Promise { const bin = getPaths().galleryDlBin; const res = await execa( bin, buildGalleryDlCaptureArgs({ id: args.id, dir: args.dir, ...args.login }), { reject: false, stdin: "ignore", timeout: CAPTURE_MEDIA_TIMEOUT_MS, cancelSignal: args.signal, env: { PYTHONUNBUFFERED: "1" }, }, ); const spawnError = `${res.code ?? ""} ${res.message ?? ""}`; if (res.failed && /ENOENT/.test(spawnError) && !res.stdout) { throw new Error(`gallery-dl not found (looked for "${bin}"; set GALLERY_DL_BIN)`); } const stderr = `${res.stderr ?? ""}`; for (const line of stderr.split("\n")) { if (line.trim()) args.onLog?.(`gallery-dl: ${line.trim()}`); } if (res.isCanceled || res.timedOut || res.exitCode !== 0) { const tail = stderr.trim().split("\n").slice(-5).join("\n"); if (looksLikeAuthFailure(tail)) { return { ok: false, needsCookies: true, error: `gallery-dl could not authenticate to X: ${tail}` }; } return { ok: false, error: res.isCanceled ? "Cancelled during the media download." : res.timedOut ? "The media download ran past its time limit." : `gallery-dl exited ${res.exitCode ?? res.signal ?? "abnormally"}: ${tail || "(no output)"}`, }; } // The printed files, plus any on disk it did not print (an older // gallery-dl): every one is in the record. const printed = parseCapturePrints(`${res.stdout ?? ""}`); const files = printed.filter((f) => existsSync(path.join(args.dir, f.name))); const names = new Set(files.map((f) => f.name)); for (const name of await listCapturedMediaFiles(args.dir)) { if (!names.has(name)) files.push({ name }); } return { ok: true, files }; } // `captureByIds` for X: the shot through the session profile, the media // through gallery-dl, the run paced as one (xPostCapture.ts). export async function captureXPostsByIds( input: PostCaptureInput, ): Promise { const paths = getPaths(); // The profile shoots the posts and opens the articles. A post is known to // link to an article before any page only by its archived text; one found // on a card is found by a shot, which already needs the profile. const articles = (input.articles ?? true) && input.ids.some((id) => xArticleLinkFromArchive(input.archived?.get(id)) !== null); if ((input.shots ?? true) || articles) { const status = await readXSessionStatus(paths); if (!status.hasProfile) { const why = "No X session profile to shoot posts or open articles with: connect an X account on /settings, then run it again."; input.onLog?.(`[auth] ${why}`); return { outcomes: [], needsCookies: true, stoppedEarly: why }; } if (!status.looksAuthenticated) { input.onLog?.( "[warn] The X session profile's exported cookies carry no login; X may show its login wall.", ); } } let login: Promise<{ cookies?: string; cookieFile?: string }> | undefined; return captureXPosts(input, { openPage: async () => { const context = await launchXProfile(paths, { onLog: input.onLog }); const page = context.pages()[0] ?? (await context.newPage()); return { page, close: () => context.close() }; }, downloadMedia: async ({ id, dir, signal, onLog }) => { login ??= galleryDlLoginFor(input); return downloadXPostMedia({ id, dir, login: await login, signal, onLog }); }, pauseMs: () => xRequestPauseMs(), pause, articleImagePauseMs: () => ARTICLE_IMAGE_PAUSE_MS + Math.floor(Math.random() * ARTICLE_IMAGE_PAUSE_MS), }); } // The gap between an article's images: 1–2 s. They are the page's own // pictures, already loaded once by the browser, so not the full post gap. const ARTICLE_IMAGE_PAUSE_MS = 1_000; // --------------------------------------------------------------------------- // Search: the older-posts backfill // --------------------------------------------------------------------------- // // X's timeline only pages back so far; below that, an account's posts are // reachable through search. gallery-dl's `twitter:search` extractor takes a // search URL (`https://x.com/search?q=`), reads the query's `from:` user // for the records' `user`, and pages with `max_id:` (its default // `search-pagination`): after each page it rewrites the query's `max_id:` to // just below the oldest tweet it has read. That pagination logs no cursor, so // this walk's resume point is its own: the oldest post id read in the window, // passed back as `max_id:` — which gallery-dl's rewrite then replaces, rather // than adding a second one. // // Search needs a logged-in session; a guest run is refused before spawning. // An X handle as search's `from:` takes it. Anything else (a space, a quote, a // colon) would change the query's meaning rather than name an account. const X_HANDLE_RE = /^[A-Za-z0-9_]{1,50}$/; export function xSearchHandle(handle: string): string { const bare = handle.trim().replace(/^@/, ""); if (!X_HANDLE_RE.test(bare)) { throw new Error( `"${handle}" is not an X handle search can take (letters, digits and "_")`, ); } return bare; } // The query for one window. `include:nativeretweets` keeps the account's // reposts, as the timeline walk does (search leaves them out by default). export function xSearchQuery( handle: string, position: Pick, ): string { const terms = [ `from:${xSearchHandle(handle)}`, `since:${position.since}`, `until:${position.until}`, "include:nativeretweets", ]; if (position.maxId) { if (!/^\d+$/.test(position.maxId)) { throw new Error(`"${position.maxId}" is not a post id`); } terms.push(`max_id:${position.maxId}`); } return terms.join(" "); } // The search URL gallery-dl's search extractor matches. The query is // percent-encoded whole (a space is %20, a colon %3A): gallery-dl turns "+" // into a space BEFORE it unquotes, so an encoded "+" survives as itself. // `f=live` is X's own "Latest" tab, which is also what gallery-dl asks for. export function xSearchUrl(query: string): string { return `https://x.com/search?q=${encodeURIComponent(query)}&f=live`; } export function buildGalleryDlSearchArgs( opts: GalleryDlLoginAndCap & { handle: string; position: Pick; }, ): string[] { const args = galleryDlCommonArgs(opts); args.push( // Newest first, paged by max_id — the order the resume point relies on. "-o", "extractor.twitter.search-results=latest", "-o", "extractor.twitter.search-pagination=max_id", ); args.push(xSearchUrl(xSearchQuery(opts.handle, opts.position))); return args; } // Inclusive: the post at the resume point is read again (and skipped as a // duplicate) rather than risk an off-by-one past it. function olderId(a: string | undefined, b: string): string { if (a === undefined) return b; try { return BigInt(b) < BigInt(a) ? b : a; } catch { return a; } } // Pause between two windows' searches. Each window is its own gallery-dl run // (a user lookup, then the search pages), so a stretch of empty windows would // otherwise be a burst of requests a second apart. RANDOM, 45–120 s, drawn per // window, for the same reason the per-request gap is: no fixed rhythm. export const OLDER_WINDOW_PAUSE_MIN_MS = 45_000; export const OLDER_WINDOW_PAUSE_MAX_MS = 120_000; export function olderWindowPauseMs(rand: () => number = Math.random): number { return Math.round( OLDER_WINDOW_PAUSE_MIN_MS + rand() * (OLDER_WINDOW_PAUSE_MAX_MS - OLDER_WINDOW_PAUSE_MIN_MS), ); } function pause(ms: number, signal: AbortSignal): Promise { if (ms <= 0 || signal.aborted) return Promise.resolve(); return new Promise((resolve) => { const t = setTimeout(done, ms); function done() { clearTimeout(t); signal.removeEventListener("abort", done); resolve(); } signal.addEventListener("abort", done, { once: true }); }); } // The timeline extractor's cursor is `[_]/`, 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: " 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; // The oldest post id read this run (any post, archived or not) — a search // walk's resume point. oldestId?: string; // The account's creation time, from a record whose `user` is the account. accountCreatedAt?: string; private previousCursor?: string; private savedCursor?: string; private pending: Post[] = []; private ids = new Set(); private knownRun = 0; private tail: string[] = []; private unparsed: string[] = []; constructor( private readonly opts: { channelSlug: string; seenIds: ReadonlySet; since?: string; stopAtKnown: boolean; // The account's handle, to recognise its `user` record. accountHandle?: string; 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); } // Says what the line was when it marks a point between two pages: a page's // cursor, or a wait on X's rate limit. pushStderr(line: string): "cursor" | "rate-limit" | undefined { const cursor = CURSOR_LINE.exec(line); if (cursor) { this.previousCursor = this.cursor; this.cursor = cursor[1]; return "cursor"; } if (DEBUG_LINE.test(line) || !line.trim()) return undefined; const waiting = RATE_LIMIT_LINE.test(line); if (waiting) this.rateLimited = true; this.tail.push(line); if (this.tail.length > STDERR_TAIL_LINES) this.tail.shift(); this.opts.onLog?.(`gallery-dl: ${line.trim()}`); return waiting ? "rate-limit" : undefined; } // 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 }; } // Every post not yet handed over, whatever the cursor — for a walk whose // resume point comes from the records themselves (oldestId), not stderr. takePending(): Post[] { const posts = this.pending; this.pending = []; return posts; } // 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++; if (/^\d+$/.test(post.id)) this.oldestId = olderId(this.oldestId, post.id); if (!this.accountCreatedAt && this.opts.accountHandle) { const user = rec.user as Record | undefined; if ( user && typeof user.name === "string" && user.name.toLowerCase() === this.opts.accountHandle.toLowerCase() && typeof user.date === "string" ) { this.accountCreatedAt = xCreatedAt({ date: user.date }); } } 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 `, always (gallery-dl reads the // operator's browser itself, fresh on every run); never the jar. // With no spec, a guest run that says why. // "profile" — the broker's jar when it holds a login, else the cookieMode // "always" spec, else a guest run. // unset — the same as "profile": a caller that resolves no source keeps // the behaviour from before the choice existed. export function galleryDlCookieChoice(opts: { source?: XCookieSource; // The spec the browser source reads (channel over global, any cookieMode). browserCookies?: string; // The broker's exported jar, when it carries an auth_token. jarFile?: string; // resolveCookiePolicy()'s "always"-mode spec. alwaysCookies?: string; }): { cookies?: string; cookieFile?: string; note: string } { if (opts.source === "browser") { const spec = opts.browserCookies?.trim(); if (spec) { return { cookies: spec, note: `X login: the browser (${spec}), read by gallery-dl.` }; } return { note: "X login: the browser source is chosen but cookiesFromBrowser is empty — " + "running as a guest.", }; } if (opts.jarFile) { return { cookieFile: opts.jarFile, note: "X login: the connected profile (the X session broker's exported cookies).", }; } if (opts.alwaysCookies) { return { cookies: opts.alwaysCookies, note: `X login: no connected profile; the cookieMode "always" browser (${opts.alwaysCookies}).`, }; } return { note: "X login: none (no connected profile) — running as a guest." }; } // X ids are 64-bit and gallery-dl emits them as UNQUOTED JSON NUMBERS // (`"tweet_id": 2085320225776427457`). That exceeds 2^53, so a plain // JSON.parse silently rounds it — 2085320225776427457 becomes // ...427500 — which would corrupt the post id, its slug, its permalink AND // the archive dedupe key that makes re-runs incremental. Verified against // real gallery-dl 1.32.9 output. // // So quote every long integer literal BEFORE parsing. The lookahead leaves the // trailing delimiter in place so adjacent numbers in an array both match, and // requiring a `:`/`[`/`,` prefix means digits inside string values are // untouched (they are already inside quotes). export function quoteBigIntegers(json: string): string { return json.replace(/([:[,]\s*)(-?\d{16,})(?=\s*[,}\]])/g, '$1"$2"'); } // gallery-dl's --dump-json emits either one object per line, or a top-level // array. Tolerate both, plus interleaved non-JSON log noise. export function parseGalleryDlOutput(stdout: string): XTweetRaw[] { const trimmed = stdout.trim(); if (!trimmed) return []; // Whole-payload array first. if (trimmed.startsWith("[")) { try { const parsed = JSON.parse(quoteBigIntegers(trimmed)) as unknown; if (Array.isArray(parsed)) return flattenRecords(parsed); } catch { // fall through to line mode } } const out: XTweetRaw[] = []; for (const line of trimmed.split("\n")) { const t = line.trim(); if (!t || (!t.startsWith("{") && !t.startsWith("["))) continue; try { out.push(...flattenRecords([JSON.parse(quoteBigIntegers(t)) as unknown])); } catch { // Not JSON (a progress line) — ignore. } } return out; } // gallery-dl's dump entries are sometimes `[, , ]` tuples // rather than bare objects. Unwrap to the metadata object either way. function flattenRecords(items: unknown[]): XTweetRaw[] { const out: XTweetRaw[] = []; for (const item of items) { if (Array.isArray(item)) { for (const part of item) { if (part && typeof part === "object" && !Array.isArray(part)) { out.push(part as XTweetRaw); } } continue; } if (item && typeof item === "object") out.push(item as XTweetRaw); } return out; } // The login source (galleryDlCookieChoice) for one run, logged: the operator's // browser, or the session broker's exported jar — kept fresh by a live browser // profile, so it survives the cookie expiry that otherwise breaks gallery-dl // within days. One resolution for the timeline walk and the search walk. async function galleryDlLoginFor( input: Pick, ): Promise<{ cookies?: string; cookieFile?: string }> { const paths = getPaths(); const status = await readXSessionStatus(paths); const choice = galleryDlCookieChoice({ source: input.cookieSource, browserCookies: input.browserCookies, jarFile: status.hasCookies && status.looksAuthenticated ? xCookieFile(paths) : undefined, alwaysCookies: input.cookies, }); input.onLog?.(choice.note); return { cookies: choice.cookies, cookieFile: choice.cookieFile }; } // 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 { if (!input) return Promise.resolve(); return new Promise((resolve) => { const rl = createInterface({ input, crlfDelay: Infinity }); rl.on("line", onLine); rl.on("close", resolve); }); } const SEARCH_NEEDS_LOGIN = "X search needs a logged-in session: connect an X account on /settings, or " + "set cookiesFromBrowser for this channel (or globally), then run it again."; // The older-posts backfill: search windows walked backwards from `position` // (olderBackfill.ts does the window arithmetic). One gallery-dl run per // window, a pause between windows, the same per-run budget as a timeline // backfill across all of them. Progress is saved mid-window (the oldest post // read so far is the resume point) and at every window boundary. export async function fetchOlderViaSearch( input: OlderPostFetchInput, ): Promise { const { channelSlug, seenIds, signal, onLog } = input; const bin = getPaths().galleryDlBin; let handle: string; try { handle = xSearchHandle(input.handle); } catch (err) { return { posts: [], complete: false, position: input.position, error: (err as Error).message, }; } const { cookies, cookieFile } = await galleryDlLoginFor(input); if (!cookies && !cookieFile) { onLog?.(`[auth] ${SEARCH_NEEDS_LOGIN}`); return { posts: [], complete: false, position: input.position, needsCookies: true, error: SEARCH_NEEDS_LOGIN, }; } const deadline = Date.now() + BACKFILL_BUDGET_MS; // A fixed pause when the caller names one (tests pass 0); otherwise a fresh // random draw before every window. const pauseFor = () => input.windowPauseMs ?? olderWindowPauseMs(); const { drain } = input; // A drain ends the pause between windows as a cancel does. const pauseSignal = drain ? AbortSignal.any([signal, drain]) : signal; // Nothing archived: a walk that finds nothing either ends sooner // (stepOlderWindow's `nothingFound`). const nothingArchived = seenIds.size === 0; let position: OlderBackfillPosition = { ...input.position }; let accountCreatedAt = input.accountCreatedAt; // Posts not handed to a checkpoint, returned for the final write. const carried: Post[] = []; let checkpointsOn = Boolean(input.onCheckpoint); let read = 0; const save = async (posts: Post[], at: OlderBackfillPosition) => { if (!checkpointsOn || !input.onCheckpoint) { carried.push(...posts); return; } try { await input.onCheckpoint({ posts, position: at, accountCreatedAt }); } catch (err) { // Keep the posts for the final write and stop checkpointing, so a saved // position never runs ahead of what is on disk. checkpointsOn = false; carried.push(...posts); onLog?.(`[warn] could not save progress mid-run: ${(err as Error).message}`); } }; const stopped = ( at: OlderBackfillPosition, extra: Partial = {}, ): OlderPostFetchResult => ({ posts: carried, complete: false, position: at, accountCreatedAt, ...extra, }); const drainedAt = (at: OlderBackfillPosition): OlderPostFetchResult => { onLog?.( `Drained; the next run resumes at ${at.since} – ${at.until}` + (at.maxId ? " below the oldest post read." : "."), ); return stopped(at, { drained: true }); }; onLog?.( `Searching @${handle}'s posts backwards, ${position.since} – ${position.until}` + (position.maxId ? " from the saved resume point" : "") + ` (up to ${Math.round(BACKFILL_BUDGET_MS / 60_000)} min this run; progress is saved as it goes)`, ); for (;;) { const floor = olderFloorDay(input.floor, accountCreatedAt); if (floor !== undefined && position.until <= floor) { return { posts: carried, complete: true, completeReason: `the archive already reaches back to ${floor}, the earliest date to walk to`, position, accountCreatedAt, }; } const remainingMs = deadline - Date.now(); if (remainingMs <= 0) { onLog?.(`Stopped at this run's ${Math.round(BACKFILL_BUDGET_MS / 60_000)}-minute budget; the next run resumes at ${position.since} – ${position.until}.`); return stopped(position); } const remainingLimit = input.limit ? input.limit - read : undefined; if (remainingLimit !== undefined && remainingLimit <= 0) { return stopped(position); } if (drain?.aborted) return drainedAt(position); const windowStart = position; const stream = new GalleryDlStream({ channelSlug, seenIds, stopAtKnown: false, accountHandle: handle, onLog, }); const args = buildGalleryDlSearchArgs({ handle, position: windowStart, cookies, cookieFile, limit: remainingLimit, }); const subprocess = execa(bin, args, { reject: false, buffer: false, stdin: "ignore", timeout: remainingMs, cancelSignal: signal, env: { PYTHONUNBUFFERED: "1" }, }); // DRAIN mid-window stops the search at once: every post read is its own // resume point (the oldest one read becomes the window's max_id), so there // is no page boundary to wait for. SIGTERM, then execa's SIGKILL if // gallery-dl has not exited a few seconds later. let drainedHere = false; const onDrain = () => { drainedHere = true; onLog?.("Drain requested: stopping this window's search."); subprocess.kill("SIGTERM"); }; drain?.addEventListener("abort", onDrain, { once: true }); // Mid-window checkpoints, one at a time and in order. The resume point is // the oldest post read so far, which every post taken here is at or above. let checkpoints: Promise = Promise.resolve(); let lastCheckpointAt = Date.now(); const maybeCheckpoint = () => { if (!checkpointsOn) return; const due = stream.pendingCount() >= CHECKPOINT_EVERY_POSTS || (stream.pendingCount() > 0 && Date.now() - lastCheckpointAt >= CHECKPOINT_EVERY_MS); if (!due || !stream.oldestId) return; lastCheckpointAt = Date.now(); const posts = stream.takePending(); const at = { ...windowStart, maxId: stream.oldestId }; checkpoints = checkpoints.then(() => save(posts, at)); }; const stdoutDone = readLines(subprocess.stdout, (line) => { stream.pushStdout(line); maybeCheckpoint(); }); const stderrDone = readLines(subprocess.stderr, (line) => { stream.pushStderr(line); maybeCheckpoint(); }); const res = await subprocess; drain?.removeEventListener("abort", onDrain); 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)`, ); } read += stream.distinct; accountCreatedAt ??= stream.accountCreatedAt; carried.push(...stream.finish()); const here: OlderBackfillPosition = { ...windowStart, maxId: stream.oldestId ?? windowStart.maxId, }; if (stream.rateLimited) { onLog?.("[rate-limit] X throttled this search — gallery-dl paused between pages."); } onLog?.( `${windowStart.since} – ${windowStart.until}: read ${stream.distinct} post(s), ${stream.newPosts} new` + (stream.known ? `, ${stream.known} already archived` : ""), ); if (res.isCanceled || signal.aborted || res.timedOut) { onLog?.( res.timedOut ? `Stopped at this run's ${Math.round(BACKFILL_BUDGET_MS / 60_000)}-minute budget; the next run resumes where this one left off.` : "Cancelled; the next run resumes where this one left off.", ); return stopped(here); } if (drainedHere) return drainedAt(here); if (res.exitCode !== 0) { const tail = stream.stderrTail(); if (looksLikeAuthFailure(tail)) { onLog?.(`[auth] gallery-dl could not authenticate to X:\n${tail}`); return stopped(here, { needsCookies: true, error: SEARCH_NEEDS_LOGIN }); } return stopped(here, { error: `gallery-dl exited ${res.exitCode ?? res.signal ?? "abnormally"}: ${tail || "(no output)"}`, }); } if (remainingLimit !== undefined && stream.distinct >= remainingLimit) { return stopped(here); } // The window is walked to its end. const step = stepOlderWindow(windowStart, { // A window resumed below a saved point held posts above it. hadPosts: stream.distinct > 0 || Boolean(windowStart.maxId), floor: olderFloorDay(input.floor, accountCreatedAt), nothingFound: nothingArchived && read === 0, }); if (step.complete) { onLog?.(`Older-posts backfill complete: ${step.reason}.`); return { posts: carried, complete: true, completeReason: step.reason, position: { since: windowStart.since, until: windowStart.until, emptyWindows: step.emptyWindows }, accountCreatedAt, }; } position = step.next; // Every window boundary is saved: the posts so far, and the next window. await save(carried.splice(0), position); if (signal.aborted) { onLog?.("Cancelled; the next run resumes where this one left off."); return stopped(position); } // Before the gap, not after it: a drain never waits out a pause. if (drain?.aborted) return drainedAt(position); await pause(pauseFor(), pauseSignal); if (signal.aborted) { onLog?.("Cancelled; the next run resumes where this one left off."); return stopped(position); } if (drain?.aborted) return drainedAt(position); } } export const xGalleryDlFetcher: SocialFetcher = { id: "x-gallery-dl", label: "X / Twitter (gallery-dl)", platform: "twitter", fields: { cookies: true, binPath: true, limit: true }, detect(url: string): boolean { try { const host = new URL(url).hostname.toLowerCase(); return ( host === "x.com" || host.endsWith(".x.com") || host === "twitter.com" || host.endsWith(".twitter.com") ); } catch { return false; } }, // A cheap liveness probe: read a single item. There is no free unauthenticated // X metadata endpoint, so unlike Bluesky this genuinely spawns the binary. async probe(url: string): Promise { const bin = getPaths().galleryDlBin; try { const res = await execa( bin, buildGalleryDlArgs({ accountUrl: url, limit: 1 }), { reject: false, timeout: 60_000 }, ); if (res.exitCode !== 0) { const err = `${res.stderr ?? ""}`.trim().split("\n").slice(-3).join(" "); return { ok: false, error: looksLikeAuthFailure(err) ? `gallery-dl needs X credentials (configure cookies-from-browser): ${err}` : err || `gallery-dl exited ${res.exitCode}`, }; } const records = parseGalleryDlOutput(res.stdout ?? ""); const first = records[0]; const author = first && typeof first === "object" ? ((first.author as Record | undefined) ?? undefined) : undefined; const handle = (typeof author?.name === "string" ? author.name : undefined) ?? (typeof first?.screen_name === "string" ? first.screen_name : undefined); const nick = typeof author?.nick === "string" ? author.nick : undefined; return { ok: true, name: nick ?? handle, handle, url: handle ? `https://x.com/${handle}` : url, }; } catch (err) { const message = (err as Error).message; return { ok: false, error: /ENOENT/.test(message) ? `gallery-dl not found (looked for "${bin}"; set GALLERY_DL_BIN)` : message, }; } }, async fetch(input: PostFetchInput): Promise { const { accountUrl, channelSlug, since, seenIds, limit, signal, onLog } = input; const bin = getPaths().galleryDlBin; const { cookies, cookieFile } = await galleryDlLoginFor(input); 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)`, ); 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 = 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}`); } }); }; // DRAIN stops gallery-dl at a page boundary — the "Cursor:" line it logs // after each page, which is a resume point — so the next run does not // repeat a page's request. Between pages already (a cursor just logged, // or a rate-limit wait), it stops at once. SIGTERM, then execa's SIGKILL // if gallery-dl has not exited a few seconds later; the resume point is // the latest cursor, safe once both pipes have drained. let stoppedAtKnown = false; let drained = false; let betweenPages = false; const stopForDrain = () => { if (drained || stoppedAtKnown) return; drained = true; onLog?.("Drain requested: stopping gallery-dl at this page boundary."); subprocess.kill("SIGTERM"); }; const onDrain = () => { if (betweenPages) stopForDrain(); }; input.drain?.addEventListener("abort", onDrain, { once: true }); const stdoutDone = readLines(subprocess.stdout, (line) => { if (line.trim()) betweenPages = false; 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(); } }); const stderrDone = readLines(subprocess.stderr, (line) => { if (stream.pushStderr(line)) { betweenPages = true; if (input.drain?.aborted) stopForDrain(); } maybeCheckpoint(); }); const res = await subprocess; input.drain?.removeEventListener("abort", onDrain); 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)`, ); } const posts = stream.finish(); if (stream.rateLimited) { onLog?.( "[rate-limit] X throttled this run — gallery-dl paused between pages. " + (cookieFile || cookies ? "Consider a longer sync interval." : "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` : ""), ); // 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 (drained) { onLog?.( resumeAt ? "Drained; the next run resumes from the last page read." : "Drained before the first page; the next run starts where this one did.", ); return { posts, complete: false, cursor: resumeAt, drained: true }; } 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 }; } return { posts, complete: false, cursor: resumeAt, error: `gallery-dl exited ${res.exitCode ?? res.signal ?? "abnormally"}: ${tail || "(no output)"}`, }; } // A clean exit walked everything it was given: gallery-dl waits out rate // limits rather than skipping pages. Only a --post-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 }; }, fetchOlder: fetchOlderViaSearch, captureByIds: captureXPostsByIds, }; registerSocialFetcher(xGalleryDlFetcher);