diff --git a/.changeset/backfill-method-and-user-filter.md b/.changeset/backfill-method-and-user-filter.md new file mode 100644 index 0000000..83fd6d7 --- /dev/null +++ b/.changeset/backfill-method-and-user-filter.md @@ -0,0 +1,8 @@ +--- +"@atmo-dev/contrail": minor +--- + +Add `backfillMethod` config (`"listRecords" | "car"`) and `userFilter` for excluding users from indexing. + +- `backfillMethod` chooses how a user's history is fetched. Default `"listRecords"` keeps the existing per-(user, collection) paginated walk. Opt into `"car"` for one streamed `com.atproto.sync.getRepo` per user covering every configured collection — dramatically faster on multi-collection configs but pulls the whole repo regardless of which collections you index. +- `userFilter?: ({ did, handle, pds }) => boolean` runs after identity resolution. Returning true marks the user excluded: `identities.excluded` flips to `1` (new column with auto-migration), pending `backfills` rows are dropped, future PDS lookups short-circuit, and jetstream drops their commits and identity events. Filter checks fire on first PDS resolution, stale-identity refresh, and `#identity` events. Use for handle-suffix exclusions, deny-lists, etc. diff --git a/.gitignore b/.gitignore index 09e8139..6e21497 100644 --- a/.gitignore +++ b/.gitignore @@ -11,3 +11,6 @@ apps/*/src/lexicon-types/ apps/*/lex.config.js stuff/ + + +apps/backfill-test/ \ No newline at end of file diff --git a/docs/01-indexing.md b/docs/01-indexing.md index 54659fd..05f0321 100644 --- a/docs/01-indexing.md +++ b/docs/01-indexing.md @@ -58,6 +58,50 @@ await contrail.backfill({ concurrency: 100 }); // fetch history for registered D `backfill()` picks up where it left off across runs — safe to re-run. +### Backfill strategies + +Two ways to fetch a user's history. Set `backfillMethod` in your top-level config: + +```ts +{ + // ... + backfillMethod: "listRecords", // default + // backfillMethod: "car", +} +``` + +| | `listRecords` (default) | `car` | +|---|---|---| +| protocol | `com.atproto.repo.listRecords`, paginated | `com.atproto.sync.getRepo`, streamed CAR | +| requests | one page per (user, collection) | one fetch per user, all collections | +| speed | slower | dramatically faster on multi-collection configs / deep histories | +| bandwidth | per-collection cap | pulls the whole repo even if you only care about one collection | +| resumable | yes (`backfills.pds_cursor`) | no — failure restarts the user | + +If you're indexing a single collection or want a strict per-collection bandwidth cap, stick with `listRecords`. If you index several collections per user, `car` is usually a big win. + +### Excluding users (`userFilter`) + +Skip users whose handle, DID, or PDS matches a predicate — they won't be backfilled, and their commits get dropped from jetstream. + +```ts +{ + userFilter: ({ did, handle, pds }) => + handle?.endsWith(".example.com") ?? false, +} +``` + +When the filter returns true: + +- `identities.excluded` is flipped to `1` for the user. +- Pending `backfills` rows for the user are dropped. +- Future `getPDS` lookups short-circuit (so on-demand backfill no-ops). +- Jetstream loads the excluded set into memory and drops their commits and identity events. + +Filter checks fire whenever an identity is freshly resolved — first PDS lookup, stale-identity refresh, and `#identity` events from jetstream. Handle changes that flip exclusion are picked up on the next identity event. + +The filter is sync. `handle` and `pds` may be `null` when called from a path that doesn't have them yet (e.g. a `#identity` event carries handle but not PDS). + ### Workers CLI For Cloudflare Workers deploys, `@atmo-dev/contrail` ships a `contrail` bin that handles the `wrangler.getPlatformProxy` dance — no script file, no package.json alias needed: @@ -200,6 +244,8 @@ const db = createPostgresDatabase(pool); | `jetstreams` | Bluesky | Jetstream URLs | | `relays` | Bluesky | Relay URLs for discovery | | `notify` | off | `true` opens `notifyOfUpdate`; a string requires `Bearer` | +| `backfillMethod` | `"listRecords"` | `"car"` opts into one-fetch-per-user CAR streaming | +| `userFilter` | — | `(id) => boolean` — `true` excludes the user from indexing | | `feeds` | — | See [Feeds](./04-feeds.md) | | `spaces` | — | See [Spaces](./06-spaces.md) | | `community` | — | See [Communities](./07-communities.md) | diff --git a/packages/contrail/package.json b/packages/contrail/package.json index 49123de..a9547bb 100644 --- a/packages/contrail/package.json +++ b/packages/contrail/package.json @@ -65,6 +65,7 @@ }, "dependencies": { "@atcute/atproto": "^3.1.10", + "@atcute/car": "^5.1.1", "@atcute/cbor": "^2.3.2", "@atcute/cid": "^2.4.1", "@atcute/client": "^4.2.1", @@ -72,6 +73,7 @@ "@atcute/identity-resolver": "^1.2.2", "@atcute/jetstream": "^1.0.2", "@atcute/lexicons": "^1.2.9", + "@atcute/repo": "^0.1.4", "@atcute/xrpc-server": "^0.1.12", "cac": "^7.0.0", "hono": "^4.12.8", diff --git a/packages/contrail/src/contrail.ts b/packages/contrail/src/contrail.ts index 7c52aa0..469eb03 100644 --- a/packages/contrail/src/contrail.ts +++ b/packages/contrail/src/contrail.ts @@ -188,17 +188,27 @@ export class Contrail { // caller hasn't supplied their own. Throttle at 2s so we don't spam in // fast/local runs; final summary always prints. let effective = options; + let lastUsersExcluded = 0; if (!options?.onProgress) { let lastLogAt = 0; effective = { ...options, - onProgress: ({ records, usersComplete, usersTotal, usersFailed }) => { + onProgress: ({ + records, + usersComplete, + usersTotal, + usersFailed, + usersExcluded, + }) => { + lastUsersExcluded = usersExcluded; const now = Date.now(); if (now - lastLogAt < 2_000) return; lastLogAt = now; const failStr = usersFailed > 0 ? `, ${usersFailed} failed` : ""; + const exclStr = + usersExcluded > 0 ? `, ${usersExcluded} excluded` : ""; logger?.log?.( - ` ${records} records | ${usersComplete}/${usersTotal} users${failStr}` + ` ${records} records | ${usersComplete}/${usersTotal} users${failStr}${exclStr}` ); }, }; @@ -207,8 +217,10 @@ export class Contrail { logger?.log?.("backfilling…"); const backfilled = await this.backfill(effective, d); const elapsedS = ((Date.now() - startedAt) / 1000).toFixed(1); + const exclSummary = + lastUsersExcluded > 0 ? ` (${lastUsersExcluded} excluded)` : ""; logger?.log?.( - ` done: ${backfilled} records across ${discovered.length} users in ${elapsedS}s` + ` done: ${backfilled} records across ${discovered.length} users in ${elapsedS}s${exclSummary}` ); return { discovered: discovered.length, backfilled }; } diff --git a/packages/contrail/src/core/backfill-car.ts b/packages/contrail/src/core/backfill-car.ts new file mode 100644 index 0000000..4094549 --- /dev/null +++ b/packages/contrail/src/core/backfill-car.ts @@ -0,0 +1,196 @@ +/** + * CAR-based backfill: fetch a user's entire repo via + * `com.atproto.sync.getRepo` and stream-decode it with `@atcute/repo`. + * + * One HTTP request per DID covers every configured collection — replaces + * paginated `com.atproto.repo.listRecords` walks (one per collection) and + * is dramatically faster on users with several configured collections or + * deep history. + */ +import { type Did } from "@atcute/lexicons"; +import { isDid } from "@atcute/lexicons/syntax"; +import * as Repo from "@atcute/repo"; + +import type { ContrailConfig, Database, IngestEvent } from "./types"; +import { getCollectionNsids } from "./types"; +import { applyEvents } from "./db"; +import { getPDS } from "./client"; +import { recordTimeUs, filterEventsBySubject } from "./backfill-shared"; + +const DEFAULT_REQUEST_TIMEOUT_MS = 60_000; +const BATCH_SIZE = 100; + +export interface BackfillCarOptions { + /** Skip replay detection in applyEvents (safe during initial backfill). */ + skipReplayDetection?: boolean; + /** Per-request timeout in ms (default: 60000). */ + requestTimeout?: number; + /** Abort the in-flight fetch (rolled into the timeout). */ + signal?: AbortSignal; +} + +export interface BackfillCarResult { + /** Records inserted across every configured collection. */ + inserted: number; + /** Streamed entries that matched a configured collection. */ + matched: number; +} + +/** + * Stream the user's repo CAR and apply records for every configured + * collection. Marks every `(did, *)` row in the `backfills` table complete + * on success. Throws on PDS resolution / fetch / parse errors so the caller + * can attribute the failure. + */ +export async function backfillUserCar( + db: Database, + did: string, + deadline: number, + config: ContrailConfig, + options?: BackfillCarOptions +): Promise { + if (Date.now() >= deadline) return { inserted: 0, matched: 0 }; + if (!isDid(did)) throw new Error(`Invalid DID: ${did}`); + + const wantedNsids = new Set(getCollectionNsids(config)); + if (wantedNsids.size === 0) { + await markRepoComplete(db, did); + return { inserted: 0, matched: 0 }; + } + + const subjectFields = collectSubjectFields(config); + + const pds = await getPDS(did as Did, db); + if (!pds) throw new Error(`PDS not found for ${did}`); + + const url = new URL("/xrpc/com.atproto.sync.getRepo", pds); + url.searchParams.set("did", did); + + const timeoutMs = options?.requestTimeout ?? DEFAULT_REQUEST_TIMEOUT_MS; + const ctrl = new AbortController(); + const onUserAbort = () => ctrl.abort(options?.signal?.reason); + options?.signal?.addEventListener("abort", onUserAbort); + const timer = setTimeout( + () => ctrl.abort(new Error(`Timeout: getRepo(${did})`)), + timeoutMs + ); + + let response: Response; + try { + response = await fetch(url.toString(), { + signal: ctrl.signal, + headers: { accept: "application/vnd.ipld.car" }, + }); + } catch (err) { + clearTimeout(timer); + options?.signal?.removeEventListener("abort", onUserAbort); + throw err; + } + + if (!response.ok || !response.body) { + clearTimeout(timer); + options?.signal?.removeEventListener("abort", onUserAbort); + throw new Error(`getRepo HTTP ${response.status}`); + } + + const nowUs = Date.now() * 1000; + let inserted = 0; + let matched = 0; + const buffers = new Map(); + + const flush = async (collection: string): Promise => { + const events = buffers.get(collection); + if (!events || events.length === 0) return; + buffers.set(collection, []); + + let toApply = events; + const subjectField = subjectFields.get(collection); + if (subjectField) { + toApply = await filterEventsBySubject(db, events, subjectField); + } + if (toApply.length > 0) { + await applyEvents(db, toApply, config, { + skipReplayDetection: options?.skipReplayDetection, + skipFeedFanout: true, + }); + inserted += toApply.length; + } + }; + + const reader = Repo.fromStream(response.body); + try { + for await (const entry of reader) { + if (Date.now() >= deadline) { + throw new Error(`Deadline exceeded mid-stream for ${did}`); + } + if (!wantedNsids.has(entry.collection)) continue; + + matched++; + const buf = buffers.get(entry.collection) ?? []; + const record = entry.record; + buf.push({ + uri: `at://${did}/${entry.collection}/${entry.rkey}`, + did, + collection: entry.collection, + rkey: entry.rkey, + operation: "create", + cid: entry.cid.$link, + record: JSON.stringify(record), + time_us: recordTimeUs(record, entry.collection, config, nowUs), + indexed_at: nowUs, + }); + buffers.set(entry.collection, buf); + + if (buf.length >= BATCH_SIZE) { + await flush(entry.collection); + } + } + + for (const collection of [...buffers.keys()]) { + await flush(collection); + } + } finally { + clearTimeout(timer); + options?.signal?.removeEventListener("abort", onUserAbort); + await reader.dispose(); + } + + await markRepoComplete(db, did); + return { inserted, matched }; +} + +/** Mark every `(did, *)` row in `backfills` complete in one statement. */ +export async function markRepoComplete( + db: Database, + did: string +): Promise { + await db + .prepare( + "UPDATE backfills SET completed = 1, last_error = NULL WHERE did = ?" + ) + .bind(did) + .run(); +} + +/** Bump `retries` and stamp `last_error` on every `(did, *)` row. */ +export async function markRepoFailed( + db: Database, + did: string, + error: string +): Promise { + await db + .prepare( + "UPDATE backfills SET retries = retries + 1, last_error = ? WHERE did = ?" + ) + .bind(error, did) + .run(); +} + +/** Build NSID → subjectField map for collections that declare one. */ +function collectSubjectFields(config: ContrailConfig): Map { + const out = new Map(); + for (const c of Object.values(config.collections)) { + if (c.subjectField) out.set(c.collection, c.subjectField); + } + return out; +} diff --git a/packages/contrail/src/core/backfill-list-records.ts b/packages/contrail/src/core/backfill-list-records.ts new file mode 100644 index 0000000..accdf93 --- /dev/null +++ b/packages/contrail/src/core/backfill-list-records.ts @@ -0,0 +1,224 @@ +/** + * Paginated `com.atproto.repo.listRecords` backfill — the default strategy. + * + * One request per (did, collection) page; cursor stored in `backfills.pds_cursor` + * so a partial walk can resume next cycle. Slower than the CAR path + * (`./backfill-car.ts`) but lets you cap bandwidth per collection. + */ +import { type Did } from "@atcute/lexicons"; +import { isDid, isNsid } from "@atcute/lexicons/syntax"; + +import type { Client } from "@atcute/client"; +import type { ContrailConfig, Database, IngestEvent } from "./types"; +import { applyEvents } from "./db"; +import { getClient } from "./client"; +import { recordTimeUs, filterEventsBySubject } from "./backfill-shared"; + +const PAGE_SIZE = 100; +const REQUEST_TIMEOUT_MS = 10_000; + +export interface BackfillListRecordsOptions { + /** Pre-resolved client — avoids redundant PDS lookups when batching by DID. */ + client?: Client; + /** Skip replay detection in applyEvents (safe during initial backfill). */ + skipReplayDetection?: boolean; + /** Max retries per request (default: 3). Set to 0 for single-attempt mode. */ + maxRetries?: number; + /** Per-request timeout in ms (default: 10000). */ + requestTimeout?: number; +} + +async function withRetry( + fn: () => Promise, + label: string, + maxRetries = 3, + timeoutMs = REQUEST_TIMEOUT_MS +): Promise { + let lastError: unknown; + for (let attempt = 0; attempt <= maxRetries; attempt++) { + try { + return await Promise.race([ + fn(), + new Promise((_, reject) => + setTimeout(() => reject(new Error(`Timeout: ${label}`)), timeoutMs) + ), + ]); + } catch (err) { + lastError = err; + if (attempt < maxRetries) { + const delay = Math.min(1000 * 2 ** attempt, 10000); + await new Promise((r) => setTimeout(r, delay)); + } + } + } + throw lastError; +} + +async function markFailed( + db: Database, + did: string, + collection: string, + error: string +): Promise { + await db + .prepare( + "UPDATE backfills SET retries = retries + 1, last_error = ? WHERE did = ? AND collection = ?" + ) + .bind(error, did, collection) + .run(); +} + +export async function backfillUserListRecords( + db: Database, + did: string, + collection: string, + deadline: number, + config: ContrailConfig | undefined, + options?: BackfillListRecordsOptions +): Promise { + if (Date.now() >= deadline) return 0; + + const status = await db + .prepare( + "SELECT completed, pds_cursor, retries FROM backfills WHERE did = ? AND collection = ?" + ) + .bind(did, collection) + .first<{ completed: number; pds_cursor: string | null; retries: number }>(); + + if (status?.completed) return 0; + + if (!status) { + await db + .prepare( + "INSERT INTO backfills (did, collection, completed) VALUES (?, ?, 0) ON CONFLICT DO NOTHING" + ) + .bind(did, collection) + .run(); + } + + let currentCursor: string | undefined = status?.pds_cursor ?? undefined; + const retries = options?.maxRetries ?? 3; + const timeout = options?.requestTimeout ?? REQUEST_TIMEOUT_MS; + + if (!isDid(did)) { + await markFailed(db, did, collection, `Invalid DID: ${did}`); + return 0; + } + + if (!isNsid(collection)) { + await markFailed(db, did, collection, `Invalid NSID: ${collection}`); + return 0; + } + + let client = options?.client; + if (!client) { + try { + client = await withRetry( + () => getClient(did as Did, db), + `getClient(${did})`, + Math.min(retries, 1), + timeout + ); + } catch (err) { + await markFailed(db, did, collection, String(err)); + return 0; + } + } + + let totalInserted = 0; + let done = false; + + // If this collection declares a subjectField, we drop records whose subject + // DID isn't already in our identities table (lookup happens once per page). + const colConfig = config + ? Object.values(config.collections).find((c) => c.collection === collection) + : undefined; + const subjectField = colConfig?.subjectField; + + try { + while (Date.now() < deadline) { + const response = await withRetry( + () => + client!.get("com.atproto.repo.listRecords", { + params: { + repo: did as Did, + collection, + limit: PAGE_SIZE, + cursor: currentCursor, + }, + }), + `listRecords(${did}/${collection})`, + retries, + timeout + ); + if (!response.ok) { + await markFailed( + db, + did, + collection, + `listRecords status ${response.status}` + ); + return totalInserted; + } + + if (response.data.records.length === 0) { + done = true; + break; + } + + const now = Date.now(); + const nowUs = now * 1000; + let events: IngestEvent[] = response.data.records.map((r) => ({ + uri: r.uri, + did, + collection, + rkey: r.uri.split("/").pop()!, + operation: "create" as const, + cid: r.cid, + record: JSON.stringify(r.value), + time_us: recordTimeUs(r.value, collection, config, nowUs), + indexed_at: nowUs, + })); + + if (subjectField) { + events = await filterEventsBySubject(db, events, subjectField); + } + + if (events.length > 0) { + await applyEvents(db, events, config, { + skipReplayDetection: options?.skipReplayDetection, + skipFeedFanout: true, + }); + } + totalInserted += events.length; + + currentCursor = response.data.cursor ?? undefined; + + await db + .prepare( + "UPDATE backfills SET pds_cursor = ? WHERE did = ? AND collection = ?" + ) + .bind(currentCursor ?? null, did, collection) + .run(); + + if (!currentCursor) { + done = true; + break; + } + } + } catch (err) { + await markFailed(db, did, collection, String(err)); + return totalInserted; + } + + if (done) { + await db + .prepare( + "UPDATE backfills SET completed = 1 WHERE did = ? AND collection = ?" + ) + .bind(did, collection) + .run(); + } + + return totalInserted; +} diff --git a/packages/contrail/src/core/backfill-shared.ts b/packages/contrail/src/core/backfill-shared.ts new file mode 100644 index 0000000..5abbb8f --- /dev/null +++ b/packages/contrail/src/core/backfill-shared.ts @@ -0,0 +1,75 @@ +/** Helpers shared between the listRecords-based and CAR-based backfill paths. */ +import { isDid } from "@atcute/lexicons/syntax"; + +import type { ContrailConfig, Database, IngestEvent } from "./types"; +import { shortNameForNsid } from "./types"; + +const DEFAULT_TIME_FIELD = "createdAt"; + +/** Parse the record's canonical time (e.g. createdAt) and return microseconds. + * Falls back to `nowUs` when missing/invalid. Clamps to nowUs to avoid + * user-controlled future timestamps pinning records at the top of feeds. */ +export function recordTimeUs( + record: unknown, + collection: string, + config: ContrailConfig | undefined, + nowUs: number +): number { + if (!config) return nowUs; + const short = shortNameForNsid(config, collection); + const colCfg = short ? config.collections[short] : undefined; + const field = colCfg?.timeField ?? DEFAULT_TIME_FIELD; + if (field === false) return nowUs; + const raw = + record && typeof record === "object" + ? (record as Record)[field] + : undefined; + if (typeof raw !== "string") return nowUs; + const ms = Date.parse(raw); + if (!Number.isFinite(ms) || ms <= 0) return nowUs; + const us = ms * 1000; + return us > nowUs ? nowUs : us; +} + +/** Drop events whose `subjectField` value is a DID we have no identity for. + * One bulk SELECT per call, suitable for use after each backfill batch. */ +export async function filterEventsBySubject( + db: Database, + events: IngestEvent[], + subjectField: string +): Promise { + const subjects = new Set(); + const eventSubjects = new Map(); + for (const e of events) { + if (!e.record) continue; + let subj: unknown; + try { + subj = JSON.parse(e.record)?.[subjectField]; + } catch { + continue; + } + if (typeof subj === "string" && isDid(subj)) { + subjects.add(subj); + eventSubjects.set(e.uri, subj); + } + } + if (subjects.size === 0) return []; + + const known = new Set(); + const list = [...subjects]; + const CHUNK = 100; + for (let i = 0; i < list.length; i += CHUNK) { + const chunk = list.slice(i, i + CHUNK); + const placeholders = chunk.map(() => "?").join(","); + const rows = await db + .prepare(`SELECT did FROM identities WHERE did IN (${placeholders})`) + .bind(...chunk) + .all<{ did: string }>(); + for (const r of rows.results ?? []) known.add(r.did); + } + + return events.filter((e) => { + const subj = eventSubjects.get(e.uri); + return subj !== undefined && known.has(subj); + }); +} diff --git a/packages/contrail/src/core/backfill.ts b/packages/contrail/src/core/backfill.ts index 559131a..418267f 100644 --- a/packages/contrail/src/core/backfill.ts +++ b/packages/contrail/src/core/backfill.ts @@ -1,49 +1,61 @@ import { type Did } from "@atcute/lexicons"; -import { isDid, isNsid } from "@atcute/lexicons/syntax"; import type { Client } from "@atcute/client"; -import type { ContrailConfig, Database, IngestEvent } from "./types"; +import type { ContrailConfig, Database } from "./types"; import { getDiscoverableNsids, getDependentNsids, DEFAULT_RELAYS, - shortNameForNsid, } from "./types"; -import { applyEvents, getLastCursor, saveCursor } from "./db"; +import { getLastCursor, saveCursor } from "./db"; import { getClient, getPDS } from "./client"; +import { isExcluded, sweepUserFilter } from "./user-filter"; +import { + backfillUserCar, + markRepoComplete, + markRepoFailed, +} from "./backfill-car"; +import type { BackfillCarOptions } from "./backfill-car"; +import { backfillUserListRecords } from "./backfill-list-records"; +import type { BackfillListRecordsOptions } from "./backfill-list-records"; -const DEFAULT_TIME_FIELD = "createdAt"; +const MAX_RETRIES = 5; +const REQUEST_TIMEOUT_MS = 10_000; -/** Parse the record's canonical time (e.g. createdAt) and return microseconds. - * Falls back to `nowUs` when missing/invalid. Clamps to nowUs to avoid - * user-controlled future timestamps pinning records at the top of feeds. */ -function recordTimeUs( - record: unknown, - collection: string, - config: ContrailConfig | undefined, - nowUs: number -): number { - if (!config) return nowUs; - const short = shortNameForNsid(config, collection); - const colCfg = short ? config.collections[short] : undefined; - const field = colCfg?.timeField ?? DEFAULT_TIME_FIELD; - if (field === false) return nowUs; - const raw = - record && typeof record === "object" - ? (record as Record)[field] - : undefined; - if (typeof raw !== "string") return nowUs; - const ms = Date.parse(raw); - if (!Number.isFinite(ms) || ms <= 0) return nowUs; - const us = ms * 1000; - return us > nowUs ? nowUs : us; +async function countExcluded(db: Database): Promise { + const r = await db + .prepare("SELECT COUNT(*) AS c FROM identities WHERE excluded = 1") + .first<{ c: number }>(); + return Number(r?.c ?? 0); } -const PAGE_SIZE = 100; -const BATCH_SIZE = 100; -const MAX_RETRIES = 5; - -const REQUEST_TIMEOUT_MS = 10_000; +/** + * Resolve identities for every DID in `dids` in parallel batches. This + * populates the `identities` table and — when `config.userFilter` is set — + * lets the filter mark exclusions (and delete the matching `backfills` rows) + * BEFORE the per-user backfill loop runs. Without this, excluded users + * still pay for a full slingshot resolve + getClient throw cycle inside + * each concurrency slot, which is the slow path when most of your population + * is filtered out. + * + * No-op when no filter is configured (the per-user loop already pre-warms + * the PDS cache in the background). + */ +async function preResolveIdentitiesForFilter( + db: Database, + config: ContrailConfig, + dids: string[], + concurrency = 200 +): Promise { + if (!config.userFilter || dids.length === 0) return; + for (let i = 0; i < dids.length; i += concurrency) { + await Promise.allSettled( + dids + .slice(i, i + concurrency) + .map((did) => getPDS(did as Did, db, config).catch(() => {})) + ); + } +} async function withRetry( fn: () => Promise, @@ -71,90 +83,58 @@ async function withRetry( throw lastError; } -/** Drop events whose `subjectField` value is a DID we have no identity for. - * One bulk SELECT per call, suitable for use after each backfill page. */ -async function filterEventsBySubject( - db: Database, - events: IngestEvent[], - subjectField: string -): Promise { - const subjects = new Set(); - const eventSubjects = new Map(); - for (const e of events) { - if (!e.record) continue; - let subj: unknown; - try { - subj = JSON.parse(e.record)?.[subjectField]; - } catch { - continue; - } - if (typeof subj === "string" && isDid(subj)) { - subjects.add(subj); - eventSubjects.set(e.uri, subj); - } - } - if (subjects.size === 0) return []; - - const known = new Set(); - const list = [...subjects]; - const CHUNK = 100; - for (let i = 0; i < list.length; i += CHUNK) { - const chunk = list.slice(i, i + CHUNK); - const placeholders = chunk.map(() => "?").join(","); - const rows = await db - .prepare(`SELECT did FROM identities WHERE did IN (${placeholders})`) - .bind(...chunk) - .all<{ did: string }>(); - for (const r of rows.results ?? []) known.add(r.did); - } +export interface BackfillOptions + extends BackfillListRecordsOptions, + BackfillCarOptions {} - return events.filter((e) => { - const subj = eventSubjects.get(e.uri); - return subj !== undefined && known.has(subj); - }); +function getBackfillMethod(config?: ContrailConfig): "car" | "listRecords" { + return config?.backfillMethod ?? "listRecords"; } -async function markFailed( +/** + * Ensure this user's records are backfilled. Dispatches on + * `config.backfillMethod`: + * - `"listRecords"` (default): per-collection paginated walk. + * - `"car"`: one CAR fetch covering every configured collection. + * + * Excluded users (`identities.excluded = 1`) short-circuit to 0. + */ +export async function backfillUser( db: Database, did: string, collection: string, - error: string -): Promise { - await db - .prepare( - "UPDATE backfills SET retries = retries + 1, last_error = ? WHERE did = ? AND collection = ?" - ) - .bind(error, did, collection) - .run(); -} + deadline: number, + config: ContrailConfig | undefined, + options?: BackfillOptions +): Promise { + if (Date.now() >= deadline) return 0; + if (await isExcluded(db, did)) return 0; -export interface BackfillOptions { - /** Pre-resolved client — avoids redundant PDS lookups when batching by DID */ - client?: Client; - /** Skip replay detection in applyEvents (safe during initial backfill) */ - skipReplayDetection?: boolean; - /** Max retries per request (default: 3). Set to 0 for single-attempt mode. */ - maxRetries?: number; - /** Per-request timeout in ms (default: 10000). */ - requestTimeout?: number; + if (getBackfillMethod(config) === "car") { + return backfillUserViaCar(db, did, collection, deadline, config, options); + } + return backfillUserListRecords(db, did, collection, deadline, config, options); } -export async function backfillUser( +/** CAR-flavored wrapper that gates on the (did, collection) row but always + * fetches the whole repo. Once the CAR succeeds, every (did, *) row gets + * marked complete in one statement. */ +async function backfillUserViaCar( db: Database, did: string, collection: string, deadline: number, - config?: ContrailConfig, + config: ContrailConfig | undefined, options?: BackfillOptions ): Promise { - if (Date.now() >= deadline) return 0; + if (!config) return 0; const status = await db .prepare( - "SELECT completed, pds_cursor, retries FROM backfills WHERE did = ? AND collection = ?" + "SELECT completed FROM backfills WHERE did = ? AND collection = ?" ) .bind(did, collection) - .first<{ completed: number; pds_cursor: string | null; retries: number }>(); + .first<{ completed: number }>(); if (status?.completed) return 0; @@ -167,142 +147,36 @@ export async function backfillUser( .run(); } - let currentCursor: string | undefined = status?.pds_cursor ?? undefined; const retries = options?.maxRetries ?? 3; const timeout = options?.requestTimeout ?? REQUEST_TIMEOUT_MS; - if (!isDid(did)) { - await markFailed(db, did, collection, `Invalid DID: ${did}`); - return 0; - } - - if (!isNsid(collection)) { - await markFailed(db, did, collection, `Invalid NSID: ${collection}`); - return 0; - } - - let client = options?.client; - if (!client) { - try { - client = await withRetry( - () => getClient(did as Did, db), - `getClient(${did})`, - Math.min(retries, 1), - timeout - ); - } catch (err) { - await markFailed(db, did, collection, String(err)); - return 0; - } - } - - let totalInserted = 0; - let done = false; - - // Lookup subject filter once: if this collection declares a subjectField, we - // drop records whose subject DID isn't already in our identities table. - const collectionShort = config - ? shortNameForNsid(config, collection) - : undefined; - const subjectField = collectionShort - ? config?.collections[collectionShort]?.subjectField - : undefined; - try { - while (Date.now() < deadline) { - const response = await withRetry( - () => - client!.get("com.atproto.repo.listRecords", { - params: { - repo: did as Did, - collection, - limit: PAGE_SIZE, - cursor: currentCursor, - }, - }), - `listRecords(${did}/${collection})`, - retries, - timeout - ); - if (!response.ok) { - await markFailed( - db, - did, - collection, - `listRecords status ${response.status}` - ); - return totalInserted; - } - - if (response.data.records.length === 0) { - done = true; - break; - } - - const now = Date.now(); - const nowUs = now * 1000; - let events: IngestEvent[] = response.data.records.map((r) => ({ - uri: r.uri, - did, - collection, - rkey: r.uri.split("/").pop()!, - operation: "create" as const, - cid: r.cid, - record: JSON.stringify(r.value), - time_us: recordTimeUs(r.value, collection, config, nowUs), - indexed_at: nowUs, - })); - - if (subjectField) { - events = await filterEventsBySubject(db, events, subjectField); - } - - if (events.length > 0) { - await applyEvents(db, events, config, { - skipReplayDetection: options?.skipReplayDetection, - skipFeedFanout: true, - }); - } - totalInserted += events.length; - - currentCursor = response.data.cursor ?? undefined; - - await db - .prepare( - "UPDATE backfills SET pds_cursor = ? WHERE did = ? AND collection = ?" - ) - .bind(currentCursor ?? null, did, collection) - .run(); - - if (!currentCursor) { - done = true; - break; - } - } + const result = await withRetry( + () => + backfillUserCar(db, did, deadline, config, { + ...options, + requestTimeout: timeout, + }), + `backfillUserCar(${did})`, + retries, + timeout + ); + return result.inserted; } catch (err) { - await markFailed(db, did, collection, String(err)); - return totalInserted; - } - - if (done) { - await db - .prepare( - "UPDATE backfills SET completed = 1 WHERE did = ? AND collection = ?" - ) - .bind(did, collection) - .run(); + await markRepoFailed(db, did, String(err)); + return 0; } - - return totalInserted; } -// --- Bulk backfill (groups by DID, resolves client once) --- +// --- Bulk backfill --- export interface BackfillProgress { records: number; usersComplete: number; usersTotal: number; usersFailed: number; + /** Identities currently flagged `excluded = 1` (matched `config.userFilter`). */ + usersExcluded: number; } export interface BackfillAllOptions { @@ -315,9 +189,6 @@ export async function backfillPending( config: ContrailConfig, options?: BackfillAllOptions ): Promise { - const concurrency = options?.concurrency ?? 100; - let totalBackfilled = 0; - // Anchor the jetstream cursor to now if it hasn't been set yet, so records // emitted during backfill are replayed once jetstream starts. if ((await getLastCursor(db)) === null) { @@ -329,17 +200,88 @@ export async function backfillPending( .prepare("UPDATE backfills SET retries = 0 WHERE completed = 0") .run(); + // Apply userFilter to all identities — catches users resolved before the + // filter was configured. Returns the running excluded count for progress. + const log = config.logger ?? console; + const sweep = await sweepUserFilter(db, config); + if (sweep.newlyExcluded > 0) { + log.log?.( + `[backfill] userFilter excluded ${sweep.newlyExcluded} new user(s) (${sweep.totalExcluded} total)` + ); + } else if (sweep.totalExcluded > 0) { + log.log?.( + `[backfill] ${sweep.totalExcluded} user(s) excluded by userFilter (skipped)` + ); + } + + return getBackfillMethod(config) === "car" + ? backfillPendingCar(db, config, options, sweep.totalExcluded) + : backfillPendingListRecords(db, config, options, sweep.totalExcluded); +} + +/** Per-(did, collection) listRecords loop. Groups by DID so the PDS lookup + * and `getClient` happen once per user. */ +async function backfillPendingListRecords( + db: Database, + config: ContrailConfig, + options?: BackfillAllOptions, + usersExcluded = 0 +): Promise { + const concurrency = options?.concurrency ?? 100; + let totalBackfilled = 0; + + const log = config.logger ?? console; + while (true) { - const pending = await db + let pending = await db .prepare( - "SELECT did, collection FROM backfills WHERE completed = 0 AND retries < ? ORDER BY did" + `SELECT b.did, b.collection + FROM backfills b + LEFT JOIN identities i ON i.did = b.did + WHERE b.completed = 0 AND b.retries < ? AND COALESCE(i.excluded, 0) = 0 + ORDER BY b.did` ) .bind(MAX_RETRIES) .all<{ did: string; collection: string }>(); - const rows = pending.results ?? []; + let rows = pending.results ?? []; if (rows.length === 0) break; + // When userFilter is set, resolve identities upfront for every pending + // DID so exclusions take effect (and matching backfills rows get + // deleted) before the slow listRecords loop spins up. Then re-query. + if (config.userFilter) { + const allDids = [...new Set(rows.map((r) => r.did))]; + const before = Date.now(); + log.log?.( + `[backfill] pre-resolving ${allDids.length} identities to apply userFilter…` + ); + await preResolveIdentitiesForFilter(db, config, allDids); + pending = await db + .prepare( + `SELECT b.did, b.collection + FROM backfills b + LEFT JOIN identities i ON i.did = b.did + WHERE b.completed = 0 AND b.retries < ? AND COALESCE(i.excluded, 0) = 0 + ORDER BY b.did` + ) + .bind(MAX_RETRIES) + .all<{ did: string; collection: string }>(); + const filteredRows = pending.results ?? []; + const droppedDids = + allDids.length - new Set(filteredRows.map((r) => r.did)).size; + log.log?.( + `[backfill] pre-resolve done in ${( + (Date.now() - before) / + 1000 + ).toFixed(1)}s — ${droppedDids} excluded, ${ + filteredRows.length + } (did, collection) rows remaining` + ); + rows = filteredRows; + if (rows.length === 0) break; + } + // Group by DID so we resolve PDS once per user const byDid = new Map(); for (const row of rows) { @@ -350,33 +292,42 @@ export async function backfillPending( const dids = [...byDid.keys()]; - // Resolve PDS endpoints in background (populates in-memory cache) - const resolvePromise = (async () => { - for (let i = 0; i < dids.length; i += 200) { - await Promise.allSettled( - dids.slice(i, i + 200).map((did) => - getPDS(did as Did, db).catch(() => {}) - ) - ); - } - })(); + // Resolve PDS endpoints in background (populates in-memory cache). + // Skipped when we already pre-resolved above for the filter sweep. + const resolvePromise = config.userFilter + ? Promise.resolve() + : (async () => { + for (let i = 0; i < dids.length; i += 200) { + await Promise.allSettled( + dids.slice(i, i + 200).map((did) => + getPDS(did as Did, db, config).catch(() => {}) + ) + ); + } + })(); let roundBackfilled = 0; let usersComplete = 0; let usersFailed = 0; - const failedDids: string[] = []; + const failedDids = new Set(); const FAST_TIMEOUT = 3_000; + const hasFilter = !!config.userFilter; - const emitProgress = () => + const emitProgress = async () => { + if (hasFilter) usersExcluded = await countExcluded(db); options?.onProgress?.({ records: totalBackfilled + roundBackfilled, usersComplete, usersTotal: dids.length, usersFailed, + usersExcluded, }); + }; - // Fast pass: single attempt per user with short timeout + // Fast pass: single attempt per user with short timeout. Only count + // toward `usersComplete` when *every* collection for the DID succeeds — + // partial failures get punted to the retry pass and accounted there. for (let i = 0; i < dids.length; i += concurrency) { const batch = dids.slice(i, i + concurrency); @@ -385,32 +336,42 @@ export async function backfillPending( let client: Client | undefined; try { client = await withRetry( - () => getClient(did as Did, db), + () => getClient(did as Did, db, config), `getClient(${did})`, 0, FAST_TIMEOUT ); } catch { - failedDids.push(did); + // Distinguish "filter excluded the user" from a real PDS failure — + // the userFilter side-channel marks identities.excluded=1, so a + // getClient failure on an excluded DID is a clean skip, not a + // failure to retry. + if (hasFilter && (await isExcluded(db, did))) { + usersComplete++; + return 0; + } + failedDids.add(did); return 0; } const cols = byDid.get(did)!; + let anyFailed = false; const counts = await Promise.all( cols.map((col) => - backfillUser(db, did, col, Infinity, config, { + backfillUserListRecords(db, did, col, Infinity, config, { client, skipReplayDetection: true, maxRetries: 0, requestTimeout: FAST_TIMEOUT, }).catch(() => { - failedDids.push(did); + anyFailed = true; return 0; }) ) ); - usersComplete++; + if (anyFailed) failedDids.add(did); + else usersComplete++; return counts.reduce((a, b) => a + b, 0); }) ); @@ -419,13 +380,12 @@ export async function backfillPending( if (r.status === "fulfilled") roundBackfilled += r.value; } - emitProgress(); + await emitProgress(); } // Retry pass: failed DIDs get retries with backoff, still in concurrent batches - if (failedDids.length > 0) { - const uniqueFailed = [...new Set(failedDids)]; - usersComplete -= uniqueFailed.length; // don't count them yet + if (failedDids.size > 0) { + const uniqueFailed = [...failedDids]; for (let i = 0; i < uniqueFailed.length; i += concurrency) { const batch = uniqueFailed.slice(i, i + concurrency); @@ -435,13 +395,23 @@ export async function backfillPending( let client: Client | undefined; try { client = await withRetry( - () => getClient(did as Did, db), + () => getClient(did as Did, db, config), `getClient(${did})`, 2 ); } catch (err) { + // Same exclusion-vs-failure distinction as the fast pass. + if (hasFilter && (await isExcluded(db, did))) { + usersComplete++; + return 0; + } for (const col of byDid.get(did)!) { - await markFailed(db, did, col, String(err)); + await db + .prepare( + "UPDATE backfills SET retries = retries + 1, last_error = ? WHERE did = ? AND collection = ?" + ) + .bind(String(err), did, col) + .run(); } usersFailed++; usersComplete++; @@ -451,7 +421,7 @@ export async function backfillPending( const cols = byDid.get(did)!; const counts = await Promise.all( cols.map((col) => - backfillUser(db, did, col, Infinity, config, { + backfillUserListRecords(db, did, col, Infinity, config, { client, skipReplayDetection: true, maxRetries: 2, @@ -467,7 +437,7 @@ export async function backfillPending( if (r.status === "fulfilled") roundBackfilled += r.value; } - emitProgress(); + await emitProgress(); } } @@ -481,6 +451,181 @@ export async function backfillPending( return totalBackfilled; } +/** Per-DID CAR loop. One `com.atproto.sync.getRepo` call per user covers + * every configured collection in one pass. */ +async function backfillPendingCar( + db: Database, + config: ContrailConfig, + options?: BackfillAllOptions, + usersExcluded = 0 +): Promise { + const concurrency = options?.concurrency ?? 50; + let totalBackfilled = 0; + + const log = config.logger ?? console; + + while (true) { + let pending = await db + .prepare( + `SELECT DISTINCT b.did + FROM backfills b + LEFT JOIN identities i ON i.did = b.did + WHERE b.completed = 0 AND b.retries < ? AND COALESCE(i.excluded, 0) = 0 + ORDER BY b.did` + ) + .bind(MAX_RETRIES) + .all<{ did: string }>(); + + let dids = (pending.results ?? []).map((r) => r.did); + if (dids.length === 0) break; + + // Pre-resolve identities so userFilter can prune the list before we + // start firing CAR fetches. + if (config.userFilter) { + const before = Date.now(); + log.log?.( + `[backfill] pre-resolving ${dids.length} identities to apply userFilter…` + ); + await preResolveIdentitiesForFilter(db, config, dids); + pending = await db + .prepare( + `SELECT DISTINCT b.did + FROM backfills b + LEFT JOIN identities i ON i.did = b.did + WHERE b.completed = 0 AND b.retries < ? AND COALESCE(i.excluded, 0) = 0 + ORDER BY b.did` + ) + .bind(MAX_RETRIES) + .all<{ did: string }>(); + const filtered = (pending.results ?? []).map((r) => r.did); + log.log?.( + `[backfill] pre-resolve done in ${( + (Date.now() - before) / + 1000 + ).toFixed(1)}s — ${dids.length - filtered.length} excluded, ${ + filtered.length + } users remaining` + ); + dids = filtered; + if (dids.length === 0) break; + } + + const resolvePromise = config.userFilter + ? Promise.resolve() + : (async () => { + for (let i = 0; i < dids.length; i += 200) { + await Promise.allSettled( + dids + .slice(i, i + 200) + .map((did) => getPDS(did as Did, db, config).catch(() => {})) + ); + } + })(); + + let roundBackfilled = 0; + let usersComplete = 0; + let usersFailed = 0; + const failedDids = new Set(); + + const FAST_TIMEOUT = 30_000; + const hasFilter = !!config.userFilter; + + const emitProgress = async () => { + if (hasFilter) usersExcluded = await countExcluded(db); + options?.onProgress?.({ + records: totalBackfilled + roundBackfilled, + usersComplete, + usersTotal: dids.length, + usersFailed, + usersExcluded, + }); + }; + + for (let i = 0; i < dids.length; i += concurrency) { + const batch = dids.slice(i, i + concurrency); + + const results = await Promise.allSettled( + batch.map(async (did) => { + try { + const r = await backfillUserCar(db, did, Infinity, config, { + skipReplayDetection: true, + requestTimeout: FAST_TIMEOUT, + }); + usersComplete++; + return r.inserted; + } catch { + // Filter exclusion looks like a "PDS not found" failure here. + if (hasFilter && (await isExcluded(db, did))) { + usersComplete++; + return 0; + } + failedDids.add(did); + return 0; + } + }) + ); + + for (const r of results) { + if (r.status === "fulfilled") roundBackfilled += r.value; + } + await emitProgress(); + } + + if (failedDids.size > 0) { + const uniqueFailed = [...failedDids]; + for (let i = 0; i < uniqueFailed.length; i += concurrency) { + const batch = uniqueFailed.slice(i, i + concurrency); + + const results = await Promise.allSettled( + batch.map(async (did) => { + try { + const r = await withRetry( + () => + backfillUserCar(db, did, Infinity, config, { + skipReplayDetection: true, + requestTimeout: FAST_TIMEOUT * 2, + }), + `backfillUserCar(${did})`, + 2, + FAST_TIMEOUT * 2 + ); + usersComplete++; + return r.inserted; + } catch (err) { + if (hasFilter && (await isExcluded(db, did))) { + usersComplete++; + return 0; + } + await markRepoFailed(db, did, String(err)); + usersFailed++; + usersComplete++; + return 0; + } + }) + ); + + for (const r of results) { + if (r.status === "fulfilled") roundBackfilled += r.value; + } + await emitProgress(); + } + } + + await resolvePromise; + totalBackfilled += roundBackfilled; + + if (roundBackfilled === 0) break; + } + + return totalBackfilled; +} + +// Re-export helpers used outside this module. +export { backfillUserCar, markRepoComplete, markRepoFailed }; +export type { BackfillCarOptions }; +export { backfillUserListRecords }; +export type { BackfillListRecordsOptions }; + // --- Discovery --- interface DiscoveryPage { diff --git a/packages/contrail/src/core/client.ts b/packages/contrail/src/core/client.ts index f37bea3..ba45a00 100644 --- a/packages/contrail/src/core/client.ts +++ b/packages/contrail/src/core/client.ts @@ -6,7 +6,8 @@ import { import { type Did } from "@atcute/lexicons"; import { Client, simpleFetchHandler } from "@atcute/client"; import type {} from "@atcute/atproto"; -import type { Database } from "./types"; +import type { ContrailConfig, Database } from "./types"; +import { checkUserFilter } from "./user-filter"; // Slingshot-first PDS resolution with fallback to DID document resolution const SLINGSHOT_URL = @@ -134,7 +135,8 @@ function pdsCacheSet(did: string, pds: string): void { export async function getPDS( did: Did, - db?: Database + db?: Database, + config?: ContrailConfig ): Promise { const mem = pdsCacheGet(did); if (mem) return mem; @@ -143,7 +145,7 @@ export async function getPDS( const inflight = pdsInflight.get(did); if (inflight) return inflight; - const promise = resolvePDSCached(did, db); + const promise = resolvePDSCached(did, db, config); pdsInflight.set(did, promise); try { return await promise; @@ -154,13 +156,15 @@ export async function getPDS( async function resolvePDSCached( did: Did, - db?: Database + db?: Database, + config?: ContrailConfig ): Promise { if (db) { const cached = await db - .prepare("SELECT pds FROM identities WHERE did = ? AND pds IS NOT NULL") + .prepare("SELECT pds, excluded FROM identities WHERE did = ?") .bind(did) - .first<{ pds: string }>(); + .first<{ pds: string | null; excluded: number }>(); + if (cached?.excluded) return undefined; if (cached?.pds) { pdsCacheSet(did, cached.pds); return cached.pds; @@ -170,8 +174,6 @@ async function resolvePDSCached( const resolved = await resolvePDS(did); if (!resolved?.pds) return undefined; - pdsCacheSet(did, resolved.pds); - // Persist to DB for future runs if (db) { await db @@ -180,13 +182,29 @@ async function resolvePDSCached( ) .bind(did, resolved.handle, resolved.pds, Date.now()) .run(); + + // Run the configured user filter — if it excludes this user, drop pending + // backfill rows and return as if we never resolved a PDS. + if (config?.userFilter) { + const excluded = await checkUserFilter( + db, + { did, handle: resolved.handle, pds: resolved.pds }, + config + ); + if (excluded) return undefined; + } } + pdsCacheSet(did, resolved.pds); return resolved.pds; } -export async function getClient(did: Did, db?: Database): Promise { - const pds = await getPDS(did, db); +export async function getClient( + did: Did, + db?: Database, + config?: ContrailConfig +): Promise { + const pds = await getPDS(did, db, config); if (!pds) throw new Error(`PDS not found for ${did}`); return new Client({ handler: simpleFetchHandler({ service: pds }), diff --git a/packages/contrail/src/core/db/schema.ts b/packages/contrail/src/core/db/schema.ts index 5a56cbd..4a889c3 100644 --- a/packages/contrail/src/core/db/schema.ts +++ b/packages/contrail/src/core/db/schema.ts @@ -44,7 +44,8 @@ CREATE TABLE IF NOT EXISTS identities ( did TEXT PRIMARY KEY, handle TEXT, pds TEXT, - resolved_at ${dialect.bigintType} NOT NULL + resolved_at ${dialect.bigintType} NOT NULL, + excluded INTEGER NOT NULL DEFAULT 0 ); CREATE INDEX IF NOT EXISTS idx_identities_handle ON identities(handle); `; @@ -258,6 +259,10 @@ const MIGRATIONS = [ "ALTER TABLE feed_backfills ADD COLUMN retries INTEGER NOT NULL DEFAULT 0", "ALTER TABLE feed_backfills ADD COLUMN last_error TEXT", "ALTER TABLE feed_backfills ADD COLUMN started_at BIGINT", + "ALTER TABLE identities ADD COLUMN excluded INTEGER NOT NULL DEFAULT 0", + // Index references `excluded`, so it has to run *after* the ALTER above. + // Both are guarded by `runMigrations`'s try/catch (already-exists is fine). + "CREATE INDEX IF NOT EXISTS idx_identities_excluded ON identities(excluded)", ]; async function runMigrations(db: Database): Promise { diff --git a/packages/contrail/src/core/identity.ts b/packages/contrail/src/core/identity.ts index fe207b3..71f9957 100644 --- a/packages/contrail/src/core/identity.ts +++ b/packages/contrail/src/core/identity.ts @@ -1,7 +1,10 @@ import type { Did } from "@atcute/lexicons"; -import type { Database, Logger } from "./types"; +import type { ContrailConfig, Database } from "./types"; import { isDid, isHandle } from "@atcute/lexicons/syntax"; import { resolvePDS } from "./client"; +import { checkUserFilter } from "./user-filter"; + +export { checkUserFilter, isExcluded } from "./user-filter"; const STALE_MS = 24 * 60 * 60 * 1000; // 24 hours @@ -10,6 +13,7 @@ export interface Identity { handle: string | null; pds: string | null; resolved_at: number; + excluded?: number; } async function saveIdentity(db: Database, identity: Identity): Promise { @@ -28,7 +32,8 @@ function isStale(resolvedAt: number): boolean { async function fetchAndSave( db: Database, identifier: string, - cached?: Identity | null + cached?: Identity | null, + config?: ContrailConfig ): Promise { const resolved = await resolvePDS(identifier); const identity: Identity = { @@ -38,21 +43,23 @@ async function fetchAndSave( resolved_at: Date.now(), }; await saveIdentity(db, identity); + await checkUserFilter(db, identity, config); return identity; } export async function resolveIdentity( db: Database, - did: Did + did: Did, + config?: ContrailConfig ): Promise { const cached = await db - .prepare("SELECT did, handle, pds, resolved_at FROM identities WHERE did = ?") + .prepare("SELECT did, handle, pds, resolved_at, excluded FROM identities WHERE did = ?") .bind(did) .first(); if (cached && !isStale(cached.resolved_at)) return cached; - return fetchAndSave(db, did, cached); + return fetchAndSave(db, did, cached, config); } export async function resolveIdentities( @@ -129,17 +136,23 @@ export async function resolveActor( export async function applyIdentityEvent( db: Database, did: string, - handle: string + handle: string, + config?: ContrailConfig ): Promise { await db .prepare("UPDATE identities SET handle = ?, resolved_at = ? WHERE did = ?") .bind(handle, Date.now(), did) .run(); + // Re-check the filter — handle changes can flip exclusion state. + if (config?.userFilter) { + await checkUserFilter(db, { did, handle, pds: null }, config); + } } export async function refreshStaleIdentities( db: Database, - dids: string[] + dids: string[], + config?: ContrailConfig ): Promise { if (dids.length === 0) return; @@ -169,7 +182,7 @@ export async function refreshStaleIdentities( for (const did of toRefresh) { try { - await fetchAndSave(db, did); + await fetchAndSave(db, did, undefined, config); } catch { // Silently skip unresolvable identities } diff --git a/packages/contrail/src/core/jetstream.ts b/packages/contrail/src/core/jetstream.ts index abc0287..561096f 100644 --- a/packages/contrail/src/core/jetstream.ts +++ b/packages/contrail/src/core/jetstream.ts @@ -16,6 +16,9 @@ const FEED_PRUNE_INTERVAL_MS = 60 * 60 * 1000; // 1 hour /** Mutable state that persists across ingest cycles within the same process. */ export interface IngestState { cachedKnownDids?: Set; + /** DIDs the configured `userFilter` has marked excluded. Loaded once per + * process, kept in sync via mutations from identity-event handlers. */ + cachedExcludedDids?: Set; schemaInitialized: boolean; lastFeedPruneMs: number; } @@ -32,7 +35,8 @@ export async function ingestEvents( config: ContrailConfig, cursor: number | null, safetyTimeoutMs: number = 25_000, - knownDids?: Set + knownDids?: Set, + excludedDids?: Set ): Promise<{ events: IngestEvent[]; lastCursor: number | null; @@ -86,6 +90,11 @@ export async function ingestEvents( const { commit } = event; totalCommits++; + // userFilter exclusion list — drop the event regardless of collection + if (excludedDids?.has(event.did)) { + continue; + } + const uri = `at://${event.did}/${commit.collection}/${commit.rkey}`; const short = shortNameForNsid(config, commit.collection); @@ -247,21 +256,36 @@ export async function runIngestCycle( }, timeout=${timeoutMs}ms, collections=${collections.join(", ")}` ); - // Load known DIDs for filtering dependent collections + // Load known DIDs for filtering dependent collections, plus the excluded + // set if the user configured a `userFilter`. Both share one query. const dependentCollections = getDependentNsids(config); let knownDids: Set | undefined; + let excludedDids: Set | undefined; - if (dependentCollections.length > 0) { - if (s.cachedKnownDids) { + if (dependentCollections.length > 0 || config.userFilter) { + if (s.cachedKnownDids && (s.cachedExcludedDids || !config.userFilter)) { knownDids = s.cachedKnownDids; - log.log(`Using cached known DIDs (${knownDids.size} users)`); + excludedDids = s.cachedExcludedDids; + log.log( + `Using cached identities (${knownDids.size} known, ${ + excludedDids?.size ?? 0 + } excluded)` + ); } else { const result = await db - .prepare("SELECT did FROM identities") - .all<{ did: string }>(); - knownDids = new Set((result.results ?? []).map((r) => r.did)); + .prepare("SELECT did, excluded FROM identities") + .all<{ did: string; excluded: number }>(); + knownDids = new Set(); + excludedDids = new Set(); + for (const r of result.results ?? []) { + knownDids.add(r.did); + if (r.excluded) excludedDids.add(r.did); + } s.cachedKnownDids = knownDids; - log.log(`Loaded ${knownDids.size} known DIDs from database`); + s.cachedExcludedDids = excludedDids; + log.log( + `Loaded ${knownDids.size} known DIDs (${excludedDids.size} excluded) from database` + ); } } @@ -269,7 +293,8 @@ export async function runIngestCycle( config, cursor, timeoutMs, - knownDids + knownDids, + excludedDids ); if (events.length > 0) { @@ -295,7 +320,17 @@ export async function runIngestCycle( if (identityUpdates.size > 0) { for (const [did, handle] of identityUpdates) { try { - await applyIdentityEvent(db, did, handle); + await applyIdentityEvent(db, did, handle, config); + // Re-mirror exclusion state into the cached set if the filter just + // flipped it. Cheap one-row read. + if (config.userFilter && excludedDids) { + const row = await db + .prepare("SELECT excluded FROM identities WHERE did = ?") + .bind(did) + .first<{ excluded: number }>(); + if (row?.excluded) excludedDids.add(did); + else excludedDids.delete(did); + } } catch (err) { log.warn(`[ingest] identity update failed for ${did}: ${err}`); } @@ -307,7 +342,7 @@ export async function runIngestCycle( const uniqueDids = [...new Set(events.map((e) => e.did))]; if (uniqueDids.length > 0) { try { - await refreshStaleIdentities(db, uniqueDids); + await refreshStaleIdentities(db, uniqueDids, config); } catch (err) { log.warn(`Identity refresh failed: ${err}`); } diff --git a/packages/contrail/src/core/types.ts b/packages/contrail/src/core/types.ts index c59f116..b8e5b0f 100644 --- a/packages/contrail/src/core/types.ts +++ b/packages/contrail/src/core/types.ts @@ -239,6 +239,30 @@ export interface ContrailConfig { * ingests synthesized rows for any follower already in our identities * table. Lets newcomers immediately appear in existing users' feeds. */ constellation?: ConstellationConfig | false; + /** Which strategy to use when backfilling a user's existing records. + * - `"listRecords"` (default): paginated `com.atproto.repo.listRecords` + * per (user, collection). Slower but lets you cap bandwidth per + * collection and resume mid-walk via `backfills.pds_cursor`. + * - `"car"`: one `com.atproto.sync.getRepo` call per user, streamed CAR. + * Dramatically faster on users with multiple configured collections, + * but pulls the whole repo even if you only care about one collection. */ + backfillMethod?: "car" | "listRecords"; + /** Predicate run after identity resolution. Returning true marks the user + * excluded — their `identities.excluded` column is flipped to 1, any + * pending `backfills` rows are deleted, jetstream drops their commits, + * and future backfill calls short-circuit. Use for e.g. excluding handles + * by suffix or DIDs by deny-list. */ + userFilter?: (identity: UserFilterInput) => boolean; +} + +export interface UserFilterInput { + did: string; + /** May be null when the filter is invoked from a path that hasn't + * resolved the handle yet (e.g. some identity event flows). */ + handle: string | null; + /** May be null when the filter is invoked from a path that doesn't carry + * PDS info (e.g. jetstream identity events). */ + pds: string | null; } export interface ConstellationConfig { diff --git a/packages/contrail/src/core/user-filter.ts b/packages/contrail/src/core/user-filter.ts new file mode 100644 index 0000000..2cf600e --- /dev/null +++ b/packages/contrail/src/core/user-filter.ts @@ -0,0 +1,92 @@ +import type { ContrailConfig, Database, UserFilterInput } from "./types"; + +/** + * Run `config.userFilter` against a freshly resolved identity. When the + * filter returns true the user is marked excluded (`identities.excluded = 1`) + * and any pending backfill rows are dropped so the bulk loop won't enumerate + * them again. Safe to call repeatedly. + * + * Returns true if the user is now excluded. + */ +export async function checkUserFilter( + db: Database, + identity: UserFilterInput, + config: ContrailConfig | undefined +): Promise { + if (!config?.userFilter) return false; + let excluded = false; + try { + excluded = !!config.userFilter(identity); + } catch (err) { + (config.logger ?? console).warn( + `[userFilter] threw for ${identity.did}: ${err}` + ); + return false; + } + if (!excluded) return false; + await db + .prepare("UPDATE identities SET excluded = 1 WHERE did = ?") + .bind(identity.did) + .run(); + await db + .prepare("DELETE FROM backfills WHERE did = ?") + .bind(identity.did) + .run(); + return true; +} + +/** True if `identities.excluded = 1` for this DID. */ +export async function isExcluded(db: Database, did: string): Promise { + const row = await db + .prepare("SELECT excluded FROM identities WHERE did = ?") + .bind(did) + .first<{ excluded: number }>(); + return !!row?.excluded; +} + +export interface UserFilterSweepResult { + /** Identities the filter newly marked excluded during this sweep. */ + newlyExcluded: number; + /** Total identities marked excluded after the sweep. */ + totalExcluded: number; +} + +/** + * Apply `config.userFilter` to every identity that isn't already excluded. + * Catches users that were resolved BEFORE the filter was added to config — + * fresh resolutions filter at write time, but cached identities would + * otherwise stay un-filtered until their next stale refresh. + * + * If no filter is configured, this still returns the current excluded count. + */ +export async function sweepUserFilter( + db: Database, + config: ContrailConfig | undefined +): Promise { + let newlyExcluded = 0; + + if (config?.userFilter) { + const BATCH = 500; + let lastDid = ""; + while (true) { + const rows = await db + .prepare( + "SELECT did, handle, pds FROM identities WHERE excluded = 0 AND did > ? ORDER BY did LIMIT ?" + ) + .bind(lastDid, BATCH) + .all<{ did: string; handle: string | null; pds: string | null }>(); + const list = rows.results ?? []; + if (list.length === 0) break; + for (const r of list) { + if (await checkUserFilter(db, r, config)) newlyExcluded++; + } + lastDid = list[list.length - 1].did; + if (list.length < BATCH) break; + } + } + + const total = await db + .prepare("SELECT COUNT(*) AS c FROM identities WHERE excluded = 1") + .first<{ c: number }>(); + return { newlyExcluded, totalExcluded: Number(total?.c ?? 0) }; +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 8e648bf..29e265b 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -18,6 +18,28 @@ importers: specifier: ^5.7.3 version: 5.9.3 + apps/backfill-test: + dependencies: + '@atmo-dev/contrail': + specifier: workspace:* + version: link:../../packages/contrail + devDependencies: + '@atcute/lex-cli': + specifier: ^2.8.1 + version: 2.8.1 + '@atmo-dev/contrail-lexicons': + specifier: workspace:* + version: link:../../packages/lexicons + '@cloudflare/workers-types': + specifier: ^4.20250124.0 + version: 4.20260424.1 + typescript: + specifier: ^5.7.3 + version: 5.9.3 + wrangler: + specifier: ^4.63.0 + version: 4.84.1(@cloudflare/workers-types@4.20260424.1) + apps/cloudflare-workers: dependencies: '@atmo-dev/contrail': @@ -349,6 +371,9 @@ importers: '@atcute/atproto': specifier: ^3.1.10 version: 3.1.11 + '@atcute/car': + specifier: ^5.1.1 + version: 5.1.1 '@atcute/cbor': specifier: ^2.3.2 version: 2.3.2 @@ -370,6 +395,9 @@ importers: '@atcute/lexicons': specifier: ^1.2.9 version: 1.3.0 + '@atcute/repo': + specifier: ^0.1.4 + version: 0.1.4 '@atcute/xrpc-server': specifier: ^0.1.12 version: 0.1.12