From c74aec5b8ce3b719d7cf2987674e6bd6a815296f Mon Sep 17 00:00:00 2001 From: Florian <45694132+flo-bit@users.noreply.github.com> Date: Wed, 5 Aug 2026 01:35:07 +0200 Subject: [PATCH] Speed up host-aware D1 backfill --- .changeset/backfill-throughput.md | 5 + apps/benchmark/README.md | 14 +- .../baselines/calendar-host-aware.json | 66 +++ .../configs/calendar-no-follow.config.json | 47 +++ apps/benchmark/src/bench.ts | 145 ++++++- docs/01-indexing.md | 14 +- .../contrail/src/cli/commands/backfill.ts | 20 +- packages/contrail/src/core/backfill.ts | 387 +++++++++++------- packages/contrail/src/core/client.ts | 36 +- packages/contrail/src/core/db/records.ts | 238 +++++++++-- packages/contrail/src/core/ingest.ts | 3 + packages/contrail/src/core/status.ts | 13 +- packages/contrail/src/workers/backfill.ts | 10 +- .../contrail/tests/backfill-status.test.ts | 298 +++++++++++++- packages/contrail/tests/helpers.ts | 1 + packages/contrail/tests/ingest.test.ts | 35 ++ packages/contrail/tests/records.test.ts | 35 ++ packages/contrail/tests/search.test.ts | 38 ++ 18 files changed, 1195 insertions(+), 210 deletions(-) create mode 100644 .changeset/backfill-throughput.md create mode 100644 apps/benchmark/baselines/calendar-host-aware.json create mode 100644 apps/benchmark/configs/calendar-no-follow.config.json diff --git a/.changeset/backfill-throughput.md b/.changeset/backfill-throughput.md new file mode 100644 index 0000000..2aa8788 --- /dev/null +++ b/.changeset/backfill-throughput.md @@ -0,0 +1,5 @@ +--- +"@atmo-dev/contrail": patch +--- + +Group historical fetches by PDS with bounded per-host concurrency, cancel timed-out requests, defer failed initial accounts to scheduled retries, and rebuild FTS and relation counts with set-based SQL after canonical bulk loading. diff --git a/apps/benchmark/README.md b/apps/benchmark/README.md index 1ef4b43..f3bfcd5 100644 --- a/apps/benchmark/README.md +++ b/apps/benchmark/README.md @@ -12,20 +12,20 @@ pnpm bench --config calendar.config.json ## Comparing concurrency -Each command runs in a new Node process and starts from a fresh local D1: +Each command runs in a new Node process and starts from a fresh local D1. Identity resolution, active PDS hosts, and accounts per PDS are separate controls: ```bash -pnpm bench --config calendar.config.json --concurrency 25 -pnpm bench --config calendar.config.json --concurrency 50 pnpm bench --config calendar.config.json --concurrency 100 -pnpm bench --config calendar.config.json --concurrency 200 +pnpm bench --config calendar.config.json --pds-concurrency 5 --dids-per-pds 3 +pnpm bench --config calendar.config.json --pds-concurrency 10 --dids-per-pds 3 +pnpm bench --config calendar.config.json --pds-concurrency 20 --dids-per-pds 3 ``` -The defaults are `--concurrency 100` and `--max-attempts 5`, matching `contrail backfill`. The baseline therefore includes the same immediate retry behavior as an ordinary Atmo RSVP re-backfill. Change only one option at a time after recording that baseline. +The current defaults are 100 concurrent identity resolutions, 20 active PDS hosts, 3 accounts per PDS, and one immediate attempt. Failures retain their cursor and move to scheduled cron retries instead of slowing the initial pass. The checked-in 774.49-second baseline records the older global-concurrency/5-attempt behavior and remains the historical comparison point. Before every run the harness recursively deletes its config/concurrency-specific `.cache` directory. It disposes and deletes the local D1 afterward as well; pass `--keep-cache` only for debugging. -Results are written to ignored JSON files under `results/`. Selected reference runs can be copied into `baselines/`; the first exact-default Atmo RSVP run is [`baselines/calendar-default.json`](baselines/calendar-default.json). +Results are written to ignored JSON files under `results/`. Selected reference runs live in `baselines/`: [`calendar-default.json`](baselines/calendar-default.json) is the original 774.49-second global-concurrency run, while [`calendar-host-aware.json`](baselines/calendar-host-aware.json) is the comparable 219.74-second host-aware run with set-based derived projection rebuilds. Each result includes: @@ -34,7 +34,7 @@ Each result includes: - records per second; - peak RSS; - account and per-collection backfill state; and -- the exact concurrency and attempt settings. +- the exact resolution, PDS-host, per-PDS account, and attempt settings. ## Adding configs diff --git a/apps/benchmark/baselines/calendar-host-aware.json b/apps/benchmark/baselines/calendar-host-aware.json new file mode 100644 index 0000000..6216fea --- /dev/null +++ b/apps/benchmark/baselines/calendar-host-aware.json @@ -0,0 +1,66 @@ +{ + "format": "contrail.backfill-benchmark-baseline", + "version": 1, + "config": "configs/calendar.config.json", + "implementation_base_commit": "c19cb96fa63392510ef621f1676d0335168de2ac", + "environment": { + "os": "Darwin 25.4.0 arm64", + "cpu": "Apple M1 Pro", + "memory_bytes": 17179869184, + "node": "22.14.0", + "backend": "wrangler-local-d1" + }, + "options": { + "concurrency": 100, + "pdsConcurrency": 20, + "didsPerPds": 3, + "maxAttempts": 1 + }, + "started_at": "2026-08-04T23:21:40.711Z", + "completed_at": "2026-08-04T23:25:20.552Z", + "timings_ms": { + "binding": 216.91, + "init": 239.37, + "discovery": 1602.71, + "backfill": 217206.93, + "total": 219738.34 + }, + "throughput": { + "accepted_records_per_second": 401.69, + "indexed_records_per_second": 397.01 + }, + "discovered_accounts": 1633, + "accepted_records": 87249, + "indexed_records": 87238, + "peak_rss_kib": 578640, + "network_max_concurrent": 100, + "accounts": { + "total": 1633, + "complete": 1590, + "pending": 0, + "retrying": 43, + "failed": 0 + }, + "collections": [ + { + "collection": "community.lexicon.calendar.event", + "records": 14610, + "unique_users": 318 + }, + { + "collection": "community.lexicon.calendar.rsvp", + "records": 6269, + "unique_users": 1430 + }, + { + "collection": "app.bsky.actor.profile", + "records": 1401, + "unique_users": 1396 + }, + { + "collection": "app.bsky.graph.follow", + "records": 64958, + "unique_users": 1283 + } + ] +} diff --git a/apps/benchmark/configs/calendar-no-follow.config.json b/apps/benchmark/configs/calendar-no-follow.config.json new file mode 100644 index 0000000..dd7af75 --- /dev/null +++ b/apps/benchmark/configs/calendar-no-follow.config.json @@ -0,0 +1,47 @@ +{ + "namespace": "rsvp.atmo", + "maintenance": { + "optimize": true + }, + "collections": { + "event": { + "collection": "community.lexicon.calendar.event", + "queryable": { + "mode": {}, + "name": {}, + "status": {}, + "description": {}, + "preferences.showInDiscovery": {}, + "startsAt": { "type": "range" }, + "endsAt": { "type": "range" }, + "createdAt": { "type": "range" } + }, + "searchable": ["mode", "name", "status", "description"], + "relations": { + "rsvps": { + "collection": "rsvp", + "groupBy": "status", + "groups": { + "going": "community.lexicon.calendar.rsvp#going", + "interested": "community.lexicon.calendar.rsvp#interested", + "notgoing": "community.lexicon.calendar.rsvp#notgoing" + } + } + } + }, + "rsvp": { + "collection": "community.lexicon.calendar.rsvp", + "queryable": { + "status": {}, + "subject.uri": {}, + "createdAt": { "type": "range" } + }, + "references": { + "event": { + "collection": "event", + "field": "subject.uri" + } + } + } + } +} diff --git a/apps/benchmark/src/bench.ts b/apps/benchmark/src/bench.ts index 0079218..3ea6b7b 100644 --- a/apps/benchmark/src/bench.ts +++ b/apps/benchmark/src/bench.ts @@ -14,8 +14,11 @@ const WRANGLER_CONFIG = resolve(APP_DIR, "wrangler.jsonc"); interface Options { config: string; concurrency: number; + pdsConcurrency: number; + didsPerPds: number; maxAttempts: number; keepCache: boolean; + excludeDids: string[]; } function positiveInteger(raw: string | undefined, name: string, fallback: number): number { @@ -30,8 +33,11 @@ function positiveInteger(raw: string | undefined, name: string, fallback: number function parseArgs(argv: string[]): Options { let config: string | undefined; let concurrencyRaw: string | undefined; + let pdsConcurrencyRaw: string | undefined; + let didsPerPdsRaw: string | undefined; let maxAttemptsRaw: string | undefined; let keepCache = false; + const excludeDids: string[] = []; for (let index = 0; index < argv.length; index++) { const arg = argv[index]; @@ -40,17 +46,29 @@ function parseArgs(argv: string[]): Options { else if (arg === "--concurrency") concurrencyRaw = argv[++index]; else if (arg.startsWith("--concurrency=")) { concurrencyRaw = arg.slice("--concurrency=".length); + } else if (arg === "--pds-concurrency") pdsConcurrencyRaw = argv[++index]; + else if (arg.startsWith("--pds-concurrency=")) { + pdsConcurrencyRaw = arg.slice("--pds-concurrency=".length); + } else if (arg === "--dids-per-pds") didsPerPdsRaw = argv[++index]; + else if (arg.startsWith("--dids-per-pds=")) { + didsPerPdsRaw = arg.slice("--dids-per-pds=".length); } else if (arg === "--max-attempts") maxAttemptsRaw = argv[++index]; else if (arg.startsWith("--max-attempts=")) { maxAttemptsRaw = arg.slice("--max-attempts=".length); + } else if (arg === "--exclude-did") excludeDids.push(argv[++index]); + else if (arg.startsWith("--exclude-did=")) { + excludeDids.push(arg.slice("--exclude-did=".length)); } else if (arg === "--keep-cache") keepCache = true; else if (arg === "--help" || arg === "-h") { console.log(`Usage: pnpm bench --config [options] Options: --config JSON config path or filename under configs/ (required) - --concurrency Concurrent accounts (default: 100) - --max-attempts Initial attempts per failed account (default: 5) + --concurrency Concurrent identity resolutions (default: 100) + --pds-concurrency Concurrent PDS hosts (default: 20) + --dids-per-pds Concurrent accounts per PDS (default: 3) + --max-attempts Immediate attempts per failed account (default: 1) + --exclude-did Skip an actor after discovery (repeatable) --keep-cache Keep the disposable local D1 after the run `); process.exit(0); @@ -63,8 +81,11 @@ Options: return { config, concurrency: positiveInteger(concurrencyRaw, "--concurrency", 100), - maxAttempts: positiveInteger(maxAttemptsRaw, "--max-attempts", 5), + pdsConcurrency: positiveInteger(pdsConcurrencyRaw, "--pds-concurrency", 20), + didsPerPds: positiveInteger(didsPerPdsRaw, "--dids-per-pds", 3), + maxAttempts: positiveInteger(maxAttemptsRaw, "--max-attempts", 1), keepCache, + excludeDids, }; } @@ -94,12 +115,83 @@ function safeName(path: string): string { .replace(/[^a-z0-9_-]+/gi, "-"); } +interface FetchMetric { + requests: number; + errors: number; + total_ms: number; + max_ms: number; + statuses: Record; +} + +function instrumentFetch(): { + metrics: Map; + restore(): void; + maxActive(): number; +} { + const original = globalThis.fetch; + const metrics = new Map(); + let active = 0; + let peakActive = 0; + + globalThis.fetch = async (input, init) => { + let key = "invalid-url"; + try { + const raw = input instanceof Request ? input.url : String(input); + const url = new URL(raw); + const operation = url.pathname.split("/").pop() || url.pathname; + key = `${url.host}/${operation}`; + } catch { + // Keep the fallback key. + } + + const metric = metrics.get(key) ?? { + requests: 0, + errors: 0, + total_ms: 0, + max_ms: 0, + statuses: {}, + }; + metrics.set(key, metric); + metric.requests++; + active++; + peakActive = Math.max(peakActive, active); + const start = performance.now(); + try { + const response = await original(input, init); + const status = String(response.status); + metric.statuses[status] = (metric.statuses[status] ?? 0) + 1; + return response; + } catch (error) { + metric.errors++; + throw error; + } finally { + const duration = performance.now() - start; + metric.total_ms += duration; + metric.max_ms = Math.max(metric.max_ms, duration); + active--; + } + }; + + return { + metrics, + restore() { + globalThis.fetch = original; + }, + maxActive() { + return peakActive; + }, + }; +} + async function main(): Promise { const options = parseArgs(process.argv.slice(2)); const configPath = await resolveConfigPath(options.config); const config = JSON.parse(await readFile(configPath, "utf8")) as ContrailConfig; const name = safeName(configPath); - const cachePath = resolve(CACHE_DIR, `${name}-c${options.concurrency}`); + const cachePath = resolve( + CACHE_DIR, + `${name}-r${options.concurrency}-h${options.pdsConcurrency}-d${options.didsPerPds}`, + ); await rm(cachePath, { recursive: true, force: true }); await mkdir(cachePath, { recursive: true }); @@ -107,8 +199,11 @@ async function main(): Promise { console.log(`config: ${configPath}`); console.log(`backend: fresh local D1`); - console.log(`concurrency: ${options.concurrency}`); + console.log(`resolution: ${options.concurrency}`); + console.log(`PDS hosts: ${options.pdsConcurrency}`); + console.log(`DIDs / PDS: ${options.didsPerPds}`); console.log(`max attempts: ${options.maxAttempts}`); + console.log(`excluded: ${options.excludeDids.length} actors`); console.log(`cache: reset ${cachePath}`); const startedAt = new Date(); @@ -121,6 +216,7 @@ async function main(): Promise { let discovered = 0; let acceptedRecords = 0; let overview: any; + let fetchInstrumentation: ReturnType | undefined; try { let phaseStart = performance.now(); @@ -131,6 +227,7 @@ async function main(): Promise { envFiles: [], }); bindingMs = elapsed(phaseStart); + fetchInstrumentation = instrumentFetch(); const db = (proxy.env as Record).DB as Database; const contrail = new Contrail({ ...config, db }); @@ -145,17 +242,25 @@ async function main(): Promise { discoveryMs = elapsed(phaseStart); console.log(`discovered: ${discovered} accounts in ${(discoveryMs / 1000).toFixed(2)}s`); + for (const did of new Set(options.excludeDids)) { + await db.prepare("DELETE FROM backfills WHERE did = ?").bind(did).run(); + } + phaseStart = performance.now(); + const backfillStart = phaseStart; let lastProgressAt = 0; acceptedRecords = await contrail.backfill({ concurrency: options.concurrency, + pdsConcurrency: options.pdsConcurrency, + didsPerPds: options.didsPerPds, maxAttempts: options.maxAttempts, onProgress(progress) { const now = Date.now(); if (now - lastProgressAt < 2_000) return; lastProgressAt = now; + const seconds = ((performance.now() - backfillStart) / 1000).toFixed(1); console.log( - `progress: ${progress.usersComplete}/${progress.usersTotal} accounts, ` + + `progress: +${seconds}s ${progress.usersComplete}/${progress.usersTotal} accounts, ` + `${progress.records} accepted, ${progress.usersFailed} failed`, ); }, @@ -168,6 +273,7 @@ async function main(): Promise { if (!response.ok) throw new Error(`Status request failed: ${response.status}`); overview = await response.json(); } finally { + fetchInstrumentation?.restore(); await proxy?.dispose(); if (!options.keepCache) { await rm(cachePath, { recursive: true, force: true }); @@ -178,6 +284,16 @@ async function main(): Promise { const completedAt = new Date(); const indexedRecords = Number(overview.total_records ?? 0); const relativeConfig = relative(APP_DIR, configPath); + const network = Object.fromEntries( + [...(fetchInstrumentation?.metrics ?? new Map())].map(([key, metric]) => [ + key, + { + ...metric, + total_ms: Math.round(metric.total_ms * 100) / 100, + max_ms: Math.round(metric.max_ms * 100) / 100, + }, + ]), + ); const result = { format: "contrail.backfill-benchmark", version: 1, @@ -185,7 +301,10 @@ async function main(): Promise { backend: "wrangler-local-d1", options: { concurrency: options.concurrency, + pdsConcurrency: options.pdsConcurrency, + didsPerPds: options.didsPerPds, maxAttempts: options.maxAttempts, + excludedDids: options.excludeDids, }, started_at: startedAt.toISOString(), completed_at: completedAt.toISOString(), @@ -206,6 +325,10 @@ async function main(): Promise { accepted_records: acceptedRecords, indexed_records: indexedRecords, peak_rss_kib: process.resourceUsage().maxRSS, + network: { + max_concurrent: fetchInstrumentation?.maxActive() ?? 0, + requests: network, + }, backfill: overview.backfill, collections: overview.collections, }; @@ -213,7 +336,7 @@ async function main(): Promise { const timestamp = completedAt.toISOString().replace(/[:.]/g, "-"); const resultPath = resolve( RESULTS_DIR, - `${name}-c${options.concurrency}-${timestamp}.json`, + `${name}-r${options.concurrency}-h${options.pdsConcurrency}-d${options.didsPerPds}-${timestamp}.json`, ); await writeFile(resultPath, `${JSON.stringify(result, null, 2)}\n`); @@ -230,6 +353,14 @@ async function main(): Promise { `${overview.backfill.accounts.retrying} retrying, ` + `${overview.backfill.accounts.failed} failed`, ); + console.log(`network max: ${result.network.max_concurrent} concurrent requests`); + for (const [key, metric] of Object.entries(network)) { + console.log( + `network: ${key} — ${metric.requests} requests, ` + + `${(metric.total_ms / 1000).toFixed(2)}s cumulative, ` + + `${(metric.max_ms / 1000).toFixed(2)}s max`, + ); + } console.log(`result: ${resultPath}`); if (overview.backfill.state !== "complete") { diff --git a/docs/01-indexing.md b/docs/01-indexing.md index 580c376..f4117ba 100644 --- a/docs/01-indexing.md +++ b/docs/01-indexing.md @@ -52,13 +52,19 @@ await contrail.backfillAll({ concurrency: 100 }); // discover + backfill, logs p Under the hood this is two steps you can call separately if you want finer control: ```ts -await contrail.discover(); // walk relays, register DIDs -await contrail.backfill({ concurrency: 100 }); // fetch history for registered DIDs +await contrail.discover(); // walk relays, register DIDs +await contrail.backfill({ + concurrency: 100, // identity resolution + pdsConcurrency: 20, // active PDS hosts + didsPerPds: 3, // accounts per active PDS +}); ``` `backfill()` picks up each account/collection at its saved PDS cursor. A row is marked complete only after the PDS listing reaches its end. Timeouts, failed identity resolution, `429`, and `5xx` responses leave the row pending with its last error. -Each initial invocation has a bounded failure budget (five attempts by default), so one dead PDS cannot hang the whole command forever. Failed account rows retain their cursors and receive an exponential `next_retry_at`. Scheduled retries start at 15 minutes, double to a maximum of 48 hours, and stop after ten failed scheduled attempts. Cloudflare's scheduled Worker retries a small due slice after each live-ingest cycle; an explicit later invocation resets exhausted rows and forces another bounded pass. `backfillAll()` returns a durable `status` summary alongside the number of discovered accounts and accepted records. +Each initial invocation attempts a failed account once by default, then gets out of the way. Failed rows retain their cursors and receive an exponential `next_retry_at`. Scheduled retries start at 15 minutes, double to a maximum of 48 hours, and stop after ten failed scheduled attempts. Cloudflare's scheduled Worker retries a small due slice after each live-ingest cycle; an explicit later invocation resets exhausted rows and forces another pass. `backfillAll()` returns a durable `status` summary alongside the number of discovered accounts and accepted records. + +Historical loading writes canonical records first, then rebuilds FTS and materialized relation counts with set-based SQL. A durable dirty marker keeps status `incomplete` if the process stops between those phases; the next manual or scheduled backfill repairs the projections before reporting readiness. Live ingestion and scheduled account retries continue maintaining both projections incrementally. ### Workers CLI @@ -69,7 +75,7 @@ pnpm contrail backfill # local D1 (wrangler dev's bindings) pnpm contrail backfill --remote # production D1 ``` -Auto-detects configs at `contrail.config.ts`, `src/contrail.config.ts`, `src/lib/contrail.config.ts`, or `app/contrail.config.ts` (first match wins). Override with `--config `. Other flags: `--binding ` (default `DB`), `--concurrency ` (default 100), and `--max-attempts ` (default 5). Once every known account has either completed or received a deferred failure, the initial pass is complete and scheduled retries continue in the background. Interrupted or undiscovered work still reports the pass as incomplete. +Auto-detects configs at `contrail.config.ts`, `src/contrail.config.ts`, `src/lib/contrail.config.ts`, or `app/contrail.config.ts` (first match wins). Override with `--config `. Other flags include `--binding ` (default `DB`), `--concurrency ` for identity resolution (default 100), `--pds-concurrency ` (default 20), `--dids-per-pds ` (default 3), and `--max-attempts ` (default 1). Once every known account has either completed or received a deferred failure, the initial pass is complete and scheduled retries continue in the background. Interrupted or undiscovered work still reports the pass as incomplete. If you'd rather embed backfill inside your own script, `@atmo-dev/contrail/workers` exports the same logic as a function: diff --git a/packages/contrail/src/cli/commands/backfill.ts b/packages/contrail/src/cli/commands/backfill.ts index 3a02d49..2d91198 100644 --- a/packages/contrail/src/cli/commands/backfill.ts +++ b/packages/contrail/src/cli/commands/backfill.ts @@ -8,6 +8,8 @@ interface BackfillOpts { remote?: boolean; binding: string; concurrency: number; + pdsConcurrency: number; + didsPerPds: number; maxAttempts: number; only?: string; } @@ -28,13 +30,23 @@ export function registerBackfill(cli: CAC): void { }) .option( "--concurrency ", - "Concurrency for record backfill (labels are per-labeler serial)", + "Concurrent identity resolutions (labels are per-labeler serial)", { default: 100 } ) + .option( + "--pds-concurrency ", + "PDS hosts fetched concurrently", + { default: 20 } + ) + .option( + "--dids-per-pds ", + "Accounts fetched concurrently from each PDS", + { default: 3 } + ) .option( "--max-attempts ", - "Failed attempts before leaving an account pending for the next run", - { default: 5 } + "Immediate attempts before deferring failures to scheduled retries", + { default: 1 } ) .option( "--only ", @@ -64,6 +76,8 @@ export function registerBackfill(cli: CAC): void { const result = await backfillAll({ ...wrangler, concurrency: Number(options.concurrency), + pdsConcurrency: Number(options.pdsConcurrency), + didsPerPds: Number(options.didsPerPds), maxAttempts: Number(options.maxAttempts), }); recordsIncomplete = result.status.state !== "complete"; diff --git a/packages/contrail/src/core/backfill.ts b/packages/contrail/src/core/backfill.ts index 86d1292..cde5bbb 100644 --- a/packages/contrail/src/core/backfill.ts +++ b/packages/contrail/src/core/backfill.ts @@ -11,8 +11,10 @@ import { shortNameForNsid, } from "./types"; import { getLastCursor, saveCursor } from "./db"; +import { getMeta, setMeta } from "./db/meta"; import { createIngestEvent, ingestRecords } from "./ingest"; -import { getClient, getPDS } from "./client"; +import { rebuildDerivedProjections } from "./db/records"; +import { createPdsClient, getClient, getPDS } from "./client"; import { finishBackfillRun, heartbeatBackfillRun, @@ -46,12 +48,15 @@ function recordTimeUs( } const PAGE_SIZE = 100; -const DEFAULT_MAX_ATTEMPTS = 5; +const DEFAULT_MAX_ATTEMPTS = 1; +const DEFAULT_PDS_CONCURRENCY = 20; +const DEFAULT_DIDS_PER_PDS = 3; const REQUEST_TIMEOUT_MS = 10_000; const BACKFILL_RETRY_BASE_MS = 15 * 60_000; const BACKFILL_RETRY_MAX_MS = 48 * 60 * 60_000; const DEFAULT_SCHEDULED_MAX_ATTEMPTS = 10; +const DERIVED_PROJECTIONS_DIRTY_KEY = "backfill_derived_projections_dirty"; function positiveInteger(value: number | undefined, fallback: number): number { return typeof value === "number" && Number.isFinite(value) && value > 0 @@ -60,26 +65,28 @@ function positiveInteger(value: number | undefined, fallback: number): number { } async function withRetry( - fn: () => Promise, + fn: (signal: AbortSignal) => Promise, label: string, maxRetries = 3, timeoutMs = REQUEST_TIMEOUT_MS ): Promise { let lastError: unknown; for (let attempt = 0; attempt <= maxRetries; attempt++) { + const controller = new AbortController(); + const timeout = setTimeout( + () => controller.abort(new Error(`Timeout: ${label}`)), + timeoutMs + ); try { - return await Promise.race([ - fn(), - new Promise((_, reject) => - setTimeout(() => reject(new Error(`Timeout: ${label}`)), timeoutMs) - ), - ]); + return await fn(controller.signal); } catch (err) { lastError = err; - if (attempt < maxRetries) { - const delay = Math.min(1000 * 2 ** attempt, 10000); - await new Promise((r) => setTimeout(r, delay)); - } + } finally { + clearTimeout(timeout); + } + if (attempt < maxRetries) { + const delay = Math.min(1000 * 2 ** attempt, 10000); + await new Promise((r) => setTimeout(r, delay)); } } throw lastError; @@ -139,6 +146,8 @@ async function markFailed( export interface BackfillOptions { /** Pre-resolved client — avoids redundant PDS lookups when batching by DID */ client?: Client; + /** Complete relay-discovered actor set for dependent-record admission. */ + knownDids?: ReadonlySet; /** Skip replay detection during initial backfill. */ skipReplayDetection?: boolean; /** Max retries per request (default: 3). Set to 0 for single-attempt mode. */ @@ -147,11 +156,16 @@ export interface BackfillOptions { requestTimeout?: number; /** Mark the row terminal when this consecutive-failure count is reached. */ exhaustAfterAttempts?: number; + /** Yield after this many successful pages without consuming a retry. */ + maxPages?: number; + /** Defer FTS and relation counts until the bulk pass finishes. */ + skipDerivedProjections?: boolean; } interface BackfillUserAttempt { records: number; completed: boolean; + failed: boolean; } async function backfillUserAttempt( @@ -162,7 +176,9 @@ async function backfillUserAttempt( config: ContrailConfig, options?: BackfillOptions ): Promise { - if (Date.now() >= deadline) return { records: 0, completed: false }; + if (Date.now() >= deadline) { + return { records: 0, completed: false, failed: false }; + } const status = await db .prepare( @@ -171,7 +187,9 @@ async function backfillUserAttempt( .bind(did, collection) .first<{ completed: number; pds_cursor: string | null; retries: number }>(); - if (status?.completed) return { records: 0, completed: true }; + if (status?.completed) { + return { records: 0, completed: true, failed: false }; + } if (!status) { await db @@ -194,7 +212,7 @@ async function backfillUserAttempt( `Invalid DID: ${did}`, options?.exhaustAfterAttempts ); - return { records: 0, completed: false }; + return { records: 0, completed: false, failed: true }; } if (!isNsid(collection)) { @@ -205,14 +223,14 @@ async function backfillUserAttempt( `Invalid NSID: ${collection}`, options?.exhaustAfterAttempts ); - return { records: 0, completed: false }; + return { records: 0, completed: false, failed: true }; } let client = options?.client; if (!client) { try { client = await withRetry( - () => getClient(did as Did, db, config), + (signal) => getClient(did as Did, db, config, signal), `getClient(${did})`, Math.min(retries, 1), timeout @@ -225,17 +243,18 @@ async function backfillUserAttempt( err, options?.exhaustAfterAttempts ); - return { records: 0, completed: false }; + return { records: 0, completed: false, failed: true }; } } let totalInserted = 0; let done = false; + let pages = 0; try { while (Date.now() < deadline) { const response = await withRetry( - () => + (signal) => client!.get("com.atproto.repo.listRecords", { params: { repo: did as Did, @@ -243,6 +262,7 @@ async function backfillUserAttempt( limit: PAGE_SIZE, cursor: currentCursor, }, + signal, }), `listRecords(${did}/${collection})`, retries, @@ -259,7 +279,7 @@ async function backfillUserAttempt( `listRecords status ${response.status} (${detail})`, options?.exhaustAfterAttempts ); - return { records: totalInserted, completed: false }; + return { records: totalInserted, completed: false, failed: true }; } if (response.data.records.length === 0) { @@ -287,6 +307,8 @@ async function backfillUserAttempt( const result = await ingestRecords(db, events, config, { skipReplayDetection: options?.skipReplayDetection, skipFeedFanout: true, + knownDids: options?.knownDids, + skipDerivedProjections: options?.skipDerivedProjections, // Let sinks bulk-flush differently from live ingestion. phase: "backfill", }); @@ -302,10 +324,12 @@ async function backfillUserAttempt( .bind(currentCursor ?? null, now, did, collection) .run(); + pages++; if (!currentCursor) { done = true; break; } + if (pages >= (options?.maxPages ?? Infinity)) break; } } catch (err) { await markFailed( @@ -315,7 +339,7 @@ async function backfillUserAttempt( err, options?.exhaustAfterAttempts ); - return { records: totalInserted, completed: false }; + return { records: totalInserted, completed: false, failed: true }; } if (done) { @@ -327,7 +351,7 @@ async function backfillUserAttempt( .run(); } - return { records: totalInserted, completed: done }; + return { records: totalInserted, completed: done, failed: false }; } export async function backfillUser( @@ -359,25 +383,95 @@ export interface BackfillProgress { } export interface BackfillAllOptions { + /** Concurrent identity resolutions before work is grouped by PDS. Default: 100. */ concurrency?: number; - /** Failed account/collection attempts allowed in one run before leaving the - * row pending for a future run. Default: 5. */ + /** PDS hosts allowed to fetch concurrently. Default: 20. */ + pdsConcurrency?: number; + /** Accounts allowed to fetch concurrently from one PDS. Default: 3. */ + didsPerPds?: number; + /** Per-request timeout in milliseconds. Default: 10000. */ + requestTimeoutMs?: number; + /** Immediate attempts per failed account. Default: 1; scheduled retries handle + * later attempts. Values above 1 are retained for explicit manual recovery. */ maxAttempts?: number; onProgress?: (progress: BackfillProgress) => void; } +/** Keep a fixed number of jobs active without batch barriers. Returning true + * from consume requeues only that item, allowing paginated repositories to + * yield fairly while completed slots refill immediately. */ +async function drainQueue( + items: TItem[], + concurrency: number, + run: (item: TItem) => Promise, + consume: (result: TResult) => boolean | void +): Promise { + if (items.length === 0) return; + const queue = [...items]; + let nextIndex = 0; + let active = 0; + let settled = false; + + await new Promise((resolve, reject) => { + const pump = () => { + if (settled) return; + while (active < concurrency && nextIndex < queue.length) { + const item = queue[nextIndex++]; + active++; + run(item).then( + (result) => { + active--; + if (consume(result) === true) queue.push(item); + if (active === 0 && nextIndex >= queue.length) { + settled = true; + resolve(); + } else { + pump(); + } + }, + (error) => { + settled = true; + reject(error); + } + ); + } + }; + pump(); + }); +} + +async function loadKnownBackfillDids(db: Database): Promise> { + const rows = await db + .prepare("SELECT DISTINCT did FROM backfills") + .all<{ did: string }>(); + return new Set((rows.results ?? []).map((row) => row.did)); +} + async function backfillPendingWork( db: Database, config: ContrailConfig, options: BackfillAllOptions | undefined, runId: string ): Promise { - const concurrency = positiveInteger(options?.concurrency, 100); + const resolutionConcurrency = positiveInteger(options?.concurrency, 100); + const pdsConcurrency = positiveInteger( + options?.pdsConcurrency, + DEFAULT_PDS_CONCURRENCY + ); + const didsPerPds = positiveInteger( + options?.didsPerPds, + DEFAULT_DIDS_PER_PDS + ); + const requestTimeout = positiveInteger( + options?.requestTimeoutMs, + REQUEST_TIMEOUT_MS + ); const maxAttempts = positiveInteger( options?.maxAttempts, DEFAULT_MAX_ATTEMPTS ); let totalBackfilled = 0; + const knownDids = await loadKnownBackfillDids(db); // Anchor the jetstream cursor to now if it hasn't been set yet, so records // emitted during backfill are replayed once jetstream starts. @@ -385,8 +479,14 @@ async function backfillPendingWork( await saveCursor(db, Date.now() * 1000); } - // The attempt cap is per invocation. Rows remain incomplete and a later run - // gets a fresh bounded retry budget. + // Mark the set-based catch-up dirty before canonical writes. A crash can leave + // search/count projections stale, so status stays incomplete and the next + // manual or scheduled pass repairs them before claiming readiness. + await setMeta(db, DERIVED_PROJECTIONS_DIRTY_KEY, "1"); + + // A manual/initial run resets scheduled exhaustion, but its default one-shot + // attempt leaves failures for cron instead of waiting on the same unhealthy + // upstream repeatedly. await db .prepare( "UPDATE backfills SET retries = 0, scheduled_retries = 0, next_retry_at = NULL, retry_exhausted = 0 WHERE completed = 0" @@ -404,32 +504,18 @@ async function backfillPendingWork( const rows = pending.results ?? []; if (rows.length === 0) break; - // Group by DID so we resolve each PDS once per pass. 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 collections = byDid.get(row.did) ?? []; + collections.push(row.collection); + byDid.set(row.did, collections); } const dids = [...byDid.keys()]; - await heartbeatBackfillRun(db, runId); - - // Warm the identity/PDS cache separately from record requests. Failures are - // deliberately ignored here: the actual pass records them per pending row. - for (let i = 0; i < dids.length; i += 200) { - await Promise.allSettled( - dids.slice(i, i + 200).map((did) => - getPDS(did as Did, db, config).catch(() => {}) - ) - ); - } - + const byPds = new Map(); let roundBackfilled = 0; let usersComplete = 0; let usersFailed = 0; - const failedDids: string[] = []; - const FAST_TIMEOUT = 3_000; const emitProgress = () => options?.onProgress?.({ @@ -439,116 +525,113 @@ async function backfillPendingWork( usersFailed, }); - // Fast pass: one short attempt. A DID is complete only when every pending - // collection reaches the end of its PDS listing. - for (let i = 0; i < dids.length; i += concurrency) { - const batch = dids.slice(i, i + concurrency); - const results = await Promise.all( - batch.map(async (did) => { - let client: Client; - try { - client = await withRetry( - () => getClient(did as Did, db, config), - `getClient(${did})`, - 0, - FAST_TIMEOUT - ); - } catch { - return { did, records: 0, completed: false }; - } + await heartbeatBackfillRun(db, runId); - const attempts = await Promise.all( - byDid.get(did)!.map((collection) => - backfillUserAttempt(db, did, collection, Infinity, config, { - client, - skipReplayDetection: true, - maxRetries: 0, - requestTimeout: FAST_TIMEOUT, - }) - ) + // Resolve first, then schedule by host. This deliberately separates cheap + // identity fan-out from PDS traffic and lets every host share one client. + await drainQueue( + dids, + resolutionConcurrency, + async (did) => { + try { + const pds = await withRetry( + (signal) => getPDS(did as Did, db, config, signal), + `getPDS(${did})`, + 0, + requestTimeout ); - return { - did, - records: attempts.reduce((sum, attempt) => sum + attempt.records, 0), - completed: attempts.every((attempt) => attempt.completed), - }; - }) - ); - - for (const result of results) { - roundBackfilled += result.records; - if (result.completed) usersComplete++; - else failedDids.push(result.did); + if (!pds) throw new Error(`PDS not found for ${did}`); + return { did, pds: pds.replace(/\/+$/, ""), failed: false }; + } catch (error) { + for (const collection of byDid.get(did)!) { + await markFailed(db, did, collection, error); + } + return { did, pds: null, failed: true }; + } + }, + (result) => { + if (result.pds) { + const hostDids = byPds.get(result.pds) ?? []; + hostDids.push(result.did); + byPds.set(result.pds, hostDids); + } else if (result.failed) { + usersFailed++; + emitProgress(); + } } - emitProgress(); - await heartbeatBackfillRun(db, runId); - } - - // Slow retry pass: use normal timeouts and request backoff for DIDs that did - // not finish. They remain incomplete if this pass also fails. - const uniqueFailed = [...new Set(failedDids)]; - const retryableRows = await db - .prepare( - "SELECT DISTINCT did FROM backfills WHERE completed = 0 AND retries < ?" - ) - .bind(maxAttempts) - .all<{ did: string }>(); - const retryableDids = new Set( - (retryableRows.results ?? []).map((row) => row.did) ); - const retryDids = uniqueFailed.filter((did) => retryableDids.has(did)); - usersFailed += uniqueFailed.length - retryDids.length; - emitProgress(); - - for (let i = 0; i < retryDids.length; i += concurrency) { - const batch = retryDids.slice(i, i + concurrency); - const results = await Promise.all( - batch.map(async (did) => { - let client: Client; - try { - client = await withRetry( - () => getClient(did as Did, db, config), - `getClient(${did})`, - 2 - ); - } catch (error) { + + // Keep only a small number of PDS hosts active. Within each host, process a + // few accounts and one collection page at a time. This preserves connection + // reuse, prevents request storms, and keeps large repositories from blocking + // unrelated hosts or accounts on the same host. + await drainQueue( + [...byPds.entries()], + pdsConcurrency, + async ([pds, hostDids]) => { + const client = createPdsClient(pds); + await drainQueue( + hostDids, + didsPerPds, + async (did) => { + const attempts: BackfillUserAttempt[] = []; for (const collection of byDid.get(did)!) { - await markFailed(db, did, collection, error); + attempts.push( + await backfillUserAttempt( + db, + did, + collection, + Infinity, + config, + { + client, + knownDids, + skipReplayDetection: true, + maxRetries: 0, + requestTimeout, + maxPages: 1, + skipDerivedProjections: true, + } + ) + ); } - return { records: 0, completed: false }; + return { + records: attempts.reduce( + (sum, attempt) => sum + attempt.records, + 0 + ), + completed: attempts.every((attempt) => attempt.completed), + failed: attempts.some((attempt) => attempt.failed), + }; + }, + (result) => { + roundBackfilled += result.records; + if (result.completed) usersComplete++; + else if (result.failed) usersFailed++; + emitProgress(); + return !result.completed && !result.failed; } - - const attempts = await Promise.all( - byDid.get(did)!.map((collection) => - backfillUserAttempt(db, did, collection, Infinity, config, { - client, - skipReplayDetection: true, - maxRetries: 2, - }) - ) - ); - return { - records: attempts.reduce((sum, attempt) => sum + attempt.records, 0), - completed: attempts.every((attempt) => attempt.completed), - }; - }) - ); - - for (const result of results) { - roundBackfilled += result.records; - if (result.completed) usersComplete++; - else usersFailed++; - } - emitProgress(); - await heartbeatBackfillRun(db, runId); - } + ); + await heartbeatBackfillRun(db, runId); + return undefined; + }, + () => false + ); totalBackfilled += roundBackfilled; - // Do not use inserted-record count as a completion signal: a reachable PDS - // may legitimately have zero records, while an unreachable PDS must consume - // its bounded retry budget and remain pending. + await heartbeatBackfillRun(db, runId); + // With the default maxAttempts=1 this ends after one failed attempt; rows + // remain retrying with persisted backoff. Explicit larger values cause + // another host-aware pass without changing the scheduled retry budget. } + // Canonical records and cursors are durable now. Rebuild expensive derived + // projections once with set-based SQL instead of hundreds of statements per + // network page. This also repairs a prior interrupted bulk pass. + await rebuildDerivedProjections(db, config); + await setMeta(db, DERIVED_PROJECTIONS_DIRTY_KEY, "0"); + await heartbeatBackfillRun(db, runId); + return totalBackfilled; } @@ -611,6 +694,13 @@ export async function retryPendingBackfills( let records = 0; try { + if ((await getMeta(db, DERIVED_PROJECTIONS_DIRTY_KEY)) === "1") { + await rebuildDerivedProjections(db, config); + await setMeta(db, DERIVED_PROJECTIONS_DIRTY_KEY, "0"); + await heartbeatBackfillRun(db, runId); + } + + const knownDids = await loadKnownBackfillDids(db); const due = await db .prepare( "SELECT did FROM backfills WHERE completed = 0 AND retry_exhausted = 0 AND (next_retry_at IS NULL OR next_retry_at <= ?) GROUP BY did ORDER BY MIN(COALESCE(next_retry_at, 0)), did LIMIT ?" @@ -633,7 +723,7 @@ export async function retryPendingBackfills( let client: Client; try { client = await withRetry( - () => getClient(did as Did, db, config), + (signal) => getClient(did as Did, db, config, signal), `getClient(${did})`, 0, Math.min(requestTimeoutMs, Math.max(1, deadline - Date.now())) @@ -659,6 +749,7 @@ export async function retryPendingBackfills( attempts.push( await backfillUserAttempt(db, did, collection, deadline, config, { client, + knownDids, skipReplayDetection: true, maxRetries: 0, requestTimeout: Math.min( @@ -710,8 +801,8 @@ async function fetchPage( } return withRetry( - async () => { - const response = await fetch(url.toString()); + async (signal) => { + const response = await fetch(url.toString(), { signal }); if (!response.ok) { throw new Error(`HTTP ${response.status}`); } diff --git a/packages/contrail/src/core/client.ts b/packages/contrail/src/core/client.ts index 10c89fa..464a257 100644 --- a/packages/contrail/src/core/client.ts +++ b/packages/contrail/src/core/client.ts @@ -58,12 +58,13 @@ export function validateExternalUrl(url: string, additionalAllowedHosts?: string async function resolveViaSlingshot( identifier: string, slingshotUrl: string, + signal?: AbortSignal, ): Promise { const url = new URL(slingshotUrl); url.searchParams.set("identifier", identifier); try { - const response = await fetch(url.toString()); + const response = await fetch(url.toString(), { signal }); if (!response.ok) return undefined; const data = (await response.json()) as { did?: string; @@ -76,7 +77,8 @@ async function resolveViaSlingshot( handle: data.handle ?? null, pds: data.pds ?? null, }; - } catch { + } catch (error) { + if (signal?.aborted) throw signal.reason ?? error; return undefined; } } @@ -91,9 +93,12 @@ const DEFAULT_DID_RESOLVER: DidDocumentResolver = new CompositeDidDocumentResolv async function getPDSViaDidDoc( did: Did, config?: ContrailConfig, + signal?: AbortSignal, ): Promise { + signal?.throwIfAborted(); const resolver = config?.networkOverrides?.resolver ?? DEFAULT_DID_RESOLVER; const doc = await resolver.resolve(did as Did<"plc"> | Did<"web">); + signal?.throwIfAborted(); return doc.service ?.find((s) => s.id === "#atproto_pds") ?.serviceEndpoint.toString(); @@ -110,10 +115,11 @@ async function getPDSViaDidDoc( export async function resolvePDS( identifier: string, config?: ContrailConfig, + signal?: AbortSignal, ): Promise { const slingshotUrl = config?.networkOverrides?.slingshotUrl ?? SLINGSHOT_URL; const allowed = config?.networkOverrides?.additionalAllowedHosts; - const result = await resolveViaSlingshot(identifier, slingshotUrl); + const result = await resolveViaSlingshot(identifier, slingshotUrl, signal); if (result?.pds) { if (!validateExternalUrl(result.pds, allowed)) return { ...result, pds: null }; return result; @@ -122,7 +128,7 @@ export async function resolvePDS( // Fall back to DID doc resolution (only works for DIDs, not handles) if (identifier.startsWith("did:")) { try { - const pds = await getPDSViaDidDoc(identifier as Did, config); + const pds = await getPDSViaDidDoc(identifier as Did, config, signal); if (pds && validateExternalUrl(pds, allowed)) { return { did: identifier, @@ -130,7 +136,8 @@ export async function resolvePDS( pds, }; } - } catch { + } catch (error) { + if (signal?.aborted) throw signal.reason ?? error; // ignore } } @@ -174,6 +181,7 @@ export async function getPDS( did: Did, db?: Database, config?: ContrailConfig, + signal?: AbortSignal, ): Promise { const mem = pdsCacheGet(did); if (mem) return mem; @@ -182,7 +190,7 @@ export async function getPDS( const inflight = pdsInflight.get(did); if (inflight) return inflight; - const promise = resolvePDSCached(did, db, config); + const promise = resolvePDSCached(did, db, config, signal); pdsInflight.set(did, promise); try { return await promise; @@ -195,6 +203,7 @@ async function resolvePDSCached( did: Did, db?: Database, config?: ContrailConfig, + signal?: AbortSignal, ): Promise { let knownPds: string | undefined; if (db) { @@ -215,7 +224,7 @@ async function resolvePDSCached( } } - const resolved = await resolvePDS(did, config); + const resolved = await resolvePDS(did, config, signal); if (!resolved?.pds) return knownPds; pdsCacheSet(did, resolved.pds); @@ -234,16 +243,21 @@ async function resolvePDSCached( return resolved.pds; } +export function createPdsClient(pds: string): Client { + return new Client({ + handler: simpleFetchHandler({ service: pds }), + }); +} + export async function getClient( did: Did, db?: Database, config?: ContrailConfig, + signal?: AbortSignal, ): Promise { - const pds = await getPDS(did, db, config); + const pds = await getPDS(did, db, config, signal); if (!pds) throw new Error(`PDS not found for ${did}`); - return new Client({ - handler: simpleFetchHandler({ service: pds }), - }); + return createPdsClient(pds); } /** Test-only: clear module-level PDS caches. Production code MUST NOT call this. diff --git a/packages/contrail/src/core/db/records.ts b/packages/contrail/src/core/db/records.ts index bb71884..2b555b6 100644 --- a/packages/contrail/src/core/db/records.ts +++ b/packages/contrail/src/core/db/records.ts @@ -594,6 +594,73 @@ export async function lookupExistingRecords( // --- Events --- +// D1 accepts at most 100 bound parameters per statement. A record upsert uses +// seven, so fourteen rows fit while leaving the same SQL valid for SQLite and +// PostgreSQL. Packing rows avoids one prepared statement per record during +// backfill without changing the admission or derived-projection path. +const MAX_STATEMENT_BINDINGS = 100; +const RECORD_UPSERT_BINDINGS = 7; +const RECORD_UPSERT_ROWS = Math.floor( + MAX_STATEMENT_BINDINGS / RECORD_UPSERT_BINDINGS +); + +interface StorageMutation { + event: IngestEvent; + table: string; +} + +function buildRecordMutationStatements( + db: Database, + mutations: Iterable +): Statement[] { + const byTable = new Map< + string, + { upserts: IngestEvent[]; deletes: IngestEvent[] } + >(); + for (const { event, table } of mutations) { + const group = byTable.get(table) ?? { upserts: [], deletes: [] }; + if (event.operation === "delete") group.deletes.push(event); + else group.upserts.push(event); + byTable.set(table, group); + } + + const statements: Statement[] = []; + for (const [table, { upserts, deletes }] of byTable) { + for (let index = 0; index < upserts.length; index += RECORD_UPSERT_ROWS) { + const chunk = upserts.slice(index, index + RECORD_UPSERT_ROWS); + const values = chunk.map(() => "(?, ?, ?, ?, ?, ?, ?)").join(", "); + statements.push( + db + .prepare( + `INSERT INTO ${table} (uri, did, rkey, cid, record, time_us, indexed_at) VALUES ${values} ON CONFLICT(uri) DO UPDATE SET cid = excluded.cid, record = excluded.record, time_us = excluded.time_us, indexed_at = excluded.indexed_at` + ) + .bind( + ...chunk.flatMap((event) => [ + event.uri, + event.did, + event.rkey, + event.cid, + event.record, + event.time_us, + event.indexed_at, + ]) + ) + ); + } + + for (let index = 0; index < deletes.length; index += MAX_STATEMENT_BINDINGS) { + const chunk = deletes.slice(index, index + MAX_STATEMENT_BINDINGS); + const placeholders = chunk.map(() => "?").join(", "); + statements.push( + db + .prepare(`DELETE FROM ${table} WHERE uri IN (${placeholders})`) + .bind(...chunk.map((event) => event.uri)) + ); + } + } + return statements; +} + export async function projectEvents( db: Database, events: IngestEvent[], @@ -601,6 +668,8 @@ export async function projectEvents( options?: { skipReplayDetection?: boolean; skipFeedFanout?: boolean; + /** Skip FTS and relation-count maintenance during canonical bulk loading. */ + skipDerivedProjections?: boolean; /** Pre-fetched existing records — skips the internal lookup when provided */ existing?: Map; /** Ingest phase forwarded to `config.sinks`. `"live"` for jetstream / @@ -611,9 +680,13 @@ export async function projectEvents( if (events.length === 0) return; const followCollections = getFeedFollowShortNames(config); - const hasCountingRelations = Object.values(config.collections).some(c => - Object.values(c.relations ?? {}).some(r => r.count !== false) - ); + const hasCountingRelations = + !options?.skipDerivedProjections && + Object.values(config.collections).some((collection) => + Object.values(collection.relations ?? {}).some( + (relation) => relation.count !== false + ) + ); const needRecordContent = followCollections.length > 0 || hasCountingRelations; // Use pre-fetched data or look up existing records @@ -627,6 +700,9 @@ export async function projectEvents( } const batch: Statement[] = []; + // Keep only the final storage mutation for a URI within this atomic batch. + // Derived projections still inspect every admitted event below. + const storageMutations = new Map(); // Build a record-content map for feed statements (needs string values) const existingRecordStrings = new Map(); @@ -647,30 +723,15 @@ export async function projectEvents( continue; } const table = recordsTableName(short); + storageMutations.set(`${table}\0${e.uri}`, { event: e, table }); - if (e.operation === "delete") { - batch.push(db.prepare(`DELETE FROM ${table} WHERE uri = ?`).bind(e.uri)); - } else { - batch.push( - db.prepare( - `INSERT INTO ${table} (uri, did, rkey, cid, record, time_us, indexed_at) VALUES (?, ?, ?, ?, ?, ?, ?) ON CONFLICT(uri) DO UPDATE SET cid = excluded.cid, record = excluded.record, time_us = excluded.time_us, indexed_at = excluded.indexed_at` - ).bind( - e.uri, - e.did, - e.rkey, - e.cid, - e.record, - e.time_us, - e.indexed_at - ) - ); + if (!options?.skipDerivedProjections) { + // Collect count targets (deduplicated across the whole batch) + const existingRecordJson = existingMap.get(e.uri)?.record ?? null; + collectCountTargets(e, config, existingRecordJson, countTargets); + collectParentCountTargets(e, config, countTargets); } - // Collect count targets (deduplicated across the whole batch) - const existingRecordJson = existingMap.get(e.uri)?.record ?? null; - collectCountTargets(e, config, existingRecordJson, countTargets); - collectParentCountTargets(e, config, countTargets); - // Feed fanout still needs replay detection const existingInfo = existingMap.get(e.uri); const isReplay = @@ -681,9 +742,15 @@ export async function projectEvents( if (!isReplay && !options?.skipFeedFanout) { batch.push(...buildFeedStatements(db, e, config, existingRecordStrings)); } - batch.push(...buildFtsStatements(db, e, config)); + if (!options?.skipDerivedProjections) { + batch.push(...buildFtsStatements(db, e, config)); + } } + // Storage runs first so FTS, feeds, and count statements in the same atomic + // batch observe the final records. + batch.unshift(...buildRecordMutationStatements(db, storageMutations.values())); + // Build deduplicated count statements — one UPDATE per unique target batch.push(...buildBatchCountStatements(db, config, countTargets)); @@ -721,6 +788,127 @@ export async function projectEvents( } } +/** Rebuild projections that are intentionally skipped during canonical bulk + * loading. Live ingestion and scheduled retries continue maintaining them + * incrementally after this set-based catch-up. */ +export async function rebuildDerivedProjections( + db: Database, + config: ContrailConfig +): Promise { + const dialect = getDialect(db); + + if (dialect.ftsStrategy === "virtual-table") { + for (const [short, collection] of Object.entries(config.collections)) { + const fields = getSearchableFields(short, collection); + if (!fields || fields.length === 0) continue; + + const ftsTable = ftsTableName(short); + const recordsTable = recordsTableName(short); + const terms = fields.map((field) => { + const path = `$.${field}`; + return `CASE WHEN json_type(record, '${path}') = 'text' THEN json_extract(record, '${path}') ELSE '' END`; + }); + const content = `trim(${terms.join(" || ' ' || ")})`; + try { + await db.batch([ + db.prepare(`DELETE FROM ${ftsTable}`), + db.prepare( + `INSERT INTO ${ftsTable} (uri, content) + SELECT uri, content FROM ( + SELECT uri, ${content} AS content FROM ${recordsTable} + ) rebuilt + WHERE content <> ''` + ), + ]); + } catch { + // FTS5 is optional in some SQLite builds, matching schema initialization. + } + } + } + + const countStatements: Statement[] = []; + for (const [parentShort, parentConfig] of Object.entries(config.collections)) { + const parentTable = recordsTableName(parentShort); + for (const [relationName, relation] of Object.entries( + parentConfig.relations ?? {} + )) { + if (relation.count === false) continue; + + const childTable = recordsTableName(relation.collection); + const childTarget = dialect.jsonExtract( + "child.record", + getRelationField(relation) + ); + const parentTarget = relation.match === "did" ? "did" : "uri"; + const distinctExpression = relation.countDistinct + ? ["uri", "did", "rkey"].includes(relation.countDistinct) + ? `child.${relation.countDistinct}` + : dialect.jsonExtract("child.record", relation.countDistinct) + : null; + const totalCount = distinctExpression + ? `COUNT(DISTINCT ${distinctExpression})` + : "COUNT(*)"; + const totalColumn = countColumnName(relation.collection); + const projections = [`${totalCount} AS ${totalColumn}`]; + const columns = [totalColumn]; + const bindings: string[] = []; + + const mapping = (config as ResolvedContrailConfig)._resolved?.relations[ + parentShort + ]?.[relationName]; + if (relation.groupBy && mapping?.groups) { + const childGroup = dialect.jsonExtract( + "child.record", + relation.groupBy + ); + for (const [groupKey, token] of Object.entries(mapping.groups)) { + const column = groupedCountColumnName( + relation.collection, + groupKey + ); + const groupedCount = distinctExpression + ? `COUNT(DISTINCT CASE WHEN ${childGroup} = ? THEN ${distinctExpression} END)` + : `SUM(CASE WHEN ${childGroup} = ? THEN 1 ELSE 0 END)`; + projections.push(`${groupedCount} AS ${column}`); + columns.push(column); + bindings.push(token); + } + } + + // Reset parents with no matching children, then update only parents that + // appear in one grouped child scan. This avoids N parents × M correlated + // COUNT subqueries during a historical rebuild. + countStatements.push( + db.prepare( + `UPDATE ${parentTable} SET ${columns + .map((column) => `${column} = 0`) + .join(", ")}` + ) + ); + const aggregateUpdate = db.prepare( + `WITH derived AS ( + SELECT ${childTarget} AS target, ${projections.join(", ")} + FROM ${childTable} child + GROUP BY ${childTarget} + ) + UPDATE ${parentTable} AS parent + SET ${columns + .map((column) => `${column} = derived.${column}`) + .join(", ")} + FROM derived + WHERE parent.${parentTarget} = derived.target` + ); + countStatements.push( + bindings.length > 0 + ? aggregateUpdate.bind(...bindings) + : aggregateUpdate + ); + } + } + + if (countStatements.length > 0) await db.batch(countStatements); +} + function safeParseJson(s: string): Record { try { const v = JSON.parse(s); diff --git a/packages/contrail/src/core/ingest.ts b/packages/contrail/src/core/ingest.ts index 1d6828c..bcc4244 100644 --- a/packages/contrail/src/core/ingest.ts +++ b/packages/contrail/src/core/ingest.ts @@ -43,6 +43,9 @@ export function createIngestEvent(input: RecordEventInput): IngestEvent { export interface IngestRecordsOptions { skipReplayDetection?: boolean; skipFeedFanout?: boolean; + /** Skip FTS and relation-count maintenance for a bulk load that will rebuild + * both projections once canonical records are durable. */ + skipDerivedProjections?: boolean; /** Pre-fetched rows, used by immediate synchronization. */ existing?: Map; phase?: "live" | "backfill"; diff --git a/packages/contrail/src/core/status.ts b/packages/contrail/src/core/status.ts index a21104d..1a40767 100644 --- a/packages/contrail/src/core/status.ts +++ b/packages/contrail/src/core/status.ts @@ -111,6 +111,7 @@ export async function getBackfillStatus( collectionRows, retryRow, runRow, + projectionRow, ] = await Promise.all([ db .prepare( @@ -195,7 +196,12 @@ export async function getBackfillStatus( run_id: string | null; heartbeat_at: number | null; finished_at: number | null; - }>() + }>(), + db + .prepare( + "SELECT value FROM _contrail_meta WHERE key = 'backfill_derived_projections_dirty'" + ) + .first<{ value: string }>() ]); const taskCounts = counts(taskRow); @@ -230,9 +236,12 @@ export async function getBackfillStatus( !!runRow?.run_id && runRow.finished_at === null && number(runRow.heartbeat_at) >= now - BACKFILL_RUN_STALE_MS; + const projectionDirty = projectionRow?.value === "1"; const state = running ? "running" - : taskCounts.total === 0 && discovery.total === 0 + : projectionDirty + ? "incomplete" + : taskCounts.total === 0 && discovery.total === 0 ? expectsDiscovery ? "not_started" : "complete" diff --git a/packages/contrail/src/workers/backfill.ts b/packages/contrail/src/workers/backfill.ts index feb9a24..2eed322 100644 --- a/packages/contrail/src/workers/backfill.ts +++ b/packages/contrail/src/workers/backfill.ts @@ -24,9 +24,13 @@ interface WranglerCommon { } export interface BackfillAllViaWranglerOptions extends WranglerCommon { - /** Passed through to `contrail.backfillAll()`. Default: 100. */ + /** Concurrent identity resolutions. Default: 100. */ concurrency?: number; - /** Failed attempts allowed in this run before a row remains pending. */ + /** PDS hosts fetched concurrently. Default: 20. */ + pdsConcurrency?: number; + /** Accounts fetched concurrently from each PDS. Default: 3. */ + didsPerPds?: number; + /** Immediate attempts before scheduled retries take over. Default: 1. */ maxAttempts?: number; /** Override the built-in progress logging. */ onProgress?: BackfillAllOptions["onProgress"]; @@ -74,6 +78,8 @@ export async function backfillAll( contrail.backfillAll( { concurrency: opts.concurrency ?? 100, + pdsConcurrency: opts.pdsConcurrency, + didsPerPds: opts.didsPerPds, maxAttempts: opts.maxAttempts, onProgress: opts.onProgress, }, diff --git a/packages/contrail/tests/backfill-status.test.ts b/packages/contrail/tests/backfill-status.test.ts index 4c9a0a0..bce4648 100644 --- a/packages/contrail/tests/backfill-status.test.ts +++ b/packages/contrail/tests/backfill-status.test.ts @@ -6,12 +6,14 @@ import { discoverDIDs, finishBackfillRun, getBackfillStatus, + initSchema, + queryRecords, resolveConfig, retryPendingBackfills, tryStartBackfillRun } from "../src/index"; import { __resetPdsCachesForTests } from "../src/core/client"; -import { TEST_CONFIG, createTestDbWithSchema, ingestRecords, makeEvent } from "./helpers"; +import { TEST_CONFIG, createTestDb, createTestDbWithSchema, ingestRecords, makeEvent } from "./helpers"; const DID = "did:plc:backfilltest"; const EVENT = "community.lexicon.calendar.event"; @@ -73,6 +75,281 @@ describe("backfill failure state", () => { expect(row?.next_retry_at).toBeGreaterThan(Date.now()); }); + it("admits dependent records using all relay-discovered actors", async () => { + const actor = "did:plc:dependent-actor"; + const subject = "did:plc:discovered-subject"; + const follow = "app.bsky.graph.follow"; + const config = resolveConfig({ + namespace: "com.example", + collections: { + event: { collection: EVENT }, + follow: { + collection: follow, + discover: false, + subjectField: "subject" + } + } + }); + const db = createTestDb(); + await initSchema(db, config); + await db + .prepare("INSERT INTO identities (did, handle, pds, resolved_at) VALUES (?, ?, ?, ?)") + .bind(actor, "actor.test", "https://pds.test", Date.now()) + .run(); + await db + .prepare("INSERT INTO backfills (did, collection, completed) VALUES (?, ?, 0)") + .bind(actor, follow) + .run(); + await db + .prepare("INSERT INTO backfills (did, collection, completed) VALUES (?, ?, 1)") + .bind(subject, EVENT) + .run(); + + fetchSpy.mockResolvedValue( + new Response( + JSON.stringify({ + records: [ + { + uri: `at://${actor}/${follow}/one`, + cid: "follow-cid", + value: { $type: follow, subject, createdAt: "2026-01-01T00:00:00Z" } + } + ] + }), + { status: 200, headers: { "content-type": "application/json" } } + ) + ); + + await backfillPending(db, config, { concurrency: 1, maxAttempts: 1 }); + + expect((await queryRecords(db, config, { collection: "follow" })).records).toHaveLength(1); + }); + + it("starts the next queued account without waiting for a slow worker", async () => { + const db = await createTestDbWithSchema(); + const dids = ["did:plc:queue-a", "did:plc:queue-b", "did:plc:queue-c"]; + for (const [index, did] of dids.entries()) { + await db + .prepare("INSERT INTO identities (did, handle, pds, resolved_at) VALUES (?, ?, ?, ?)") + .bind(did, `${did}.test`, `https://pds-${index}.test`, Date.now()) + .run(); + await db + .prepare("INSERT INTO backfills (did, collection, completed) VALUES (?, ?, 0)") + .bind(did, EVENT) + .run(); + } + + let releaseSlow!: () => void; + const slow = new Promise((resolve) => { + releaseSlow = resolve; + }); + let observedThird!: () => void; + const thirdStarted = new Promise((resolve) => { + observedThird = resolve; + }); + fetchSpy.mockImplementation(async (input) => { + const repo = new URL(input instanceof Request ? input.url : String(input)).searchParams.get("repo"); + if (repo === dids[0]) await slow; + if (repo === dids[2]) observedThird(); + return new Response(JSON.stringify({ records: [] }), { + status: 200, + headers: { "content-type": "application/json" } + }); + }); + + const backfill = backfillPending(db, TEST_CONFIG, { + concurrency: 2, + pdsConcurrency: 2, + didsPerPds: 1, + maxAttempts: 1 + }); + try { + await Promise.race([ + thirdStarted, + new Promise((_, reject) => + setTimeout(() => reject(new Error("third account stayed behind the slow batch")), 500) + ) + ]); + } finally { + releaseSlow(); + } + await backfill; + + expect((await getBackfillStatus(db, TEST_CONFIG)).accounts.complete).toBe(3); + }); + + it("requeues a paginated account behind other waiting accounts", async () => { + const db = await createTestDbWithSchema(); + const first = "did:plc:paged-a"; + const second = "did:plc:paged-b"; + for (const did of [first, second]) { + await db + .prepare("INSERT INTO identities (did, handle, pds, resolved_at) VALUES (?, ?, ?, ?)") + .bind(did, `${did}.test`, "https://pds.test", Date.now()) + .run(); + await db + .prepare("INSERT INTO backfills (did, collection, completed) VALUES (?, ?, 0)") + .bind(did, EVENT) + .run(); + } + + const order: string[] = []; + fetchSpy.mockImplementation(async (input) => { + const url = new URL(input instanceof Request ? input.url : String(input)); + const repo = url.searchParams.get("repo")!; + const cursor = url.searchParams.get("cursor"); + order.push(repo); + if (repo === first && cursor === null) { + return new Response( + JSON.stringify({ + records: [ + { + uri: `at://${first}/${EVENT}/one`, + cid: "event-cid", + value: { name: "Queued", startsAt: "2026-01-01T00:00:00Z" } + } + ], + cursor: "next" + }), + { status: 200, headers: { "content-type": "application/json" } } + ); + } + return new Response(JSON.stringify({ records: [] }), { + status: 200, + headers: { "content-type": "application/json" } + }); + }); + + await backfillPending(db, TEST_CONFIG, { + concurrency: 1, + pdsConcurrency: 1, + didsPerPds: 1, + maxAttempts: 1 + }); + + expect(order).toEqual([first, second, first]); + }); + + it("bounds active PDS hosts and accounts per host", async () => { + const db = await createTestDbWithSchema(); + const dids = Array.from({ length: 9 }, (_, index) => + `did:plc:host-limit-${index}` + ); + for (const [index, did] of dids.entries()) { + await db + .prepare("INSERT INTO identities (did, handle, pds, resolved_at) VALUES (?, ?, ?, ?)") + .bind( + did, + `${did}.test`, + `https://pds-${Math.floor(index / 3)}.test`, + Date.now() + ) + .run(); + await db + .prepare("INSERT INTO backfills (did, collection, completed) VALUES (?, ?, 0)") + .bind(did, EVENT) + .run(); + } + + let activeRequests = 0; + let maxRequests = 0; + let maxHosts = 0; + const activeByHost = new Map(); + fetchSpy.mockImplementation(async (input) => { + const host = new URL(input instanceof Request ? input.url : String(input)).host; + activeRequests++; + maxRequests = Math.max(maxRequests, activeRequests); + activeByHost.set(host, (activeByHost.get(host) ?? 0) + 1); + maxHosts = Math.max( + maxHosts, + [...activeByHost.values()].filter((count) => count > 0).length + ); + expect(activeByHost.get(host)).toBeLessThanOrEqual(2); + await new Promise((resolve) => setTimeout(resolve, 10)); + activeRequests--; + activeByHost.set(host, activeByHost.get(host)! - 1); + return new Response(JSON.stringify({ records: [] }), { + status: 200, + headers: { "content-type": "application/json" } + }); + }); + + await backfillPending(db, TEST_CONFIG, { + concurrency: 9, + pdsConcurrency: 2, + didsPerPds: 2 + }); + + expect(maxRequests).toBeLessThanOrEqual(4); + expect(maxHosts).toBeLessThanOrEqual(2); + expect((await getBackfillStatus(db)).accounts.complete).toBe(9); + }); + + it("defers a failed initial account without retrying it immediately", async () => { + const db = await createTestDbWithSchema(); + await db + .prepare("INSERT INTO identities (did, handle, pds, resolved_at) VALUES (?, ?, ?, ?)") + .bind(DID, "backfill.test", "https://pds.test", Date.now()) + .run(); + await db + .prepare("INSERT INTO backfills (did, collection, completed) VALUES (?, ?, 0)") + .bind(DID, EVENT) + .run(); + fetchSpy.mockResolvedValue(new Response("unavailable", { status: 503 })); + + await backfillPending(db, TEST_CONFIG, { + pdsConcurrency: 1, + didsPerPds: 1 + }); + + expect(fetchSpy).toHaveBeenCalledTimes(1); + const row = await db + .prepare("SELECT retries, next_retry_at FROM backfills WHERE did = ?") + .bind(DID) + .first<{ retries: number; next_retry_at: number | null }>(); + expect(row?.retries).toBe(1); + expect(row?.next_retry_at).toBeGreaterThan(Date.now()); + expect((await getBackfillStatus(db)).accounts.retrying).toBe(1); + }); + + it("aborts a timed-out PDS request instead of leaving it running", async () => { + const db = await createTestDbWithSchema(); + await db + .prepare("INSERT INTO identities (did, handle, pds, resolved_at) VALUES (?, ?, ?, ?)") + .bind(DID, "backfill.test", "https://pds.test", Date.now()) + .run(); + await db + .prepare("INSERT INTO backfills (did, collection, completed) VALUES (?, ?, 0)") + .bind(DID, EVENT) + .run(); + + let aborted = false; + fetchSpy.mockImplementation((_input, init) => + new Promise((_resolve, reject) => { + const signal = init?.signal; + if (!signal) return reject(new Error("missing abort signal")); + signal.addEventListener( + "abort", + () => { + aborted = true; + reject(signal.reason); + }, + { once: true } + ); + }) + ); + + await backfillPending(db, TEST_CONFIG, { + pdsConcurrency: 1, + didsPerPds: 1, + requestTimeoutMs: 10 + }); + + expect(aborted).toBe(true); + expect(fetchSpy).toHaveBeenCalledTimes(1); + expect((await getBackfillStatus(db)).accounts.retrying).toBe(1); + }); + it("retries pending rows on the next run and only then completes them", async () => { const db = await createTestDbWithSchema(); await db @@ -140,6 +417,25 @@ describe("scheduled backfill retries", () => { __resetPdsCachesForTests(); }); + it("repairs an interrupted derived rebuild before scheduled retry work", async () => { + const db = await createTestDbWithSchema(); + await db + .prepare("INSERT INTO backfills (did, collection, completed) VALUES (?, ?, 1)") + .bind(DID, EVENT) + .run(); + await db + .prepare( + "INSERT INTO _contrail_meta (key, value) VALUES ('backfill_derived_projections_dirty', '1') ON CONFLICT(key) DO UPDATE SET value = '1'" + ) + .run(); + + expect((await getBackfillStatus(db, TEST_CONFIG)).state).toBe("incomplete"); + expect( + await retryPendingBackfills(db, TEST_CONFIG, { maxAccounts: 1 }) + ).toMatchObject({ attempted: 0, skipped: false }); + expect((await getBackfillStatus(db, TEST_CONFIG)).state).toBe("complete"); + }); + it("waits for persisted backoff and then completes a due account", async () => { const db = await createTestDbWithSchema(); await db diff --git a/packages/contrail/tests/helpers.ts b/packages/contrail/tests/helpers.ts index e4ba767..b4a6105 100644 --- a/packages/contrail/tests/helpers.ts +++ b/packages/contrail/tests/helpers.ts @@ -55,6 +55,7 @@ export function ingestRecords( options?: { skipReplayDetection?: boolean; skipFeedFanout?: boolean; + skipDerivedProjections?: boolean; existing?: Map; } ): Promise { diff --git a/packages/contrail/tests/ingest.test.ts b/packages/contrail/tests/ingest.test.ts index aa66038..38dafea 100644 --- a/packages/contrail/tests/ingest.test.ts +++ b/packages/contrail/tests/ingest.test.ts @@ -7,6 +7,7 @@ import { queryRecords, resolveConfig, type Database, + type Statement, } from "../src/index"; const logger = { log() {}, warn() {}, error() {} }; @@ -71,6 +72,40 @@ describe("ingestRecords", () => { expect(stored.records.map((record) => record.rkey)).toEqual(["kept"]); }); + it("packs record writes into statements within D1's binding limit", async () => { + const sql: string[] = []; + const observedDb: Database = { + prepare(statement) { + sql.push(statement); + return db.prepare(statement); + }, + batch(statements: Statement[]) { + return db.batch(statements); + }, + dialect: db.dialect, + }; + const events = Array.from({ length: 30 }, (_, index) => + createIngestEvent({ + did: "did:plc:alice", + collection: "com.example.event", + rkey: `bulk-${index}`, + operation: "create", + cid: `cid-${index}`, + value: { keep: true }, + timeUs: index + 1, + }), + ); + + await ingestRecords(observedDb, events, config); + + const inserts = sql.filter((statement) => + statement.startsWith("INSERT INTO records_event"), + ); + expect(inserts).toHaveLength(3); + expect(inserts.every((statement) => (statement.match(/\?/g) ?? []).length <= 100)).toBe(true); + expect((await queryRecords(db, config, { collection: "event" })).records).toHaveLength(30); + }); + it("drops malformed records and untracked collections before projection", async () => { const malformed = createIngestEvent({ did: "did:plc:alice", diff --git a/packages/contrail/tests/records.test.ts b/packages/contrail/tests/records.test.ts index b415af0..604decb 100644 --- a/packages/contrail/tests/records.test.ts +++ b/packages/contrail/tests/records.test.ts @@ -2,6 +2,7 @@ import { describe, it, expect, beforeEach, vi } from "vitest"; import type { Database } from "../src/index"; import { ingestRecords, createTestDbWithSchema, makeEvent, TEST_CONFIG } from "./helpers"; import { queryRecords, getLastCursor, saveCursor } from "../src/index"; +import { rebuildDerivedProjections } from "../src/core/db/records"; let db: Database; @@ -100,6 +101,40 @@ describe("ingestRecords", () => { expect(result.records[0].counts!["rsvp"]).toBe(1); }); + it("rebuilds deferred relation counts after bulk canonical writes", async () => { + const eventUri = "at://did:plc:test/community.lexicon.calendar.event/deferred"; + await ingestRecords( + db, + [ + makeEvent({ uri: eventUri, rkey: "deferred" }), + makeEvent({ + uri: "at://did:plc:user1/community.lexicon.calendar.rsvp/deferred", + did: "did:plc:user1", + collection: "community.lexicon.calendar.rsvp", + rkey: "deferred", + record: { + subject: { uri: eventUri }, + status: "community.lexicon.calendar.rsvp#going", + }, + }), + ], + { skipDerivedProjections: true } + ); + + let result = await queryRecords(db, TEST_CONFIG, { collection: "event" }); + expect(result.records[0].counts?.rsvp).toBeUndefined(); + + await rebuildDerivedProjections(db, TEST_CONFIG); + + result = await queryRecords(db, TEST_CONFIG, { collection: "event" }); + expect(result.records[0].counts?.rsvp).toBe(1); + expect( + result.records[0].counts?.[ + "community.lexicon.calendar.rsvp#going" + ] + ).toBe(1); + }); + it("recounts children that arrived before their parent", async () => { const eventUri = "at://did:plc:test/community.lexicon.calendar.event/late"; diff --git a/packages/contrail/tests/search.test.ts b/packages/contrail/tests/search.test.ts index 3a0ed69..3bc2c2e 100644 --- a/packages/contrail/tests/search.test.ts +++ b/packages/contrail/tests/search.test.ts @@ -4,6 +4,7 @@ import { resolveConfig } from "../src/index"; import { createTestDb, makeEvent } from "./helpers"; import { initSchema } from "../src/index"; import { ingestRecords, queryRecords } from "../src/index"; +import { rebuildDerivedProjections } from "../src/core/db/records"; // Detect FTS5 support at module level (node:sqlite doesn't include it) let hasFts = false; @@ -86,6 +87,43 @@ describe.skipIf(!hasFts)("FTS with explicit searchable fields", () => { ); }); + it("rebuilds deferred FTS rows after bulk canonical writes", async () => { + await ingestRecords( + db, + [ + makeEvent({ + uri: `at://did:plc:deferred/${collection}/deferred`, + did: "did:plc:deferred", + collection, + rkey: "deferred", + record: { + name: "DeferredNeedle", + mode: "online", + description: "bulk loaded", + }, + }), + ], + SEARCH_CONFIG, + { skipDerivedProjections: true } + ); + + expect( + (await queryRecords(db, SEARCH_CONFIG, { + collection, + search: "DeferredNeedle", + })).records + ).toHaveLength(0); + + await rebuildDerivedProjections(db, SEARCH_CONFIG); + + expect( + (await queryRecords(db, SEARCH_CONFIG, { + collection, + search: "DeferredNeedle", + })).records + ).toHaveLength(1); + }); + it("finds records matching a search term", async () => { const result = await queryRecords(db, SEARCH_CONFIG, { collection, search: "Rust" }); expect(result.records).toHaveLength(2); -- 2.51.2