diff --git a/package.json b/package.json index 070788b..752b9fb 100644 --- a/package.json +++ b/package.json @@ -12,6 +12,7 @@ "build:viz": "bun build public/viz/viz-hydrate.ts --outfile public/viz/dist.js", "build:ui": "bun build public/ui.ts --outfile public/ui.js", "build:htmx": "cp node_modules/htmx.org/dist/htmx.min.js public/htmx.min.js", + "backfill": "bun run scripts/backfill.ts", "test": "bun test", "lint": "biome check src/ public/ tests/", "lint:fix": "biome check --write src/ public/ tests/", diff --git a/scripts/backfill.ts b/scripts/backfill.ts new file mode 100644 index 0000000..152198f --- /dev/null +++ b/scripts/backfill.ts @@ -0,0 +1,80 @@ +#!/usr/bin/env bun + +// PDS backfill: crawl wiki.lichen.* records from one or more repos and feed them +// through the firehose handler. The firehose is live-only; this is how we pull in +// records we missed (joined late, another client wrote them, targeted recovery). +// +// bun run scripts/backfill.ts [...] [--dry-run] [--collection ] +// +// Reads/writes the DB at DB_PATH (default lichen.db) — point it at the prod DB on +// the VPS. listRecords is public, so no auth is needed. + +import { runBackfill } from "../src/lib/backfill/crawler.ts"; +import { formatReport, hasFatal } from "../src/lib/backfill/report.ts"; +import { COLLECTIONS } from "../src/lib/collections.ts"; + +interface ParsedArgs { + dids: string[]; + dryRun: boolean; + collection?: string; +} + +function parseArgs(argv: string[]): ParsedArgs { + const dids: string[] = []; + let dryRun = false; + let collection: string | undefined; + + for (let i = 0; i < argv.length; i++) { + const arg = argv[i]; + if (arg === "--dry-run") { + dryRun = true; + } else if (arg === "--collection") { + collection = argv[++i]; + if (!collection) + throw new Error("--collection requires an NSID argument"); + } else if (arg?.startsWith("did:")) { + dids.push(arg); + } else { + throw new Error(`unexpected argument: ${arg}`); + } + } + + return { dids, dryRun, ...(collection ? { collection } : {}) }; +} + +const USAGE = + "Usage: bun run scripts/backfill.ts [...] [--dry-run] [--collection ]"; + +async function main(): Promise { + let parsed: ParsedArgs; + try { + parsed = parseArgs(process.argv.slice(2)); + } catch (err) { + console.error(`${String(err)}\n${USAGE}`); + return 2; + } + + if (parsed.dids.length === 0) { + console.error(`No DIDs given.\n${USAGE}`); + return 2; + } + + if ( + parsed.collection && + !Object.values(COLLECTIONS).includes(parsed.collection as never) + ) { + console.warn( + `Warning: --collection ${parsed.collection} is not a wiki.lichen collection; it will match nothing.`, + ); + } + + const report = await runBackfill(parsed.dids, { + dryRun: parsed.dryRun, + ...(parsed.collection ? { collections: [parsed.collection] } : {}), + }); + + console.log(`\n${formatReport(report, { dryRun: parsed.dryRun })}`); + return hasFatal(report) ? 1 : 0; +} + +process.exit(await main()); diff --git a/src/lib/backfill/crawler.ts b/src/lib/backfill/crawler.ts new file mode 100644 index 0000000..ba2a6cf --- /dev/null +++ b/src/lib/backfill/crawler.ts @@ -0,0 +1,241 @@ +import { type Client, ClientResponseError, ok } from "@atcute/client"; +import type { Did, Nsid } from "@atcute/lexicons/syntax"; +import { + type FirehoseCommit, + handleCommitEvent, +} from "../../firehose/handlers.ts"; +import { getNoteByAtUri } from "../../server/db/queries/index.ts"; +import { COLLECTIONS } from "../collections.ts"; +import { resolvePdsClient } from "./pds-client.ts"; +import { + addError, + type BackfillReport, + countIngest, + countSkip, + createReport, + setEndpoint, +} from "./report.ts"; + +interface BackoffConfig { + baseMs: number; + maxMs: number; + maxRetries: number; +} + +const DEFAULT_BACKOFF: BackoffConfig = { + baseMs: 1000, + maxMs: 30_000, + maxRetries: 5, +}; +const DEFAULT_THROTTLE_MS = 100; +const PAGE_LIMIT = 100; + +export interface BackfillOptions { + /** Log/count what would be ingested without dispatching to the handler. */ + dryRun?: boolean; + /** Restrict the crawl to these collection NSIDs (the --collection flag). */ + collections?: string[]; + /** Politeness delay between paginated calls (ms). */ + throttleMs?: number; + backoff?: BackoffConfig; + // --- Injection seams (tests) --- + resolveClient?: ( + did: string, + ) => Promise<{ client: Client; endpoint: string } | null>; + dispatch?: (commit: FirehoseCommit) => void; + noteExists?: (uri: string) => boolean; + sleep?: (ms: number) => Promise; +} + +const realSleep = (ms: number): Promise => + new Promise((resolve) => setTimeout(resolve, ms)); + +const defaultNoteExists = (uri: string): boolean => + getNoteByAtUri(uri) !== null; + +interface ListRecordsPage { + records: { uri: string; cid: string; value: unknown }[]; + cursor?: string; +} + +// One listRecords page with retry. Retries 429/5xx/network with exponential +// backoff; rethrows other 4xx immediately (caller logs and skips the pair). +async function listPage( + client: Client, + did: string, + collection: string, + cursor: string | undefined, + reverse: boolean, + options: BackfillOptions, +): Promise { + const cfg = options.backoff ?? DEFAULT_BACKOFF; + const sleep = options.sleep ?? realSleep; + let attempt = 0; + + for (;;) { + try { + const res = await ok( + client.get("com.atproto.repo.listRecords", { + params: { + repo: did as Did, + collection: collection as Nsid, + limit: PAGE_LIMIT, + reverse, + ...(cursor ? { cursor } : {}), + }, + }), + ); + return res as ListRecordsPage; + } catch (err) { + const status = err instanceof ClientResponseError ? err.status : 0; // 0 = network/timeout + // Non-429 client errors aren't worth retrying (e.g. deactivated account → 400). + const retryable = status === 0 || status === 429 || status >= 500; + if (!retryable || attempt >= cfg.maxRetries) throw err; + + const delay = Math.min(cfg.baseMs * 2 ** attempt, cfg.maxMs); + await sleep(delay); + attempt++; + } + } +} + +function rkeyOf(uri: string): string { + return uri.split("/").pop() ?? ""; +} + +// Feed one listRecords row into the firehose handler, classifying the outcome. +// Revisions get a parent-missing pre-check (option (a) from the design): the +// handler would silently no-op on an unknown note, so detect and count it here +// rather than reimplement any handler logic. +function dispatchRow( + did: string, + collection: string, + row: { uri: string; value: unknown }, + report: BackfillReport, + options: BackfillOptions, +): void { + if (collection === COLLECTIONS.noteRevision) { + const noteRef = (row.value as { noteRef?: unknown }).noteRef; + const noteExists = options.noteExists ?? defaultNoteExists; + if (typeof noteRef !== "string" || !noteExists(noteRef)) { + countSkip(report, did, collection, "parent-missing"); + return; + } + } + + if (!options.dryRun) { + const dispatch = options.dispatch ?? handleCommitEvent; + dispatch({ + did, + collection, + rkey: rkeyOf(row.uri), + operation: "create", + record: row.value, + }); + } + countIngest(report, did, collection); +} + +// Paginate one (DID, collection) pair to exhaustion, dispatching each record. +async function crawlCollection( + client: Client, + did: string, + collection: string, + reverse: boolean, + report: BackfillReport, + options: BackfillOptions, +): Promise { + const sleep = options.sleep ?? realSleep; + const throttleMs = options.throttleMs ?? DEFAULT_THROTTLE_MS; + let cursor: string | undefined; + + do { + let page: ListRecordsPage; + try { + page = await listPage(client, did, collection, cursor, reverse, options); + } catch (err) { + const status = err instanceof ClientResponseError ? err.status : 0; + // 4xx (other than 429): the PDS reachable but rejected — skip this pair, + // not a fatal DID failure. Network/5xx exhaustion is fatal. + const fatal = status === 0 || status === 429 || status >= 500; + addError(report, did, `${collection}: ${String(err)}`, fatal); + return; + } + + for (const row of page.records) { + dispatchRow(did, collection, row, report, options); + } + + cursor = page.cursor; + // Stop on an absent cursor or an empty page (defends against PDSes that + // echo a stable cursor forever). + if (cursor && page.records.length > 0) { + await sleep(throttleMs); + } else { + cursor = undefined; + } + } while (cursor); +} + +interface Stage { + collections: string[]; + reverse: boolean; +} + +// Cross-collection ordering matters: revisions reference notes, notes/memberships +// reference wikis. Crawl every DID through each stage before advancing. +const STAGES: Stage[] = [ + { collections: [COLLECTIONS.wiki], reverse: false }, + { collections: [COLLECTIONS.note], reverse: false }, + // reverse: true → ascending TID (oldest first) so the diff chain applies in order. + { collections: [COLLECTIONS.noteRevision], reverse: true }, + { + collections: [ + COLLECTIONS.membership, + COLLECTIONS.memberRequest, + COLLECTIONS.bookmark, + ], + reverse: false, + }, +]; + +export async function runBackfill( + dids: string[], + options: BackfillOptions = {}, +): Promise { + const report = createReport(); + const resolveClient = options.resolveClient ?? resolvePdsClient; + const allowed = options.collections ? new Set(options.collections) : null; + + // Resolve every PDS up front; a DID with no reachable PDS is fatal but doesn't + // abort the others. + const clients = new Map(); + for (const did of dids) { + const resolved = await resolveClient(did); + if (!resolved) { + setEndpoint(report, did, null); + addError(report, did, "could not resolve PDS endpoint", true); + continue; + } + setEndpoint(report, did, resolved.endpoint); + clients.set(did, resolved.client); + } + + for (const stage of STAGES) { + for (const [did, client] of clients) { + for (const collection of stage.collections) { + if (allowed && !allowed.has(collection)) continue; + await crawlCollection( + client, + did, + collection, + stage.reverse, + report, + options, + ); + } + } + } + + return report; +} diff --git a/src/lib/backfill/pds-client.ts b/src/lib/backfill/pds-client.ts new file mode 100644 index 0000000..c887edb --- /dev/null +++ b/src/lib/backfill/pds-client.ts @@ -0,0 +1,32 @@ +import { Client, simpleFetchHandler } from "@atcute/client"; +import { resolvePdsEndpoint } from "../identity.ts"; + +// `com.atproto.repo.listRecords` is a public, unauthenticated query, so the +// backfill crawler talks to each PDS with a bare fetch handler — no OAuth session, +// no app password. `fetchImpl` is an injection seam for tests (a canned fetch). +export function makePdsClient( + endpoint: string, + fetchImpl?: typeof fetch, +): Client { + return new Client({ + handler: simpleFetchHandler({ + service: endpoint, + ...(fetchImpl ? { fetch: fetchImpl } : {}), + }), + }); +} + +export interface ResolvedPdsClient { + client: Client; + endpoint: string; +} + +// Resolve a DID's PDS endpoint and build an unauthenticated client for it. +// Returns null when the DID document can't be resolved or exposes no PDS. +export async function resolvePdsClient( + did: string, +): Promise { + const endpoint = await resolvePdsEndpoint(did); + if (!endpoint) return null; + return { client: makePdsClient(endpoint), endpoint }; +} diff --git a/src/lib/backfill/report.ts b/src/lib/backfill/report.ts new file mode 100644 index 0000000..aa103b4 --- /dev/null +++ b/src/lib/backfill/report.ts @@ -0,0 +1,179 @@ +import { COLLECTIONS } from "../collections.ts"; + +interface CollectionCounts { + ingested: number; + skipped: number; + // Skip reasons the crawler can observe directly (e.g. "parent-missing"), keyed by reason. + reasons: Record; +} + +interface DidReport { + did: string; + endpoint: string | null; + collections: Map; + // Soft notes (4xx skips) and fatal failures (unreachable PDS, retry exhaustion). + errors: string[]; + // True when the DID couldn't be crawled to completion — drives the CLI exit code. + fatal: boolean; +} + +export interface BackfillReport { + dids: Map; +} + +export function createReport(): BackfillReport { + return { dids: new Map() }; +} + +function ensureDid(report: BackfillReport, did: string): DidReport { + let d = report.dids.get(did); + if (!d) { + d = { + did, + endpoint: null, + collections: new Map(), + errors: [], + fatal: false, + }; + report.dids.set(did, d); + } + return d; +} + +function ensureCollection(d: DidReport, collection: string): CollectionCounts { + let c = d.collections.get(collection); + if (!c) { + c = { ingested: 0, skipped: 0, reasons: {} }; + d.collections.set(collection, c); + } + return c; +} + +export function setEndpoint( + report: BackfillReport, + did: string, + endpoint: string | null, +): void { + ensureDid(report, did).endpoint = endpoint; +} + +export function countIngest( + report: BackfillReport, + did: string, + collection: string, +): void { + ensureCollection(ensureDid(report, did), collection).ingested++; +} + +export function countSkip( + report: BackfillReport, + did: string, + collection: string, + reason: string, +): void { + const c = ensureCollection(ensureDid(report, did), collection); + c.skipped++; + c.reasons[reason] = (c.reasons[reason] ?? 0) + 1; +} + +export function addError( + report: BackfillReport, + did: string, + message: string, + fatal = false, +): void { + const d = ensureDid(report, did); + d.errors.push(message); + if (fatal) d.fatal = true; +} + +// Non-zero CLI exit when any DID failed to crawl to completion. +export function hasFatal(report: BackfillReport): boolean { + for (const d of report.dids.values()) { + if (d.fatal) return true; + } + return false; +} + +// Stable display order; anything unexpected sorts to the end alphabetically. +const COLLECTION_ORDER: string[] = [ + COLLECTIONS.wiki, + COLLECTIONS.note, + COLLECTIONS.noteRevision, + COLLECTIONS.membership, + COLLECTIONS.memberRequest, + COLLECTIONS.bookmark, +]; + +function shortName(collection: string): string { + return collection.startsWith("wiki.lichen.") + ? collection.slice("wiki.lichen.".length) + : collection; +} + +function orderedCollections(d: DidReport): string[] { + return [...d.collections.keys()].sort((a, b) => { + const ia = COLLECTION_ORDER.indexOf(a); + const ib = COLLECTION_ORDER.indexOf(b); + if (ia !== -1 && ib !== -1) return ia - ib; + if (ia !== -1) return -1; + if (ib !== -1) return 1; + return a.localeCompare(b); + }); +} + +function formatReasons(reasons: Record): string { + const parts = Object.entries(reasons) + .sort(([a], [b]) => a.localeCompare(b)) + .map(([reason, n]) => `${n} ${reason}`); + return parts.length > 0 ? ` (${parts.join(", ")})` : ""; +} + +export function formatReport( + report: BackfillReport, + opts: { dryRun?: boolean } = {}, +): string { + const lines: string[] = []; + if (opts.dryRun) lines.push("DRY RUN — no records dispatched\n"); + + const totals = { ingested: 0, skipped: 0 }; + + for (const did of [...report.dids.keys()].sort()) { + const d = report.dids.get(did); + if (!d) continue; + lines.push(`DID ${did}`); + if (d.endpoint) lines.push(` pds: ${d.endpoint}`); + + const cols = orderedCollections(d); + const labelWidth = Math.max( + 0, + ...cols.map((c) => `${shortName(c)}:`.length), + ); + for (const collection of cols) { + const c = d.collections.get(collection); + if (!c) continue; + totals.ingested += c.ingested; + totals.skipped += c.skipped; + const label = `${shortName(collection)}:`.padEnd(labelWidth); + lines.push( + ` ${label} ${String(c.ingested).padStart(5)} ingested, ` + + `${String(c.skipped).padStart(4)} skipped${formatReasons(c.reasons)}`, + ); + } + if (cols.length === 0) lines.push(" (no records)"); + for (const err of d.errors) lines.push(` ! ${err}`); + lines.push(""); + } + + lines.push( + `TOTAL: ${totals.ingested} ingested, ${totals.skipped} skipped` + + (hasFatal(report) ? " — completed with errors" : ""), + ); + // "ingested" counts records dispatched to the firehose handler; the handler may + // still drop malformed/oversized/unauthorized records (see warnings above). + lines.push( + 'Note: "ingested" = dispatched to the handler; check logs for handler-side drops.', + ); + + return lines.join("\n"); +} diff --git a/src/types.d.ts b/src/types.d.ts index 938f079..cf4fbb3 100644 --- a/src/types.d.ts +++ b/src/types.d.ts @@ -1,4 +1,5 @@ declare const process: { + argv: string[]; env: Record; exit(code?: number): never; on(event: string, listener: (...args: unknown[]) => void): void; diff --git a/tests/integration/backfill.test.ts b/tests/integration/backfill.test.ts new file mode 100644 index 0000000..70e733d --- /dev/null +++ b/tests/integration/backfill.test.ts @@ -0,0 +1,216 @@ +import { afterAll, beforeAll, describe, expect, test } from "bun:test"; +import * as TID from "@atcute/tid"; +import { runBackfill } from "../../src/lib/backfill/crawler.ts"; +import { makePdsClient } from "../../src/lib/backfill/pds-client.ts"; +import { COLLECTIONS } from "../../src/lib/collections.ts"; +import { createDiff } from "../../src/lib/diff.ts"; +import { + getCurrentNoteContent, + getNoteByAtUri, + getWikiByAtUri, + listRevisions, +} from "../../src/server/db/queries/index.ts"; +import { cleanupWikiAndDependents } from "../helpers/cleanup.ts"; + +const DID = "did:plc:backfillitest000000000000"; +const WIKI_SLUG = "backfill-itest"; +const WIKI_URI = `at://${DID}/${COLLECTIONS.wiki}/${WIKI_SLUG}`; +const ISO = "2026-01-01T00:00:00.000Z"; + +type Row = { uri: string; cid: string; value: unknown }; + +function jsonResponse(body: unknown): Response { + return new Response(JSON.stringify(body), { + status: 200, + headers: { "content-type": "application/json" }, + }); +} + +// One-page listRecords shim keyed by the `collection` query param. +function fixtureFetch(byCollection: Record): typeof fetch { + const impl = async (url: string | URL): Promise => { + const collection = + new URL(String(url)).searchParams.get("collection") ?? ""; + return jsonResponse({ records: byCollection[collection] ?? [] }); + }; + return impl as unknown as typeof fetch; +} + +const NOTE_SPECS = [ + { + slug: "alpha", + contents: [ + "# Alpha\n\none", + "# Alpha\n\none two", + "# Alpha\n\none two three", + ], + }, + { slug: "beta", contents: ["beta v1", "beta v2"] }, + { slug: "gamma", contents: ["gamma only"] }, +]; + +interface FixtureNote { + slug: string; + uri: string; + finalContent: string; + revCount: number; +} + +interface Fixture { + byCollection: Record; + notes: FixtureNote[]; +} + +function buildFixture(): Fixture { + const wiki: Row = { + uri: WIKI_URI, + cid: "wikicid", + value: { + $type: COLLECTIONS.wiki, + name: "Backfill ITest", + visibility: "public", + language: "en", + createdAt: ISO, + }, + }; + const noteRows: Row[] = []; + const revisions: Row[] = []; + const notes: FixtureNote[] = []; + + for (const spec of NOTE_SPECS) { + const noteUri = `at://${DID}/${COLLECTIONS.note}/${TID.now()}`; + noteRows.push({ + uri: noteUri, + cid: "notecid", + value: { + $type: COLLECTIONS.note, + slug: spec.slug, + title: spec.slug, + wikiRef: WIKI_URI, + createdAt: ISO, + }, + }); + + let parent: string | null = null; + let prevContent = ""; + for (const content of spec.contents) { + const revUri = `at://${DID}/${COLLECTIONS.noteRevision}/${TID.now()}`; + const value: Record = { + $type: COLLECTIONS.noteRevision, + noteRef: noteUri, + diff: createDiff(prevContent, content), + diffFormat: "diff-match-patch", + createdAt: ISO, + }; + if (parent) value["parentRevision"] = parent; + revisions.push({ uri: revUri, cid: "revcid", value }); + parent = revUri; + prevContent = content; + } + + notes.push({ + slug: spec.slug, + uri: noteUri, + finalContent: spec.contents[spec.contents.length - 1] ?? "", + revCount: spec.contents.length, + }); + } + + return { + byCollection: { + [COLLECTIONS.wiki]: [wiki], + [COLLECTIONS.note]: noteRows, + [COLLECTIONS.noteRevision]: revisions, + }, + notes, + }; +} + +function crawl(byCollection: Record, dids = [DID]) { + return runBackfill(dids, { + throttleMs: 0, + sleep: async () => {}, + resolveClient: async () => ({ + client: makePdsClient("http://pds.fake", fixtureFetch(byCollection)), + endpoint: "http://pds.fake", + }), + }); +} + +describe("backfill end-to-end against a fake PDS", () => { + const fx = buildFixture(); + + beforeAll(() => { + cleanupWikiAndDependents(WIKI_SLUG); + }); + afterAll(() => { + cleanupWikiAndDependents(WIKI_SLUG); + }); + + test("ingests wiki, notes, and revision chains into the DB", async () => { + const report = await crawl(fx.byCollection); + + expect(getWikiByAtUri(WIKI_URI)?.name).toBe("Backfill ITest"); + + for (const note of fx.notes) { + expect(getNoteByAtUri(note.uri)?.slug).toBe(note.slug); + expect(getCurrentNoteContent(note.uri)).toBe(note.finalContent); + expect(listRevisions(note.uri).length).toBe(note.revCount); + } + + // No fatal errors; revision counts add up. + expect(report.dids.get(DID)?.fatal).toBe(false); + const revIngested = report.dids + .get(DID) + ?.collections.get(COLLECTIONS.noteRevision)?.ingested; + expect(revIngested).toBe(fx.notes.reduce((a, n) => a + n.revCount, 0)); + }); + + test("re-running the crawl is a no-op", async () => { + await crawl(fx.byCollection); + for (const note of fx.notes) { + expect(listRevisions(note.uri).length).toBe(note.revCount); + expect(getCurrentNoteContent(note.uri)).toBe(note.finalContent); + } + }); +}); + +describe("backfill parent-missing handling", () => { + test("revision pointing at an uncrawled note is counted, not ingested", async () => { + const missingNoteUri = `at://${DID}/${COLLECTIONS.note}/nonexistent`; + const revUri = `at://${DID}/${COLLECTIONS.noteRevision}/orphan`; + const report = await runBackfill([DID], { + throttleMs: 0, + sleep: async () => {}, + collections: [COLLECTIONS.noteRevision], + resolveClient: async () => ({ + client: makePdsClient( + "http://pds.fake", + fixtureFetch({ + [COLLECTIONS.noteRevision]: [ + { + uri: revUri, + cid: "c", + value: { + $type: COLLECTIONS.noteRevision, + noteRef: missingNoteUri, + diff: createDiff("", "orphan"), + diffFormat: "diff-match-patch", + createdAt: ISO, + }, + }, + ], + }), + ), + endpoint: "http://pds.fake", + }), + }); + + expect(getNoteByAtUri(missingNoteUri)).toBeNull(); + const counts = report.dids + .get(DID) + ?.collections.get(COLLECTIONS.noteRevision); + expect(counts?.ingested).toBe(0); + expect(counts?.reasons["parent-missing"]).toBe(1); + }); +}); diff --git a/tests/lib/backfill/crawler.test.ts b/tests/lib/backfill/crawler.test.ts new file mode 100644 index 0000000..d6e404d --- /dev/null +++ b/tests/lib/backfill/crawler.test.ts @@ -0,0 +1,290 @@ +import { describe, expect, test } from "bun:test"; +import type { FirehoseCommit } from "../../../src/firehose/handlers.ts"; +import { + type BackfillOptions, + runBackfill, +} from "../../../src/lib/backfill/crawler.ts"; +import { makePdsClient } from "../../../src/lib/backfill/pds-client.ts"; +import { COLLECTIONS } from "../../../src/lib/collections.ts"; + +interface Page { + records: { uri: string; cid: string; value: unknown }[]; + cursor?: string; +} + +function rec( + collection: string, + n: number, +): { + uri: string; + cid: string; + value: unknown; +} { + return { + uri: `at://did:plc:x/${collection}/rk${n}`, + cid: `cid${n}`, + value: { $type: collection, n }, + }; +} + +function jsonResponse(body: unknown, status = 200): Response { + return new Response(JSON.stringify(body), { + status, + headers: { "content-type": "application/json" }, + }); +} + +// A fetch shim returning a fixed sequence of pages, recording each request URL. +function sequenceFetch(pages: Page[]): { + fetch: typeof fetch; + urls: string[]; +} { + const urls: string[] = []; + let i = 0; + const fetchImpl = async (url: string | URL): Promise => { + urls.push(String(url)); + return jsonResponse(pages[i++] ?? { records: [] }); + }; + return { fetch: fetchImpl as unknown as typeof fetch, urls }; +} + +// A fetch shim returning a status (and empty records on 200) keyed per call, +// recording delays would-be slept. +function statusFetch(statuses: number[]): { + fetch: typeof fetch; + calls: number; +} { + let i = 0; + const state = { calls: 0 }; + const fetchImpl = async (): Promise => { + state.calls++; + const status = statuses[i++] ?? 200; + if (status === 200) return jsonResponse({ records: [] }); + return jsonResponse({ error: "Err", message: "nope" }, status); + }; + return { + fetch: fetchImpl as unknown as typeof fetch, + get calls() { + return state.calls; + }, + }; +} + +const noWait: Pick = { + sleep: async () => {}, + throttleMs: 0, +}; + +function spyDispatch(): { + dispatch: (c: FirehoseCommit) => void; + commits: FirehoseCommit[]; +} { + const commits: FirehoseCommit[] = []; + return { dispatch: (c) => commits.push(c), commits }; +} + +describe("crawler pagination", () => { + test("follows cursor hops and stops when cursor absent", async () => { + const { fetch, urls } = sequenceFetch([ + { records: [rec(COLLECTIONS.wiki, 1)], cursor: "c1" }, + { records: [rec(COLLECTIONS.wiki, 2)], cursor: "c2" }, + { records: [rec(COLLECTIONS.wiki, 3)] }, // no cursor → terminate + ]); + const spy = spyDispatch(); + await runBackfill(["did:plc:x"], { + ...noWait, + collections: [COLLECTIONS.wiki], + dispatch: spy.dispatch, + resolveClient: async () => ({ + client: makePdsClient("http://pds.fake", fetch), + endpoint: "http://pds.fake", + }), + }); + + expect(spy.commits.length).toBe(3); + expect(urls.length).toBe(3); + expect(urls[1]).toContain("cursor=c1"); + expect(urls[2]).toContain("cursor=c2"); + }); + + test("stops on an empty page even if a cursor is echoed", async () => { + const { fetch, urls } = sequenceFetch([ + { records: [rec(COLLECTIONS.wiki, 1)], cursor: "c1" }, + { records: [], cursor: "c2" }, // empty → terminate despite cursor + ]); + const spy = spyDispatch(); + await runBackfill(["did:plc:x"], { + ...noWait, + collections: [COLLECTIONS.wiki], + dispatch: spy.dispatch, + resolveClient: async () => ({ + client: makePdsClient("http://pds.fake", fetch), + endpoint: "http://pds.fake", + }), + }); + + expect(spy.commits.length).toBe(1); + expect(urls.length).toBe(2); + }); +}); + +describe("crawler backoff", () => { + test("retries 429 with exponential backoff then succeeds", async () => { + const sf = statusFetch([429, 429, 200]); + const delays: number[] = []; + await runBackfill(["did:plc:x"], { + collections: [COLLECTIONS.wiki], + throttleMs: 0, + sleep: async (ms) => { + delays.push(ms); + }, + backoff: { baseMs: 10, maxMs: 1000, maxRetries: 5 }, + dispatch: () => {}, + resolveClient: async () => ({ + client: makePdsClient("http://pds.fake", sf.fetch), + endpoint: "http://pds.fake", + }), + }); + + expect(sf.calls).toBe(3); + expect(delays).toEqual([10, 20]); + }); + + test("gives up after maxRetries and marks the DID fatal", async () => { + const sf = statusFetch([429, 429, 429, 429]); + const spy = spyDispatch(); + const report = await runBackfill(["did:plc:x"], { + collections: [COLLECTIONS.wiki], + throttleMs: 0, + sleep: async () => {}, + backoff: { baseMs: 1, maxMs: 10, maxRetries: 2 }, + dispatch: spy.dispatch, + resolveClient: async () => ({ + client: makePdsClient("http://pds.fake", sf.fetch), + endpoint: "http://pds.fake", + }), + }); + + expect(spy.commits.length).toBe(0); + expect(report.dids.get("did:plc:x")?.fatal).toBe(true); + }); + + test("non-429 4xx skips the pair without a fatal failure", async () => { + const sf = statusFetch([400]); + const report = await runBackfill(["did:plc:x"], { + ...noWait, + collections: [COLLECTIONS.wiki], + backoff: { baseMs: 1, maxMs: 10, maxRetries: 3 }, + dispatch: () => {}, + resolveClient: async () => ({ + client: makePdsClient("http://pds.fake", sf.fetch), + endpoint: "http://pds.fake", + }), + }); + + expect(sf.calls).toBe(1); // no retries on 400 + expect(report.dids.get("did:plc:x")?.fatal).toBe(false); + expect(report.dids.get("did:plc:x")?.errors.length).toBe(1); + }); +}); + +describe("crawler staging", () => { + test("finishes every DID's wikis before any notes (stage ordering)", async () => { + // One empty page per (did, collection); record the collection request order. + const order: string[] = []; + const fetchImpl = async (url: string | URL): Promise => { + const collection = + new URL(String(url)).searchParams.get("collection") ?? ""; + order.push(collection); + return jsonResponse({ records: [] }); + }; + await runBackfill(["did:plc:a", "did:plc:b"], { + ...noWait, + dispatch: () => {}, + resolveClient: async () => ({ + client: makePdsClient( + "http://pds.fake", + fetchImpl as unknown as typeof fetch, + ), + endpoint: "http://pds.fake", + }), + }); + + const lastWiki = order.lastIndexOf(COLLECTIONS.wiki); + const firstNote = order.indexOf(COLLECTIONS.note); + const lastNote = order.lastIndexOf(COLLECTIONS.note); + const firstRev = order.indexOf(COLLECTIONS.noteRevision); + expect(lastWiki).toBeLessThan(firstNote); + expect(lastNote).toBeLessThan(firstRev); + // Both DIDs crawled in each stage. + expect(order.filter((c) => c === COLLECTIONS.wiki).length).toBe(2); + }); +}); + +describe("crawler dispatch classification", () => { + test("revision with missing note is counted parent-missing, not dispatched", async () => { + const fetchImpl = async (): Promise => + jsonResponse({ + records: [ + { + uri: `at://did:plc:x/${COLLECTIONS.noteRevision}/rk1`, + cid: "c1", + value: { + $type: COLLECTIONS.noteRevision, + noteRef: "at://did:plc:x/wiki.lichen.note/missing", + }, + }, + ], + }); + const spy = spyDispatch(); + const report = await runBackfill(["did:plc:x"], { + ...noWait, + collections: [COLLECTIONS.noteRevision], + dispatch: spy.dispatch, + noteExists: () => false, + resolveClient: async () => ({ + client: makePdsClient( + "http://pds.fake", + fetchImpl as unknown as typeof fetch, + ), + endpoint: "http://pds.fake", + }), + }); + + expect(spy.commits.length).toBe(0); + const counts = report.dids + .get("did:plc:x") + ?.collections.get(COLLECTIONS.noteRevision); + expect(counts?.ingested).toBe(0); + expect(counts?.reasons["parent-missing"]).toBe(1); + }); + + test("dry run counts but does not dispatch", async () => { + const { fetch } = sequenceFetch([{ records: [rec(COLLECTIONS.wiki, 1)] }]); + const spy = spyDispatch(); + const report = await runBackfill(["did:plc:x"], { + ...noWait, + dryRun: true, + collections: [COLLECTIONS.wiki], + dispatch: spy.dispatch, + resolveClient: async () => ({ + client: makePdsClient("http://pds.fake", fetch), + endpoint: "http://pds.fake", + }), + }); + + expect(spy.commits.length).toBe(0); + expect( + report.dids.get("did:plc:x")?.collections.get(COLLECTIONS.wiki)?.ingested, + ).toBe(1); + }); + + test("unresolvable DID is fatal", async () => { + const report = await runBackfill(["did:plc:gone"], { + ...noWait, + dispatch: () => {}, + resolveClient: async () => null, + }); + expect(report.dids.get("did:plc:gone")?.fatal).toBe(true); + }); +}); diff --git a/tests/lib/backfill/report.test.ts b/tests/lib/backfill/report.test.ts new file mode 100644 index 0000000..830cd7c --- /dev/null +++ b/tests/lib/backfill/report.test.ts @@ -0,0 +1,61 @@ +import { describe, expect, test } from "bun:test"; +import { + addError, + countIngest, + countSkip, + createReport, + formatReport, + hasFatal, + setEndpoint, +} from "../../../src/lib/backfill/report.ts"; +import { COLLECTIONS } from "../../../src/lib/collections.ts"; + +const DID = "did:plc:reporttest"; + +describe("report aggregation", () => { + test("counts ingested and skipped per collection", () => { + const r = createReport(); + countIngest(r, DID, COLLECTIONS.note); + countIngest(r, DID, COLLECTIONS.note); + countSkip(r, DID, COLLECTIONS.noteRevision, "parent-missing"); + countSkip(r, DID, COLLECTIONS.noteRevision, "parent-missing"); + countIngest(r, DID, COLLECTIONS.noteRevision); + + const note = r.dids.get(DID)?.collections.get(COLLECTIONS.note); + const rev = r.dids.get(DID)?.collections.get(COLLECTIONS.noteRevision); + expect(note?.ingested).toBe(2); + expect(rev?.ingested).toBe(1); + expect(rev?.skipped).toBe(2); + expect(rev?.reasons["parent-missing"]).toBe(2); + }); + + test("hasFatal is false for soft errors, true for fatal", () => { + const r = createReport(); + addError(r, DID, "400 skip", false); + expect(hasFatal(r)).toBe(false); + addError(r, DID, "unreachable", true); + expect(hasFatal(r)).toBe(true); + }); + + test("formatReport renders totals, reasons, and endpoint", () => { + const r = createReport(); + setEndpoint(r, DID, "https://pds.example"); + countIngest(r, DID, COLLECTIONS.wiki); + countIngest(r, DID, COLLECTIONS.note); + countSkip(r, DID, COLLECTIONS.noteRevision, "parent-missing"); + + const out = formatReport(r); + expect(out).toContain(DID); + expect(out).toContain("https://pds.example"); + expect(out).toContain("parent-missing"); + expect(out).toContain("TOTAL: 2 ingested, 1 skipped"); + }); + + test("formatReport flags dry run and fatal completion", () => { + const r = createReport(); + addError(r, DID, "unreachable", true); + const out = formatReport(r, { dryRun: true }); + expect(out).toContain("DRY RUN"); + expect(out).toContain("completed with errors"); + }); +});