diff --git a/src/lib/backfill/crawler.ts b/src/lib/backfill/crawler.ts index ba2a6cf..f46e23d 100644 --- a/src/lib/backfill/crawler.ts +++ b/src/lib/backfill/crawler.ts @@ -4,7 +4,7 @@ import { type FirehoseCommit, handleCommitEvent, } from "../../firehose/handlers.ts"; -import { getNoteByAtUri } from "../../server/db/queries/index.ts"; +import { getNoteByAtUri, recordLanded } from "../../server/db/queries/index.ts"; import { COLLECTIONS } from "../collections.ts"; import { resolvePdsClient } from "./pds-client.ts"; import { @@ -12,6 +12,7 @@ import { type BackfillReport, countIngest, countSkip, + countUnchanged, createReport, setEndpoint, } from "./report.ts"; @@ -103,10 +104,11 @@ 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. +// Feed one listRecords row into the firehose handler and classify the *outcome* +// by checking the DB before and after dispatch — the handler silently drops +// records it can't place, so "dispatched" must not be counted as "ingested". +// Revisions keep a parent-missing pre-check for a precise reason; other drops are +// reported generically as "rejected". function dispatchRow( did: string, collection: string, @@ -114,6 +116,14 @@ function dispatchRow( report: BackfillReport, options: BackfillOptions, ): void { + // Dry-run can't observe outcomes (nothing is written), so it only inventories + // what the PDS holds. Skipping the dispatch-dependent parent-missing check here + // also avoids false positives (the parent notes were never written). + if (options.dryRun) { + countIngest(report, did, collection); + return; + } + if (collection === COLLECTIONS.noteRevision) { const noteRef = (row.value as { noteRef?: unknown }).noteRef; const noteExists = options.noteExists ?? defaultNoteExists; @@ -123,17 +133,24 @@ function dispatchRow( } } - if (!options.dryRun) { - const dispatch = options.dispatch ?? handleCommitEvent; - dispatch({ - did, - collection, - rkey: rkeyOf(row.uri), - operation: "create", - record: row.value, - }); + const existedBefore = recordLanded(collection, did, row.uri, row.value); + const dispatch = options.dispatch ?? handleCommitEvent; + dispatch({ + did, + collection, + rkey: rkeyOf(row.uri), + operation: "create", + record: row.value, + }); + + if (!recordLanded(collection, did, row.uri, row.value)) { + // Dispatched but not present afterwards: the handler declined it. + countSkip(report, did, collection, "rejected"); + } else if (existedBefore) { + countUnchanged(report, did, collection); + } else { + countIngest(report, did, collection); } - countIngest(report, did, collection); } // Paginate one (DID, collection) pair to exhaustion, dispatching each record. diff --git a/src/lib/backfill/report.ts b/src/lib/backfill/report.ts index aa103b4..74f8d9f 100644 --- a/src/lib/backfill/report.ts +++ b/src/lib/backfill/report.ts @@ -1,9 +1,14 @@ import { COLLECTIONS } from "../collections.ts"; interface CollectionCounts { + // Newly written rows (existed after dispatch, not before). In dry-run this + // holds the inventory count (records found on the PDS), shown as "found". ingested: number; + // Already present before dispatch — re-run / idempotent no-ops. + unchanged: number; + // Dispatched but not written: the handler dropped it. `reasons` breaks these + // down where the crawler can tell ("parent-missing"); "rejected" otherwise. skipped: number; - // Skip reasons the crawler can observe directly (e.g. "parent-missing"), keyed by reason. reasons: Record; } @@ -43,7 +48,7 @@ function ensureDid(report: BackfillReport, did: string): DidReport { function ensureCollection(d: DidReport, collection: string): CollectionCounts { let c = d.collections.get(collection); if (!c) { - c = { ingested: 0, skipped: 0, reasons: {} }; + c = { ingested: 0, unchanged: 0, skipped: 0, reasons: {} }; d.collections.set(collection, c); } return c; @@ -65,6 +70,14 @@ export function countIngest( ensureCollection(ensureDid(report, did), collection).ingested++; } +export function countUnchanged( + report: BackfillReport, + did: string, + collection: string, +): void { + ensureCollection(ensureDid(report, did), collection).unchanged++; +} + export function countSkip( report: BackfillReport, did: string, @@ -133,10 +146,11 @@ export function formatReport( report: BackfillReport, opts: { dryRun?: boolean } = {}, ): string { + const dryRun = opts.dryRun ?? false; const lines: string[] = []; - if (opts.dryRun) lines.push("DRY RUN — no records dispatched\n"); + if (dryRun) lines.push("DRY RUN — inventory only, nothing written\n"); - const totals = { ingested: 0, skipped: 0 }; + const totals = { ingested: 0, unchanged: 0, skipped: 0 }; for (const did of [...report.dids.keys()].sort()) { const d = report.dids.get(did); @@ -153,11 +167,15 @@ export function formatReport( const c = d.collections.get(collection); if (!c) continue; totals.ingested += c.ingested; + totals.unchanged += c.unchanged; 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)}`, + dryRun + ? ` ${label} ${String(c.ingested).padStart(5)} found` + : ` ${label} ${String(c.ingested).padStart(5)} ingested, ` + + `${String(c.unchanged).padStart(4)} unchanged, ` + + `${String(c.skipped).padStart(4)} skipped${formatReasons(c.reasons)}`, ); } if (cols.length === 0) lines.push(" (no records)"); @@ -165,14 +183,25 @@ export function formatReport( lines.push(""); } + if (dryRun) { + lines.push( + `TOTAL: ${totals.ingested} found` + + (hasFatal(report) ? " — completed with errors" : ""), + ); + lines.push("Run without --dry-run to ingest these records."); + return lines.join("\n"); + } + lines.push( - `TOTAL: ${totals.ingested} ingested, ${totals.skipped} skipped` + + `TOTAL: ${totals.ingested} ingested, ${totals.unchanged} unchanged, ` + + `${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). + // ingested = newly written; unchanged = already present; skipped = the handler + // declined to write it (foreign/unknown wiki, unauthorized, invalid, or + // out-of-order). parent-missing = the referenced note wasn't in the crawl set. lines.push( - 'Note: "ingested" = dispatched to the handler; check logs for handler-side drops.', + 'Note: "skipped" records were rejected by the handler — see warnings above for details.', ); return lines.join("\n"); diff --git a/src/server/db/queries/exists.ts b/src/server/db/queries/exists.ts new file mode 100644 index 0000000..7684c19 --- /dev/null +++ b/src/server/db/queries/exists.ts @@ -0,0 +1,66 @@ +import { COLLECTIONS } from "../../../lib/collections.ts"; +import { getDb } from "../index.ts"; + +// Collections whose row is uniquely identified by its AT-URI (no dedup on write). +const TABLE_BY_ATURI: Record = { + [COLLECTIONS.wiki]: "wikis", + [COLLECTIONS.note]: "notes", + [COLLECTIONS.noteRevision]: "revisions", +}; + +function existsByAtUri(table: string, atUri: string): boolean { + return ( + getDb() + .query(`SELECT 1 FROM ${table} WHERE at_uri = ? LIMIT 1`) + .get(atUri) !== null + ); +} + +function existsByDidWiki( + table: string, + did: string, + wikiAtUri: string, +): boolean { + return ( + getDb() + .query(`SELECT 1 FROM ${table} WHERE did = ? AND wiki_at_uri = ? LIMIT 1`) + .get(did, wikiAtUri) !== null + ); +} + +// Did a dispatched record actually land in the DB? The firehose handler silently +// drops records it can't place or authorize, so "dispatched" never equals +// "ingested" on its own. +// +// AT-URI is the right key for wiki/note/revision (unique, no replace). The +// peripheral collections instead dedup by their logical identity — memberships by +// (wiki, member), requests/bookmarks by (wiki, author) — and rewrite at_uri on +// conflict, so they must be checked by that identity, not by at_uri. +export function recordLanded( + collection: string, + did: string, + atUri: string, + value: unknown, +): boolean { + const byAtUri = TABLE_BY_ATURI[collection]; + if (byAtUri) return existsByAtUri(byAtUri, atUri); + + const v = (value ?? {}) as { wikiRef?: unknown; memberDid?: unknown }; + const wikiRef = typeof v.wikiRef === "string" ? v.wikiRef : null; + if (!wikiRef) return false; + + switch (collection) { + case COLLECTIONS.membership: { + const memberDid = typeof v.memberDid === "string" ? v.memberDid : null; + return ( + memberDid !== null && existsByDidWiki("memberships", memberDid, wikiRef) + ); + } + case COLLECTIONS.memberRequest: + return existsByDidWiki("requests", did, wikiRef); + case COLLECTIONS.bookmark: + return existsByDidWiki("bookmarks", did, wikiRef); + default: + return false; + } +} diff --git a/src/server/db/queries/index.ts b/src/server/db/queries/index.ts index b1597cd..0ae6cd9 100644 --- a/src/server/db/queries/index.ts +++ b/src/server/db/queries/index.ts @@ -14,6 +14,7 @@ export { loadDraft, saveDraftState, } from "./draft.ts"; +export { recordLanded } from "./exists.ts"; export { deleteMembership, deleteMembershipByUri, diff --git a/tests/integration/backfill.test.ts b/tests/integration/backfill.test.ts index 70e733d..cefe32b 100644 --- a/tests/integration/backfill.test.ts +++ b/tests/integration/backfill.test.ts @@ -166,12 +166,58 @@ describe("backfill end-to-end against a fake PDS", () => { 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); + test("re-running the crawl is a no-op and reports unchanged", async () => { + const report = 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); } + // Nothing new landed; everything already present is reported as unchanged. + const rev = report.dids.get(DID)?.collections.get(COLLECTIONS.noteRevision); + const total = fx.notes.reduce((a, n) => a + n.revCount, 0); + expect(rev?.ingested).toBe(0); + expect(rev?.unchanged).toBe(total); + }); +}); + +describe("backfill reports handler drops accurately", () => { + const DROP_DID = "did:plc:backfilldrop00000000000000"; + + test("a note whose wiki was not crawled is reported skipped, not ingested", async () => { + // Crawl only the note collection: its wikiRef points at a wiki that never + // gets ingested, so the handler drops the note. The report must not claim it. + const noteUri = `at://${DROP_DID}/${COLLECTIONS.note}/orphan`; + const report = await runBackfill([DROP_DID], { + throttleMs: 0, + sleep: async () => {}, + collections: [COLLECTIONS.note], + resolveClient: async () => ({ + client: makePdsClient( + "http://pds.fake", + fixtureFetch({ + [COLLECTIONS.note]: [ + { + uri: noteUri, + cid: "c", + value: { + $type: COLLECTIONS.note, + slug: "orphan", + title: "Orphan", + wikiRef: `at://${DROP_DID}/${COLLECTIONS.wiki}/uncrawled`, + createdAt: ISO, + }, + }, + ], + }), + ), + endpoint: "http://pds.fake", + }), + }); + + expect(getNoteByAtUri(noteUri)).toBeNull(); + const counts = report.dids.get(DROP_DID)?.collections.get(COLLECTIONS.note); + expect(counts?.ingested).toBe(0); + expect(counts?.reasons["rejected"]).toBe(1); }); }); diff --git a/tests/lib/backfill/report.test.ts b/tests/lib/backfill/report.test.ts index 830cd7c..6106ecb 100644 --- a/tests/lib/backfill/report.test.ts +++ b/tests/lib/backfill/report.test.ts @@ -3,6 +3,7 @@ import { addError, countIngest, countSkip, + countUnchanged, createReport, formatReport, hasFatal, @@ -13,10 +14,11 @@ import { COLLECTIONS } from "../../../src/lib/collections.ts"; const DID = "did:plc:reporttest"; describe("report aggregation", () => { - test("counts ingested and skipped per collection", () => { + test("counts ingested, unchanged, and skipped per collection", () => { const r = createReport(); countIngest(r, DID, COLLECTIONS.note); countIngest(r, DID, COLLECTIONS.note); + countUnchanged(r, DID, COLLECTIONS.note); countSkip(r, DID, COLLECTIONS.noteRevision, "parent-missing"); countSkip(r, DID, COLLECTIONS.noteRevision, "parent-missing"); countIngest(r, DID, COLLECTIONS.noteRevision); @@ -24,6 +26,7 @@ describe("report aggregation", () => { 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(note?.unchanged).toBe(1); expect(rev?.ingested).toBe(1); expect(rev?.skipped).toBe(2); expect(rev?.reasons["parent-missing"]).toBe(2); @@ -42,20 +45,31 @@ describe("report aggregation", () => { setEndpoint(r, DID, "https://pds.example"); countIngest(r, DID, COLLECTIONS.wiki); countIngest(r, DID, COLLECTIONS.note); + countUnchanged(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"); + expect(out).toContain("TOTAL: 2 ingested, 1 unchanged, 1 skipped"); }); - test("formatReport flags dry run and fatal completion", () => { + test("dry-run report shows an inventory, not ingest/skip counts", () => { const r = createReport(); - addError(r, DID, "unreachable", true); + countIngest(r, DID, COLLECTIONS.note); + countIngest(r, DID, COLLECTIONS.note); const out = formatReport(r, { dryRun: true }); expect(out).toContain("DRY RUN"); + expect(out).toContain("2 found"); + expect(out).toContain("TOTAL: 2 found"); + expect(out).not.toContain("ingested"); + }); + + test("formatReport flags fatal completion", () => { + const r = createReport(); + addError(r, DID, "unreachable", true); + const out = formatReport(r); expect(out).toContain("completed with errors"); }); });