commit df0ab16ec786e040a1b4eefbad20fa1f3a5fab80
parent c3f3f4e94f69ae7f10f49269366147b7d270da02
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sun, 4 Oct 2026 22:04:12 -0400
common: X fetches honour Drain at a resume point; an older walk over an empty archive that finds nothing stops after two empty windows
The timeline fetch stops gallery-dl at the next page boundary (its cursor
line), or at once between pages, and returns drained with the cursor kept.
The search walk stops a window's search at once (every post read is its
resume point), and checks drain at each window boundary before the pause.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Diffstat:
3 files changed, 100 insertions(+), 9 deletions(-)
diff --git a/common/social/fetchers.ts b/common/social/fetchers.ts
@@ -63,6 +63,10 @@ export type PostFetchInput = {
// again in the result. A fetcher that pages in one shot never calls it.
onCheckpoint?: (checkpoint: { posts: Post[]; cursor: string }) => Promise<void>;
signal: AbortSignal;
+ // A soft stop (the job's Drain): the run ends at its next resume point —
+ // between two pages — keeping that point, rather than walking on. `signal`
+ // cancels outright. A fetcher that pages in one shot may ignore it.
+ drain?: AbortSignal;
// Progress/diagnostic sink, wired to the job log.
onLog?: (line: string) => void;
};
@@ -80,6 +84,8 @@ export type PostFetchResult = {
// The run failed part-way. `posts` and `cursor` still hold what it got, so
// the controller saves them before reporting the failure.
error?: string;
+ // The run stopped on `drain` at a resume point (`cursor`). Not a failure.
+ drained?: boolean;
};
// Input for the older-posts backfill (`SocialFetcher.fetchOlder`): walk the
@@ -97,6 +103,7 @@ export type OlderPostFetchInput = Pick<
| "browserCookies"
| "limit"
| "signal"
+ | "drain"
| "onLog"
> & {
// Where to start: the stored position of an unfinished walk, or the first
@@ -128,6 +135,8 @@ export type OlderPostFetchResult = {
accountCreatedAt?: string;
needsCookies?: boolean;
error?: string;
+ // The walk stopped on `drain`; `position` is where the next run resumes.
+ drained?: boolean;
};
export type SocialFetcherProbe = {
diff --git a/common/social/olderBackfill.ts b/common/social/olderBackfill.ts
@@ -21,6 +21,11 @@ export const OLDER_WINDOW_MONTHS = 3;
// three-month windows, a year with nothing.
export const OLDER_MAX_EMPTY_WINDOWS = 4;
+// The same, for a walk over an account with nothing archived that has found
+// nothing in this run either: two windows. Searching an empty account window
+// after window only repeats empty searches.
+export const OLDER_MAX_EMPTY_WINDOWS_NOTHING_ARCHIVED = 2;
+
const DAY_RE = /^\d{4}-\d{2}-\d{2}$/;
// A real YYYY-MM-DD calendar day (2021-02-30 is not one).
@@ -102,7 +107,9 @@ export type OlderWindowStep =
// What follows a window walked to its end. `hadPosts` counts any post the
// window held, archived already or not — a window full of known posts is not
-// an empty one.
+// an empty one. `nothingFound`: nothing is archived and this run has read no
+// post — the walk ends after OLDER_MAX_EMPTY_WINDOWS_NOTHING_ARCHIVED empty
+// windows, not a year of them.
export function stepOlderWindow(
position: OlderBackfillPosition,
opts: {
@@ -110,10 +117,13 @@ export function stepOlderWindow(
floor?: string;
months?: number;
maxEmptyWindows?: number;
+ nothingFound?: boolean;
},
): OlderWindowStep {
const months = opts.months ?? OLDER_WINDOW_MONTHS;
- const maxEmpty = opts.maxEmptyWindows ?? OLDER_MAX_EMPTY_WINDOWS;
+ const maxEmpty =
+ opts.maxEmptyWindows ??
+ (opts.nothingFound ? OLDER_MAX_EMPTY_WINDOWS_NOTHING_ARCHIVED : OLDER_MAX_EMPTY_WINDOWS);
const emptyWindows = opts.hadPosts ? 0 : position.emptyWindows + 1;
if (isAtFloor(position, opts.floor)) {
return {
@@ -125,7 +135,11 @@ export function stepOlderWindow(
if (emptyWindows >= maxEmpty) {
return {
complete: true,
- reason: `${emptyWindows} windows of ${months} months in a row held no posts, back to ${position.since}`,
+ reason:
+ `${emptyWindows} windows of ${months} months in a row held no posts, back to ${position.since}` +
+ (opts.nothingFound
+ ? ", with nothing archived and none found — an empty account is not searched further"
+ : ""),
emptyWindows,
};
}
diff --git a/common/social/xGalleryDlFetcher.ts b/common/social/xGalleryDlFetcher.ts
@@ -574,18 +574,22 @@ export class GalleryDlStream {
for (const rec of records) this.addRecord(rec);
}
- pushStderr(line: string): void {
+ // 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;
+ return "cursor";
}
- if (DEBUG_LINE.test(line) || !line.trim()) return;
- if (RATE_LIMIT_LINE.test(line)) this.rateLimited = true;
+ 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.
@@ -854,6 +858,12 @@ export async function fetchOlderViaSearch(
// 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.
@@ -886,6 +896,13 @@ export async function fetchOlderViaSearch(
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}` +
@@ -913,6 +930,7 @@ export async function fetchOlderViaSearch(
if (remainingLimit !== undefined && remainingLimit <= 0) {
return stopped(position);
}
+ if (drain?.aborted) return drainedAt(position);
const windowStart = position;
const stream = new GalleryDlStream({
@@ -937,6 +955,17 @@ export async function fetchOlderViaSearch(
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.
@@ -963,6 +992,7 @@ export async function fetchOlderViaSearch(
maybeCheckpoint();
});
const res = await subprocess;
+ drain?.removeEventListener("abort", onDrain);
await Promise.all([stdoutDone, stderrDone]);
await checkpoints;
@@ -996,6 +1026,7 @@ export async function fetchOlderViaSearch(
);
return stopped(here);
}
+ if (drainedHere) return drainedAt(here);
if (res.exitCode !== 0) {
const tail = stream.stderrTail();
if (looksLikeAuthFailure(tail)) {
@@ -1015,6 +1046,7 @@ export async function fetchOlderViaSearch(
// 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}.`);
@@ -1033,11 +1065,14 @@ export async function fetchOlderViaSearch(
onLog?.("Cancelled; the next run resumes where this one left off.");
return stopped(position);
}
- await pause(pauseFor(), signal);
+ // 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);
}
}
@@ -1182,8 +1217,28 @@ export const xGalleryDlFetcher: SocialFetcher = {
});
};
+ // 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;
@@ -1194,11 +1249,15 @@ export const xGalleryDlFetcher: SocialFetcher = {
}
});
const stderrDone = readLines(subprocess.stderr, (line) => {
- stream.pushStderr(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;
@@ -1246,6 +1305,15 @@ export const xGalleryDlFetcher: SocialFetcher = {
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)) {