diff --git a/.changeset/backfill-status.md b/.changeset/backfill-status.md index c0c16da..b37ebd3 100644 --- a/.changeset/backfill-status.md +++ b/.changeset/backfill-status.md @@ -2,4 +2,4 @@ "@atmo-dev/contrail": patch --- -Keep failed PDS and relay work pending with bounded retries instead of reporting partial backfills as complete. Add durable backfill and discovery state to the JSON overview and `/status` endpoint, including pending tasks and unreachable-account counts. +Keep failed PDS and relay work pending with bounded retries instead of marking account rows complete. Retry due PDS accounts automatically in small scheduled slices with persisted exponential backoff and an overlap lease. Add durable backfill and discovery state to the JSON overview and `/status` endpoint, including running/complete state, pending tasks, retry timing, and unreachable-account counts. diff --git a/README.md b/README.md index f9f1eca..793a4b8 100644 --- a/README.md +++ b/README.md @@ -79,7 +79,7 @@ GET /xrpc/com.example.event.listRecords?startsAtMin=2026-01-01&limit=10 GET /status ``` -The JSON status response reports live cursor lag, indexed records, known backfill progress, and currently unreachable accounts. Failed PDS work remains pending and is retried on the next backfill run. +The JSON status response reports live cursor lag, indexed records, known backfill progress, and currently unreachable accounts. Failed PDS work remains pending and is retried automatically in small, backed-off slices after scheduled live ingestion. For ordinary Lexicon parsing, validation, pulling, and TypeScript generation, use [Atcute](https://github.com/mary-ext/atcute) directly. Contrail no longer ships a separate Lexicon toolchain. diff --git a/apps/cloudflare-workers/README.md b/apps/cloudflare-workers/README.md index 8a2a15b..f229da9 100644 --- a/apps/cloudflare-workers/README.md +++ b/apps/cloudflare-workers/README.md @@ -39,13 +39,13 @@ pnpm contrail backfill # backfill against the local D1 created by wrangler `pnpm dev` runs `contrail dev` under the hood. On start it inspects the local D1 and offers to start or resume backfill whenever known work remains, then runs `wrangler dev --test-scheduled` with a 60s timer hitting `/__scheduled` so the cron actually fires locally. -A failed PDS stays pending rather than being marked complete. The backfill command retries failures with a bounded attempt budget and exits unsuccessfully if work remains; rerun it later to give unreachable accounts a fresh retry budget. +A failed PDS account stays pending rather than being marked complete. The initial command uses a bounded attempt budget, then records an exponential retry time. The normal one-minute Worker cron retries a small due slice after live ingestion, so dead PDSes cannot monopolize the tick. Rerunning the command forces an immediate bounded pass. ```bash curl http://localhost:8787/status ``` -This reports the durable state without exposing account DIDs or raw upstream errors. +This reports the durable state without exposing account DIDs or raw upstream errors. `state: "running"` means a manual or scheduled backfill slice currently owns the database lease. `state: "complete"` means discovery and the initial pass finished; `accounts.unreachable` and `retries` may still show deferred account work that cron will continue retrying. ## extending diff --git a/docs/01-indexing.md b/docs/01-indexing.md index b20b409..b4de6d6 100644 --- a/docs/01-indexing.md +++ b/docs/01-indexing.md @@ -58,7 +58,7 @@ await contrail.backfill({ concurrency: 100 }); // fetch history for registered D `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 invocation has a bounded failure budget (five attempts by default), so one dead PDS cannot hang the whole command forever. A later invocation resets that budget and retries the pending rows. `backfillAll()` returns a durable `status` summary alongside the number of discovered accounts and accepted records. +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, receive an exponential `next_retry_at`, and remain incomplete. Cloudflare's scheduled Worker retries a small due slice after each live-ingest cycle; an explicit later invocation forces another bounded pass. `backfillAll()` returns a durable `status` summary alongside the number of discovered accounts and accepted records. ### Workers CLI @@ -69,7 +69,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). The command exits unsuccessfully when known work remains; rerun it later rather than treating the partial result as complete. +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. If you'd rather embed backfill inside your own script, `@atmo-dev/contrail/workers` exports the same logic as a function: @@ -93,7 +93,10 @@ Workers can't hold long-lived connections, so run one catch-up cycle per cron fi ```ts // wrangler.jsonc: "triggers": { "crons": ["*/1 * * * *"] } async scheduled(_ev, env, ctx) { - ctx.waitUntil(contrail.ingest({}, env.DB)); + ctx.waitUntil((async () => { + await contrail.ingest({}, env.DB); + await contrail.retryBackfill({}, env.DB); // small due slice + })()); } ``` diff --git a/docs/02-querying.md b/docs/02-querying.md index b9f5647..e54643a 100644 --- a/docs/02-querying.md +++ b/docs/02-querying.md @@ -35,8 +35,11 @@ Dotted field names become camelCase params — `queryable: { "subject.uri": {} } - discovery sources complete, pending, or failed; - known account totals, including pending and currently unreachable accounts; - account/collection tasks complete, pending, or failed; -- known-task completion percentage; and -- per-collection backfill progress. +- known-task completion percentage; +- per-collection backfill progress; and +- scheduled/due account retries plus the next retry time. + +`state` is `running` while a manual or scheduled slice holds the backfill lease. It becomes `complete` after discovery and the initial pass finish, even when `accounts.unreachable` and `tasks.failed` are non-zero; those account rows are still incomplete and remain scheduled for background retry. `incomplete` means discovery or an initial account attempt has not finished. "Known" is deliberate: while relay discovery is incomplete, Contrail cannot honestly claim how many accounts remain undiscovered. `/health` remains a lightweight liveness response and does not claim that historical backfill is complete. diff --git a/packages/contrail/src/contrail.ts b/packages/contrail/src/contrail.ts index d078302..184f066 100644 --- a/packages/contrail/src/contrail.ts +++ b/packages/contrail/src/contrail.ts @@ -18,8 +18,12 @@ import { } from "./core/jetstream"; import { backfillPending, + discoverAndBackfill, discoverDIDs, + retryPendingBackfills, type BackfillAllOptions, + type BackfillRetryOptions, + type BackfillRetryResult, } from "./core/backfill"; import { getBackfillStatus, type BackfillStatus } from "./core/status"; import { @@ -166,6 +170,15 @@ export class Contrail { return backfillPending(this.getDb(db), this.config, options); } + /** Retry a bounded slice of due or interrupted account backfills. Intended + * for scheduled runtimes; persisted backoff prevents hammering failures. */ + async retryBackfill( + options?: BackfillRetryOptions, + db?: Database + ): Promise { + return retryPendingBackfills(this.getDb(db), this.config, options); + } + /** Discover every DID with records in the configured collections, then * backfill their history. Logs progress via `config.logger` — supply * `onProgress` to take over output, or pass a no-op logger in the config @@ -182,10 +195,6 @@ export class Contrail { const logger = this.config.logger; const startedAt = Date.now(); - logger?.log?.("discovering users…"); - const discovered = await this.discover(d); - logger?.log?.(` discovered ${discovered.length} users`); - // Wrap the call with a throttled default progress logger when the // caller hasn't supplied their own. Throttle at 2s so we don't spam in // fast/local runs; final summary always prints. @@ -206,17 +215,32 @@ export class Contrail { }; } - logger?.log?.("backfilling…"); - const backfilled = await this.backfill(effective, d); + logger?.log?.("discovering users…"); + const result = await discoverAndBackfill( + d, + this.config, + effective, + (count) => { + logger?.log?.(` discovered ${count} users`); + logger?.log?.("backfilling…"); + } + ); + const discovered = result.discovered; + const backfilled = result.backfilled; + const status = await getBackfillStatus(d, this.config); const elapsedS = ((Date.now() - startedAt) / 1000).toFixed(1); - if (status.state === "complete") { + if (status.state === "complete" && status.accounts.unreachable > 0) { + logger?.warn?.( + ` done with deferred failures: ${status.accounts.complete}/${status.accounts.total} known accounts complete; ${status.accounts.unreachable} unavailable accounts scheduled for retry (${elapsedS}s)` + ); + } else if (status.state === "complete") { logger?.log?.( ` done: ${backfilled} records; ${status.accounts.complete}/${status.accounts.total} known accounts complete in ${elapsedS}s` ); } else { logger?.warn?.( - ` incomplete: ${status.accounts.pending} known accounts remain; ${status.accounts.unreachable} currently unreachable (${elapsedS}s)` + ` interrupted: ${status.accounts.pending} known accounts still need an initial attempt (${elapsedS}s)` ); } return { discovered: discovered.length, backfilled, status }; diff --git a/packages/contrail/src/core/backfill.ts b/packages/contrail/src/core/backfill.ts index 925d1b8..368ae89 100644 --- a/packages/contrail/src/core/backfill.ts +++ b/packages/contrail/src/core/backfill.ts @@ -13,6 +13,11 @@ import { import { getLastCursor, saveCursor } from "./db"; import { createIngestEvent, ingestRecords } from "./ingest"; import { getClient, getPDS } from "./client"; +import { + finishBackfillRun, + heartbeatBackfillRun, + tryStartBackfillRun, +} from "./status"; const DEFAULT_TIME_FIELD = "createdAt"; @@ -44,6 +49,8 @@ const PAGE_SIZE = 100; const DEFAULT_MAX_ATTEMPTS = 5; const REQUEST_TIMEOUT_MS = 10_000; +const BACKFILL_RETRY_BASE_MS = 60_000; +const BACKFILL_RETRY_MAX_MS = 60 * 60_000; function positiveInteger(value: number | undefined, fallback: number): number { return typeof value === "number" && Number.isFinite(value) && value > 0 @@ -87,11 +94,30 @@ async function markFailed( collection: string, error: unknown ): Promise { + const row = await db + .prepare( + "SELECT retries FROM backfills WHERE did = ? AND collection = ?" + ) + .bind(did, collection) + .first<{ retries: number }>(); + const retries = (row?.retries ?? 0) + 1; + const now = Date.now(); + const retryDelay = Math.min( + BACKFILL_RETRY_BASE_MS * 2 ** Math.min(retries - 1, 6), + BACKFILL_RETRY_MAX_MS + ); await db .prepare( - "UPDATE backfills SET retries = retries + 1, last_error = ?, last_attempt_at = ? WHERE did = ? AND collection = ?" + "UPDATE backfills SET retries = ?, last_error = ?, last_attempt_at = ?, next_retry_at = ? WHERE did = ? AND collection = ?" + ) + .bind( + retries, + errorMessage(error), + now, + now + retryDelay, + did, + collection ) - .bind(errorMessage(error), Date.now(), did, collection) .run(); } @@ -235,7 +261,7 @@ async function backfillUserAttempt( await db .prepare( - "UPDATE backfills SET pds_cursor = ?, last_error = NULL, last_attempt_at = ? WHERE did = ? AND collection = ?" + "UPDATE backfills SET pds_cursor = ?, retries = 0, last_error = NULL, last_attempt_at = ?, next_retry_at = NULL WHERE did = ? AND collection = ?" ) .bind(currentCursor ?? null, now, did, collection) .run(); @@ -253,7 +279,7 @@ async function backfillUserAttempt( if (done) { await db .prepare( - "UPDATE backfills SET completed = 1, last_error = NULL, last_attempt_at = ? WHERE did = ? AND collection = ?" + "UPDATE backfills SET completed = 1, retries = 0, last_error = NULL, last_attempt_at = ?, next_retry_at = NULL WHERE did = ? AND collection = ?" ) .bind(Date.now(), did, collection) .run(); @@ -298,10 +324,11 @@ export interface BackfillAllOptions { onProgress?: (progress: BackfillProgress) => void; } -export async function backfillPending( +async function backfillPendingWork( db: Database, config: ContrailConfig, - options?: BackfillAllOptions + options: BackfillAllOptions | undefined, + runId: string ): Promise { const concurrency = positiveInteger(options?.concurrency, 100); const maxAttempts = positiveInteger( @@ -342,6 +369,7 @@ export async function backfillPending( } 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. @@ -409,6 +437,7 @@ export async function backfillPending( else failedDids.push(result.did); } emitProgress(); + await heartbeatBackfillRun(db, runId); } // Slow retry pass: use normal timeouts and request backoff for DIDs that did @@ -467,6 +496,7 @@ export async function backfillPending( else usersFailed++; } emitProgress(); + await heartbeatBackfillRun(db, runId); } totalBackfilled += roundBackfilled; @@ -478,6 +508,128 @@ export async function backfillPending( return totalBackfilled; } +export async function backfillPending( + db: Database, + config: ContrailConfig, + options?: BackfillAllOptions +): Promise { + const runId = await tryStartBackfillRun(db); + if (!runId) throw new Error("A backfill is already running"); + try { + return await backfillPendingWork(db, config, options, runId); + } finally { + await finishBackfillRun(db, runId); + } +} + +export interface BackfillRetryOptions { + /** Maximum accounts to attempt in one scheduled slice. Default: 5. */ + maxAccounts?: number; + /** Total wall-clock budget for the slice. Default: 10000ms. */ + timeoutMs?: number; + /** Deadline for each PDS request within the slice. Default: 3000ms. */ + requestTimeoutMs?: number; +} + +export interface BackfillRetryResult { + attempted: number; + completed: number; + failed: number; + records: number; + skipped: boolean; +} + +/** Retry a small due slice without resetting persisted failure backoff. Safe for + * scheduled runtimes: one database-backed run lease prevents overlap. */ +export async function retryPendingBackfills( + db: Database, + config: ContrailConfig, + options?: BackfillRetryOptions +): Promise { + const runId = await tryStartBackfillRun(db); + if (!runId) { + return { attempted: 0, completed: 0, failed: 0, records: 0, skipped: true }; + } + + const maxAccounts = positiveInteger(options?.maxAccounts, 5); + const timeoutMs = positiveInteger(options?.timeoutMs, 10_000); + const requestTimeoutMs = positiveInteger(options?.requestTimeoutMs, 3_000); + const deadline = Date.now() + timeoutMs; + let attempted = 0; + let completed = 0; + let failed = 0; + let records = 0; + + try { + const due = await db + .prepare( + "SELECT did FROM backfills WHERE completed = 0 AND (next_retry_at IS NULL OR next_retry_at <= ?) GROUP BY did ORDER BY MIN(COALESCE(next_retry_at, 0)), did LIMIT ?" + ) + .bind(Date.now(), maxAccounts) + .all<{ did: string }>(); + + for (const { did } of due.results ?? []) { + if (Date.now() >= deadline) break; + attempted++; + const pending = await db + .prepare( + "SELECT collection FROM backfills WHERE did = ? AND completed = 0 AND (next_retry_at IS NULL OR next_retry_at <= ?) ORDER BY collection" + ) + .bind(did, Date.now()) + .all<{ collection: string }>(); + const collections = (pending.results ?? []).map((row) => row.collection); + if (collections.length === 0) continue; + + let client: Client; + try { + client = await withRetry( + () => getClient(did as Did, db, config), + `getClient(${did})`, + 0, + Math.min(requestTimeoutMs, Math.max(1, deadline - Date.now())) + ); + } catch (error) { + for (const collection of collections) { + await markFailed(db, did, collection, error); + } + failed++; + await heartbeatBackfillRun(db, runId); + continue; + } + + const attempts: BackfillUserAttempt[] = []; + for (const collection of collections) { + if (Date.now() >= deadline) break; + attempts.push( + await backfillUserAttempt(db, did, collection, deadline, config, { + client, + skipReplayDetection: true, + maxRetries: 0, + requestTimeout: Math.min( + requestTimeoutMs, + Math.max(1, deadline - Date.now()) + ), + }) + ); + } + records += attempts.reduce((sum, attempt) => sum + attempt.records, 0); + if ( + attempts.length === collections.length && + attempts.every((attempt) => attempt.completed) + ) { + completed++; + } else { + failed++; + } + await heartbeatBackfillRun(db, runId); + } + } finally { + await finishBackfillRun(db, runId); + } + + return { attempted, completed, failed, records, skipped: false }; +} + // --- Discovery --- interface DiscoveryPage { @@ -678,3 +830,34 @@ export async function discoverDIDs( return discovered; } + +export interface DiscoverAndBackfillResult { + discovered: string[]; + backfilled: number; +} + +/** Hold one lease across relay discovery and the complete initial PDS pass. */ +export async function discoverAndBackfill( + db: Database, + config: ContrailConfig, + options?: BackfillAllOptions, + onDiscovered?: (count: number) => void +): Promise { + const runId = await tryStartBackfillRun(db); + if (!runId) throw new Error("A backfill is already running"); + + const allDiscovered = new Set(); + try { + while (true) { + const dids = await discoverDIDs(db, config, Infinity); + await heartbeatBackfillRun(db, runId); + if (dids.length === 0) break; + for (const did of dids) allDiscovered.add(did); + } + onDiscovered?.(allDiscovered.size); + const backfilled = await backfillPendingWork(db, config, options, runId); + return { discovered: [...allDiscovered], backfilled }; + } finally { + await finishBackfillRun(db, runId); + } +} diff --git a/packages/contrail/src/core/db/schema.ts b/packages/contrail/src/core/db/schema.ts index 46a04f7..c8bdbf3 100644 --- a/packages/contrail/src/core/db/schema.ts +++ b/packages/contrail/src/core/db/schema.ts @@ -17,7 +17,7 @@ import { getSearchableFields } from "../search"; import { buildLabelsSchema } from "../labels/schema"; import { getMeta, setMeta } from "./meta"; -export const CONTRAIL_SCHEMA_VERSION = 3; +export const CONTRAIL_SCHEMA_VERSION = 4; const SCHEMA_FINGERPRINT_KEY = "schema_fingerprint"; function getResolved(config: ContrailConfig): ResolvedMaps { @@ -40,6 +40,7 @@ CREATE TABLE IF NOT EXISTS backfills ( retries INTEGER NOT NULL DEFAULT 0, last_error TEXT, last_attempt_at ${dialect.bigintType}, + next_retry_at ${dialect.bigintType}, PRIMARY KEY (did, collection) ); CREATE TABLE IF NOT EXISTS discovery ( @@ -57,6 +58,13 @@ CREATE TABLE IF NOT EXISTS cursor ( id INTEGER PRIMARY KEY CHECK (id = 1), time_us ${dialect.bigintType} NOT NULL ); +CREATE TABLE IF NOT EXISTS backfill_state ( + id INTEGER PRIMARY KEY CHECK (id = 1), + run_id TEXT, + started_at ${dialect.bigintType}, + heartbeat_at ${dialect.bigintType}, + finished_at ${dialect.bigintType} +); CREATE TABLE IF NOT EXISTS identities ( did TEXT PRIMARY KEY, handle TEXT, @@ -309,6 +317,7 @@ const MIGRATIONS: MigrationOp[] = [ }, { table: "backfills", column: "last_error", columnDef: "TEXT" }, { table: "backfills", column: "last_attempt_at", columnDef: "BIGINT" }, + { table: "backfills", column: "next_retry_at", columnDef: "BIGINT" }, { table: "discovery", column: "retries", diff --git a/packages/contrail/src/core/status.ts b/packages/contrail/src/core/status.ts index dec09d6..ef19064 100644 --- a/packages/contrail/src/core/status.ts +++ b/packages/contrail/src/core/status.ts @@ -26,12 +26,20 @@ export interface BackfillCollectionStatus extends BackfillCounts { collection: string; } +export interface BackfillRetryStatus { + scheduled_accounts: number; + due_accounts: number; + next_retry_at: number | null; + next_retry_date: string | null; +} + export interface BackfillStatus { - state: "not_started" | "incomplete" | "complete"; + state: "not_started" | "running" | "incomplete" | "complete"; known_progress_percent: number; accounts: BackfillAccountCounts; tasks: BackfillCounts; discovery: BackfillCounts; + retries: BackfillRetryStatus; collections: BackfillCollectionStatus[]; } @@ -48,11 +56,67 @@ function counts(row: AggregateRow | null): BackfillCounts { }; } +const BACKFILL_RUN_STALE_MS = 2 * 60_000; + +function createRunId(): string { + return `${Date.now().toString(36)}-${Math.random().toString(36).slice(2)}`; +} + +export async function tryStartBackfillRun( + db: Database +): Promise { + const runId = createRunId(); + const now = Date.now(); + await db + .prepare( + "INSERT INTO backfill_state (id, run_id, started_at, heartbeat_at, finished_at) VALUES (1, ?, ?, ?, NULL) ON CONFLICT(id) DO UPDATE SET run_id = excluded.run_id, started_at = excluded.started_at, heartbeat_at = excluded.heartbeat_at, finished_at = NULL WHERE backfill_state.finished_at IS NOT NULL OR backfill_state.heartbeat_at IS NULL OR backfill_state.heartbeat_at < ?" + ) + .bind(runId, now, now, now - BACKFILL_RUN_STALE_MS) + .run(); + const row = await db + .prepare("SELECT run_id FROM backfill_state WHERE id = 1") + .first<{ run_id: string | null }>(); + return row?.run_id === runId ? runId : null; +} + +export async function heartbeatBackfillRun( + db: Database, + runId: string +): Promise { + await db + .prepare( + "UPDATE backfill_state SET heartbeat_at = ? WHERE id = 1 AND run_id = ? AND finished_at IS NULL" + ) + .bind(Date.now(), runId) + .run(); +} + +export async function finishBackfillRun( + db: Database, + runId: string +): Promise { + const now = Date.now(); + await db + .prepare( + "UPDATE backfill_state SET heartbeat_at = ?, finished_at = ? WHERE id = 1 AND run_id = ?" + ) + .bind(now, now, runId) + .run(); +} + export async function getBackfillStatus( db: Database, config?: ContrailConfig ): Promise { - const [taskRow, accountRow, discoveryRow, collectionRows] = await Promise.all([ + const now = Date.now(); + const [ + taskRow, + accountRow, + discoveryRow, + collectionRows, + retryRow, + runRow, + ] = await Promise.all([ db .prepare( `SELECT @@ -102,7 +166,35 @@ export async function getBackfillStatus( GROUP BY collection ORDER BY collection` ) - .all() + .all(), + db + .prepare( + `SELECT + COUNT(*) AS scheduled_accounts, + COALESCE(SUM(CASE WHEN next_retry_at IS NULL OR next_retry_at <= ? THEN 1 ELSE 0 END), 0) AS due_accounts, + MIN(next_retry_at) AS next_retry_at + FROM ( + SELECT did, MIN(next_retry_at) AS next_retry_at + FROM backfills + WHERE completed = 0 AND last_error IS NOT NULL + GROUP BY did + ) AS retry_accounts` + ) + .bind(now) + .first<{ + scheduled_accounts: number | string | null; + due_accounts: number | string | null; + next_retry_at: number | string | null; + }>(), + db + .prepare( + "SELECT run_id, heartbeat_at, finished_at FROM backfill_state WHERE id = 1" + ) + .first<{ + run_id: string | null; + heartbeat_at: number | null; + finished_at: number | null; + }>() ]); const tasks = counts(taskRow); @@ -119,18 +211,40 @@ export async function getBackfillStatus( ...counts(row) })); + const scheduledAccounts = number(retryRow?.scheduled_accounts); + const dueAccounts = number(retryRow?.due_accounts); + const nextRetryValue = retryRow?.next_retry_at; + const nextRetryAt = + dueAccounts > 0 || + nextRetryValue === null || + nextRetryValue === undefined + ? null + : number(nextRetryValue); + const retries: BackfillRetryStatus = { + scheduled_accounts: scheduledAccounts, + due_accounts: dueAccounts, + next_retry_at: nextRetryAt, + next_retry_date: nextRetryAt ? new Date(nextRetryAt).toISOString() : null + }; const expectsDiscovery = config ? getDiscoverableNsids(config).length > 0 && (config.relays ?? DEFAULT_RELAYS).length > 0 : true; - const state = - tasks.total === 0 && discovery.total === 0 + const running = + !!runRow?.run_id && + runRow.finished_at === null && + number(runRow.heartbeat_at) >= now - BACKFILL_RUN_STALE_MS; + const state = running + ? "running" + : tasks.total === 0 && discovery.total === 0 ? expectsDiscovery ? "not_started" : "complete" - : tasks.pending === 0 && discovery.pending === 0 - ? "complete" - : "incomplete"; + : discovery.pending > 0 + ? "incomplete" + : tasks.pending === 0 || tasks.failed === tasks.pending + ? "complete" + : "incomplete"; const knownProgressPercent = tasks.total === 0 ? (state === "complete" ? 100 : 0) : Math.floor((tasks.complete / tasks.total) * 10_000) / 100; @@ -140,6 +254,7 @@ export async function getBackfillStatus( accounts, tasks, discovery, + retries, collections }; } diff --git a/packages/contrail/src/worker/index.ts b/packages/contrail/src/worker/index.ts index 4c1e72d..2d01d7a 100644 --- a/packages/contrail/src/worker/index.ts +++ b/packages/contrail/src/worker/index.ts @@ -10,19 +10,23 @@ * export default createWorker(config, { lexicons }); * * The handler lazily inits the DB schema on first request per isolate, - * registers every XRPC route, and runs `contrail.ingest()` on the - * `scheduled` event. Pass `binding` if your D1 binding isn't named `DB`; + * registers every XRPC route, and runs live ingestion plus a bounded due + * backfill-retry slice on the `scheduled` event. Pass `binding` if your D1 binding isn't named `DB`; * pass `onInit` for app-specific one-shot setup that needs the DB. */ import { Contrail } from "../contrail.js"; import { createHandler } from "../server.js"; import type { ContrailConfig, Database } from "../core/types.js"; +import type { BackfillRetryOptions } from "../core/backfill.js"; export interface CreateWorkerOptions { /** D1 binding name in wrangler env. Default: `"DB"`. */ binding?: string; /** Bundled Lexicon documents to expose for application type generation. */ lexicons?: object[]; + /** Bounded pending-account retry slice after each scheduled ingest. Enabled + * by default; pass `false` to disable or options to tune its budget. */ + backfillRetries?: BackfillRetryOptions | false; /** Runs once per isolate, after schema init, before handling the first * request. Use for app-specific setup that needs a live DB handle. */ onInit?: (env: Record, db: Database) => void | Promise; @@ -61,9 +65,16 @@ export function createWorker( ): Promise { const db = env[binding] as Database; await ensureReady(env, db); - // ingest() drives both record and label ingestion in parallel when - // labels are configured — single waitUntil covers the whole cron tick. - ctx.waitUntil(contrail.ingest({}, db)); + // Keep live catch-up first, then spend a bounded slice on due historical + // failures. A database lease prevents overlap with a manual backfill. + ctx.waitUntil( + (async () => { + await contrail.ingest({}, db); + if (options.backfillRetries !== false) { + await contrail.retryBackfill(options.backfillRetries, db); + } + })() + ); }, }; } diff --git a/packages/contrail/tests/backfill-status.test.ts b/packages/contrail/tests/backfill-status.test.ts index 2ab84e4..7657d07 100644 --- a/packages/contrail/tests/backfill-status.test.ts +++ b/packages/contrail/tests/backfill-status.test.ts @@ -1,5 +1,15 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; -import { backfillPending, backfillUser, createApp, discoverDIDs, getBackfillStatus, resolveConfig } from "../src/index"; +import { + backfillPending, + backfillUser, + createApp, + discoverDIDs, + finishBackfillRun, + getBackfillStatus, + resolveConfig, + retryPendingBackfills, + tryStartBackfillRun +} from "../src/index"; import { __resetPdsCachesForTests } from "../src/core/client"; import { TEST_CONFIG, createTestDbWithSchema, ingestRecords, makeEvent } from "./helpers"; @@ -46,19 +56,21 @@ describe("backfill failure state", () => { expect(inserted).toBe(0); const row = await db - .prepare("SELECT completed, retries, last_error, last_attempt_at FROM backfills WHERE did = ? AND collection = ?") + .prepare("SELECT completed, retries, last_error, last_attempt_at, next_retry_at FROM backfills WHERE did = ? AND collection = ?") .bind(DID, EVENT) .first<{ completed: number; retries: number; last_error: string | null; last_attempt_at: number | null; + next_retry_at: number | null; }>(); expect(row?.completed).toBe(0); expect(row?.retries).toBe(1); expect(row?.last_error).toContain("503"); expect(row?.last_error).toContain("UpstreamFailure: PDS is unavailable"); expect(row?.last_attempt_at).toBeTypeOf("number"); + expect(row?.next_retry_at).toBeGreaterThan(Date.now()); }); it("retries pending rows on the next run and only then completes them", async () => { @@ -113,6 +125,86 @@ describe("backfill failure state", () => { }); }); +describe("scheduled backfill retries", () => { + let fetchSpy: ReturnType; + + beforeEach(() => { + __resetPdsCachesForTests(); + fetchSpy = vi.spyOn(global, "fetch"); + }); + + afterEach(() => { + fetchSpy.mockRestore(); + __resetPdsCachesForTests(); + }); + + it("waits for persisted backoff and then completes a due account", 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, retries, last_error, next_retry_at) VALUES (?, ?, 0, 1, 'unavailable', ?)" + ) + .bind(DID, EVENT, Date.now() + 60_000) + .run(); + fetchSpy.mockResolvedValue( + new Response(JSON.stringify({ records: [] }), { + status: 200, + headers: { "content-type": "application/json" } + }) + ); + + expect( + await retryPendingBackfills(db, TEST_CONFIG, { maxAccounts: 1 }) + ).toMatchObject({ attempted: 0, completed: 0, skipped: false }); + + await db + .prepare("UPDATE backfills SET next_retry_at = ? WHERE did = ?") + .bind(Date.now() - 1, DID) + .run(); + expect( + await retryPendingBackfills(db, TEST_CONFIG, { maxAccounts: 1 }) + ).toMatchObject({ attempted: 1, completed: 1, failed: 0, skipped: false }); + + const row = await db + .prepare("SELECT completed, retries, last_error, next_retry_at FROM backfills WHERE did = ?") + .bind(DID) + .first<{ + completed: number; + retries: number; + last_error: string | null; + next_retry_at: number | null; + }>(); + expect(row).toEqual({ + completed: 1, + retries: 0, + last_error: null, + next_retry_at: null + }); + }); + + it("reports an active run and prevents overlapping retry work", async () => { + const db = await createTestDbWithSchema(); + const runId = await tryStartBackfillRun(db); + expect(runId).not.toBeNull(); + expect((await getBackfillStatus(db, TEST_CONFIG)).state).toBe("running"); + + expect(await retryPendingBackfills(db, TEST_CONFIG)).toEqual({ + attempted: 0, + completed: 0, + failed: 0, + records: 0, + skipped: true + }); + + await finishBackfillRun(db, runId!); + expect((await getBackfillStatus(db, TEST_CONFIG)).state).toBe("not_started"); + }); +}); + describe("discovery failure state", () => { it("keeps a failed relay pending and falls back to another relay", async () => { vi.useFakeTimers(); @@ -172,6 +264,33 @@ describe("discovery failure state", () => { }); describe("backfill status JSON", () => { + it("reports the initial pass complete while failed accounts remain scheduled", async () => { + const db = await createTestDbWithSchema(); + const retryAt = Date.now() + 60_000; + await db + .prepare( + "INSERT INTO backfills (did, collection, completed, retries, last_error, next_retry_at) VALUES (?, ?, 0, 1, 'unavailable', ?)" + ) + .bind(DID, EVENT, retryAt) + .run(); + await db + .prepare( + "INSERT INTO discovery (collection, relay, completed) VALUES (?, 'https://relay.test', 1)" + ) + .bind(EVENT) + .run(); + + const status = await getBackfillStatus(db, TEST_CONFIG); + expect(status.state).toBe("complete"); + expect(status.accounts.unreachable).toBe(1); + expect(status.retries).toEqual({ + scheduled_accounts: 1, + due_accounts: 0, + next_retry_at: retryAt, + next_retry_date: new Date(retryAt).toISOString() + }); + }); + it("reports known work, unreachable accounts, discovery, records, and cursor", async () => { const db = await createTestDbWithSchema(); await ingestRecords(db, [makeEvent()]); @@ -217,7 +336,13 @@ describe("backfill status JSON", () => { unreachable: 1 }, tasks: { total: 5, complete: 2, pending: 3, failed: 1 }, - discovery: { total: 2, complete: 1, pending: 1, failed: 1 } + discovery: { total: 2, complete: 1, pending: 1, failed: 1 }, + retries: { + scheduled_accounts: 1, + due_accounts: 1, + next_retry_at: null, + next_retry_date: null + } }); expect(JSON.stringify(overview)).not.toContain("did:plc:b"); expect(JSON.stringify(overview)).not.toContain("PDS unavailable"); diff --git a/packages/contrail/tests/schema.test.ts b/packages/contrail/tests/schema.test.ts index dbcb694..48a03ea 100644 --- a/packages/contrail/tests/schema.test.ts +++ b/packages/contrail/tests/schema.test.ts @@ -15,6 +15,7 @@ describe("initSchema", () => { expect(names).toContain("records_event"); expect(names).toContain("records_rsvp"); expect(names).toContain("backfills"); + expect(names).toContain("backfill_state"); expect(names).toContain("discovery"); expect(names).toContain("cursor"); expect(names).toContain("identities"); @@ -22,8 +23,8 @@ describe("initSchema", () => { const backfillColumns = await db .prepare("PRAGMA table_info(backfills)") .all<{ name: string }>(); - expect(backfillColumns.results.map((column) => column.name)).toContain( - "last_attempt_at" + expect(backfillColumns.results.map((column) => column.name)).toEqual( + expect.arrayContaining(["last_attempt_at", "next_retry_at"]) ); const discoveryColumns = await db @@ -82,6 +83,13 @@ describe("initSchema", () => { await initSchema(db, TEST_CONFIG); + const backfillColumns = await db + .prepare("PRAGMA table_info(backfills)") + .all<{ name: string }>(); + expect(backfillColumns.results.map((column) => column.name)).toEqual( + expect.arrayContaining(["last_attempt_at", "next_retry_at"]) + ); + const discoveryColumns = await db .prepare("PRAGMA table_info(discovery)") .all<{ name: string }>(); diff --git a/packages/contrail/tests/worker.test.ts b/packages/contrail/tests/worker.test.ts index 98ce8e9..b595a82 100644 --- a/packages/contrail/tests/worker.test.ts +++ b/packages/contrail/tests/worker.test.ts @@ -1,5 +1,6 @@ import { describe, it, expect, vi } from "vitest"; import { createWorker } from "../src/worker"; +import { Contrail } from "../src/contrail"; import { createSqliteDatabase } from "../src/adapters/sqlite"; import type { ContrailConfig } from "../src/index"; @@ -84,20 +85,46 @@ describe("createWorker", () => { expect(res.status).toBe(404); }); - it("scheduled handler hands the ingest promise to ctx.waitUntil", async () => { - const db = createSqliteDatabase(":memory:"); - const worker = createWorker(MINIMAL_CONFIG); - const env = { DB: db }; - - const waitUntil = vi.fn(); - const ctx = { waitUntil, passThroughOnException: vi.fn() } as unknown as ExecutionContext; - - // Schedule returns once init + waitUntil have been called. The actual - // ingest is a long-running promise that would try to connect to a real - // Jetstream — we don't drain it; we just verify the wire-up. - await worker.scheduled({} as ScheduledEvent, env, ctx); - - expect(waitUntil).toHaveBeenCalledTimes(1); - expect(waitUntil.mock.calls[0][0]).toBeInstanceOf(Promise); + it("scheduled handler runs live ingest then a bounded backfill retry slice", async () => { + const order: string[] = []; + const ingest = vi + .spyOn(Contrail.prototype, "ingest") + .mockImplementation(async () => { + order.push("ingest"); + }); + const retry = vi + .spyOn(Contrail.prototype, "retryBackfill") + .mockImplementation(async () => { + order.push("retry"); + return { + attempted: 0, + completed: 0, + failed: 0, + records: 0, + skipped: false, + }; + }); + + try { + const db = createSqliteDatabase(":memory:"); + const worker = createWorker(MINIMAL_CONFIG); + const env = { DB: db }; + const waitUntil = vi.fn(); + const ctx = { + waitUntil, + passThroughOnException: vi.fn(), + } as unknown as ExecutionContext; + + await worker.scheduled({} as ScheduledEvent, env, ctx); + + expect(waitUntil).toHaveBeenCalledTimes(1); + const scheduled = waitUntil.mock.calls[0][0]; + expect(scheduled).toBeInstanceOf(Promise); + await scheduled; + expect(order).toEqual(["ingest", "retry"]); + } finally { + ingest.mockRestore(); + retry.mockRestore(); + } }); });