commit 098ffba4698aecf07af3270a220aff358251847a
parent 78e5b101e8570e4156cb51435536a06fbf30dd43
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Mon, 5 Oct 2026 15:07:35 -0400
posts: a forum walk with a cursor catches up from the last page first; pages over ops; e2e
A run that resumes a stored cursor now reads the newest pages first (until it
meets archived posts), then continues the older pages at the cursor, so "the
latest N pages" are always the newest and an unfinished history is kept.
`pages` reaches POST /api/ops/fetch-posts and `pnpm ops` help; the social
channel e2e creates a forum-thread channel and imports a synthetic saved page.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Diffstat:
6 files changed, 177 insertions(+), 38 deletions(-)
diff --git a/common/social/xenforoFetcher.test.ts b/common/social/xenforoFetcher.test.ts
@@ -139,13 +139,35 @@ test("a pages cap reads the latest N pages and leaves the cursor at the next one
assert.equal(loader.urls.length, 2);
});
-test("resume: a stored cursor continues backwards, and an archived post there does not end it", async () => {
+test("resume: a run with a cursor catches up from the last page, then continues the older pages there", async () => {
const loader = fakeLoader();
- // Post 9 (page 3) shifted in from a page already read: archived.
- const res = await walkXenforoThread(input({ cursor: "3", seenIds: new Set(["9009", "9013"]) }), deps(loader).deps);
+ // An earlier run read pages 5 and 4 and stopped (cursor 3). Post 9 (page 3)
+ // shifted in from a page already read: archived, and it does not end the
+ // backfill.
+ const seen = new Set(["9009", "9010", "9011", "9012", "9013"]);
+ const res = await walkXenforoThread(input({ cursor: "3", seenIds: seen }), deps(loader).deps);
assert.equal(res.complete, true);
- assert.deepEqual(loader.urls, [`${THREAD_URL}page-3`, `${THREAD_URL}page-2`, THREAD_URL]);
- assert.deepEqual(ids(res.posts), [7, 8, 4, 5, 6, 1, 2, 3]);
+ assert.deepEqual(loader.urls, [
+ `${THREAD_URL}page-${LAST_PAGE_PROBE}`,
+ `${THREAD_URL}page-3`,
+ `${THREAD_URL}page-2`,
+ THREAD_URL,
+ ]);
+ // The newest first (posts 14, 15 arrived since), then the older pages.
+ assert.deepEqual(ids(res.posts), [14, 15, 7, 8, 4, 5, 6, 1, 2, 3]);
+});
+
+test("resume with a pages cap: the latest page is always read first, the cap spans both", async () => {
+ const seen = new Set(["9007", "9008", "9009", "9010", "9011", "9012", "9013", "9014", "9015"]);
+ const one = fakeLoader();
+ const r1 = await walkXenforoThread(input({ cursor: "2", seenIds: seen, pages: 1 }), deps(one).deps);
+ assert.deepEqual(one.urls, [`${THREAD_URL}page-${LAST_PAGE_PROBE}`]);
+ assert.equal(r1.cursor, "2", "nothing new, and the older pages still owed");
+ const two = fakeLoader();
+ const r2 = await walkXenforoThread(input({ cursor: "2", seenIds: seen, pages: 2 }), deps(two).deps);
+ assert.deepEqual(two.urls, [`${THREAD_URL}page-${LAST_PAGE_PROBE}`, `${THREAD_URL}page-2`]);
+ assert.equal(r2.cursor, "1");
+ assert.deepEqual(ids(r2.posts), [4, 5, 6]);
});
test("incremental: stops on the page holding an already-archived post, keeping its new ones", async () => {
@@ -249,6 +271,7 @@ test("a captcha on the very first load keeps no cursor; a 429 is not a session p
input({ cursor: "2" }),
deps(fakeLoader(() => ({ html: "<html><body>slow down</body></html>", status: 429 }))).deps,
);
+ // Refused on the catch-up load: the stored cursor is kept, not lost.
assert.equal(limited.needsCookies, undefined);
assert.equal(limited.cursor, "2");
assert.match(limited.error ?? "", /HTTP 429/);
diff --git a/common/social/xenforoFetcher.ts b/common/social/xenforoFetcher.ts
@@ -9,8 +9,10 @@
// watermark, at page 1, or at a cap — `pages` ("the latest N pages") or
// `limit` (posts), checked after a whole page so a page is never half-read.
// The cursor is the next page to read, so a capped or drained run resumes
-// there; a resumed run walks back to page 1 (pages shift as posts are
-// deleted, so meeting an archived post there does not end it).
+// there — after catching up from the last page first, so a run always starts
+// with the newest posts; past the cursor the walk goes on to page 1 (pages
+// shift as posts are deleted, so meeting an archived post there does not end
+// it).
//
// PACED, SERIAL, ONE BROWSER. One page at a time with a jittered pause between
// loads (about 10–20 s by default; `pagePauseMs` from the channel), in one
@@ -146,6 +148,14 @@ export type ThreadWalkDeps = {
};
// THE WALK. Pure over its deps: the loader hands back HTML, the parser reads it.
+//
+// Two modes. "new": from the last page down, stopping at the first archived
+// post (or the watermark) — what a run is for. "backfill": from a stored
+// cursor down to page 1, where an archived post means nothing (pages shift as
+// posts are deleted). A run with a cursor does both, newest first: it catches
+// up from the last page, then continues the older pages at the cursor — so
+// "the latest N pages" are always the newest, and an unfinished history is
+// never forgotten.
export async function walkXenforoThread(
input: PostFetchInput,
deps: ThreadWalkDeps,
@@ -155,13 +165,18 @@ export async function walkXenforoThread(
if (!thread) {
return { posts: [], complete: false, error: `${input.accountUrl} is not a XenForo thread URL (…/threads/<title>.<id>/).` };
}
- const resuming = Boolean(input.cursor);
- const startPage = resuming ? Number(input.cursor) : undefined;
- if (resuming && (!Number.isInteger(startPage) || startPage! < 1)) {
- return { posts: [], complete: false, error: `The stored resume point "${input.cursor}" is not a page number.` };
+ let resumeAt: number | undefined;
+ if (input.cursor) {
+ resumeAt = Number(input.cursor);
+ if (!Number.isInteger(resumeAt) || resumeAt < 1) {
+ return { posts: [], complete: false, error: `The stored resume point "${input.cursor}" is not a page number.` };
+ }
}
- const watermark = resuming ? undefined : input.since;
- const stopAtKnown = (input.stopAtKnown ?? true) && !resuming;
+ const stopAtKnown = input.stopAtKnown ?? true;
+ // A full re-walk (stopAtKnown false) reads every page from the end; a
+ // cursor means nothing to it.
+ if (!stopAtKnown) resumeAt = undefined;
+ const watermark = input.since;
const posts: Post[] = [];
const taken = new Set<string>();
let returnedCount = 0;
@@ -173,37 +188,40 @@ export async function walkXenforoThread(
return deps.loader.load(url, signal);
};
- let page = startPage ?? LAST_PAGE_PROBE;
- if (resuming) onLog?.(`Resuming the thread walk at page ${page}.`);
+ let mode: "new" | "backfill" = "new";
+ let page = LAST_PAGE_PROBE;
+ if (resumeAt) onLog?.(`An earlier run stopped at page ${resumeAt}: catching up from the last page first, then continuing there.`);
let lastPage: number | undefined;
const visited = new Set<number>();
+ const cursorOf = (n: number) => (n === LAST_PAGE_PROBE ? (resumeAt ? String(resumeAt) : undefined) : String(n));
try {
for (;;) {
if (signal.aborted) {
- return { posts, complete: false, cursor: String(page) };
+ const cursor = cursorOf(page);
+ return { posts, complete: false, ...(cursor ? { cursor } : {}) };
}
const url = xenforoPageUrl(thread, page);
let load;
try {
load = await fetchPage(url);
} catch (err) {
- return {
- posts,
- complete: false,
- cursor: page === LAST_PAGE_PROBE ? undefined : String(page),
- error: (err as Error).message,
- };
+ const cursor = cursorOf(page);
+ return { posts, complete: false, ...(cursor ? { cursor } : {}), error: (err as Error).message };
+ }
+ if (signal.aborted) {
+ const cursor = cursorOf(page);
+ return { posts, complete: false, ...(cursor ? { cursor } : {}) };
}
- if (signal.aborted) return { posts, complete: false, cursor: page === LAST_PAGE_PROBE ? undefined : String(page) };
const block = classifyForumPage(load.html, load.status);
if (block) {
const why = forumBlockFailure(block, thread.host);
+ const cursor = cursorOf(page);
onLog?.(`Page ${page === LAST_PAGE_PROBE ? "(last)" : page}: ${block.kind} — ${block.detail}.`);
return {
posts,
complete: false,
- ...(page === LAST_PAGE_PROBE ? {} : { cursor: String(page) }),
+ ...(cursor ? { cursor } : {}),
error: why.error,
...(why.needsCookies ? { needsCookies: true } : {}),
};
@@ -213,7 +231,7 @@ export async function walkXenforoThread(
return {
posts,
complete: false,
- cursor: String(page),
+ ...(cursorOf(page) ? { cursor: cursorOf(page) } : {}),
error: `Asked for page ${page} and was shown page ${parsed.page} again; stopped rather than loop.`,
};
}
@@ -221,6 +239,8 @@ export async function walkXenforoThread(
page = parsed.page;
lastPage = Math.max(lastPage ?? 0, parsed.lastPage);
pagesRead++;
+ // Catching up has reached the pages the backfill has yet to read.
+ if (mode === "new" && resumeAt && page <= resumeAt) mode = "backfill";
let sawKnown = false;
let sawOld = false;
@@ -239,10 +259,10 @@ export async function walkXenforoThread(
continue;
}
if (seenIds.has(post.id)) {
- if (stopAtKnown) sawKnown = true;
+ if (mode === "new" && stopAtKnown) sawKnown = true;
continue;
}
- if (watermark && post.createdAt <= watermark) {
+ if (mode === "new" && watermark && post.createdAt <= watermark) {
sawOld = true;
continue;
}
@@ -269,7 +289,18 @@ export async function walkXenforoThread(
};
}
- const next = page - 1;
+ let next = page - 1;
+ if (sawKnown || sawOld) {
+ if (!resumeAt) {
+ if (input.onCheckpoint && fresh.length > 0) await input.onCheckpoint({ posts: fresh, cursor: String(Math.max(1, next)) });
+ else posts.push(...fresh);
+ return { posts, complete: true };
+ }
+ // Caught up: on to the older pages an earlier run left.
+ mode = "backfill";
+ next = Math.min(resumeAt, next);
+ onLog?.(`Caught up with the archive; continuing the older pages at page ${next}.`);
+ }
returnedCount += fresh.length;
if (input.onCheckpoint && fresh.length > 0) {
await input.onCheckpoint({ posts: fresh, cursor: String(Math.max(1, next)) });
@@ -277,9 +308,7 @@ export async function walkXenforoThread(
posts.push(...fresh);
}
- if (sawKnown || sawOld || next < 1) {
- return { posts, complete: true };
- }
+ if (next < 1) return { posts, complete: true };
if (input.pages && pagesRead >= input.pages) {
onLog?.(`Read ${pagesRead} page(s), the cap for this run; the next run continues at page ${next}.`);
return { posts, complete: false, cursor: String(next) };
diff --git a/editor/app/api/ops/fetch-posts/route.test.ts b/editor/app/api/ops/fetch-posts/route.test.ts
@@ -73,3 +73,9 @@ test("an older walk over an account that shows no posts is refused before any jo
// No job was ever written.
assert.deepEqual(await readdir(path.join(ROOT, ".jobs")).catch(() => []), []);
});
+
+test("pages must be a positive whole number", async () => {
+ const res = await post({ slug: SLUG, pages: 0 });
+ assert.equal(res.status, 400);
+ assert.match(res.error, /pages/);
+});
diff --git a/editor/app/api/ops/fetch-posts/route.ts b/editor/app/api/ops/fetch-posts/route.ts
@@ -10,7 +10,10 @@ import {
export const dynamic = "force-dynamic";
-// POST { slug, full?, older?, floor?, force?, limit?, queueKey? } -> { ok: true, jobId }
+// POST { slug, full?, older?, floor?, force?, limit?, pages?, queueKey? } -> { ok: true, jobId }
+//
+// `pages` caps how many pages a run reads, for a channel read page by page (a
+// forum thread: its latest N pages).
//
// A social channel's "Fetch posts" button, over HTTP — and its "Re-fetch full
// history" (`full`) and "Fetch older posts" (`older`, with an optional `floor`
@@ -25,7 +28,7 @@ export const dynamic = "force-dynamic";
export async function POST(request: Request) {
return ops(
request,
- ["slug", "full", "older", "floor", "force", "limit", "queueKey"],
+ ["slug", "full", "older", "floor", "force", "limit", "pages", "queueKey"],
async (body) => {
const slug = reqSlug(body, "slug");
return jobResponse(
@@ -37,6 +40,7 @@ export async function POST(request: Request) {
optBool(body, "older"),
optString(body, "floor"),
optBool(body, "force"),
+ optPositiveInt(body, "pages"),
),
);
},
diff --git a/editor/e2e/social-channel.spec.ts b/editor/e2e/social-channel.spec.ts
@@ -7,8 +7,14 @@
// and no test run can risk a real account.
import { test, expect } from "@playwright/test";
-import { readFile } from "node:fs/promises";
+import { mkdir, mkdtemp, readFile, writeFile } from "node:fs/promises";
+import os from "node:os";
+import path from "node:path";
import { resetData, readJson, pathExists, resolvePath, generateReport } from "./helpers";
+import {
+ THREAD_URL as FORUM_THREAD_URL,
+ threadPage,
+} from "yt-dlp-transcript-common/social/__fixtures__/xenforoPages";
type ChannelConfig = {
handling: string;
@@ -189,3 +195,72 @@ test("fetch-posts writes month-sharded JSONL + a posts-archive, and re-running i
)
.toEqual({ june: 2, archive: 3 });
});
+
+// A FORUM THREAD (XenForo): the URL alone makes a forum-thread channel. The
+// forum is invented (forum.example) and never contacted: the create skips its
+// first fetch, and the posts come in through Import saved pages, from a
+// SYNTHETIC page written to a temp dir.
+test("a forum thread URL makes a forum-thread channel, and saved pages import into it", async ({
+ page,
+}) => {
+ await resetData("empty");
+ await generateReport(page, "new");
+ await page.goto("/channels/new");
+
+ await page.locator('input[name="url"]').fill(`${FORUM_THREAD_URL}page-3`);
+ await expect(page.locator('select[name="platform"]')).toHaveValue("xenforo");
+ await expect(page.locator('input[name="socialHandle"]')).toHaveValue(
+ "the-teapot-collectors-thread.4242",
+ );
+ await expect(page.locator('input[name="slug"]')).toHaveValue("the-teapot-collectors-thread");
+ await expect(page.locator('select[name="audioFormat"]')).toHaveCount(0);
+
+ await page.locator('input[name="name"]').fill("Teapot thread");
+ await page.locator('input[name="fetchPostsNow"]').uncheck();
+ await page.getByRole("button", { name: /create|add channel|save/i }).first().click();
+
+ const configPath = "test-transcripts/channels/the-teapot-collectors-thread/config.json";
+ await expect
+ .poll(async () => ((await pathExists(configPath)) ? await readJson<ChannelConfig>(configPath) : null), {
+ timeout: 15_000,
+ })
+ .toMatchObject({
+ sourceKind: "social",
+ platform: "xenforo",
+ postFetcher: "xenforo-thread",
+ socialHandle: "the-teapot-collectors-thread.4242",
+ });
+
+ const saves = await mkdtemp(path.join(os.tmpdir(), "forum-saves-"));
+ await mkdir(path.join(saves, "The Teapot Collectors Thread_files"), { recursive: true });
+ const post = (id: number, position: number) => ({
+ id,
+ author: `Member${id % 3}`,
+ userId: 10 + (id % 3),
+ ts: 1_760_000_000 + position * 60,
+ position,
+ body: `Words of post ${position}.`,
+ });
+ await writeFile(
+ path.join(saves, "The Teapot Collectors Thread _ Page 2.html"),
+ threadPage({ page: 2, last: 2, saved: true, posts: [post(21, 21), post(22, 22), post(23, 23)] }),
+ );
+
+ await generateReport(page, "the-teapot-collectors-thread");
+ await page.goto("/channels/the-teapot-collectors-thread");
+ await expect(page.locator("[data-forum-session]")).toBeVisible();
+ await expect(page.getByLabel("pages to fetch")).toBeVisible();
+ await page.getByLabel("saved pages path").fill(saves);
+ await page.getByRole("button", { name: "import saved pages" }).click();
+
+ const archive = "test-transcripts/channels/the-teapot-collectors-thread/posts-archive";
+ await expect
+ .poll(
+ async () =>
+ (await pathExists(archive))
+ ? (await readFile(resolvePath(archive), "utf8")).split("\n").filter(Boolean)
+ : [],
+ { timeout: 30_000 },
+ )
+ .toEqual(["xenforo 21", "xenforo 22", "xenforo 23"]);
+});
diff --git a/scripts/archilyzer-ops.mjs b/scripts/archilyzer-ops.mjs
@@ -342,14 +342,16 @@ export function usage() {
' re-walks the whole timeline; "older": true walks back from the oldest',
" archived post through search (X; needs a login), saving its place for",
' the next run, down to "floor": "YYYY-MM-DD" when given. "limit": N caps',
- ' the posts one run reads. "full" and "older" together are refused. An',
+ ' the posts one run reads; "pages": N caps the pages (a forum thread: its',
+ ' latest N pages). "full" and "older" together are refused. An',
' older walk over an account that shows no posts (nothing archived, and',
' the last timeline fetch read none) is refused unless "force": true. A',
" drained fetch stops at its next resume point and the next run resumes.",
"",
- 'capture-posts captures archived posts of a social channel (X): a',
- ' screenshot of each through the connected X profile, and its attached',
- ' media through gallery-dl, into the channel\'s posts-media/<id>/:',
+ 'capture-posts captures archived posts of a social channel (X, forum): a',
+ ' screenshot of each through the connected X profile (a forum thread: its',
+ ' host\'s forum profile), and its attached media through gallery-dl (a',
+ ' forum thread: the same profile), into the channel\'s posts-media/<id>/:',
' {"slug", "ids": [...]}. Every id must be in the channel\'s posts archive.',
' "shots": false or "media": false skips that half; posts already captured',
' are skipped unless "force": true. A post that links to an X Article also',