From 8c5cea638f5c5a2ca30ba7672fa098d215fbae83 Mon Sep 17 00:00:00 2001 From: Florian <45694132+flo-bit@users.noreply.github.com> Date: Tue, 4 Aug 2026 02:18:12 +0200 Subject: [PATCH] Unify record ingestion --- .changeset/lean-contrail.md | 5 + packages/contrail/src/core/backfill.ts | 94 ++------ packages/contrail/src/core/db/index.ts | 2 +- packages/contrail/src/core/db/records.ts | 53 ++--- packages/contrail/src/core/ingest.ts | 207 ++++++++++++++++++ packages/contrail/src/core/jetstream.ts | 97 +++----- packages/contrail/src/core/persistent.ts | 83 +++---- packages/contrail/src/core/refresh.ts | 32 +-- packages/contrail/src/core/router/notify.ts | 68 +++--- packages/contrail/src/core/types.ts | 2 +- packages/contrail/src/index.ts | 19 +- .../tests/feed-prune-guardrail.test.ts | 6 +- .../contrail/tests/get-record-handle.test.ts | 4 +- packages/contrail/tests/helpers.ts | 8 +- packages/contrail/tests/hydrate.test.ts | 26 +-- packages/contrail/tests/ingest.test.ts | 129 +++++++++++ packages/contrail/tests/labels-router.test.ts | 6 +- packages/contrail/tests/notify.test.ts | 20 +- packages/contrail/tests/postgres-e2e.test.ts | 30 +-- packages/contrail/tests/postgres.test.ts | 12 +- packages/contrail/tests/records.test.ts | 28 +-- packages/contrail/tests/refresh.test.ts | 8 +- packages/contrail/tests/search.test.ts | 30 +-- packages/contrail/tests/sinks.test.ts | 22 +- 24 files changed, 630 insertions(+), 361 deletions(-) create mode 100644 .changeset/lean-contrail.md create mode 100644 packages/contrail/src/core/ingest.ts create mode 100644 packages/contrail/tests/ingest.test.ts diff --git a/.changeset/lean-contrail.md b/.changeset/lean-contrail.md new file mode 100644 index 0000000..4f5a7c9 --- /dev/null +++ b/.changeset/lean-contrail.md @@ -0,0 +1,5 @@ +--- +"@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. diff --git a/packages/contrail/src/core/backfill.ts b/packages/contrail/src/core/backfill.ts index 3a7f931..d0d2836 100644 --- a/packages/contrail/src/core/backfill.ts +++ b/packages/contrail/src/core/backfill.ts @@ -10,7 +10,8 @@ import { DEFAULT_RELAYS, shortNameForNsid, } from "./types"; -import { applyEvents, getLastCursor, saveCursor } from "./db"; +import { getLastCursor, saveCursor } from "./db"; +import { createIngestEvent, ingestRecords } from "./ingest"; import { getClient, getPDS } from "./client"; const DEFAULT_TIME_FIELD = "createdAt"; @@ -21,10 +22,9 @@ const DEFAULT_TIME_FIELD = "createdAt"; function recordTimeUs( record: unknown, collection: string, - config: ContrailConfig | undefined, + config: ContrailConfig, 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; @@ -72,49 +72,6 @@ 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); - } - - return events.filter((e) => { - const subj = eventSubjects.get(e.uri); - return subj !== undefined && known.has(subj); - }); -} - async function markFailed( db: Database, did: string, @@ -132,7 +89,7 @@ async function markFailed( 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) */ + /** Skip replay detection during initial backfill. */ skipReplayDetection?: boolean; /** Max retries per request (default: 3). Set to 0 for single-attempt mode. */ maxRetries?: number; @@ -145,7 +102,7 @@ export async function backfillUser( did: string, collection: string, deadline: number, - config?: ContrailConfig, + config: ContrailConfig, options?: BackfillOptions ): Promise { if (Date.now() >= deadline) return 0; @@ -200,15 +157,6 @@ export async function backfillUser( 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( @@ -242,31 +190,29 @@ export async function backfillUser( 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); - } + const events: IngestEvent[] = response.data.records.map((record) => + createIngestEvent({ + uri: record.uri, + did, + collection, + rkey: record.uri.split("/").pop()!, + operation: "create", + cid: record.cid, + value: record.value, + timeUs: recordTimeUs(record.value, collection, config, nowUs), + indexedAt: nowUs, + }), + ); if (events.length > 0) { - await applyEvents(db, events, config, { + const result = await ingestRecords(db, events, config, { skipReplayDetection: options?.skipReplayDetection, skipFeedFanout: true, // Let sinks bulk-flush differently from live ingestion. phase: "backfill", }); + totalInserted += result.accepted.length; } - totalInserted += events.length; currentCursor = response.data.cursor ?? undefined; diff --git a/packages/contrail/src/core/db/index.ts b/packages/contrail/src/core/db/index.ts index b7a04cb..4c2acc4 100644 --- a/packages/contrail/src/core/db/index.ts +++ b/packages/contrail/src/core/db/index.ts @@ -1,6 +1,6 @@ export { initSchema, CONTRAIL_SCHEMA_VERSION } from "./schema"; export { getMeta, setMeta, getMetaNumber } from "./meta"; export { optimizeDatabase } from "./optimize"; -export { getLastCursor, saveCursor, applyEvents, lookupExistingRecords, queryRecords, pruneFeedItems, pruneActorFeed, sweepFeedItems, getFeedPruneCursor, saveFeedPruneCursor } from "./records"; +export { getLastCursor, saveCursor, lookupExistingRecords, queryRecords, pruneFeedItems, pruneActorFeed, sweepFeedItems, getFeedPruneCursor, saveFeedPruneCursor } from "./records"; export type { QueryOptions, SortOption, ExistingRecordInfo, FeedSweepResult } from "./records"; export type { RecordSource } from "../types"; diff --git a/packages/contrail/src/core/db/records.ts b/packages/contrail/src/core/db/records.ts index 3b1a380..627afd1 100644 --- a/packages/contrail/src/core/db/records.ts +++ b/packages/contrail/src/core/db/records.ts @@ -554,10 +554,10 @@ export async function lookupExistingRecords( // --- Events --- -export async function applyEvents( +export async function projectEvents( db: Database, events: IngestEvent[], - config?: ContrailConfig, + config: ContrailConfig, options?: { skipReplayDetection?: boolean; skipFeedFanout?: boolean; @@ -570,17 +570,17 @@ export async function applyEvents( ): Promise { if (events.length === 0) return; - const followCollections = config ? getFeedFollowShortNames(config) : []; - const hasCountingRelations = config ? Object.values(config.collections).some(c => + const followCollections = getFeedFollowShortNames(config); + const hasCountingRelations = Object.values(config.collections).some(c => Object.values(c.relations ?? {}).some(r => r.count !== false) - ) : false; + ); const needRecordContent = followCollections.length > 0 || hasCountingRelations; // Use pre-fetched data or look up existing records let existingMap: Map; if (options?.existing) { existingMap = options.existing; - } else if (config && !options?.skipReplayDetection) { + } else if (!options?.skipReplayDetection) { existingMap = await lookupExistingRecords(db, events, needRecordContent, config); } else { existingMap = new Map(); @@ -599,11 +599,10 @@ export async function applyEvents( for (const e of events) { // Event's collection is an NSID. Resolve its storage key from config. - // If no config, treat collection string as-is (for tests that pre-populate tables). - const short = config ? resolveCollectionKey(config, e.collection) : e.collection; + const short = resolveCollectionKey(config, e.collection); if (!short) { - (config?.logger ?? console).warn( - `[ingest] drop (unknown collection in applyEvents): ${e.operation} ${e.uri} collection=${e.collection}` + (config.logger ?? console).warn( + `[ingest] drop (unknown collection in projection): ${e.operation} ${e.uri} collection=${e.collection}` ); continue; } @@ -627,29 +626,25 @@ export async function applyEvents( ); } - if (config) { - // Collect count targets (deduplicated across the whole batch) - const existingRecordJson = existingMap.get(e.uri)?.record ?? null; - collectCountTargets(e, config, existingRecordJson, countTargets); + // Collect count targets (deduplicated across the whole batch) + const existingRecordJson = existingMap.get(e.uri)?.record ?? null; + collectCountTargets(e, config, existingRecordJson, countTargets); - // Feed fanout still needs replay detection - const existingInfo = existingMap.get(e.uri); - const isReplay = - e.operation === "delete" - ? existingInfo === undefined - : existingInfo?.cid === e.cid; + // Feed fanout still needs replay detection + const existingInfo = existingMap.get(e.uri); + const isReplay = + e.operation === "delete" + ? existingInfo === undefined + : existingInfo?.cid === e.cid; - if (!isReplay && !options?.skipFeedFanout) { - batch.push(...buildFeedStatements(db, e, config, existingRecordStrings)); - } - batch.push(...buildFtsStatements(db, e, config)); + if (!isReplay && !options?.skipFeedFanout) { + batch.push(...buildFeedStatements(db, e, config, existingRecordStrings)); } + batch.push(...buildFtsStatements(db, e, config)); } // Build deduplicated count statements — one UPDATE per unique target - if (config) { - batch.push(...buildBatchCountStatements(db, config, countTargets)); - } + batch.push(...buildBatchCountStatements(db, config, countTargets)); await db.batch(batch); @@ -657,7 +652,7 @@ export async function applyEvents( // This fires on both the live and backfill // paths (driven by `options.phase`), carries one deduplicated event per // record, and isolates failures so a throwing sink never blocks ingestion. - const sinks = config?.sinks; + const sinks = config.sinks; if (sinks && sinks.length > 0) { const records: RecordEvent[] = events.map((e) => e.operation === "delete" @@ -674,7 +669,7 @@ export async function applyEvents( } ); const ctx = { phase: options?.phase ?? "live" } as const; - const logger = config?.logger ?? console; + const logger = config.logger ?? console; for (const sink of sinks) { try { await sink.onRecords(records, ctx); diff --git a/packages/contrail/src/core/ingest.ts b/packages/contrail/src/core/ingest.ts new file mode 100644 index 0000000..fb5a6d7 --- /dev/null +++ b/packages/contrail/src/core/ingest.ts @@ -0,0 +1,207 @@ +import { isDid } from "@atcute/lexicons/syntax"; +import type { ContrailConfig, Database, IngestEvent } from "./types"; +import { + getNestedValue, + resolveCollectionKey, +} from "./types"; +import { + projectEvents, + type ExistingRecordInfo, +} from "./db/records"; + +export interface RecordEventInput { + uri?: string; + did: string; + collection: string; + rkey: string; + operation: "create" | "update" | "delete"; + cid?: string | null; + value?: unknown; + timeUs: number; + indexedAt?: number; +} + +/** Normalize a source record into Contrail's canonical mutation shape. */ +export function createIngestEvent(input: RecordEventInput): IngestEvent { + const deleted = input.operation === "delete"; + const serialized = deleted ? null : JSON.stringify(input.value); + return { + uri: + input.uri ?? + `at://${input.did}/${input.collection}/${input.rkey}`, + did: input.did, + collection: input.collection, + rkey: input.rkey, + operation: input.operation, + cid: deleted ? null : (input.cid ?? null), + record: serialized ?? null, + time_us: input.timeUs, + indexed_at: input.indexedAt ?? Date.now() * 1000, + }; +} + +export interface IngestRecordsOptions { + skipReplayDetection?: boolean; + skipFeedFanout?: boolean; + /** Pre-fetched rows, used by immediate synchronization. */ + existing?: Map; + phase?: "live" | "backfill"; + /** Known actors used to filter dependent records without another DB read. */ + knownDids?: ReadonlySet; +} + +export interface IngestDropCounts { + unknownCollection: number; + invalidRecord: number; + recordFilter: number; + unknownSubject: number; +} + +export interface IngestRecordsResult { + accepted: IngestEvent[]; + dropped: IngestDropCounts; +} + +/** + * The single admission and projection path for records from every source. + * + * Jetstream, persistent subscriptions, PDS backfill, refresh, and immediate + * synchronization all produce the same IngestEvent shape and enter here. + * Source connection and checkpoint handling remain outside this function. + */ +export async function ingestRecords( + db: Database, + events: IngestEvent[], + config: ContrailConfig, + options: IngestRecordsOptions = {}, +): Promise { + const accepted: IngestEvent[] = []; + const dropped: IngestDropCounts = { + unknownCollection: 0, + invalidRecord: 0, + recordFilter: 0, + unknownSubject: 0, + }; + const logger = config.logger ?? console; + + for (const event of events) { + const shortName = resolveCollectionKey(config, event.collection); + if (!shortName) { + dropped.unknownCollection++; + logger.warn( + `[ingest] drop unknown collection: ${event.operation} ${event.uri} collection=${event.collection}`, + ); + continue; + } + + if (event.operation !== "delete") { + const record = parseRecord(event.record); + if (!record) { + dropped.invalidRecord++; + logger.warn(`[ingest] drop invalid record: ${event.uri}`); + continue; + } + + const filter = config.collections[shortName]?.recordFilter; + if (filter) { + let keep = false; + try { + keep = filter(record); + } catch (error) { + logger.warn(`[ingest] recordFilter threw for ${event.uri}: ${error}`); + } + if (!keep) { + dropped.recordFilter++; + continue; + } + } + } + + accepted.push(event); + } + + const subjectFiltered = await filterUnknownSubjects( + db, + config, + accepted, + options.knownDids, + dropped, + ); + + if (subjectFiltered.length > 0) { + await projectEvents(db, subjectFiltered, config, options); + } + + return { accepted: subjectFiltered, dropped }; +} + +async function filterUnknownSubjects( + db: Database, + config: ContrailConfig, + events: IngestEvent[], + knownDids: ReadonlySet | undefined, + dropped: IngestDropCounts, +): Promise { + const subjectsByUri = new Map(); + const subjects = new Set(); + + for (const event of events) { + if (event.operation === "delete") continue; + const shortName = resolveCollectionKey(config, event.collection); + const subjectField = shortName + ? config.collections[shortName]?.subjectField + : undefined; + if (!subjectField) continue; + + const record = parseRecord(event.record); + const subject = record ? getNestedValue(record, subjectField) : undefined; + if (typeof subject !== "string" || !isDid(subject)) { + dropped.unknownSubject++; + subjectsByUri.set(event.uri, ""); + continue; + } + subjectsByUri.set(event.uri, subject); + subjects.add(subject); + } + + if (subjectsByUri.size === 0) return events; + + const known = knownDids ? new Set(knownDids) : await loadKnownDids(db, subjects); + return events.filter((event) => { + const subject = subjectsByUri.get(event.uri); + if (subject === undefined) return true; + if (subject !== "" && known.has(subject)) return true; + if (subject !== "") dropped.unknownSubject++; + return false; + }); +} + +async function loadKnownDids( + db: Database, + dids: Set, +): Promise> { + const known = new Set(); + const list = [...dids]; + for (let index = 0; index < list.length; index += 100) { + const chunk = list.slice(index, index + 100); + 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 row of rows.results ?? []) known.add(row.did); + } + return known; +} + +function parseRecord(value: string | null): Record | null { + if (value === null) return null; + try { + const parsed: unknown = JSON.parse(value); + return parsed !== null && typeof parsed === "object" && !Array.isArray(parsed) + ? (parsed as Record) + : null; + } catch { + return null; + } +} diff --git a/packages/contrail/src/core/jetstream.ts b/packages/contrail/src/core/jetstream.ts index 93c9a2b..9913518 100644 --- a/packages/contrail/src/core/jetstream.ts +++ b/packages/contrail/src/core/jetstream.ts @@ -4,14 +4,14 @@ import { getCollectionNsids, getDependentNsids, jetstreamUrlOption, - shortNameForNsid, buildFeedTargetCaps, getFeedMutatingNsids, optimizeEnabled, optimizeIntervalMs, optimizeAnalysisLimit, } from "./types"; -import { initSchema, getLastCursor, saveCursor, applyEvents, sweepFeedItems, getFeedPruneCursor, saveFeedPruneCursor, getMetaNumber, setMeta, optimizeDatabase } from "./db"; +import { initSchema, getLastCursor, saveCursor, sweepFeedItems, getFeedPruneCursor, saveFeedPruneCursor, getMetaNumber, setMeta, optimizeDatabase } from "./db"; +import { createIngestEvent, ingestRecords } from "./ingest"; import { refreshStaleIdentities, applyIdentityEvent } from "./identity"; import { backfillFollowersFromConstellation } from "./constellation"; @@ -187,7 +187,6 @@ export async function ingestEvents( ): Promise<{ events: IngestEvent[]; lastCursor: number | null; - newlyKnownDids: string[]; identityUpdates: Map; }> { const log = getLogger(config); @@ -207,7 +206,6 @@ export async function ingestEvents( let connectCount = 0; const seenUris = new Map(); // uri -> time_us of first occurrence const duplicateUris: string[] = []; - const newlyKnownDids = new Set(); const identityUpdates = new Map(); const subscription = new JetstreamSubscription({ @@ -246,39 +244,12 @@ export async function ingestEvents( const uri = `at://${event.did}/${commit.collection}/${commit.rkey}`; - const short = shortNameForNsid(config, commit.collection); - const collectionCfg = short ? config.collections[short] : undefined; - if (dependentCollections.has(commit.collection) && knownDids) { if (!knownDids.has(event.did)) { filteredUnknownDid++; if (filteredDidSamples.size < 10) filteredDidSamples.add(event.did); return; } - // Subject filter: for collections with subjectField (e.g. follows - // pointing at a `subject` DID), drop records whose subject isn't a - // DID we care about. Trims network-wide social graph to the - // subjects our discoverable users overlap with. - const subjectField = collectionCfg?.subjectField; - if (subjectField && commit.operation !== "delete") { - const subj = (commit.record as Record | undefined)?.[ - subjectField - ]; - if (typeof subj === "string" && !knownDids.has(subj)) { - return; - } - } - } - - if (collectionCfg?.recordFilter && commit.operation !== "delete") { - const rec = commit.record as Record | undefined; - let keep = false; - try { - keep = !!(rec && collectionCfg.recordFilter(rec)); - } catch (err) { - log.warn(`[ingest] recordFilter threw for ${uri}: ${err}`); - } - if (!keep) return; } const prev = seenUris.get(uri); @@ -291,33 +262,22 @@ export async function ingestEvents( seenUris.set(uri, event.time_us); } - const now = Date.now(); - - collected.push({ - uri, - did: event.did, - time_us: event.time_us, - collection: commit.collection, - operation: commit.operation as "create" | "update" | "delete", - rkey: commit.rkey, - cid: commit.operation === "delete" ? null : commit.cid, - record: - commit.operation === "delete" - ? null - : JSON.stringify(commit.record), - indexed_at: now * 1000, - }); + collected.push( + createIngestEvent({ + did: event.did, + timeUs: event.time_us, + collection: commit.collection, + operation: commit.operation, + rkey: commit.rkey, + cid: commit.operation === "delete" ? null : commit.cid, + value: commit.operation === "delete" ? undefined : commit.record, + }), + ); log.log( - `[ingest] keep: ${commit.operation} ${uri} time_us=${event.time_us}` + `[ingest] candidate: ${commit.operation} ${uri} time_us=${event.time_us}` ); - if (knownDids && !dependentCollections.has(commit.collection)) { - if (!knownDids.has(event.did)) { - knownDids.add(event.did); - newlyKnownDids.add(event.did); - } - } } else if (event.kind === "identity") { identityUpdates.set(event.did, event.identity.handle); } @@ -382,7 +342,7 @@ export async function ingestEvents( : 0; log.log( - `[ingest] jetstream loop done. commits_seen=${totalCommits}, filtered=${filteredUnknownDid}, kept=${collected.length}, dupes=${duplicateUris.length}, connects=${connectCount}, first_yielded=${firstYieldedTimeUs ?? "none"}, last_yielded=${lastYieldedTimeUs ?? "none"}, subscription_cursor=${lastCursor ?? "none"}, cursor_gap=${cursorGap ?? "n/a"}us, rolled_back=${rolledBackUs}us` + `[ingest] jetstream loop done. commits_seen=${totalCommits}, filtered=${filteredUnknownDid}, candidates=${collected.length}, dupes=${duplicateUris.length}, connects=${connectCount}, first_yielded=${firstYieldedTimeUs ?? "none"}, last_yielded=${lastYieldedTimeUs ?? "none"}, subscription_cursor=${lastCursor ?? "none"}, cursor_gap=${cursorGap ?? "n/a"}us, rolled_back=${rolledBackUs}us` ); if (cursorGap !== null && cursorGap > 1000) { @@ -409,7 +369,7 @@ export async function ingestEvents( } } - return { events: collected, lastCursor, newlyKnownDids: [...newlyKnownDids], identityUpdates }; + return { events: collected, lastCursor, identityUpdates }; } // Run a full ingest cycle: init schema, load cursor, ingest, apply, save cursor @@ -456,7 +416,7 @@ export async function runIngestCycle( } } - const { events, lastCursor, newlyKnownDids, identityUpdates } = await ingestEvents( + const { events, lastCursor, identityUpdates } = await ingestEvents( config, cursor, timeoutMs, @@ -476,9 +436,22 @@ export async function runIngestCycle( log.log(`[ingest] received 0 events from Jetstream`); } + const accepted: IngestEvent[] = []; for (let i = 0; i < events.length; i += BATCH_SIZE) { const batch = events.slice(i, i + BATCH_SIZE); - await applyEvents(db, batch, config); + const result = await ingestRecords(db, batch, config, { knownDids }); + accepted.push(...result.accepted); + } + + const newlyKnownDids: string[] = []; + if (knownDids) { + const dependent = new Set(dependentCollections); + for (const event of accepted) { + if (!dependent.has(event.collection) && !knownDids.has(event.did)) { + knownDids.add(event.did); + newlyKnownDids.push(event.did); + } + } } // Apply handle changes from #identity events. UPDATE-only, so unknown @@ -495,7 +468,7 @@ export async function runIngestCycle( } // Persist the cursor BEFORE the best-effort enrichment tail below. Records are - // already durably applied (applyEvents) and handle changes recorded, so the + // already durably projected and handle changes recorded, so the // cursor's forward progress is real and must be committed now. The steps that // follow — refreshStaleIdentities especially — make per-DID network calls and // can run long; if the cron isolate is aborted (e.g. a scheduled-invocation @@ -515,7 +488,7 @@ export async function runIngestCycle( // Refresh stale/missing identities for DIDs in this batch (best-effort; runs // after the cursor save so its network latency can't strand forward progress). - const uniqueDids = [...new Set(events.map((e) => e.did))]; + const uniqueDids = [...new Set(accepted.map((e) => e.did))]; if (uniqueDids.length > 0) { try { await refreshStaleIdentities(db, uniqueDids, config); @@ -556,7 +529,7 @@ export async function runIngestCycle( // for many intervals (one slice per interval) as a per-slice clock would. if (config.feeds) { const feedMutatingNsids = getFeedMutatingNsids(config); - const feedTouched = events.some((e) => feedMutatingNsids.has(e.collection)); + const feedTouched = accepted.some((e) => feedMutatingNsids.has(e.collection)); await runGatedFeedPrune(db, config, feedTouched); } @@ -564,5 +537,5 @@ export async function runIngestCycle( // config.maintenance.optimize is set). await maybeOptimize(db, config, log); - log.log(`[ingest] cycle complete. stored=${events.length}`); + log.log(`[ingest] cycle complete. stored=${accepted.length}`); } diff --git a/packages/contrail/src/core/persistent.ts b/packages/contrail/src/core/persistent.ts index bfd69dc..05fa072 100644 --- a/packages/contrail/src/core/persistent.ts +++ b/packages/contrail/src/core/persistent.ts @@ -7,9 +7,9 @@ import { getFeedMutatingNsids, jetstreamUrlOption, resolveConfig, - shortNameForNsid, } from "./types"; -import { initSchema, getLastCursor, saveCursor, applyEvents } from "./db"; +import { initSchema, getLastCursor, saveCursor } from "./db"; +import { createIngestEvent, ingestRecords } from "./ingest"; import { refreshStaleIdentities, applyIdentityEvent } from "./identity"; import { backfillFollowersFromConstellation } from "./constellation"; import { @@ -43,7 +43,7 @@ export async function runPersistent( config: ContrailConfig, options?: PersistentIngestOptions, ): Promise { - // Internals (applyEvents, count updates, query planning) read `_resolved` + // Ingestion, count updates, and query planning read `_resolved` // and silently skip features when it's missing. The Contrail class resolves // in its constructor; callers using this raw export must also get a resolved // config, so do it defensively here. resolveConfig is idempotent. @@ -157,14 +157,28 @@ async function streamAndFlush( try { if (buffer.length > 0) { const batch = buffer.splice(0); - await applyEvents(db, batch, config); + const { accepted } = await ingestRecords(db, batch, config, { + knownDids, + }); + + if (knownDids) { + for (const event of accepted) { + if ( + !dependentCollections.has(event.collection) && + !knownDids.has(event.did) + ) { + knownDids.add(event.did); + opts.newlyKnownDids?.add(event.did); + } + } + } // A feed can only go over cap right after a feed-mutating record is // applied, so remember whether this batch had one. The sweep below uses // it to prune promptly (see the cron path in jetstream.ts). if (config.feeds) { const feedMutatingNsids = getFeedMutatingNsids(config); - if (batch.some((e) => feedMutatingNsids.has(e.collection))) { + if (accepted.some((e) => feedMutatingNsids.has(e.collection))) { state.feedDirty = true; } } @@ -172,7 +186,7 @@ async function streamAndFlush( const lastTimeUs = Math.max(...batch.map((e) => e.time_us)); await saveCursor(db, lastTimeUs); - const uniqueDids = [...new Set(batch.map((e) => e.did))]; + const uniqueDids = [...new Set(accepted.map((e) => e.did))]; if (uniqueDids.length > 0) { try { await refreshStaleIdentities(db, uniqueDids, config); @@ -197,7 +211,9 @@ async function streamAndFlush( // Opt-in planner-stat maintenance (gated + persisted cadence). await maybeOptimize(db, config, log); - log.log(`Flushed ${batch.length} events. Cursor: ${lastTimeUs}`); + log.log( + `Flushed ${accepted.length}/${batch.length} records. Cursor: ${lastTimeUs}`, + ); } // Bounded, cursored feed prune (see sweepFeedItems / runFeedPruneSlice). @@ -276,53 +292,22 @@ async function streamAndFlush( if (event.kind === "commit") { const { commit } = event; - const short = shortNameForNsid(config, commit.collection); - const collectionCfg = short ? config.collections[short] : undefined; - if (dependentCollections.has(commit.collection) && knownDids) { if (!knownDids.has(event.did)) continue; - // Subject filter: skip records whose subject DID isn't known. - const subjectField = collectionCfg?.subjectField; - if (subjectField && commit.operation !== "delete") { - const subj = (commit.record as Record | undefined)?.[ - subjectField - ]; - if (typeof subj === "string" && !knownDids.has(subj)) continue; - } } - if (collectionCfg?.recordFilter && commit.operation !== "delete") { - const rec = commit.record as Record | undefined; - let keep = false; - try { - keep = !!(rec && collectionCfg.recordFilter(rec)); - } catch (err) { - log.warn(`recordFilter threw for ${commit.collection}/${commit.rkey}: ${err}`); - } - if (!keep) continue; - } + buffer.push( + createIngestEvent({ + did: event.did, + timeUs: event.time_us, + collection: commit.collection, + operation: commit.operation, + rkey: commit.rkey, + cid: commit.operation === "delete" ? null : commit.cid, + value: commit.operation === "delete" ? undefined : commit.record, + }), + ); - const now = Date.now(); - const uri = `at://${event.did}/${commit.collection}/${commit.rkey}`; - - buffer.push({ - uri, - did: event.did, - time_us: event.time_us, - collection: commit.collection, - operation: commit.operation as "create" | "update" | "delete", - rkey: commit.rkey, - cid: commit.operation === "delete" ? null : commit.cid, - record: commit.operation === "delete" ? null : JSON.stringify(commit.record), - indexed_at: now * 1000, - }); - - if (knownDids && !dependentCollections.has(commit.collection)) { - if (!knownDids.has(event.did)) { - knownDids.add(event.did); - opts.newlyKnownDids?.add(event.did); - } - } } else if (event.kind === "identity") { try { await applyIdentityEvent(db, event.did, event.identity.handle); diff --git a/packages/contrail/src/core/refresh.ts b/packages/contrail/src/core/refresh.ts index 2c8a343..c14ec37 100644 --- a/packages/contrail/src/core/refresh.ts +++ b/packages/contrail/src/core/refresh.ts @@ -25,7 +25,8 @@ import { isDid, isNsid } from "@atcute/lexicons/syntax"; import type { Client } from "@atcute/client"; import type { ContrailConfig, Database, IngestEvent } from "./types.js"; -import { applyEvents, lookupExistingRecords } from "./db/records.js"; +import { lookupExistingRecords } from "./db/records.js"; +import { createIngestEvent, ingestRecords } from "./ingest.js"; import { getClient } from "./client.js"; const PAGE_SIZE = 100; @@ -190,17 +191,19 @@ export async function refresh( if (pageRecords.length === 0) break; const now = Date.now(); - const events: IngestEvent[] = pageRecords.map((r) => ({ - uri: r.uri, - did, - collection: nsid, - rkey: r.uri.split("/").pop()!, - operation: "create" as const, - cid: r.cid, - record: JSON.stringify(r.value), - time_us: now * 1000, - indexed_at: now * 1000, - })); + 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, @@ -234,7 +237,10 @@ export async function refresh( // might genuinely have a new CID; we just don't count them as a // miss-signal. Skip feed fanout since this is a catch-up, not a // user-visible write. - await applyEvents(db, events, config, { skipFeedFanout: true }); + await ingestRecords(db, events, config, { + existing, + skipFeedFanout: true, + }); recordsScanned += events.length; cursor = nextCursor; diff --git a/packages/contrail/src/core/router/notify.ts b/packages/contrail/src/core/router/notify.ts index f6c51c3..c2b39f0 100644 --- a/packages/contrail/src/core/router/notify.ts +++ b/packages/contrail/src/core/router/notify.ts @@ -1,7 +1,8 @@ import type { Hono } from "hono"; import type { Database, ContrailConfig, IngestEvent } from "../types"; import { shortNameForNsid, getFeedMutatingNsids } from "../types"; -import { applyEvents, lookupExistingRecords } from "../db/records"; +import { lookupExistingRecords } from "../db/records"; +import { createIngestEvent, ingestRecords } from "../ingest"; import { runGatedFeedPrune } from "../jetstream"; import { getPDS } from "../client"; import type { Did } from "@atcute/lexicons"; @@ -105,39 +106,40 @@ export async function processNotifyUris( continue; } - events.push({ - uri, - did: parsed.did, - collection: parsed.collection, - rkey: parsed.rkey, - operation: existingInfo ? "update" : "create", - cid: result.cid, - record: JSON.stringify(result.value), - time_us: now, - indexed_at: now, - }); + events.push( + createIngestEvent({ + uri, + did: parsed.did, + collection: parsed.collection, + rkey: parsed.rkey, + operation: existingInfo ? "update" : "create", + cid: result.cid, + value: result.value, + timeUs: now, + indexedAt: now, + }), + ); } else if (existingInfo) { // Record gone from PDS but exists locally — delete it. - events.push({ - uri, - did: parsed.did, - collection: parsed.collection, - rkey: parsed.rkey, - operation: "delete", - cid: null, - record: existingInfo.record, - time_us: now, - indexed_at: now, - }); + events.push( + createIngestEvent({ + uri, + did: parsed.did, + collection: parsed.collection, + rkey: parsed.rkey, + operation: "delete", + timeUs: now, + indexedAt: now, + }), + ); } } - if (events.length > 0) { - // Pass pre-fetched existing records so applyEvents skips re-querying - await applyEvents(db, events, config, { existing }); - } + const appliedEvents = events.length > 0 + ? (await ingestRecords(db, events, config, { existing })).accepted + : []; - // applyEvents fans these records into feed_items exactly like the cron and + // The shared ingest path fans these records into feed_items exactly like the cron and // persistent ingest paths, so prune here too — otherwise a notify-only // deployment (no jetstream loop) would never sweep. Run the recovery-aware // gate on every call, not only when records changed: a notify-only deployment @@ -146,13 +148,17 @@ export async function processNotifyUris( // true only when this call actually applied a feed-mutating record. if (config.feeds) { const feedMutatingNsids = getFeedMutatingNsids(config); - const feedTouched = events.some((e) => feedMutatingNsids.has(e.collection)); + const feedTouched = appliedEvents.some((e) => + feedMutatingNsids.has(e.collection), + ); await runGatedFeedPrune(db, config, feedTouched); } return { - indexed: events.filter((e) => e.operation === "create" || e.operation === "update").length, - deleted: events.filter((e) => e.operation === "delete").length, + indexed: appliedEvents.filter( + (event) => event.operation === "create" || event.operation === "update", + ).length, + deleted: appliedEvents.filter((event) => event.operation === "delete").length, errors: errors.length > 0 ? errors : undefined, }; } diff --git a/packages/contrail/src/core/types.ts b/packages/contrail/src/core/types.ts index b131715..ef49ec1 100644 --- a/packages/contrail/src/core/types.ts +++ b/packages/contrail/src/core/types.ts @@ -247,7 +247,7 @@ export interface ContrailConfig { * Set to `true` for open access, or a string to require `Authorization: Bearer `. */ notify?: boolean | string; /** Write-only, post-commit observers of applied records — derived indexes, - * audit logs, webhook fan-outs. Each fires after every `applyEvents()` commit + * audit logs, webhook fan-outs. Each fires after every record-ingest commit * on both live and backfill paths, with failures isolated so a throwing * sink never blocks ingestion. */ sinks?: import("./sinks/types").Sink[]; diff --git a/packages/contrail/src/index.ts b/packages/contrail/src/index.ts index 464dfb5..a4b9c15 100644 --- a/packages/contrail/src/index.ts +++ b/packages/contrail/src/index.ts @@ -10,6 +10,7 @@ export * from "./core/client"; export * from "./core/sinks/types"; // Ingestion and maintenance. +export * from "./core/ingest"; export * from "./core/jetstream"; export * from "./core/persistent"; export * from "./core/backfill"; @@ -19,7 +20,23 @@ export * from "./core/constellation"; // Database. export * from "./core/db/schema"; -export * from "./core/db/records"; +export { + getFeedPruneCursor, + getLastCursor, + lookupExistingRecords, + pruneActorFeed, + pruneFeedItems, + queryRecords, + saveCursor, + saveFeedPruneCursor, + sweepFeedItems, +} from "./core/db/records"; +export type { + ExistingRecordInfo, + FeedSweepResult, + QueryOptions, + SortOption, +} from "./core/db/records"; export * from "./core/db/meta"; export * from "./core/db/optimize"; diff --git a/packages/contrail/tests/feed-prune-guardrail.test.ts b/packages/contrail/tests/feed-prune-guardrail.test.ts index 2eaf7b7..cce9b5c 100644 --- a/packages/contrail/tests/feed-prune-guardrail.test.ts +++ b/packages/contrail/tests/feed-prune-guardrail.test.ts @@ -9,7 +9,7 @@ import type { import { resolveConfig } from "../src/index"; import { initSchema } from "../src/index"; import { - applyEvents, + ingestRecords, sweepFeedItems, pruneActorFeed, pruneFeedItems, @@ -181,7 +181,7 @@ describe("feed maintenance stays index-bounded", () => { indexed_at: 5000, operation: "create", }; - await applyEvents(db, [event], CONFIG); + await ingestRecords(db, [event], CONFIG); // The fan-out INSERT must have been issued and it must hit the table. const fanout = sqls.find((s) => /INSERT.*feed_items/is.test(s)); @@ -195,7 +195,7 @@ describe("feed maintenance stays index-bounded", () => { await makeFollowRow(real, "did:plc:follower", "did:plc:author"); const { db, sqls } = recordingDb(real); - await applyEvents( + await ingestRecords( db, [ { diff --git a/packages/contrail/tests/get-record-handle.test.ts b/packages/contrail/tests/get-record-handle.test.ts index 1aac656..8bbc977 100644 --- a/packages/contrail/tests/get-record-handle.test.ts +++ b/packages/contrail/tests/get-record-handle.test.ts @@ -4,7 +4,7 @@ import { describe, it, expect, vi, beforeEach, afterEach } from "vitest"; import { Contrail } from "../src/contrail"; import { createSqliteDatabase } from "../src/adapters/sqlite"; -import { applyEvents } from "../src/index"; +import { ingestRecords } from "../src/index"; import { __resetPdsCachesForTests } from "../src/index"; import type { Database, IngestEvent } from "../src/index"; @@ -37,7 +37,7 @@ async function setup() { db, }); await contrail.init(); - await applyEvents(db, [ev()], contrail.config); + await ingestRecords(db, [ev()], contrail.config); return { db, contrail }; } diff --git a/packages/contrail/tests/helpers.ts b/packages/contrail/tests/helpers.ts index ea81dd7..e4ba767 100644 --- a/packages/contrail/tests/helpers.ts +++ b/packages/contrail/tests/helpers.ts @@ -2,7 +2,7 @@ import { createSqliteDatabase } from "../src/adapters/sqlite"; import type { Database, IngestEvent, ResolvedContrailConfig } from "../src/index"; import { resolveConfig } from "../src/index"; import { initSchema } from "../src/index"; -import { applyEvents as coreApplyEvents, type ExistingRecordInfo } from "../src/index"; +import { ingestRecords as coreIngestRecords, type ExistingRecordInfo } from "../src/index"; export function createTestDb(): Database { return createSqliteDatabase(":memory:"); @@ -48,8 +48,8 @@ export async function createTestDbWithSchema(): Promise { return db; } -/** Apply events with TEST_CONFIG baked in — avoids each test having to pass it. */ -export function applyEvents( +/** Ingest records with TEST_CONFIG baked in. */ +export function ingestRecords( db: Database, events: IngestEvent[], options?: { @@ -58,7 +58,7 @@ export function applyEvents( existing?: Map; } ): Promise { - return coreApplyEvents(db, events, TEST_CONFIG, options); + return coreIngestRecords(db, events, TEST_CONFIG, options).then(() => undefined); } export function makeEvent(overrides: Partial<{ diff --git a/packages/contrail/tests/hydrate.test.ts b/packages/contrail/tests/hydrate.test.ts index fb41a1a..b225213 100644 --- a/packages/contrail/tests/hydrate.test.ts +++ b/packages/contrail/tests/hydrate.test.ts @@ -3,7 +3,7 @@ import type { Database, RecordRow, RelationConfig, ReferenceConfig } from "../sr import { recordsTableName } from "../src/index"; import { parseHydrateParams, resolveHydrates, resolveReferences } from "../src/index"; import { createTestDbWithSchema, makeEvent, TEST_CONFIG } from "./helpers"; -import { applyEvents } from "./helpers"; +import { ingestRecords } from "./helpers"; describe("parseHydrateParams", () => { const relations: Record = { @@ -83,11 +83,11 @@ describe("resolveHydrates", () => { const eventUri = "at://did:plc:test/community.lexicon.calendar.event/evt1"; // Insert event - await applyEvents(db, [makeEvent({ uri: eventUri, rkey: "evt1", time_us: 1000 })]); + await ingestRecords(db, [makeEvent({ uri: eventUri, rkey: "evt1", time_us: 1000 })]); // Insert RSVPs for (let i = 0; i < 3; i++) { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: `at://did:plc:user${i}/community.lexicon.calendar.rsvp/r${i}`, did: `did:plc:user${i}`, @@ -114,10 +114,10 @@ describe("resolveHydrates", () => { it("respects hydrate limit", async () => { const eventUri = "at://did:plc:test/community.lexicon.calendar.event/evt1"; - await applyEvents(db, [makeEvent({ uri: eventUri, rkey: "evt1", time_us: 1000 })]); + await ingestRecords(db, [makeEvent({ uri: eventUri, rkey: "evt1", time_us: 1000 })]); for (let i = 0; i < 5; i++) { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: `at://did:plc:user${i}/community.lexicon.calendar.rsvp/r${i}`, did: `did:plc:user${i}`, @@ -145,10 +145,10 @@ describe("resolveHydrates", () => { const eventUri = `at://${did}/community.lexicon.calendar.event/evt1`; // Insert event owned by this DID - await applyEvents(db, [makeEvent({ uri: eventUri, did, rkey: "evt1", time_us: 1000 })]); + await ingestRecords(db, [makeEvent({ uri: eventUri, did, rkey: "evt1", time_us: 1000 })]); // Insert a related record whose "author" field contains the parent's DID - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: `at://did:plc:other/community.lexicon.calendar.rsvp/r1`, did: "did:plc:other", @@ -179,13 +179,13 @@ describe("resolveHydrates", () => { const eventUri = "at://did:plc:test/community.lexicon.calendar.event/evt1"; // Insert event - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: eventUri, rkey: "evt1", record: { name: "My Event" }, time_us: 1000 }), ]); // Insert RSVP pointing at the event const rsvpUri = "at://did:plc:user1/community.lexicon.calendar.rsvp/r1"; - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: rsvpUri, did: "did:plc:user1", @@ -215,7 +215,7 @@ describe("resolveHydrates", () => { const eventUri2 = "at://did:plc:test/community.lexicon.calendar.event/evt2"; // Insert two events - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: eventUri1, rkey: "evt1", record: { name: "Event 1" }, time_us: 1000 }), makeEvent({ uri: eventUri2, rkey: "evt2", record: { name: "Event 2" }, time_us: 1001 }), ]); @@ -223,7 +223,7 @@ describe("resolveHydrates", () => { // Insert RSVPs pointing at different events const rsvpUri1 = "at://did:plc:user1/community.lexicon.calendar.rsvp/r1"; const rsvpUri2 = "at://did:plc:user2/community.lexicon.calendar.rsvp/r2"; - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: rsvpUri1, did: "did:plc:user1", @@ -259,10 +259,10 @@ describe("resolveHydrates", () => { it("groups into 'other' when groupBy value is null", async () => { const eventUri = "at://did:plc:test/community.lexicon.calendar.event/evt1"; - await applyEvents(db, [makeEvent({ uri: eventUri, rkey: "evt1", time_us: 1000 })]); + await ingestRecords(db, [makeEvent({ uri: eventUri, rkey: "evt1", time_us: 1000 })]); // Insert RSVP without a status field — groupBy "status" should fall back to "other" - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:user1/community.lexicon.calendar.rsvp/r1", did: "did:plc:user1", diff --git a/packages/contrail/tests/ingest.test.ts b/packages/contrail/tests/ingest.test.ts new file mode 100644 index 0000000..9fa063d --- /dev/null +++ b/packages/contrail/tests/ingest.test.ts @@ -0,0 +1,129 @@ +import { beforeEach, describe, expect, it } from "vitest"; +import { createSqliteDatabase } from "../src/adapters/sqlite"; +import { + createIngestEvent, + ingestRecords, + initSchema, + queryRecords, + resolveConfig, + type Database, +} from "../src/index"; + +const logger = { log() {}, warn() {}, error() {} }; + +const config = resolveConfig({ + namespace: "com.example", + logger, + collections: { + event: { + collection: "com.example.event", + recordFilter: (record) => record.keep === true, + }, + follow: { + collection: "app.bsky.graph.follow", + discover: false, + subjectField: "subject", + }, + }, +}); + +let db: Database; + +beforeEach(async () => { + db = createSqliteDatabase(":memory:"); + await initSchema(db, config); +}); + +describe("ingestRecords", () => { + it("is the shared recordFilter admission path", async () => { + const result = await ingestRecords( + db, + [ + createIngestEvent({ + did: "did:plc:alice", + collection: "com.example.event", + rkey: "kept", + operation: "create", + cid: "cid-kept", + value: { keep: true }, + timeUs: 1, + }), + createIngestEvent({ + did: "did:plc:alice", + collection: "com.example.event", + rkey: "dropped", + operation: "create", + cid: "cid-dropped", + value: { keep: false }, + timeUs: 2, + }), + ], + config, + ); + + expect(result.accepted).toHaveLength(1); + expect(result.dropped.recordFilter).toBe(1); + const stored = await queryRecords(db, config, { collection: "event" }); + expect(stored.records.map((record) => record.rkey)).toEqual(["kept"]); + }); + + it("drops malformed records and untracked collections before projection", async () => { + const malformed = createIngestEvent({ + did: "did:plc:alice", + collection: "com.example.event", + rkey: "bad", + operation: "create", + cid: "cid-bad", + value: { keep: true }, + timeUs: 1, + }); + malformed.record = "not json"; + + const result = await ingestRecords( + db, + [ + malformed, + createIngestEvent({ + did: "did:plc:alice", + collection: "com.example.unknown", + rkey: "unknown", + operation: "create", + cid: "cid-unknown", + value: {}, + timeUs: 2, + }), + ], + config, + ); + + expect(result.accepted).toHaveLength(0); + expect(result.dropped.invalidRecord).toBe(1); + expect(result.dropped.unknownCollection).toBe(1); + }); + + it("applies dependent subject filtering in the same admission path", async () => { + await db + .prepare( + "INSERT INTO identities (did, handle, pds, resolved_at) VALUES (?, ?, ?, ?)", + ) + .bind("did:plc:known", null, null, 1) + .run(); + + const records = ["did:plc:known", "did:plc:unknown"].map( + (subject, index) => + createIngestEvent({ + did: "did:plc:alice", + collection: "app.bsky.graph.follow", + rkey: `f${index}`, + operation: "create", + cid: `cid-${index}`, + value: { subject }, + timeUs: index + 1, + }), + ); + + const result = await ingestRecords(db, records, config); + expect(result.accepted.map((record) => record.rkey)).toEqual(["f0"]); + expect(result.dropped.unknownSubject).toBe(1); + }); +}); diff --git a/packages/contrail/tests/labels-router.test.ts b/packages/contrail/tests/labels-router.test.ts index 863b77c..3ad7ef4 100644 --- a/packages/contrail/tests/labels-router.test.ts +++ b/packages/contrail/tests/labels-router.test.ts @@ -6,7 +6,7 @@ import { describe, it, expect } from "vitest"; import { Contrail } from "../src/contrail"; import { createSqliteDatabase } from "../src/adapters/sqlite"; -import { applyEvents } from "../src/index"; +import { ingestRecords } from "../src/index"; import { applyLabels } from "../src/index"; import type { IngestEvent } from "../src/index"; @@ -44,7 +44,7 @@ async function setup() { db, }); await contrail.init(); - await applyEvents(db, [ev()], contrail.config); + await ingestRecords(db, [ev()], contrail.config); await applyLabels(db, [ { src: SRC_A, uri: URI, val: "spam", cts: new Date().toISOString() }, { src: SRC_B, uri: URI, val: "porn", cts: new Date().toISOString() }, @@ -117,7 +117,7 @@ describe("labels router integration", () => { db, }); await contrail.init(); - await applyEvents(db, [ev()], contrail.config); + await ingestRecords(db, [ev()], contrail.config); const app = contrail.app(); const res = await app.fetch(new Request(`http://localhost/xrpc/ex.event.listRecords`)); diff --git a/packages/contrail/tests/notify.test.ts b/packages/contrail/tests/notify.test.ts index 671e04d..4a3178a 100644 --- a/packages/contrail/tests/notify.test.ts +++ b/packages/contrail/tests/notify.test.ts @@ -1,7 +1,7 @@ import { describe, it, expect, beforeEach, vi, afterEach } from "vitest"; import type { Database } from "../src/index"; import { resolveConfig } from "../src/index"; -import { applyEvents, createTestDb, createTestDbWithSchema, makeEvent, TEST_CONFIG } from "./helpers"; +import { ingestRecords, createTestDb, createTestDbWithSchema, makeEvent, TEST_CONFIG } from "./helpers"; import { initSchema } from "../src/index"; import { parseAtUri } from "../src/index"; import { createApp } from "../src/index"; @@ -175,7 +175,7 @@ describe("POST notifyOfUpdate", () => { const uri = `at://${did}/community.lexicon.calendar.event/evt1`; // Pre-populate a record - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri, did, rkey: "evt1", record: { name: "Old" } }), ], TEST_CONFIG); @@ -232,7 +232,7 @@ describe("POST notifyOfUpdate", () => { const uri = `at://${did}/community.lexicon.calendar.event/evt1`; // Insert original - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri, did, rkey: "evt1", record: { name: "V1" } }), ]); @@ -265,7 +265,7 @@ describe("POST notifyOfUpdate", () => { const rsvpUri = `at://${did}/community.lexicon.calendar.rsvp/r1`; // Insert parent event - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: eventUri, did, rkey: "evt1", record: { name: "Event" } }), ]); @@ -300,10 +300,10 @@ describe("POST notifyOfUpdate", () => { const rsvpRecord = { subject: { uri: eventUri }, status: "community.lexicon.calendar.rsvp#going" }; // Insert parent event and RSVP via normal ingestion - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: eventUri, did, rkey: "evt1", record: { name: "Event" } }), ]); - await applyEvents( + await ingestRecords( db, [ makeEvent({ @@ -354,10 +354,10 @@ describe("POST notifyOfUpdate", () => { const rsvpUri = `at://${did}/community.lexicon.calendar.rsvp/r1`; // Insert parent event and RSVP via normal ingestion - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: eventUri, did, rkey: "evt1", record: { name: "Event" } }), ]); - await applyEvents( + await ingestRecords( db, [ makeEvent({ @@ -476,10 +476,10 @@ describe("POST notifyOfUpdate", () => { const rsvpUri = `at://${did}/community.lexicon.calendar.rsvp/r1`; // Insert event + RSVP - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: eventUri, did, rkey: "evt1", record: { name: "Event" } }), ]); - await applyEvents( + await ingestRecords( db, [ makeEvent({ diff --git a/packages/contrail/tests/postgres-e2e.test.ts b/packages/contrail/tests/postgres-e2e.test.ts index 967a873..5510f02 100644 --- a/packages/contrail/tests/postgres-e2e.test.ts +++ b/packages/contrail/tests/postgres-e2e.test.ts @@ -12,7 +12,7 @@ import pg from "pg"; import { createPostgresDatabase } from "../src/adapters/postgres"; import { initSchema } from "../src/index"; import { - applyEvents, + ingestRecords, queryRecords, getLastCursor, saveCursor, @@ -132,7 +132,7 @@ if (!PG_URL) { describe("ingestion", () => { it("inserts create events", async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ record: { name: "Event 1", mode: "online" } }), ]); const result = await queryRecords(db, TEST_CONFIG, { @@ -144,10 +144,10 @@ if (!PG_URL) { it("upserts on conflict", async () => { const uri = "at://did:plc:test/community.lexicon.calendar.event/1"; - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri, rkey: "1", cid: "cid1", record: { name: "V1", mode: "online" } }), ]); - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri, rkey: "1", cid: "cid2", record: { name: "V2", mode: "online" } }), ]); const result = await queryRecords(db, TEST_CONFIG, { @@ -159,10 +159,10 @@ if (!PG_URL) { it("deletes events", async () => { const uri = "at://did:plc:test/community.lexicon.calendar.event/1"; - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri, rkey: "1", record: { name: "Gone", mode: "online" } }), ]); - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri, rkey: "1", operation: "delete", record: null, cid: null }), ]); const result = await queryRecords(db, TEST_CONFIG, { @@ -172,7 +172,7 @@ if (!PG_URL) { }); it("does nothing for empty events", async () => { - await applyEvents(db, []); + await ingestRecords(db, []); const cursor = await getLastCursor(db); expect(cursor).toBeNull(); }); @@ -184,13 +184,13 @@ if (!PG_URL) { const eventUri = "at://did:plc:host/community.lexicon.calendar.event/evt1"; beforeEach(async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: eventUri, rkey: "evt1", record: { name: "Counted Event", mode: "online" } }), ]); }); it("increments counts on RSVP create", async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:u1/community.lexicon.calendar.rsvp/r1", did: "did:plc:u1", @@ -227,7 +227,7 @@ if (!PG_URL) { it("decrements counts on RSVP delete", async () => { const rsvpUri = "at://did:plc:u1/community.lexicon.calendar.rsvp/r1"; - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: rsvpUri, did: "did:plc:u1", @@ -245,7 +245,7 @@ if (!PG_URL) { expect(result.records[0].counts?.["community.lexicon.calendar.rsvp"]).toBe(1); // Delete the RSVP - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: rsvpUri, did: "did:plc:u1", @@ -270,7 +270,7 @@ if (!PG_URL) { describe("queries", () => { beforeEach(async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://a/community.lexicon.calendar.event/1", rkey: "1", did: "did:plc:alice", record: { name: "Alpha", mode: "online", startsAt: "2026-04-01T10:00:00Z" }, time_us: 1000 }), makeEvent({ uri: "at://a/community.lexicon.calendar.event/2", rkey: "2", did: "did:plc:bob", record: { name: "Beta", mode: "in-person", startsAt: "2026-05-01T10:00:00Z" }, time_us: 2000 }), makeEvent({ uri: "at://a/community.lexicon.calendar.event/3", rkey: "3", did: "did:plc:alice", record: { name: "Gamma", mode: "online", startsAt: "2026-06-01T10:00:00Z" }, time_us: 3000 }), @@ -374,7 +374,7 @@ if (!PG_URL) { describe("full-text search", () => { beforeEach(async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://a/community.lexicon.calendar.event/1", rkey: "1", record: { name: "Rust Meetup", description: "Systems programming", mode: "online" }, time_us: 1000 }), makeEvent({ uri: "at://a/community.lexicon.calendar.event/2", rkey: "2", record: { name: "TypeScript Workshop", description: "Learn web dev", mode: "online" }, time_us: 2000 }), makeEvent({ uri: "at://a/community.lexicon.calendar.event/3", rkey: "3", record: { name: "Go Conference", description: "Concurrency and systems", mode: "in-person" }, time_us: 3000 }), @@ -425,11 +425,11 @@ if (!PG_URL) { const eventUri2 = "at://did:plc:host/community.lexicon.calendar.event/evt2"; beforeEach(async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: eventUri1, rkey: "evt1", record: { name: "Event 1", mode: "online" }, time_us: 1000 }), makeEvent({ uri: eventUri2, rkey: "evt2", record: { name: "Event 2", mode: "online" }, time_us: 2000 }), ]); - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:u1/community.lexicon.calendar.rsvp/r1", did: "did:plc:u1", diff --git a/packages/contrail/tests/postgres.test.ts b/packages/contrail/tests/postgres.test.ts index cacfde1..7bc29cf 100644 --- a/packages/contrail/tests/postgres.test.ts +++ b/packages/contrail/tests/postgres.test.ts @@ -2,7 +2,7 @@ import { describe, it, expect, beforeAll, afterAll, beforeEach } from "vitest"; import pg from "pg"; import { createPostgresDatabase } from "../src/adapters/postgres"; import { initSchema } from "../src/index"; -import { applyEvents, queryRecords, getLastCursor, saveCursor } from "../src/index"; +import { ingestRecords, queryRecords, getLastCursor, saveCursor } from "../src/index"; import { resolveConfig } from "../src/index"; import { makeEvent } from "./helpers"; @@ -79,7 +79,7 @@ if (!PG_URL) { }); it("inserts and queries records", async () => { - await applyEvents(db, [makeEvent({ record: { name: "PG Test", mode: "online" } })]); + await ingestRecords(db, [makeEvent({ record: { name: "PG Test", mode: "online" } })]); const result = await queryRecords(db, TEST_CONFIG, { collection: "community.lexicon.calendar.event", }); @@ -87,7 +87,7 @@ if (!PG_URL) { }); it("filters by json field", async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://a/community.lexicon.calendar.event/1", rkey: "1", record: { name: "A", mode: "online" }, time_us: 2000 }), makeEvent({ uri: "at://a/community.lexicon.calendar.event/2", rkey: "2", record: { name: "B", mode: "in-person" }, time_us: 1000 }), ]); @@ -100,8 +100,8 @@ if (!PG_URL) { it("counts relations", async () => { const eventUri = "at://did:plc:test/community.lexicon.calendar.event/evt1"; - await applyEvents(db, [makeEvent({ uri: eventUri, rkey: "evt1" })]); - await applyEvents(db, [ + await ingestRecords(db, [makeEvent({ uri: eventUri, rkey: "evt1" })]); + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:user1/community.lexicon.calendar.rsvp/r1", did: "did:plc:user1", @@ -119,7 +119,7 @@ if (!PG_URL) { }); it("full-text search works via tsvector", async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://a/community.lexicon.calendar.event/1", rkey: "1", diff --git a/packages/contrail/tests/records.test.ts b/packages/contrail/tests/records.test.ts index 16e86ac..3ad9570 100644 --- a/packages/contrail/tests/records.test.ts +++ b/packages/contrail/tests/records.test.ts @@ -1,6 +1,6 @@ import { describe, it, expect, beforeEach } from "vitest"; import type { Database } from "../src/index"; -import { applyEvents, createTestDbWithSchema, makeEvent, TEST_CONFIG } from "./helpers"; +import { ingestRecords, createTestDbWithSchema, makeEvent, TEST_CONFIG } from "./helpers"; import { queryRecords, getLastCursor, saveCursor } from "../src/index"; let db: Database; @@ -26,9 +26,9 @@ describe("cursor", () => { }); }); -describe("applyEvents", () => { +describe("ingestRecords", () => { it("does nothing for empty events", async () => { - await applyEvents(db, []); + await ingestRecords(db, []); const result = await queryRecords(db, TEST_CONFIG, { collection: "community.lexicon.calendar.event", }); @@ -36,7 +36,7 @@ describe("applyEvents", () => { }); it("inserts create events", async () => { - await applyEvents(db, [makeEvent()]); + await ingestRecords(db, [makeEvent()]); const result = await queryRecords(db, TEST_CONFIG, { collection: "community.lexicon.calendar.event", }); @@ -46,10 +46,10 @@ describe("applyEvents", () => { }); it("upserts on conflict", async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ record: { name: "V1" }, time_us: 100 }), ]); - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ record: { name: "V2" }, time_us: 200, operation: "update" }), ]); const result = await queryRecords(db, TEST_CONFIG, { @@ -61,8 +61,8 @@ describe("applyEvents", () => { }); it("deletes events", async () => { - await applyEvents(db, [makeEvent()]); - await applyEvents(db, [makeEvent({ operation: "delete" })]); + await ingestRecords(db, [makeEvent()]); + await ingestRecords(db, [makeEvent({ operation: "delete" })]); const result = await queryRecords(db, TEST_CONFIG, { collection: "community.lexicon.calendar.event", }); @@ -73,10 +73,10 @@ describe("applyEvents", () => { const eventUri = "at://did:plc:test/community.lexicon.calendar.event/evt1"; // Insert parent event - await applyEvents(db, [makeEvent({ uri: eventUri, rkey: "evt1" })]); + await ingestRecords(db, [makeEvent({ uri: eventUri, rkey: "evt1" })]); // Insert RSVP pointing at the event - await applyEvents( + await ingestRecords( db, [ makeEvent({ @@ -102,10 +102,10 @@ describe("applyEvents", () => { it("decrements counts on delete", async () => { const eventUri = "at://did:plc:test/community.lexicon.calendar.event/evt1"; - await applyEvents(db, [makeEvent({ uri: eventUri, rkey: "evt1" })]); + await ingestRecords(db, [makeEvent({ uri: eventUri, rkey: "evt1" })]); const rsvpRecord = { subject: { uri: eventUri }, status: "community.lexicon.calendar.rsvp#going" }; - await applyEvents( + await ingestRecords( db, [ makeEvent({ @@ -121,7 +121,7 @@ describe("applyEvents", () => { ); // Delete the RSVP - await applyEvents( + await ingestRecords( db, [ makeEvent({ @@ -148,7 +148,7 @@ describe("applyEvents", () => { describe("queryRecords", () => { beforeEach(async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:a/community.lexicon.calendar.event/1", did: "did:plc:a", diff --git a/packages/contrail/tests/refresh.test.ts b/packages/contrail/tests/refresh.test.ts index 3ee34fb..b84f016 100644 --- a/packages/contrail/tests/refresh.test.ts +++ b/packages/contrail/tests/refresh.test.ts @@ -1,6 +1,6 @@ import { describe, it, expect, vi, beforeEach } from "vitest"; import { refresh } from "../src/index"; -import { applyEvents, createTestDbWithSchema, makeEvent, TEST_CONFIG } from "./helpers"; +import { ingestRecords, 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. @@ -74,7 +74,7 @@ describe("refresh", () => { await registerKnownDid(db, ALICE); // Seed DB with a record at this URI. - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ did: ALICE, rkey: "k1", @@ -101,7 +101,7 @@ describe("refresh", () => { // Seed DB with an OLD record (indexed_at way in the past). const dayAgoUs = (Date.now() - 86_400_000) * 1000; - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ did: ALICE, rkey: "k2", @@ -130,7 +130,7 @@ describe("refresh", () => { // Seed with a record indexed JUST NOW. const nowUs = Date.now() * 1000; - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ did: ALICE, rkey: "k3", diff --git a/packages/contrail/tests/search.test.ts b/packages/contrail/tests/search.test.ts index f7f645b..f744d0c 100644 --- a/packages/contrail/tests/search.test.ts +++ b/packages/contrail/tests/search.test.ts @@ -3,7 +3,7 @@ import type { Database } from "../src/index"; import { resolveConfig } from "../src/index"; import { createTestDb, makeEvent } from "./helpers"; import { initSchema } from "../src/index"; -import { applyEvents, queryRecords } from "../src/index"; +import { ingestRecords, queryRecords } from "../src/index"; // Detect FTS5 support at module level (node:sqlite doesn't include it) let hasFts = false; @@ -54,7 +54,7 @@ describe.skipIf(!hasFts)("FTS with explicit searchable fields", () => { const collection = "community.lexicon.calendar.event"; beforeEach(async () => { - await applyEvents( + await ingestRecords( db, [ makeEvent({ @@ -117,7 +117,7 @@ describe.skipIf(!hasFts)("FTS with explicit searchable fields", () => { }); it("does not search range fields (startsAt)", async () => { - await applyEvents( + await ingestRecords( db, [ makeEvent({ @@ -140,10 +140,10 @@ describe.skipIf(!hasFts)("FTS sync", () => { const collection = "community.lexicon.calendar.event"; it("updates FTS on record update", async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:a/community.lexicon.calendar.event/1", collection, rkey: "1", record: { name: "Old Name", mode: "online", description: "test" }, time_us: 1000 }), ], SEARCH_CONFIG); - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:a/community.lexicon.calendar.event/1", collection, rkey: "1", record: { name: "New Name", mode: "online", description: "test" }, operation: "update", time_us: 2000 }), ], SEARCH_CONFIG); expect((await queryRecords(db, SEARCH_CONFIG, { collection, search: "Old" })).records).toHaveLength(0); @@ -151,17 +151,17 @@ describe.skipIf(!hasFts)("FTS sync", () => { }); it("removes from FTS on delete", async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:a/community.lexicon.calendar.event/1", collection, rkey: "1", record: { name: "Deletable", mode: "online", description: "test" }, time_us: 1000 }), ], SEARCH_CONFIG); - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:a/community.lexicon.calendar.event/1", collection, rkey: "1", operation: "delete", record: { name: "Deletable", mode: "online", description: "test" }, time_us: 2000 }), ], SEARCH_CONFIG); expect((await queryRecords(db, SEARCH_CONFIG, { collection, search: "Deletable" })).records).toHaveLength(0); }); it("does not duplicate FTS rows when the same record is re-applied during backfill", async () => { - // Backfill passes call applyEvents with skipReplayDetection: true, which leaves + // Backfill passes call ingestRecords with skipReplayDetection: true, which leaves // existingMap empty so every record looks brand-new. Re-backfilling the same // record must not append a second FTS row; otherwise the search JOIN fans out // and returns the event more than once (which crashes keyed lists downstream). @@ -172,8 +172,8 @@ describe.skipIf(!hasFts)("FTS sync", () => { record: { name: "Backfilled Meetup", mode: "online", description: "test" }, time_us: 1000, }); - await applyEvents(db, [event], SEARCH_CONFIG, { skipReplayDetection: true }); - await applyEvents(db, [event], SEARCH_CONFIG, { skipReplayDetection: true }); + await ingestRecords(db, [event], SEARCH_CONFIG, { skipReplayDetection: true }); + await ingestRecords(db, [event], SEARCH_CONFIG, { skipReplayDetection: true }); const result = await queryRecords(db, SEARCH_CONFIG, { collection, search: "Meetup" }); expect(result.records).toHaveLength(1); @@ -183,12 +183,12 @@ describe.skipIf(!hasFts)("FTS sync", () => { // The delete must run unconditionally. If an update leaves every searchable // field empty, buildFtsContent returns null and there is nothing to re-insert, // but the prior FTS row must still be removed so old terms stop matching. - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:a/community.lexicon.calendar.event/1", collection, rkey: "1", record: { name: "Searchable Title", mode: "online", description: "find me" }, time_us: 1000 }), ], SEARCH_CONFIG); expect((await queryRecords(db, SEARCH_CONFIG, { collection, search: "Searchable" })).records).toHaveLength(1); - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:a/community.lexicon.calendar.event/1", collection, rkey: "1", record: { startsAt: "2026-01-01T00:00:00Z" }, operation: "update", time_us: 2000 }), ], SEARCH_CONFIG); expect((await queryRecords(db, SEARCH_CONFIG, { collection, search: "Searchable" })).records).toHaveLength(0); @@ -199,7 +199,7 @@ describe.skipIf(!hasFts)("explicit searchable fields", () => { const collection = "test.explicit.collection"; beforeEach(async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:a/test.explicit.collection/1", did: "did:plc:a", collection, rkey: "1", record: { title: "Interesting Article", body: "Some content here", category: "tech" }, time_us: 1000 }), ], SEARCH_CONFIG); }); @@ -218,7 +218,7 @@ describe.skipIf(!hasFts)("searchable: false", () => { const collection = "test.disabled.collection"; it("search param is ignored when FTS is disabled", async () => { - await applyEvents(db, [ + await ingestRecords(db, [ makeEvent({ uri: "at://did:plc:a/test.disabled.collection/1", did: "did:plc:a", collection, rkey: "1", record: { name: "Should Not Be Searchable" }, time_us: 1000 }), ], SEARCH_CONFIG); const result = await queryRecords(db, SEARCH_CONFIG, { collection, search: "Searchable" }); @@ -240,7 +240,7 @@ describe.skipIf(!hasFts)("search pagination", () => { time_us: (i + 1) * 1000, }) ); - await applyEvents(db, events, SEARCH_CONFIG); + await ingestRecords(db, events, SEARCH_CONFIG); }); it("paginates search results", async () => { diff --git a/packages/contrail/tests/sinks.test.ts b/packages/contrail/tests/sinks.test.ts index e9534b4..9b57e63 100644 --- a/packages/contrail/tests/sinks.test.ts +++ b/packages/contrail/tests/sinks.test.ts @@ -1,16 +1,16 @@ -/** Sinks — write-only, post-commit observers fanned out by `applyEvents`. +/** Sinks — write-only, post-commit observers fanned out by `ingestRecords`. * * - One deduplicated event per record (not the realtime collection:/actor: pair). * - Fire on the live path by default and on backfill when `phase` says so. * - Failures are isolated: a throwing sink blocks neither the DB commit nor * the other sinks, and is logged. * - Public-record scope: space-scoped records publish via the publishing - * adapter and never reach `applyEvents`, so a sink cannot observe them. */ + * adapter and never reach `ingestRecords`, so a sink cannot observe them. */ import { describe, it, expect } from "vitest"; import { createSqliteDatabase } from "../src/adapters/sqlite"; import { initSchema } from "../src/index"; -import { applyEvents, queryRecords } from "../src/index"; +import { ingestRecords, queryRecords } from "../src/index"; import { resolveConfig } from "../src/index"; import type { ContrailConfig, @@ -61,13 +61,13 @@ async function freshDb(config: ContrailConfig) { return db; } -describe("sinks — applyEvents fan-out", () => { +describe("sinks — ingestRecords fan-out", () => { it("delivers one deduplicated created event per record, phase=live by default", async () => { const sink = new RecordingSink(); const config = configWithSinks([sink]); const db = await freshDb(config); - await applyEvents(db, [createEvent()], config); + await ingestRecords(db, [createEvent()], config); expect(sink.calls).toHaveLength(1); expect(sink.calls[0].ctx.phase).toBe("live"); @@ -90,8 +90,8 @@ describe("sinks — applyEvents fan-out", () => { const config = configWithSinks([sink]); const db = await freshDb(config); - await applyEvents(db, [createEvent()], config); - await applyEvents(db, [createEvent({ operation: "delete" })], config); + await ingestRecords(db, [createEvent()], config); + await ingestRecords(db, [createEvent({ operation: "delete" })], config); expect(sink.calls).toHaveLength(2); expect(sink.calls[1].events[0]).toEqual({ @@ -108,7 +108,7 @@ describe("sinks — applyEvents fan-out", () => { const config = configWithSinks([sink]); const db = await freshDb(config); - await applyEvents(db, [createEvent()], config, { phase: "backfill" }); + await ingestRecords(db, [createEvent()], config, { phase: "backfill" }); expect(sink.calls[0].ctx.phase).toBe("backfill"); }); @@ -119,7 +119,7 @@ describe("sinks — applyEvents fan-out", () => { const config = configWithSinks([a, b]); const db = await freshDb(config); - await applyEvents(db, [createEvent()], config); + await ingestRecords(db, [createEvent()], config); expect(a.calls).toHaveLength(1); expect(b.calls).toHaveLength(1); @@ -143,7 +143,7 @@ describe("sinks — applyEvents fan-out", () => { }; const db = await freshDb(config); - await applyEvents(db, [createEvent()], config); + await ingestRecords(db, [createEvent()], config); // The throw blocked neither the later sink... expect(after.calls).toHaveLength(1); @@ -159,7 +159,7 @@ describe("sinks — applyEvents fan-out", () => { const config = configWithSinks([sink]); const db = await freshDb(config); - await applyEvents(db, [], config); + await ingestRecords(db, [], config); expect(sink.calls).toHaveLength(0); }); -- 2.51.2