diff --git a/.changeset/lean-contrail.md b/.changeset/lean-contrail.md index b11067c..b277dd1 100644 --- a/.changeset/lean-contrail.md +++ b/.changeset/lean-contrail.md @@ -2,4 +2,4 @@ "@atmo-dev/contrail": minor --- -Collapse Contrail into one public package and one AppView implementation. Remove the spaces, authority, record-host, community, realtime, sync, and custom Lexicon-tooling products. Route Jetstream, persistent, backfill, refresh, and immediate synchronization records through the shared `ingestRecords` admission and projection path. Make materialized relation counts converge when children arrive before parents or move during refresh, and prevent transient PDS failures from being interpreted as authoritative deletions. Keep dependent-subject filtering scoped to dependent collections and restore typed example XRPC clients with Atcute's generator. Preserve `node:sqlite` in the published adapter, make all query and search cursors stable across tied and typed/null rows, use Worker-safe cursor encoding, and bound the complete notify resolution/fetch/body operation. Admit newly discovered actors and their dependent mutations as one batch, keep subject decisions mutation-local so deletes always pass, and classify refresh statistics after admission. Existing pagination cursors from 0.12 are intentionally invalidated by the new stable cursor format; clients should discard persisted cursor tokens when upgrading. +Collapse Contrail into one public package and one AppView implementation. Remove the spaces, authority, record-host, community, realtime, sync, and custom Lexicon-tooling products. Route Jetstream, persistent, backfill, and immediate synchronization records through the shared `ingestRecords` admission and projection path. Make materialized relation counts converge when children arrive before parents, and prevent transient PDS failures from being interpreted as authoritative deletions. Keep dependent-subject filtering scoped to dependent collections and restore typed example XRPC clients with Atcute's generator. Preserve `node:sqlite` in the published adapter, make all query and search cursors stable across tied and typed/null rows, use Worker-safe cursor encoding, and bound the complete notify resolution/fetch/body operation. Admit newly discovered actors and their dependent mutations as one batch, and keep subject decisions mutation-local so deletes always pass. Existing pagination cursors from 0.12 are intentionally invalidated by the new stable cursor format; clients should discard persisted cursor tokens when upgrading. diff --git a/.changeset/remove-refresh.md b/.changeset/remove-refresh.md new file mode 100644 index 0000000..cfc656a --- /dev/null +++ b/.changeset/remove-refresh.md @@ -0,0 +1,5 @@ +--- +"@atmo-dev/contrail": minor +--- + +Remove the incomplete `refresh` PDS sweep from the public API, CLI, Wrangler helpers, examples, and documentation. Normal ingestion resumes from its saved source cursor; deployments whose source history has expired should rebuild into a fresh database rather than rely on partial reconciliation that cannot safely infer deletions. diff --git a/apps/cloudflare-workers/README.md b/apps/cloudflare-workers/README.md index 68b428f..ce7d522 100644 --- a/apps/cloudflare-workers/README.md +++ b/apps/cloudflare-workers/README.md @@ -30,11 +30,11 @@ GET https://.workers.dev/xrpc/com.example.event.listRecords?startsA ## local dev ```bash -pnpm dev # wraps wrangler dev + auto-fires cron + prompts for backfill/refresh if needed +pnpm dev # wraps wrangler dev + auto-fires cron + prompts for initial backfill 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 prompts to run backfill (if nothing's indexed yet) or refresh (if the ingest cursor is >1h old), then runs `wrangler dev --test-scheduled` with a 60s timer hitting `/__scheduled` so the cron actually fires locally. +`pnpm dev` runs `contrail dev` under the hood. On start it inspects the local D1 and offers to run backfill when nothing has been indexed yet, then runs `wrangler dev --test-scheduled` with a 60s timer hitting `/__scheduled` so the cron actually fires locally. ## extending diff --git a/apps/sveltekit-cloudflare-workers/package.json b/apps/sveltekit-cloudflare-workers/package.json index a01b724..b01adb5 100644 --- a/apps/sveltekit-cloudflare-workers/package.json +++ b/apps/sveltekit-cloudflare-workers/package.json @@ -9,8 +9,6 @@ "generate:types": "lex-cli generate", "backfill": "contrail backfill", "backfill:remote": "contrail backfill --remote", - "refresh": "contrail refresh", - "refresh:remote": "contrail refresh --remote", "preview": "vite preview", "prepare": "svelte-kit sync || echo ''", "check": "pnpm generate:types && svelte-kit sync && svelte-check --tsconfig ./tsconfig.json", diff --git a/docs/01-indexing.md b/docs/01-indexing.md index 2f0c295..193d965 100644 --- a/docs/01-indexing.md +++ b/docs/01-indexing.md @@ -97,7 +97,7 @@ async scheduled(_ev, env, ctx) { `ingest()` connects to Jetstream, streams events since the saved cursor, stops when caught up. Running every minute is fine — the next fire resumes where this one left off. Each cycle is bounded, so it can't blow past the Worker time limit. -**Local dev:** wrangler's cron scheduler only runs in deployed production. For local dev use `pnpm contrail dev` — it runs `wrangler dev --test-scheduled`, fires `/__scheduled` on your configured cron interval, and prompts you to run backfill or refresh if the local DB looks stale on start. +**Local dev:** wrangler's cron scheduler only runs in deployed production. For local dev use `pnpm contrail dev` — it runs `wrangler dev --test-scheduled`, fires `/__scheduled` on your configured cron interval, and offers to run the initial backfill when the local DB is empty. ### Persistent (node / any long-lived server) @@ -139,42 +139,11 @@ Typical combos: - **workers app:** `backfillAll()` once + `ingest()` on cron + optional `notify()` for self-writes - **node server:** `backfillAll()` once + `runPersistent()` forever + optional `notify()` for self-writes -## Refresh (catch-up after outages / dev idle) +## Recovery after an outage -When Jetstream drops events — you went offline for a few days in dev, there was an outage, or you just want reassurance nothing was lost — `refresh` walks every known DID's PDS and reconciles against your DB: +Normal ingestion resumes from its saved Jetstream cursor, so a short outage needs no special command: restart `ingest()` or `runPersistent()` and let it catch up. -```bash -pnpm contrail refresh # totals only -pnpm contrail refresh --by-collection # totals + per-collection breakdown -pnpm contrail refresh --ignore-window 30 # grace seconds (default: 60) -``` - -Each record is classified as: - -- **missing** — PDS has it, DB doesn't. Inserted. -- **stale update** — DB has it with a different CID, *and* the DB row is older than the ignore window. Upserted. -- **in sync** — same CID, or DB row is within the ignore window (jetstream probably just hadn't caught up yet). - -The ignore window is there so a refresh run seconds after a normal jetstream cycle doesn't double-count records that are about to sync anyway. Records inside the window are still written if they differ; they just don't show up in the stats. - -Report shape (`--by-collection`): - -``` -by collection: - community.lexicon.calendar.event - 3 missing, 1 stale updates, 842 in sync - community.lexicon.calendar.rsvp - 12 missing, 0 stale updates, 4108 in sync - -total: - 15 missing, 1 stale updates, 4950 in sync - 234 users scanned, 1 failed in 85.3s - (ignore window: 60s) -``` - -Safe to run repeatedly — each pass converges toward zero missing / stale. Programmatic equivalent: `contrail.refresh({ ignoreWindowMs, concurrency })` returns the same structure. - -Refresh is **not** a replacement for `ingest`/`runPersistent` — it walks every user's full history, which is expensive. Use it after outages or during dev idle, not as a continuous freshness mechanism. +Contrail does not perform a full PDS sweep as a repair mechanism. Such a sweep is expensive, cannot discover repositories it never knew about, and cannot safely infer remote deletions after partial failures. If the saved cursor is older than the source's retained history, rebuild into a fresh database with `backfillAll()` rather than trusting a partial reconciliation. A replay-capable source and first-class projection rebuild command are planned follow-up work. Reading the indexed data — filters, sorts, hydration, search, pagination — has its own doc: [Querying](./02-querying.md). diff --git a/docs/frameworks/sveltekit-cloudflare.md b/docs/frameworks/sveltekit-cloudflare.md index 64270bd..e2b32a3 100644 --- a/docs/frameworks/sveltekit-cloudflare.md +++ b/docs/frameworks/sveltekit-cloudflare.md @@ -194,7 +194,7 @@ From now on: - Pages and XRPC endpoints are served under your domain. - The cron fires every minute, hitting `/api/cron`, which runs `contrail.ingest()`. - Loaders that need live data use `getServerClient()` for zero-overhead typed calls. -- Need to reconcile after an outage? `pnpm contrail refresh --remote`. +- After a short outage, ingestion resumes from its saved cursor. If source history has expired, rebuild into a fresh database with `pnpm backfill:remote`. ## Where to go next diff --git a/packages/contrail/src/cli.ts b/packages/contrail/src/cli.ts index 8b1978f..4dfffb6 100644 --- a/packages/contrail/src/cli.ts +++ b/packages/contrail/src/cli.ts @@ -7,14 +7,12 @@ */ import { cac } from "cac"; import { registerBackfill } from "./cli/commands/backfill.js"; -import { registerRefresh } from "./cli/commands/refresh.js"; import { registerDev } from "./cli/commands/dev.js"; import { registerAppendScheduled } from "./cli/commands/append-scheduled.js"; const cli = cac("contrail"); registerBackfill(cli); -registerRefresh(cli); registerDev(cli); registerAppendScheduled(cli); diff --git a/packages/contrail/src/cli/commands/dev.ts b/packages/contrail/src/cli/commands/dev.ts index 6b5a9ee..f0bc783 100644 --- a/packages/contrail/src/cli/commands/dev.ts +++ b/packages/contrail/src/cli/commands/dev.ts @@ -10,8 +10,6 @@ interface DevOpts { binding: string; cron: string; concurrency: number; - ignoreWindow?: number; - staleAfter: number; yes?: boolean; } @@ -19,7 +17,7 @@ export function registerDev(cli: CAC): void { cli .command( "dev", - "Local wrangler dev + auto-trigger cron + backfill/refresh prompts" + "Local wrangler dev + auto-trigger cron + optional backfill prompt" ) .option("--config ", "Path to Contrail config file (TS or JS)") .option("--root ", "Project root for auto-detection (default: CWD)", { @@ -38,15 +36,6 @@ export function registerDev(cli: CAC): void { "Concurrency passed to backfill if prompted (default: 100)", { default: 100 } ) - .option( - "--ignore-window ", - "Refresh ignore-window in seconds, if prompted (default: server default)" - ) - .option( - "--stale-after ", - "Prompt to refresh if the ingest cursor is older than this (default: 60)", - { default: 60 } - ) .option("--yes, -y", "Accept all prompts without asking (CI-friendly)") .action(async (options: DevOpts) => { const config = await resolveAndLoadConfig(options); @@ -82,27 +71,6 @@ export function registerDev(cli: CAC): void { db ); } - } else { - const staleAfterMs = Number(options.staleAfter) * 60_000; - const row = await db - .prepare("SELECT time_us FROM cursor WHERE id = 1") - .first<{ time_us: number }>(); - if (row?.time_us) { - const ageMs = Date.now() - Math.floor(row.time_us / 1000); - if (ageMs > staleAfterMs) { - const hrs = (ageMs / 3_600_000).toFixed(1); - console.log( - `ingest cursor is ${hrs}h old — you may have missed events.` - ); - if (await promptYesNo("run refresh first?", true, !!options.yes)) { - const ignoreWindowMs = - options.ignoreWindow !== undefined - ? Number(options.ignoreWindow) * 1000 - : undefined; - await contrail.refresh({ ignoreWindowMs }, db); - } - } - } } } diff --git a/packages/contrail/src/cli/commands/refresh.ts b/packages/contrail/src/cli/commands/refresh.ts deleted file mode 100644 index 3c5764e..0000000 --- a/packages/contrail/src/cli/commands/refresh.ts +++ /dev/null @@ -1,50 +0,0 @@ -import type { CAC } from "cac"; -import { refresh } from "../../workers/backfill.js"; -import { printRefreshReport, resolveAndLoadConfig } from "../shared.js"; - -interface RefreshOpts { - config?: string; - root?: string; - remote?: boolean; - binding: string; - concurrency: number; - ignoreWindow?: number; - byCollection?: boolean; -} - -export function registerRefresh(cli: CAC): void { - cli - .command( - "refresh", - "Fresh sweep: reconcile PDS vs DB, report missing + stale" - ) - .option("--config ", "Path to Contrail config file (TS or JS)") - .option("--root ", "Project root for auto-detection (default: CWD)") - .option("--remote", "Use production D1 bindings") - .option("--binding ", "D1 binding name in wrangler.jsonc", { - default: "DB", - }) - .option("--concurrency ", "Passed to contrail.refresh()", { - default: 50, - }) - .option( - "--ignore-window ", - "Seconds of grace for stale-update detection (default: 60)" - ) - .option("--by-collection", "Print per-collection stats, not just totals") - .action(async (options: RefreshOpts) => { - const config = await resolveAndLoadConfig(options); - const ignoreWindowMs = - options.ignoreWindow !== undefined - ? Number(options.ignoreWindow) * 1000 - : undefined; - const result = await refresh({ - config, - remote: !!options.remote, - binding: options.binding, - concurrency: Number(options.concurrency), - ignoreWindowMs, - }); - printRefreshReport(result, !!options.byCollection); - }); -} diff --git a/packages/contrail/src/cli/shared.ts b/packages/contrail/src/cli/shared.ts index 17fd8c5..51335ce 100644 --- a/packages/contrail/src/cli/shared.ts +++ b/packages/contrail/src/cli/shared.ts @@ -9,10 +9,6 @@ import { CONFIG_CANDIDATES_MESSAGE, } from "../cli-config.js"; import type { ContrailConfig } from "../core/types.js"; -import type { - CollectionStats, - RefreshResult, -} from "../core/refresh.js"; export interface ConfigOpts { config?: string; @@ -59,35 +55,3 @@ export async function promptYesNo( rl.close(); } } - -export function formatStats(s: CollectionStats): string { - return `${s.missing} missing, ${s.staleUpdates} stale updates, ${s.inSync} in sync`; -} - -export function printRefreshReport( - result: RefreshResult, - byCollection: boolean -): void { - console.log(""); - if (byCollection) { - console.log("by collection:"); - const entries = Object.entries(result.byCollection).sort(([a], [b]) => - a.localeCompare(b) - ); - for (const [nsid, stats] of entries) { - if (stats.missing === 0 && stats.staleUpdates === 0 && stats.inSync === 0) - continue; - console.log(` ${nsid}`); - console.log(` ${formatStats(stats)}`); - } - console.log(""); - } - console.log("total:"); - console.log(` ${formatStats(result.total)}`); - console.log( - ` ${result.usersScanned} users scanned` + - (result.usersFailed ? `, ${result.usersFailed} failed` : "") + - ` in ${(result.elapsedMs / 1000).toFixed(1)}s` - ); - console.log(` (ignore window: ${(result.ignoreWindowMs / 1000).toFixed(0)}s)`); -} diff --git a/packages/contrail/src/contrail.ts b/packages/contrail/src/contrail.ts index 1c8ccf7..66ba020 100644 --- a/packages/contrail/src/contrail.ts +++ b/packages/contrail/src/contrail.ts @@ -21,11 +21,6 @@ import { discoverDIDs, type BackfillAllOptions, } from "./core/backfill"; -import { - refresh as runRefresh, - type RefreshOptions, - type RefreshResult, -} from "./core/refresh"; import { processNotifyUris, type NotifyResult, @@ -215,46 +210,6 @@ export class Contrail { return { discovered: discovered.length, backfilled }; } - /** Fresh sweep: re-walk every known DID's PDS, compare each record against - * our DB, count anything missing or stale (outside the ignore window). - * Apply the deltas. Returns stats per-collection + totals. - * - * Progress logs via `config.logger` unless `onProgress` is supplied. */ - async refresh( - options?: RefreshOptions, - db?: Database - ): Promise { - const d = this.getDb(db); - const logger = this.config.logger; - - let effective = options; - if (!options?.onProgress) { - let lastLogAt = 0; - effective = { - ...options, - onProgress: ({ usersComplete, usersTotal, usersFailed, recordsScanned }) => { - const now = Date.now(); - if (now - lastLogAt < 2_000) return; - lastLogAt = now; - const failStr = usersFailed > 0 ? `, ${usersFailed} failed` : ""; - logger?.log?.( - ` ${recordsScanned} records scanned | ${usersComplete}/${usersTotal} users${failStr}` - ); - }, - }; - } - - logger?.log?.("refreshing…"); - const result = await runRefresh(d, this.config, effective); - const elapsedS = (result.elapsedMs / 1000).toFixed(1); - logger?.log?.( - ` done: ${result.total.missing} missing, ${result.total.staleUpdates} stale updates ` + - `across ${result.usersScanned} users in ${elapsedS}s` + - (result.usersFailed > 0 ? ` (${result.usersFailed} failed)` : "") - ); - return result; - } - /** Immediately fetch and index specific records from their PDS. */ async notify( uris: string | string[], diff --git a/packages/contrail/src/core/ingest.ts b/packages/contrail/src/core/ingest.ts index 546f98b..1d6828c 100644 --- a/packages/contrail/src/core/ingest.ts +++ b/packages/contrail/src/core/ingest.ts @@ -68,7 +68,7 @@ export interface IngestRecordsResult { /** * The single admission and projection path for records from every source. * - * Jetstream, persistent subscriptions, PDS backfill, refresh, and immediate + * Jetstream, persistent subscriptions, PDS backfill, and immediate * synchronization all produce the same IngestEvent shape and enter here. * Source connection and checkpoint handling remain outside this function. */ diff --git a/packages/contrail/src/core/refresh.ts b/packages/contrail/src/core/refresh.ts deleted file mode 100644 index aef737c..0000000 --- a/packages/contrail/src/core/refresh.ts +++ /dev/null @@ -1,273 +0,0 @@ -import type {} from "@atcute/atproto"; -/** - * Fresh refresh: re-walk every known DID's PDS for every configured collection - * and reconcile against what's in our DB. Unlike `backfillPending`, this - * ignores the `backfills` state machine — it's a "check what we might have - * missed" pass, not a resumable bulk load. - * - * Two categories of delta are counted: - * - missing — the PDS has a record we don't - * - staleUpdates — we have the same URI but a different CID, *and* our - * copy's `indexed_at` is older than `ignoreWindowMs` - * - * The ignore window exists because Jetstream can run ~seconds behind the - * PDS; without the window, "in-sync but racy" writes would show up as - * misses every run. Records inside the window are still applied (they - * might be legit updates), just not counted toward stats. - * - * Typical uses: - * - dev: "I ran backfillAll on Monday, haven't touched it for a week, - * how much did jetstream miss?" - * - prod: "we had jetstream outage yesterday, what did we drop?" - */ -import { type Did, type Nsid } from "@atcute/lexicons"; -import { isDid, isNsid } from "@atcute/lexicons/syntax"; - -import type { Client } from "@atcute/client"; -import type { ContrailConfig, Database, IngestEvent } from "./types.js"; -import { lookupExistingRecords } from "./db/records.js"; -import { createIngestEvent, ingestRecords } from "./ingest.js"; -import { getClient } from "./client.js"; - -const PAGE_SIZE = 100; -const REQUEST_TIMEOUT_MS = 10_000; - -async function withTimeout(fn: () => Promise, ms: number): Promise { - return Promise.race([ - fn(), - new Promise((_, rej) => - setTimeout(() => rej(new Error(`timeout after ${ms}ms`)), ms) - ), - ]); -} - -export interface CollectionStats { - /** Record exists on PDS but was absent from our DB. */ - missing: number; - /** Record exists in our DB with a different CID than the PDS, and our - * copy was written before the ignore window. */ - staleUpdates: number; - /** Record is present and matches (same CID, or within ignore window). */ - inSync: number; -} - -export interface RefreshProgress { - usersComplete: number; - usersTotal: number; - usersFailed: number; - recordsScanned: number; -} - -export interface RefreshResult { - /** Per-NSID stats. */ - byCollection: Record; - /** Sum across every NSID. */ - total: CollectionStats; - usersScanned: number; - usersFailed: number; - /** Effective ignore window used for classification, in ms. */ - ignoreWindowMs: number; - /** Wall-clock runtime, in ms. */ - elapsedMs: number; -} - -export interface RefreshOptions { - /** How many DIDs to fan out against in parallel. Default: 50. */ - concurrency?: number; - /** Records whose local `indexed_at` is within this window of `now` are - * still upserted but excluded from `staleUpdates` counts — guards - * against jetstream being briefly behind the PDS. Default: 60_000 ms. */ - ignoreWindowMs?: number; - /** Override which NSIDs to walk. Default: every `config.collections[*].collection`. */ - nsids?: string[]; - /** Optional progress callback (fires per completed DID). */ - onProgress?: (p: RefreshProgress) => void; - /** Max attempts per listRecords request. Default: 3. */ - maxRetries?: number; - /** Per-request timeout in ms. Default: 10000. */ - requestTimeout?: number; -} - -function emptyStats(): CollectionStats { - return { missing: 0, staleUpdates: 0, inSync: 0 }; -} - -export async function refresh( - db: Database, - config: ContrailConfig, - options?: RefreshOptions -): Promise { - const concurrency = options?.concurrency ?? 50; - const ignoreWindowMs = options?.ignoreWindowMs ?? 60_000; - const requestTimeout = options?.requestTimeout ?? REQUEST_TIMEOUT_MS; - const maxRetries = options?.maxRetries ?? 3; - const startedAt = Date.now(); - - // Default to every configured collection NSID. Profiles are already - // included because `resolveConfig` adds them to `config.collections`. - const nsids = - options?.nsids ?? - Object.entries(config.collections).map(([short, c]) => c.collection ?? short); - - const byCollection: Record = {}; - for (const nsid of nsids) byCollection[nsid] = emptyStats(); - const total: CollectionStats = emptyStats(); - - // Known DIDs = every author we've ever written for. `backfills` is a - // superset (it also includes failed/pending users that we never got - // records from), which is actually what we want — if we tried and - // failed before, we might succeed now. - const didRows = await db - .prepare("SELECT DISTINCT did FROM backfills") - .all<{ did: string }>(); - const dids = (didRows.results ?? []) - .map((r) => r.did) - .filter((d) => isDid(d)); - - const usersTotal = dids.length; - let usersComplete = 0; - let usersFailed = 0; - let recordsScanned = 0; - - const ignoreBeforeUs = (Date.now() - ignoreWindowMs) * 1000; - - const processDid = async (did: string): Promise => { - let client: Client; - try { - client = await withTimeout( - () => getClient(did as Did, db, config), - requestTimeout - ); - } catch { - usersFailed++; - return; - } - - for (const nsid of nsids) { - if (!isNsid(nsid)) continue; - let cursor: string | undefined; - while (true) { - let pageRecords: Array<{ uri: string; cid: string; value: unknown }>; - let nextCursor: string | undefined; - try { - // Retry listRecords: transient PDS failures are expected during refresh - let attempt = 0; - // eslint-disable-next-line no-constant-condition - while (true) { - try { - const res = await withTimeout( - () => - client.get("com.atproto.repo.listRecords", { - params: { - repo: did as Did, - collection: nsid as Nsid, - limit: PAGE_SIZE, - cursor, - }, - }), - requestTimeout - ); - if (!res.ok) { - // 400s on a collection the user doesn't have are fine; stop - // paging this collection for this user. - pageRecords = []; - nextCursor = undefined; - break; - } - pageRecords = res.data.records; - nextCursor = res.data.cursor ?? undefined; - break; - } catch (err) { - if (attempt >= maxRetries) throw err; - attempt++; - await new Promise((r) => setTimeout(r, 500 * 2 ** attempt)); - } - } - } catch { - // Give up on this collection for this user; keep going. - break; - } - - if (pageRecords.length === 0) break; - - const now = Date.now(); - const events: IngestEvent[] = pageRecords.map((record) => - createIngestEvent({ - uri: record.uri, - did, - collection: nsid, - rkey: record.uri.split("/").pop()!, - operation: "create", - cid: record.cid, - value: record.value, - timeUs: now * 1000, - indexedAt: now * 1000, - }), - ); - - const existing = await lookupExistingRecords( - db, - events.map((e) => ({ uri: e.uri, collection: e.collection })), - true, // old body is required to recount a relation's previous target - config - ); - - // Admission owns which remote records belong in the projection. Classify - // only accepted events so an intentional filter is not reported missing - // on every refresh forever. - const { accepted } = await ingestRecords(db, events, config, { - existing, - skipFeedFanout: true, - }); - - for (const event of accepted) { - const previous = existing.get(event.uri); - if (!previous) { - byCollection[nsid].missing++; - total.missing++; - } else if (previous.cid !== event.cid) { - const inWindow = - previous.indexed_at !== null && - previous.indexed_at >= ignoreBeforeUs; - if (inWindow) { - byCollection[nsid].inSync++; - total.inSync++; - } else { - byCollection[nsid].staleUpdates++; - total.staleUpdates++; - } - } else { - byCollection[nsid].inSync++; - total.inSync++; - } - } - recordsScanned += events.length; - - cursor = nextCursor; - if (!cursor) break; - } - } - - usersComplete++; - options?.onProgress?.({ - usersComplete, - usersTotal, - usersFailed, - recordsScanned, - }); - }; - - for (let i = 0; i < dids.length; i += concurrency) { - const batch = dids.slice(i, i + concurrency); - await Promise.allSettled(batch.map(processDid)); - } - - return { - byCollection, - total, - usersScanned: usersComplete, - usersFailed, - ignoreWindowMs, - elapsedMs: Date.now() - startedAt, - }; -} diff --git a/packages/contrail/src/index.ts b/packages/contrail/src/index.ts index 391d623..e598f03 100644 --- a/packages/contrail/src/index.ts +++ b/packages/contrail/src/index.ts @@ -20,7 +20,6 @@ export * from "./core/ingest"; export * from "./core/jetstream"; export * from "./core/persistent"; export * from "./core/backfill"; -export * from "./core/refresh"; export * from "./core/search"; export * from "./core/constellation"; diff --git a/packages/contrail/src/workers/backfill.ts b/packages/contrail/src/workers/backfill.ts index 9fe4dff..35e7356 100644 --- a/packages/contrail/src/workers/backfill.ts +++ b/packages/contrail/src/workers/backfill.ts @@ -5,8 +5,6 @@ * `backfills` state table to resume across runs). * - `labelsBackfillAll` — drain pending events per configured labeler in * repeated cycles until each labeler's cursor stops advancing. - * - `refresh` — reconcile every known DID's PDS against our DB, report - * what's missing or stale. Use after outages or long idle periods. * * All dynamically import `wrangler` (optional peer dep) and wire it to * the user's D1 binding, then dispose the proxy on exit. @@ -14,10 +12,6 @@ import { Contrail } from "../contrail.js"; import type { ContrailConfig, Database } from "../core/types.js"; import type { BackfillAllOptions } from "../core/backfill.js"; -import type { - RefreshOptions, - RefreshResult, -} from "../core/refresh.js"; interface WranglerCommon { config: ContrailConfig; @@ -35,15 +29,6 @@ export interface BackfillAllViaWranglerOptions extends WranglerCommon { onProgress?: BackfillAllOptions["onProgress"]; } -export interface RefreshViaWranglerOptions extends WranglerCommon { - /** Passed through to `contrail.refresh()`. Default: 50. */ - concurrency?: number; - /** Ignore-window for stale-update classification, in ms. Default: 60_000. */ - ignoreWindowMs?: number; - /** Override the built-in progress logging. */ - onProgress?: RefreshOptions["onProgress"]; -} - async function withWrangler( opts: WranglerCommon, fn: (contrail: Contrail, db: Database) => Promise @@ -86,21 +71,6 @@ export async function backfillAll( ); } -export async function refresh( - opts: RefreshViaWranglerOptions -): Promise { - return withWrangler(opts, (contrail, db) => - contrail.refresh( - { - concurrency: opts.concurrency ?? 50, - ignoreWindowMs: opts.ignoreWindowMs, - onProgress: opts.onProgress, - }, - db - ) - ); -} - export interface LabelsBackfillAllViaWranglerOptions extends WranglerCommon { /** Per-cycle subscribe timeout passed to `contrail.ingestLabels()`. Default: 60s. */ cycleTimeoutMs?: number; diff --git a/packages/contrail/tests/built-sqlite.mjs b/packages/contrail/tests/built-sqlite.mjs index ec7f4b3..72b56dc 100644 --- a/packages/contrail/tests/built-sqlite.mjs +++ b/packages/contrail/tests/built-sqlite.mjs @@ -1,5 +1,10 @@ import assert from "node:assert/strict"; +import { Contrail } from "@atmo-dev/contrail"; import { createSqliteDatabase } from "@atmo-dev/contrail/sqlite"; +import * as workers from "@atmo-dev/contrail/workers"; + +assert.equal("refresh" in Contrail.prototype, false); +assert.equal("refresh" in workers, false); const db = createSqliteDatabase(":memory:"); const row = await db.prepare("SELECT 1 AS ok").first(); diff --git a/packages/contrail/tests/refresh.test.ts b/packages/contrail/tests/refresh.test.ts deleted file mode 100644 index 21a94bd..0000000 --- a/packages/contrail/tests/refresh.test.ts +++ /dev/null @@ -1,291 +0,0 @@ -import { describe, it, expect, vi, beforeEach } from "vitest"; -import { initSchema, queryRecords, refresh, resolveConfig } from "../src/index"; -import { - ingestRecords, - createTestDb, - createTestDbWithSchema, - makeEvent, - TEST_CONFIG, -} from "./helpers"; -import type { Database } from "../src/index"; - -// We mock the PDS client so refresh() can be exercised without network IO. -// Each test sets the desired pageRecords for a given (did, collection) via -// the `pages` map below. -const pages = new Map>(); - -vi.mock("../src/core/client", async (importOriginal) => { - const actual = await importOriginal(); - return { - ...actual, - getClient: vi.fn(async (_did: string) => ({ - get: async ( - _method: string, - opts: { params: { repo: string; collection: string; cursor?: string } } - ) => { - const key = `${opts.params.repo}|${opts.params.collection}`; - if (opts.params.cursor) return { ok: true, data: { records: [], cursor: undefined } }; - const records = pages.get(key) ?? []; - return { ok: true, data: { records, cursor: undefined } }; - }, - })), - getPDS: vi.fn(), - }; -}); - -const ALICE = "did:plc:alice"; -const BOB = "did:plc:bob"; -const EVENT_NSID = "community.lexicon.calendar.event"; -const RSVP_NSID = "community.lexicon.calendar.rsvp"; - -function aliceEventUri(rkey: string): string { - return `at://${ALICE}/${EVENT_NSID}/${rkey}`; -} - -async function registerKnownDid(db: Database, did: string): Promise { - // refresh() enumerates DIDs from the `backfills` table. - await db - .prepare( - "INSERT INTO backfills (did, collection, completed) VALUES (?, ?, 1) ON CONFLICT DO NOTHING" - ) - .bind(did, EVENT_NSID) - .run(); -} - -describe("refresh", () => { - beforeEach(() => { - pages.clear(); - }); - - it("classifies an unseen-by-DB record as missing", async () => { - const db = await createTestDbWithSchema(); - await registerKnownDid(db, ALICE); - - pages.set(`${ALICE}|${EVENT_NSID}`, [ - { - uri: aliceEventUri("new1"), - cid: "bafy-new", - value: { name: "Brand new", startsAt: "2026-04-01T10:00:00Z" }, - }, - ]); - - const result = await refresh(db, TEST_CONFIG, { ignoreWindowMs: 0 }); - expect(result.total.missing).toBe(1); - expect(result.total.staleUpdates).toBe(0); - expect(result.total.inSync).toBe(0); - expect(result.usersScanned).toBe(1); - }); - - it("classifies a same-CID record as in-sync", async () => { - const db = await createTestDbWithSchema(); - await registerKnownDid(db, ALICE); - - // Seed DB with a record at this URI. - await ingestRecords(db, [ - makeEvent({ - did: ALICE, - rkey: "k1", - uri: aliceEventUri("k1"), - cid: "bafy-same", - record: { name: "Already here" }, - }), - ]); - - // PDS returns same URI + same CID. - pages.set(`${ALICE}|${EVENT_NSID}`, [ - { uri: aliceEventUri("k1"), cid: "bafy-same", value: { name: "Already here" } }, - ]); - - const result = await refresh(db, TEST_CONFIG, { ignoreWindowMs: 0 }); - expect(result.total.missing).toBe(0); - expect(result.total.staleUpdates).toBe(0); - expect(result.total.inSync).toBe(1); - }); - - it("classifies a different-CID record as a stale update (outside ignore window)", async () => { - const db = await createTestDbWithSchema(); - await registerKnownDid(db, ALICE); - - // Seed DB with an OLD record (indexed_at way in the past). - const dayAgoUs = (Date.now() - 86_400_000) * 1000; - await ingestRecords(db, [ - makeEvent({ - did: ALICE, - rkey: "k2", - uri: aliceEventUri("k2"), - cid: "bafy-old", - time_us: dayAgoUs, - indexed_at: dayAgoUs, - record: { name: "Was online" }, - }), - ]); - - // PDS returns same URI but different CID. - pages.set(`${ALICE}|${EVENT_NSID}`, [ - { uri: aliceEventUri("k2"), cid: "bafy-NEW", value: { name: "Now in-person" } }, - ]); - - const result = await refresh(db, TEST_CONFIG, { ignoreWindowMs: 60_000 }); - expect(result.total.staleUpdates).toBe(1); - expect(result.total.missing).toBe(0); - expect(result.total.inSync).toBe(0); - }); - - it("skips stale-update classification when DB row is within the ignore window", async () => { - const db = await createTestDbWithSchema(); - await registerKnownDid(db, ALICE); - - // Seed with a record indexed JUST NOW. - const nowUs = Date.now() * 1000; - await ingestRecords(db, [ - makeEvent({ - did: ALICE, - rkey: "k3", - uri: aliceEventUri("k3"), - cid: "bafy-recent", - time_us: nowUs, - indexed_at: nowUs, - record: { name: "Recent" }, - }), - ]); - - // PDS returns different CID — but DB row is fresh, so it counts as in-sync. - pages.set(`${ALICE}|${EVENT_NSID}`, [ - { uri: aliceEventUri("k3"), cid: "bafy-different", value: { name: "Recent v2" } }, - ]); - - const result = await refresh(db, TEST_CONFIG, { ignoreWindowMs: 60_000 }); - expect(result.total.staleUpdates).toBe(0); - expect(result.total.inSync).toBe(1); - }); - - it("aggregates stats per collection", async () => { - const db = await createTestDbWithSchema(); - await registerKnownDid(db, ALICE); - await registerKnownDid(db, BOB); - - pages.set(`${ALICE}|${EVENT_NSID}`, [ - { uri: aliceEventUri("a1"), cid: "x", value: {} }, - { uri: aliceEventUri("a2"), cid: "y", value: {} }, - ]); - pages.set(`${BOB}|${EVENT_NSID}`, [ - { uri: `at://${BOB}/${EVENT_NSID}/b1`, cid: "z", value: {} }, - ]); - - const result = await refresh(db, TEST_CONFIG, { ignoreWindowMs: 0 }); - expect(result.byCollection[EVENT_NSID]).toBeDefined(); - expect(result.byCollection[EVENT_NSID].missing).toBe(3); - expect(result.usersScanned).toBe(2); - }); - - it("counts a user as failed when getClient throws, and continues with others", async () => { - const { getClient } = await import("../src/core/client"); - (getClient as unknown as ReturnType).mockImplementationOnce( - async () => { - throw new Error("PDS unreachable"); - } - ); - - const db = await createTestDbWithSchema(); - await registerKnownDid(db, ALICE); - await registerKnownDid(db, BOB); - pages.set(`${BOB}|${EVENT_NSID}`, [ - { uri: `at://${BOB}/${EVENT_NSID}/x`, cid: "c", value: {} }, - ]); - - const result = await refresh(db, TEST_CONFIG, { ignoreWindowMs: 0, concurrency: 1 }); - expect(result.usersFailed).toBe(1); - expect(result.usersScanned).toBe(1); - // The non-failing user's records still classified. - expect(result.total.missing).toBeGreaterThanOrEqual(1); - }); - - it("does not report records rejected by admission as missing", async () => { - const filteredConfig = resolveConfig({ - namespace: "com.example", - profiles: [], - collections: { - event: { - collection: EVENT_NSID, - recordFilter: (record) => record.keep === true, - }, - }, - }); - const db = createTestDb(); - await initSchema(db, filteredConfig); - await registerKnownDid(db, ALICE); - pages.set(`${ALICE}|${EVENT_NSID}`, [ - { - uri: aliceEventUri("filtered"), - cid: "filtered-cid", - value: { keep: false }, - }, - ]); - - const result = await refresh(db, filteredConfig, { ignoreWindowMs: 0 }); - expect(result.total).toEqual({ missing: 0, staleUpdates: 0, inSync: 0 }); - const stored = await queryRecords(db, filteredConfig, { - collection: "event", - }); - expect(stored.records).toHaveLength(0); - }); - - it("recounts both relation targets when a refreshed child moves", async () => { - const db = await createTestDbWithSchema(); - await registerKnownDid(db, ALICE); - const eventA = aliceEventUri("event-a"); - const eventB = aliceEventUri("event-b"); - const rsvpUri = `at://${ALICE}/${RSVP_NSID}/rsvp`; - - await ingestRecords( - db, - [ - makeEvent({ uri: eventA, did: ALICE, rkey: "event-a" }), - makeEvent({ uri: eventB, did: ALICE, rkey: "event-b" }), - makeEvent({ - uri: rsvpUri, - did: ALICE, - collection: RSVP_NSID, - rkey: "rsvp", - cid: "old-cid", - record: { - subject: { uri: eventA }, - status: "community.lexicon.calendar.rsvp#going", - }, - }), - ], - TEST_CONFIG, - ); - - pages.set(`${ALICE}|${RSVP_NSID}`, [ - { - uri: rsvpUri, - cid: "new-cid", - value: { - subject: { uri: eventB }, - status: "community.lexicon.calendar.rsvp#going", - }, - }, - ]); - - await refresh(db, TEST_CONFIG, { ignoreWindowMs: 0 }); - const rows = await db - .prepare( - "SELECT uri, count_rsvp FROM records_event WHERE uri IN (?, ?) ORDER BY uri", - ) - .bind(eventA, eventB) - .all<{ uri: string; count_rsvp: number }>(); - - expect(rows.results).toEqual([ - { uri: eventA, count_rsvp: 0 }, - { uri: eventB, count_rsvp: 1 }, - ]); - }); - - it("returns elapsed time and the configured ignore window", async () => { - const db = await createTestDbWithSchema(); - const result = await refresh(db, TEST_CONFIG, { ignoreWindowMs: 30_000 }); - expect(result.elapsedMs).toBeGreaterThanOrEqual(0); - expect(result.ignoreWindowMs).toBe(30_000); - }); -});