// Drive a SocialFetcher for one social channel and land the result on disk. // // This is the posts-corpus analogue of sync(): read the channel's already-seen // ids, ask the fetcher for anything newer, append to the month-sharded JSONL // and the posts-archive, and record a small state sidecar for the channel's // snapshot buckets. Incremental by construction — a re-run over overlapping // pages writes nothing new. import path from "node:path"; import { patchChannelConfig, readChannelConfig, } from "./channels"; import type { Paths } from "../lib/paths"; import type { Post } from "../lib/posts"; import { isSocialChannel } from "../lib/channelConfig"; import { alwaysCookies, resolveCookiePolicy, type CookiePolicyInputs, } from "../lib/cookiePolicy"; import { latestPostCreatedAt, oldestPostCreatedAt, readPostFetchState, readSeenPostIds, writePostFetchState, writePosts, type OlderBackfillPosition, type OlderBackfillState, type PostFetchState, } from "../lib/posts-server"; import { handleFromAccountUrl, resolveSocialFetcher, type SocialFetcher, } from "../social/fetchers"; import { firstOlderWindow, isUtcDay, olderFloorDay, } from "../social/olderBackfill"; // Registering the built-in fetchers is a side effect of importing them. Keep // this list here (rather than inside fetchers.ts) so the registry module stays // free of imports from the heavier fetchers. import "../social/blueskyFetcher"; // Registration order encodes preference: gallery-dl is the primary X path, so // it is registered before the Playwright fallback and wins URL detection. import "../social/xGalleryDlFetcher"; // The fallback is registered LAST and never claims a URL by detection, so it is // only ever used when a channel opts into it via postFetcher: "x-playwright". import "../social/xPlaywrightFetcher"; import "../social/xNitterFetcher"; // A forum thread (platform "xenforo"): claims only thread URLs. import "../social/xenforoFetcher"; import type { XCookieSource } from "../social/xCookieSource"; import { resolveXCookieSourceFor, type XLoginSettings, } from "../social/xBrowserLogin"; export type FetchPostsOptions = { paths: Paths; slug: string; // Global settings, for the cookie policy and the X login source. Structural // so settings.ts need not be imported here. settings: CookiePolicyInputs & XLoginSettings; // Ignore the stored watermark and re-walk the account's full history. New // posts are still deduped against the archive, so this is a safe repair // operation rather than a duplicate-maker. full?: boolean; // Walk the account's history BACKWARDS from the oldest archived post, with // the fetcher's `fetchOlder` (X: search windows), instead of fetching new // posts. Its position is kept in posts-state.json's `older`, apart from the // timeline cursor. Not with `full`. older?: boolean; // The older walk's floor date (YYYY-MM-DD): it stops there. Remembered for // later runs of the walk. Default: none (it stops after a run of empty // windows, or at the account's creation date). floor?: string; // Start the older walk afresh, backwards from this date (YYYY-MM-DD), in // place of its saved position — and of a "complete" it reached. The walk // otherwise starts below the OLDEST archived post, so one surviving old post // (a timeline that returns a handful of 2018 posts beside its 3,200 recent // ones) puts years of history below a gap no window ever covers. Only with // `older`. from?: string; // The older walk's pause between two windows. Default: the fetcher's own. olderWindowPauseMs?: number; // Walk older posts even on an account that shows none (see // emptyAccountOlderProblem). Only with `older`. force?: boolean; limit?: number; // Cap on pages read this run, for a fetcher that walks pages (a forum // thread's "latest N pages"). pages?: number; onLog?: (line: string) => void; signal?: AbortSignal; // The job's Drain: the fetch stops at its next resume point and keeps it, // and the run ends as a success, not a failure. drain?: AbortSignal; }; // The refusal for both walks at once, shared with the server action so the // button, the ops route and this controller say the same thing. export const OLDER_FROM_REFUSAL = "A start date (\"from\") applies only to an older-posts fetch."; export const FULL_AND_OLDER_REFUSAL = "A full re-fetch and an older-posts fetch are different walks — run one at a time."; // Why this fetcher cannot fetch older posts, or null when it can. export function olderPostsProblem( fetcher: Pick | undefined, ): string | null { if (!fetcher) return "This channel has no post fetcher."; return typeof fetcher.fetchOlder === "function" ? null : `${fetcher.label} cannot fetch older posts.`; } // An older walk over an account that shows no posts is refused unless forced: // nothing is archived, and the last timeline fetch finished cleanly having read // nothing. Searching such an account window after window is nothing but empty // searches. A timeline run that stopped part-way, failed or needed a login // says nothing about the account, so it is not grounds to refuse. export const EMPTY_ACCOUNT_OLDER_REFUSAL = "The account shows no posts: nothing is archived and the last timeline fetch read none. " + 'A search walk would only repeat empty searches — run it with "force" to walk anyway.'; export function emptyAccountOlderProblem( archived: number, state: Pick | null, ): string | null { if (archived > 0 || !state) return null; if (state.lastFetchedCount !== 0) return null; if (state.cursor || state.lastError || state.needsCookies) return null; return EMPTY_ACCOUNT_OLDER_REFUSAL; } export type FetchPostsResult = { ok: boolean; written: number; skipped: number; complete: boolean; error?: string; needsCookies?: boolean; // Stopped on the job's Drain, at a resume point the next run continues from. drained?: boolean; }; export async function fetchPosts( opts: FetchPostsOptions, ): Promise { const { paths, slug, settings, onLog, signal } = opts; const log = (line: string) => onLog?.(line); const channelRoot = path.join(paths.channelsDir, slug); if (opts.full && opts.older) { return { ok: false, written: 0, skipped: 0, complete: false, error: FULL_AND_OLDER_REFUSAL }; } for (const [name, day] of [["from", opts.from]] as const) { if (day !== undefined && !isUtcDay(day)) { return { ok: false, written: 0, skipped: 0, complete: false, error: `"${day}" is not a date (YYYY-MM-DD) for "${name}"` }; } } if (opts.from !== undefined && !opts.older) { return { ok: false, written: 0, skipped: 0, complete: false, error: OLDER_FROM_REFUSAL }; } if (opts.floor !== undefined && !isUtcDay(opts.floor)) { return { ok: false, written: 0, skipped: 0, complete: false, error: `"${opts.floor}" is not a date (YYYY-MM-DD)`, }; } const config = await readChannelConfig(paths, slug); if (!config) { return { ok: false, written: 0, skipped: 0, complete: false, error: `No such channel: ${slug}` }; } if (!isSocialChannel(config)) { return { ok: false, written: 0, skipped: 0, complete: false, error: `Channel ${slug} is not a social source (sourceKind=${config.sourceKind ?? "video"})`, }; } const accountUrl = config.url ?? ""; const fetcher = resolveSocialFetcher(config.postFetcher, accountUrl); if (!fetcher) { return { ok: false, written: 0, skipped: 0, complete: false, error: `No post fetcher for ${accountUrl || slug} (configured: ${config.postFetcher ?? "auto-detect"})`, }; } const handle = config.socialHandle ?? handleFromAccountUrl(accountUrl) ?? ""; if (!handle) { return { ok: false, written: 0, skipped: 0, complete: false, error: `Could not determine an account handle for ${slug}`, }; } if (opts.drain?.aborted) { log("Drained before the fetch started; nothing was read."); return { ok: true, written: 0, skipped: 0, complete: false, drained: true }; } const seenIds = await readSeenPostIds(channelRoot); const priorState = await readPostFetchState(channelRoot); if (opts.older) { const problem = olderPostsProblem(fetcher); if (problem) { return { ok: false, written: 0, skipped: 0, complete: false, error: problem }; } const policy = resolveCookiePolicy(settings, config); const xLogin = fetcher.platform === "twitter" ? await resolveXCookieSourceFor(paths, settings, policy.cookies) : undefined; return fetchOlderPosts({ opts, fetcher, channelRoot, handle, accountUrl, seenIds, priorState, cookies: alwaysCookies(policy), cookieSource: xLogin?.source, browserCookies: xLogin?.browserSpec, }); } // A stored cursor means the previous run stopped early (hit its limit, or was // cancelled). Resume the backfill from there rather than applying the // watermark, which would otherwise leave that history permanently missing. // --full discards the cursor and re-walks from the top. const resumeCursor = opts.full ? undefined : priorState?.cursor; const since = opts.full || resumeCursor ? undefined : ((await latestPostCreatedAt(channelRoot)) ?? undefined); const policy = resolveCookiePolicy(settings, config); const cookies = alwaysCookies(policy); // An X fetcher's login source (social.x.cookieSource, xCookieSource.ts), // resolved against THIS channel's browser spec — its own cookiesFromBrowser // over the global one, whatever the cookieMode. const xLogin = fetcher.platform === "twitter" ? await resolveXCookieSourceFor(paths, settings, policy.cookies) : undefined; log( `Fetching posts for ${slug} via ${fetcher.label} (@${handle})` + (resumeCursor ? " — resuming an unfinished backfill" : since ? ` since ${since}` : " — full history") + `; ${seenIds.size} already archived.`, ); // A managed job's cancel/drain arrives as an AbortSignal; the fetcher // interface requires one, so synthesize a never-aborting signal when the // caller has none (a direct CLI invocation). const controller = new AbortController(); const effectiveSignal = signal ?? controller.signal; const state: PostFetchState = { lastFetchedAt: new Date().toISOString(), }; // The older walk's position rides along untouched: a normal fetch neither // reads it nor drops it. if (priorState?.older) state.older = priorState.older; // 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(); 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({ accountUrl, handle, channelSlug: slug, since, cursor: resumeCursor, seenIds, cookies, cookieSource: xLogin?.source, browserCookies: xLogin?.browserSpec, limit: opts.limit, ...(opts.pages ? { pages: opts.pages } : {}), ...(config.postPagePauseSeconds ? { pagePauseMs: config.postPagePauseSeconds * 1000 } : {}), stopAtKnown: !opts.full, onCheckpoint, signal: effectiveSignal, drain: opts.drain, onLog: log, }); } catch (err) { 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: checkpointWritten, skipped: checkpointSkipped, complete: false, error: message, }; } 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` + (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 // (common/jobs/syncScheduler.ts) paces social channels with zero changes — // it keys only off url / syncIntervalMinutes / lastSyncedAt, plus the // channel-priority document's `sync` tier (`isChannelPaused(model, slug, // "sync")`), which is where the retired `excludeFromSync` flag went. // // A PATCH, re-read at write time. It used to spread the `config` this fetch // read at its START, so any edit made during a fetch (minutes, for a long // history) was silently reverted when the stamp landed. await patchChannelConfig(paths, slug, { lastSyncedAt: new Date().toISOString(), }); return { ok: !result.error, written: written.written, skipped: written.skipped, complete: result.complete, needsCookies: result.needsCookies, error: result.error, ...(result.drained ? { drained: true } : {}), }; } // The older-posts walk (`FetchPostsOptions.older`). Its own position in // posts-state.json's `older`; every other field of the state is carried as it // was, so the timeline cursor and watermark are exactly what the last normal // fetch left. It does not stamp the channel's lastSyncedAt: it is not a sync, // and the scheduler's next new-posts fetch stays where it was. async function fetchOlderPosts(ctx: { opts: FetchPostsOptions; fetcher: SocialFetcher; channelRoot: string; handle: string; accountUrl: string; seenIds: ReadonlySet; priorState: PostFetchState | null; cookies?: string; cookieSource?: XCookieSource; browserCookies?: string; }): Promise { const { opts, fetcher, channelRoot, handle, priorState } = ctx; const { slug } = opts; const log = (line: string) => opts.onLog?.(line); const fetchOlder = fetcher.fetchOlder!; // A start date replaces the saved walk, complete or not. const prior = opts.from === undefined ? priorState?.older : undefined; if (prior?.complete) { log( `Older posts for ${slug} are already complete` + (prior.completedAt ? ` (since ${prior.completedAt})` : "") + (prior.completeReason ? `: ${prior.completeReason}` : "") + ". Nothing to walk.", ); return { ok: true, written: 0, skipped: 0, complete: true }; } const empty = emptyAccountOlderProblem(ctx.seenIds.size, priorState); if (empty && !opts.force) { log(empty); return { ok: false, written: 0, skipped: 0, complete: false, error: empty }; } if (empty) log("[warn] The account shows no posts; walking anyway, as forced."); const floor = opts.floor ?? priorState?.older?.floor; let accountCreatedAt = priorState?.older?.accountCreatedAt; const startedAt = new Date().toISOString(); let position: OlderBackfillPosition; if (prior) { position = { since: prior.since, until: prior.until, emptyWindows: prior.emptyWindows ?? 0, ...(prior.maxId ? { maxId: prior.maxId } : {}), }; } else { const oldest = opts.from ?? (await oldestPostCreatedAt(channelRoot)); position = firstOlderWindow({ oldestArchivedAt: oldest ?? undefined, floor: olderFloorDay(floor, accountCreatedAt), }); } log( `Fetching older posts for ${slug} via ${fetcher.label} (@${handle}): ` + (prior ? "resuming at " : opts.from ? `starting afresh from ${opts.from} at ` : "starting at ") + `${position.since} – ${position.until}` + (floor ? `, down to ${floor}` : "") + `; ${ctx.seenIds.size} already archived.`, ); let written = 0; let skipped = 0; const shards = new Set(); const olderState = ( at: OlderBackfillPosition, extra: Partial = {}, ): OlderBackfillState => ({ since: at.since, until: at.until, emptyWindows: at.emptyWindows, ...(at.maxId ? { maxId: at.maxId } : {}), ...(floor ? { floor } : {}), ...(accountCreatedAt ? { accountCreatedAt } : {}), lastRunAt: startedAt, lastWritten: written, ...extra, }); // Everything but `older` is carried from the state as it was. const save = (older: OlderBackfillState, needsCookies?: boolean) => writePostFetchState(channelRoot, { ...(priorState ?? {}), ...(needsCookies ? { needsCookies: true } : {}), older, }); const onCheckpoint = async (cp: { posts: Post[]; position: OlderBackfillPosition; accountCreatedAt?: string; }) => { const w = await writePosts(channelRoot, cp.posts); written += w.written; skipped += w.skipped; for (const shard of w.shards) shards.add(shard); accountCreatedAt ??= cp.accountCreatedAt; await save(olderState(cp.position)); if (w.written > 0) { log(`Saved ${w.written} older post(s) (${written} so far this run); resume point recorded.`); } }; const controller = new AbortController(); let result; try { result = await fetchOlder({ accountUrl: ctx.accountUrl, handle, channelSlug: slug, seenIds: ctx.seenIds, cookies: ctx.cookies, cookieSource: ctx.cookieSource, browserCookies: ctx.browserCookies, limit: opts.limit, signal: opts.signal ?? controller.signal, drain: opts.drain, onLog: log, position, floor, accountCreatedAt, windowPauseMs: opts.olderWindowPauseMs, onCheckpoint, }); } catch (err) { const message = (err as Error).message; log(`[error] ${message}`); // Whatever was checkpointed is on disk with its position; keep that one. const now = (await readPostFetchState(channelRoot))?.older; await save({ ...(now ?? olderState(position)), lastError: message }); return { ok: false, written, skipped, complete: false, error: message }; } accountCreatedAt ??= result.accountCreatedAt; const final = await writePosts(channelRoot, result.posts); written += final.written; skipped += final.skipped; for (const shard of final.shards) shards.add(shard); log( `Wrote ${written} older post(s), skipped ${skipped} already archived` + (shards.size > 0 ? ` (shards: ${[...shards].join(", ")})` : "") + (result.complete ? " — older posts complete." : " — more history remains."), ); if (result.error) log(`[error] ${result.error}`); await save( olderState( result.position, result.complete ? { complete: true, completedAt: new Date().toISOString(), ...(result.completeReason ? { completeReason: result.completeReason } : {}), } : result.error ? { lastError: result.error } : {}, ), result.needsCookies, ); return { ok: !result.error, written, skipped, complete: result.complete, needsCookies: result.needsCookies, error: result.error, ...(result.drained ? { drained: true } : {}), }; }