From 87d7526e89700af0952b6e3a00b497875fc6ab0d Mon Sep 17 00:00:00 2001 From: Florian <45694132+flo-bit@users.noreply.github.com> Date: Mon, 23 Mar 2026 02:22:33 +0100 Subject: [PATCH] improve backfilling --- README.md | 49 ++++---- scripts/sync.ts | 152 ++++++++++++++++------- src/core/backfill.ts | 252 +++++++++++++++++++++++++++++++++------ src/core/client.ts | 46 ++++++- src/core/db/records.ts | 6 +- src/core/router/admin.ts | 77 ------------ wrangler.jsonc | 2 +- 7 files changed, 393 insertions(+), 191 deletions(-) diff --git a/README.md b/README.md index 673a0a3..a783bf5 100644 --- a/README.md +++ b/README.md @@ -7,28 +7,6 @@ Define collections — get automatic Jetstream ingestion, PDS backfill, user dis ## Quickstart -### Dev - -```bash -pnpm install -# Edit src/config.ts with your collections -pnpm generate:pull # pull lexicons from network, auto-detect fields, generate types -pnpm dev:auto # start wrangler dev with auto-ingestion, leave running while you sync -pnpm sync # in a different terminal, discover users + backfill records from PDS -``` - -### Production - -```bash -npx wrangler d1 create contrail -# Add database_id to wrangler.toml -pnpm deploy -# to sync in production, run it locally but set your d1 to remote, then run -pnpm sync -``` - -Ingestion runs automatically via cron (`*/1 * * * *`). Schema is auto-initialized. - ## Config Edit `src/config.ts` — this is the only file you need to touch: @@ -46,13 +24,32 @@ export const config: ContrailConfig = { }, }, "community.lexicon.calendar.rsvp": {}, - }, - // profiles: ["app.bsky.actor.profile"], ← default - // jetstreams: [...] ← default: 4 Bluesky jetstream endpoints - // relays: [...] ← default: 2 Bluesky relay endpoints + } }; ``` +### Dev + +```bash +pnpm install +pnpm generate:pull # pull lexicons from network, auto-detect fields, generate types +pnpm sync # discover users and backfill records from PDS +pnpm dev:auto # start wrangler dev with auto-ingestion +``` + +### Production + +```bash +npx wrangler d1 create contrail +# Add database_id to wrangler.toml +pnpm deploy +# to sync in production, run it locally but set your d1 to remote, then run +pnpm sync +``` + +Ingestion runs automatically via cron (`*/1 * * * *`). Schema is auto-initialized. + + ### What's auto-detected from lexicons When you run `pnpm generate`, queryable fields are derived from each collection's lexicon: diff --git a/scripts/sync.ts b/scripts/sync.ts index 3d17831..9806ad7 100644 --- a/scripts/sync.ts +++ b/scripts/sync.ts @@ -1,62 +1,124 @@ /** * Sync: discover users from relays and backfill their records from PDS. - * Calls the sync endpoint in a loop until done. + * Runs directly against D1 via wrangler bindings — no dev server needed. * - * Usage: npx tsx scripts/sync.ts [base_url] [admin_secret] + * Usage: + * npx tsx scripts/sync.ts # local D1 + * npx tsx scripts/sync.ts --remote # prod D1 */ -import { config } from "../src/config"; +import { getPlatformProxy } from "wrangler"; +import { config as rawConfig } from "../src/config"; +import { resolveConfig, validateConfig, getCollectionNames } from "../src/core/types"; +import { initSchema } from "../src/core/db"; +import { discoverDIDs, backfillAll } from "../src/core/backfill"; -async function main() { - const base = process.argv[2] || "http://localhost:8787"; - const secret = process.argv[3]; - const ns = config.namespace; +const config = resolveConfig(rawConfig); +validateConfig(config); - const headers: Record = {}; - if (secret) { - headers["Authorization"] = `Bearer ${secret}`; - } +function elapsed(start: number): string { + const ms = Date.now() - start; + if (ms < 1000) return `${ms}ms`; + if (ms < 60_000) return `${(ms / 1000).toFixed(1)}s`; + const mins = Math.floor(ms / 60_000); + const secs = ((ms % 60_000) / 1000).toFixed(0); + return `${mins}m ${secs}s`; +} - console.log("=== Syncing (discover + backfill) ==="); +async function main() { + const remote = process.argv.includes("--remote"); + const syncStart = Date.now(); - let staleCount = 0; - let lastRemaining = -1; + console.log(`=== Sync (${remote ? "remote/prod" : "local"} D1) ===\n`); - while (true) { - const res = await fetch(`${base}/xrpc/${ns}.admin.sync`, { - headers, - signal: AbortSignal.timeout(60_000), - }); - const result = (await res.json()) as { - discovered: number; - backfilled: number; - remaining: number; - done: boolean; - }; - - console.log( - ` Discovered ${result.discovered}, backfilled ${result.backfilled} records, ${result.remaining} remaining` - ); + const { env, dispose } = await getPlatformProxy<{ DB: D1Database }>({ + environment: remote ? "production" : undefined, + }); + const db = env.DB; + + try { + await initSchema(db, config); - if (result.done) { - console.log(" Sync complete."); - break; + // Phase 1: Discover all DIDs from relays + console.log("--- Discovery ---"); + const discoveryStart = Date.now(); + const allDiscovered = new Set(); + while (true) { + const dids = await discoverDIDs(db, config, Infinity); + if (dids.length === 0) break; + for (const did of dids) allDiscovered.add(did); + console.log(` Found ${allDiscovered.size} unique users so far`); } + console.log(` Done: ${allDiscovered.size} users in ${elapsed(discoveryStart)}\n`); - // Detect stuck state — no progress after several attempts - if (result.remaining === lastRemaining && result.discovered === 0 && result.backfilled === 0) { - staleCount++; - if (staleCount >= 3) { - console.log(" No progress after 3 attempts, stopping. Run again to retry."); - break; - } - } else { - staleCount = 0; + // Ensure dependent collections have backfill entries for all known DIDs + const dependentCollections = getCollectionNames(config).filter( + (col) => config.collections[col]?.discover === false + ); + for (const depCol of dependentCollections) { + await db + .prepare( + `INSERT OR IGNORE INTO backfills (did, collection, completed) + SELECT i.did, ?, 0 FROM identities i + LEFT JOIN backfills b ON b.did = i.did AND b.collection = ? + WHERE b.did IS NULL` + ) + .bind(depCol, depCol) + .run(); } - lastRemaining = result.remaining; - } - console.log("\n=== Done ==="); + // Phase 2: Backfill all pending records + console.log("--- Backfill ---"); + const backfillStart = Date.now(); + + const pending = await db + .prepare("SELECT COUNT(*) as count FROM backfills WHERE completed = 0") + .first<{ count: number }>(); + const pendingCount = pending?.count ?? 0; + + const uniqueUsers = await db + .prepare("SELECT COUNT(DISTINCT did) as count FROM backfills WHERE completed = 0") + .first<{ count: number }>(); + + const userCount = uniqueUsers?.count ?? 0; + console.log(` ${pendingCount} pending collection backfills for ${userCount} users`); + + const total = await backfillAll(db, config, { + concurrency: 100, + onProgress: ({ records, usersComplete, usersTotal, usersFailed }) => { + const secs = (Date.now() - backfillStart) / 1000; + const rate = secs > 0 ? Math.round(records / secs) : 0; + const failStr = usersFailed > 0 ? ` | ${usersFailed} failed` : ""; + process.stdout.write( + `\r ${records} records | ${usersComplete}/${usersTotal} users | ${rate}/s | ${elapsed(backfillStart)}${failStr} ` + ); + }, + }); + process.stdout.write("\n"); + + console.log(` Done: ${total} records in ${elapsed(backfillStart)}\n`); + + // Summary + const finalRemaining = await db + .prepare("SELECT COUNT(*) as count FROM backfills WHERE completed = 0") + .first<{ count: number }>(); + const failed = await db + .prepare("SELECT COUNT(*) as count FROM backfills WHERE completed = 1 AND retries > 0") + .first<{ count: number }>(); + + console.log(`=== Finished in ${elapsed(syncStart)} ===`); + console.log(` Discovered: ${allDiscovered.size} users`); + console.log(` Backfilled: ${total} records`); + if ((finalRemaining?.count ?? 0) > 0) + console.log(` Remaining: ${finalRemaining!.count} backfills`); + if ((failed?.count ?? 0) > 0) + console.log(` Failed: ${failed!.count} (exceeded retries)`); + } finally { + await dispose(); + } } -main(); +main().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/src/core/backfill.ts b/src/core/backfill.ts index cfbce07..86b4d2b 100644 --- a/src/core/backfill.ts +++ b/src/core/backfill.ts @@ -1,13 +1,14 @@ import { type Did } from "@atcute/lexicons"; import { isDid, isNsid } from "@atcute/lexicons/syntax"; +import type { Client } from "@atcute/client"; import type { ContrailConfig, Database, IngestEvent } from "./types"; import { getDiscoverableCollections, getDependentCollections, DEFAULT_RELAYS } from "./types"; import { applyEvents } from "./db"; -import { getClient } from "./client"; +import { getClient, getPDS } from "./client"; const PAGE_SIZE = 100; -const BATCH_SIZE = 50; +const BATCH_SIZE = 100; const MAX_RETRIES = 5; const REQUEST_TIMEOUT_MS = 10_000; @@ -15,7 +16,8 @@ const REQUEST_TIMEOUT_MS = 10_000; async function withRetry( fn: () => Promise, label: string, - maxRetries = 3 + maxRetries = 3, + timeoutMs = REQUEST_TIMEOUT_MS ): Promise { let lastError: unknown; for (let attempt = 0; attempt <= maxRetries; attempt++) { @@ -23,16 +25,13 @@ async function withRetry( return await Promise.race([ fn(), new Promise((_, reject) => - setTimeout(() => reject(new Error(`Timeout: ${label}`)), REQUEST_TIMEOUT_MS) + setTimeout(() => reject(new Error(`Timeout: ${label}`)), timeoutMs) ), ]); } catch (err) { lastError = err; if (attempt < maxRetries) { const delay = Math.min(1000 * 2 ** attempt, 10000); - console.warn( - `[retry] ${label} failed (attempt ${attempt + 1}/${maxRetries + 1}), retrying in ${delay}ms: ${err}` - ); await new Promise((r) => setTimeout(r, delay)); } } @@ -54,12 +53,24 @@ async function markFailed( .run(); } +export interface BackfillOptions { + /** Pre-resolved client — avoids redundant PDS lookups when batching by DID */ + client?: Client; + /** Skip replay detection in applyEvents (safe during initial backfill) */ + skipReplayDetection?: boolean; + /** Max retries per request (default: 3). Set to 0 for single-attempt mode. */ + maxRetries?: number; + /** Per-request timeout in ms (default: 10000). */ + requestTimeout?: number; +} + export async function backfillUser( db: Database, did: string, collection: string, deadline: number, - config?: ContrailConfig + config?: ContrailConfig, + options?: BackfillOptions ): Promise { if (Date.now() >= deadline) return 0; @@ -82,10 +93,8 @@ export async function backfillUser( } let currentCursor: string | undefined = status?.pds_cursor ?? undefined; - - console.log( - `Backfilling ${collection} for ${did} (cursor: ${currentCursor ?? "start"}, retries: ${status?.retries ?? 0})` - ); + const retries = options?.maxRetries ?? 3; + const timeout = options?.requestTimeout ?? REQUEST_TIMEOUT_MS; if (!isDid(did)) { await markFailed(db, did, collection, `Invalid DID: ${did}`); @@ -97,16 +106,19 @@ export async function backfillUser( return 0; } - let client; - try { - client = await withRetry( - () => getClient(did as Did, db), - `getClient(${did})`, - 1 - ); - } catch (err) { - await markFailed(db, did, collection, String(err)); - return 0; + let client = options?.client; + if (!client) { + try { + client = await withRetry( + () => getClient(did as Did, db), + `getClient(${did})`, + Math.min(retries, 1), + timeout + ); + } catch (err) { + await markFailed(db, did, collection, String(err)); + return 0; + } } let totalInserted = 0; @@ -116,7 +128,7 @@ export async function backfillUser( while (Date.now() < deadline) { const response = await withRetry( () => - client.get("com.atproto.repo.listRecords", { + client!.get("com.atproto.repo.listRecords", { params: { repo: did as Did, collection, @@ -124,9 +136,10 @@ export async function backfillUser( cursor: currentCursor, }, }), - `listRecords(${did}/${collection})` + `listRecords(${did}/${collection})`, + retries, + timeout ); - if (!response.ok) { await markFailed( db, @@ -155,11 +168,10 @@ export async function backfillUser( indexed_at: now * 1000, })); - for (let i = 0; i < events.length; i += BATCH_SIZE) { - const batch = events.slice(i, i + BATCH_SIZE); - await applyEvents(db, batch, config); - totalInserted += batch.length; - } + await applyEvents(db, events, config, { + skipReplayDetection: options?.skipReplayDetection, + }); + totalInserted += events.length; currentCursor = response.data.cursor ?? undefined; @@ -187,18 +199,184 @@ export async function backfillUser( ) .bind(did, collection) .run(); - console.log( - `Backfill complete: ${totalInserted} records for ${did}/${collection}` - ); - } else { - console.log( - `Backfill paused: ${totalInserted} records for ${did}/${collection}, will resume` - ); } return totalInserted; } +// --- Bulk backfill (groups by DID, resolves client once) --- + +export interface BackfillProgress { + records: number; + usersComplete: number; + usersTotal: number; + usersFailed: number; +} + +export interface BackfillAllOptions { + concurrency?: number; + onProgress?: (progress: BackfillProgress) => void; +} + +export async function backfillAll( + db: Database, + config: ContrailConfig, + options?: BackfillAllOptions +): Promise { + const concurrency = options?.concurrency ?? 100; + let totalBackfilled = 0; + + while (true) { + const pending = await db + .prepare( + "SELECT did, collection FROM backfills WHERE completed = 0 ORDER BY did" + ) + .all<{ did: string; collection: string }>(); + + const rows = pending.results ?? []; + if (rows.length === 0) break; + + // Group by DID so we resolve PDS once per user + const byDid = new Map(); + for (const row of rows) { + const cols = byDid.get(row.did) ?? []; + cols.push(row.collection); + byDid.set(row.did, cols); + } + + const dids = [...byDid.keys()]; + + // Resolve PDS endpoints in background (populates in-memory cache) + const resolvePromise = (async () => { + for (let i = 0; i < dids.length; i += 200) { + await Promise.allSettled( + dids.slice(i, i + 200).map((did) => + getPDS(did as Did, db).catch(() => {}) + ) + ); + } + })(); + + let roundBackfilled = 0; + let usersComplete = 0; + let usersFailed = 0; + const failedDids: string[] = []; + + const FAST_TIMEOUT = 3_000; + + const emitProgress = () => + options?.onProgress?.({ + records: totalBackfilled + roundBackfilled, + usersComplete, + usersTotal: dids.length, + usersFailed, + }); + + // Fast pass: single attempt per user with short timeout + for (let i = 0; i < dids.length; i += concurrency) { + const batch = dids.slice(i, i + concurrency); + + const results = await Promise.allSettled( + batch.map(async (did) => { + let client: Client | undefined; + try { + client = await withRetry( + () => getClient(did as Did, db), + `getClient(${did})`, + 0, + FAST_TIMEOUT + ); + } catch { + failedDids.push(did); + return 0; + } + + const cols = byDid.get(did)!; + const counts = await Promise.all( + cols.map((col) => + backfillUser(db, did, col, Infinity, config, { + client, + skipReplayDetection: true, + maxRetries: 0, + requestTimeout: FAST_TIMEOUT, + }).catch(() => { + failedDids.push(did); + return 0; + }) + ) + ); + + usersComplete++; + return counts.reduce((a, b) => a + b, 0); + }) + ); + + for (const r of results) { + if (r.status === "fulfilled") roundBackfilled += r.value; + } + + emitProgress(); + } + + // Retry pass: failed DIDs get retries with backoff, still in concurrent batches + if (failedDids.length > 0) { + const uniqueFailed = [...new Set(failedDids)]; + usersComplete -= uniqueFailed.length; // don't count them yet + + for (let i = 0; i < uniqueFailed.length; i += concurrency) { + const batch = uniqueFailed.slice(i, i + concurrency); + + const results = await Promise.allSettled( + batch.map(async (did) => { + let client: Client | undefined; + try { + client = await withRetry( + () => getClient(did as Did, db), + `getClient(${did})`, + 2 + ); + } catch (err) { + for (const col of byDid.get(did)!) { + await markFailed(db, did, col, String(err)); + } + usersFailed++; + usersComplete++; + return 0; + } + + const cols = byDid.get(did)!; + const counts = await Promise.all( + cols.map((col) => + backfillUser(db, did, col, Infinity, config, { + client, + skipReplayDetection: true, + maxRetries: 2, + }) + ) + ); + usersComplete++; + return counts.reduce((a, b) => a + b, 0); + }) + ); + + for (const r of results) { + if (r.status === "fulfilled") roundBackfilled += r.value; + } + + emitProgress(); + } + } + + await resolvePromise; + totalBackfilled += roundBackfilled; + + // If nothing was backfilled this round, we're stuck + if (roundBackfilled === 0) break; + } + + return totalBackfilled; +} + // --- Discovery --- interface DiscoveryPage { diff --git a/src/core/client.ts b/src/core/client.ts index f14e242..54a012b 100644 --- a/src/core/client.ts +++ b/src/core/client.ts @@ -86,21 +86,61 @@ export async function resolvePDS( return result; } +// In-memory PDS cache + in-flight deduplication +const pdsCache = new Map(); +const pdsInflight = new Map>(); + export async function getPDS( did: Did, db?: Database ): Promise { - // Check cached PDS first + const mem = pdsCache.get(did); + if (mem) return mem; + + // Deduplicate concurrent calls for the same DID + const inflight = pdsInflight.get(did); + if (inflight) return inflight; + + const promise = resolvePDSCached(did, db); + pdsInflight.set(did, promise); + try { + return await promise; + } finally { + pdsInflight.delete(did); + } +} + +async function resolvePDSCached( + did: Did, + db?: Database +): Promise { if (db) { const cached = await db .prepare("SELECT pds FROM identities WHERE did = ? AND pds IS NOT NULL") .bind(did) .first<{ pds: string }>(); - if (cached?.pds) return cached.pds; + if (cached?.pds) { + pdsCache.set(did, cached.pds); + return cached.pds; + } } const resolved = await resolvePDS(did); - return resolved?.pds ?? undefined; + if (!resolved?.pds) return undefined; + + pdsCache.set(did, resolved.pds); + + // Persist to DB for future runs + if (db) { + await db + .prepare( + "INSERT INTO identities (did, handle, pds, resolved_at) VALUES (?, ?, ?, ?) ON CONFLICT(did) DO UPDATE SET pds = excluded.pds, handle = COALESCE(excluded.handle, identities.handle), resolved_at = excluded.resolved_at" + ) + .bind(did, resolved.handle, resolved.pds, Date.now()) + .run(); + } + + return resolved.pds; } export async function getClient(did: Did, db?: Database): Promise { diff --git a/src/core/db/records.ts b/src/core/db/records.ts index f49bead..39b691e 100644 --- a/src/core/db/records.ts +++ b/src/core/db/records.ts @@ -130,14 +130,16 @@ export async function saveCursor( export async function applyEvents( db: Database, events: IngestEvent[], - config?: ContrailConfig + config?: ContrailConfig, + options?: { skipReplayDetection?: boolean } ): Promise { if (events.length === 0) return; // Look up existing records so we can skip duplicate count updates on replayed events. // A create/update with the same CID is a replay; a delete for a missing URI is a replay. + // Can be skipped during backfill where records are known to be fresh inserts. const existingCids = new Map(); - if (config) { + if (config && !options?.skipReplayDetection) { const uris = events.map((e) => e.uri); for (let i = 0; i < uris.length; i += 50) { const chunk = uris.slice(i, i + 50); diff --git a/src/core/router/admin.ts b/src/core/router/admin.ts index 78229bd..f76c24c 100644 --- a/src/core/router/admin.ts +++ b/src/core/router/admin.ts @@ -1,10 +1,7 @@ import type { Hono, Context, Next } from "hono"; import type { ContrailConfig, Database } from "../types"; -import { getCollectionNames } from "../types"; import { getLastCursor } from "../db"; import { initSchema } from "../db/schema"; -import { backfillUser, discoverDIDs } from "../backfill"; -import { parseIntParam } from "./helpers"; export function registerAdminRoutes( app: Hono, @@ -53,80 +50,6 @@ export function registerAdminRoutes( }); }); - app.get(`/xrpc/${ns}.admin.sync`, requireAdmin, async (c) => { - const deadline = Date.now() + 25_000; - const concurrency = parseIntParam(c.req.query("concurrency"), 25) ?? 25; - - // Phase 1: Discover DIDs from relays - const dids = await discoverDIDs(db, config, deadline); - - const discoverableCount = getCollectionNames(config).filter( - (col) => config.collections[col]?.discover !== false - ).length; - const discoveryRows = await db - .prepare("SELECT COUNT(*) as count FROM discovery") - .first<{ count: number }>(); - const pendingDiscovery = await db - .prepare("SELECT COUNT(*) as count FROM discovery WHERE completed = 0") - .first<{ count: number }>(); - const discoveryDone = - (discoveryRows?.count ?? 0) >= discoverableCount && - (pendingDiscovery?.count ?? 0) === 0; - - // Ensure dependent collections have backfill entries for all known DIDs - const dependentCollections = getCollectionNames(config).filter( - (col) => config.collections[col]?.discover === false - ); - if (dependentCollections.length > 0 && Date.now() < deadline) { - for (const depCol of dependentCollections) { - await db - .prepare( - `INSERT OR IGNORE INTO backfills (did, collection, completed) - SELECT i.did, ?, 0 FROM identities i - LEFT JOIN backfills b ON b.did = i.did AND b.collection = ? - WHERE b.did IS NULL` - ) - .bind(depCol, depCol) - .run(); - } - } - - // Phase 2: Backfill records from PDS - let backfilled = 0; - if (Date.now() < deadline) { - const pending = await db - .prepare( - "SELECT did, collection FROM backfills WHERE completed = 0 LIMIT 500" - ) - .all<{ did: string; collection: string }>(); - - const rows = pending.results ?? []; - for (let i = 0; i < rows.length; i += concurrency) { - if (Date.now() >= deadline) break; - const batch = rows.slice(i, i + concurrency); - const results = await Promise.allSettled( - batch.map((row) => - backfillUser(db, row.did, row.collection, deadline, config) - ) - ); - for (const r of results) { - if (r.status === "fulfilled") backfilled += r.value; - } - } - } - - const remaining = await db - .prepare("SELECT COUNT(*) as count FROM backfills WHERE completed = 0") - .first<{ count: number }>(); - - return c.json({ - discovered: dids.length, - backfilled, - remaining: remaining?.count ?? 0, - done: discoveryDone && (remaining?.count ?? 0) === 0, - }); - }); - app.get(`/xrpc/${ns}.admin.reset`, requireAdmin, async (c) => { const tables = ["records", "backfills", "discovery", "cursor", "identities"]; await db.batch(tables.map((t) => db.prepare(`DELETE FROM ${t}`))); diff --git a/wrangler.jsonc b/wrangler.jsonc index 2e305a6..f62e354 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -11,7 +11,7 @@ "binding": "DB", "database_name": "contrail", "database_id": "29cf6e32-d0b9-4646-ac52-bfbeef085c1a", - "remote": true + "remote": false } ], "triggers": { -- 2.51.2