diff --git a/.changeset/remove-sinks.md b/.changeset/remove-sinks.md new file mode 100644 index 0000000..6592144 --- /dev/null +++ b/.changeset/remove-sinks.md @@ -0,0 +1,5 @@ +--- +"@atmo-dev/contrail": minor +--- + +Remove the best-effort post-commit sink API. Contrail now limits core ingestion to transactional SQL projection instead of invoking external callbacks after live or historical commits. diff --git a/apps/benchmark/README.md b/apps/benchmark/README.md index e744cca..cf6e90b 100644 --- a/apps/benchmark/README.md +++ b/apps/benchmark/README.md @@ -14,7 +14,7 @@ D1 is the default backend. Use `--backend sqlite` for a fresh in-memory native S pnpm bench --backend sqlite --config calendar-records-only.config.json ``` -`calendar.config.json` mirrors the retained indexing shape from `11-atproto/02-atmo-rsvp`: calendar events, RSVPs, profiles, follows/feeds, query indexes, relation counts, and event full-text search. Removed product modules and read-only pipeline handlers are intentionally absent because they do not participate in indexing. External sinks are also omitted so the benchmark measures Contrail and D1 rather than Meilisearch latency. +`calendar.config.json` mirrors the retained indexing shape from `11-atproto/02-atmo-rsvp`: calendar events, RSVPs, profiles, follows/feeds, query indexes, relation counts, and event full-text search. Removed product modules and read-only pipeline handlers are intentionally absent because they do not participate in indexing. To run the same full workload with strict runtime Lexicon validation and canonical DAG-CBOR CID verification enabled: @@ -69,4 +69,4 @@ Each result includes: ## Adding configs -Put portable JSON `ContrailConfig` files in `configs/`. The harness also accepts absolute paths or paths relative to the current directory. JSON configs cannot contain callbacks, custom query functions, or sinks; those should be omitted or represented by a benchmark-specific code harness when they materially affect ingestion. +Put portable JSON `ContrailConfig` files in `configs/`. The harness also accepts absolute paths or paths relative to the current directory. JSON configs cannot contain callbacks or custom query functions; those should be omitted or represented by a benchmark-specific code harness when they materially affect ingestion. diff --git a/packages/contrail/src/core/backfill.ts b/packages/contrail/src/core/backfill.ts index a019068..b10fa14 100644 --- a/packages/contrail/src/core/backfill.ts +++ b/packages/contrail/src/core/backfill.ts @@ -401,8 +401,6 @@ async function backfillUserAttempt( aggregateDiagnostics: options?.aggregateDiagnostics, // Canonical projection and cursor acknowledgement commit atomically. trailingStatements: [checkpoint], - // Let sinks bulk-flush differently from live ingestion. - phase: "backfill", }); totalInserted += result.accepted.length; if (metrics) { diff --git a/packages/contrail/src/core/constellation.ts b/packages/contrail/src/core/constellation.ts index dcae8f8..221bd88 100644 --- a/packages/contrail/src/core/constellation.ts +++ b/packages/contrail/src/core/constellation.ts @@ -167,7 +167,7 @@ export async function backfillFollowersFromConstellation( timeUs: timestamp, indexedAt: nowUs, // The backlink API has no record CID. TID-derived source metadata - // is stable across repeated acquisitions, keeping sinks/feed fanout + // is stable across repeated acquisitions, keeping feed fanout // idempotent while the shared validator treats this as synthetic. source: { id: "constellation", diff --git a/packages/contrail/src/core/db/records.ts b/packages/contrail/src/core/db/records.ts index 74e0b94..e60680f 100644 --- a/packages/contrail/src/core/db/records.ts +++ b/packages/contrail/src/core/db/records.ts @@ -8,7 +8,6 @@ import type { RecordRow, RecordSource, } from "../types"; -import type { RecordEvent } from "../sinks/types"; import { getNestedValue, getRelationField, @@ -887,9 +886,6 @@ export async function projectEvents( skipDerivedProjections?: boolean; /** Pre-fetched existing records — skips the internal lookup when provided */ existing?: Map; - /** Ingest phase forwarded to `config.sinks`. `"live"` for jetstream / - * persistent ingest (default), `"backfill"` for replay / rebuild. */ - phase?: "live" | "backfill"; /** Statements committed after projection in the same database batch. */ trailingStatements?: Statement[]; /** Internal: ingestRecords already checked durable source order. */ @@ -991,53 +987,6 @@ export async function projectEvents( } await db.batch(batch); - - // Fan out to write-only sinks (derived indexes, audit logs, webhooks). - // 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; - if (sinks && sinks.length > 0) { - // An absent-row tombstone is canonical ordering state, not an externally - // visible deletion. Publish deletes only when this transaction removed a - // previously visible record. - const records: RecordEvent[] = events - .filter( - (event) => - event.operation !== "delete" || existingMap.has(event.uri), - ) - .map((event) => - event.operation === "delete" - ? { - kind: "deleted", - uri: event.uri, - did: event.did, - collection: event.collection, - rkey: event.rkey, - } - : { - kind: "created", - uri: event.uri, - did: event.did, - collection: event.collection, - rkey: event.rkey, - cid: event.cid, - record: event.record ? safeParseJson(event.record) : {}, - time_us: event.time_us, - }, - ); - if (records.length === 0) return selection; - const ctx = { phase: options?.phase ?? "live" } as const; - const logger = config.logger ?? console; - for (const sink of sinks) { - try { - await sink.onRecords(records, ctx); - } catch (err) { - logger.error("[sink] onRecords failed", err); - } - } - } - return selection; } @@ -1162,15 +1111,6 @@ export async function rebuildDerivedProjections( if (countStatements.length > 0) await db.batch(countStatements); } -function safeParseJson(s: string): Record { - try { - const v = JSON.parse(s); - return v && typeof v === "object" && !Array.isArray(v) ? (v as Record) : {}; - } catch { - return {}; - } -} - // --- Count columns --- /** Count column descriptor. `type` is the identifier returned in API responses and diff --git a/packages/contrail/src/core/ingest.ts b/packages/contrail/src/core/ingest.ts index bd7ca67..d665358 100644 --- a/packages/contrail/src/core/ingest.ts +++ b/packages/contrail/src/core/ingest.ts @@ -104,7 +104,6 @@ export interface IngestRecordsOptions { skipDerivedProjections?: 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; /** Statements committed after projection in the same database batch. */ diff --git a/packages/contrail/src/core/sinks/types.ts b/packages/contrail/src/core/sinks/types.ts deleted file mode 100644 index 59e5a15..0000000 --- a/packages/contrail/src/core/sinks/types.ts +++ /dev/null @@ -1,31 +0,0 @@ -/** Write-only, post-commit observers of applied public records. */ - -export interface SinkContext { - /** `live` for normal ingestion, `backfill` for replay or rebuild work. */ - phase: "live" | "backfill"; -} - -/** One upsert or deletion event per applied record. */ -export type RecordEvent = - | { - kind: "created"; - uri: string; - did: string; - collection: string; - rkey: string; - cid: string | null; - record: Record; - time_us: number; - } - | { - kind: "deleted"; - uri: string; - did: string; - collection: string; - rkey: string; - }; - -export interface Sink { - /** Runs after a committed batch. Errors are logged without stopping ingestion. */ - onRecords(events: RecordEvent[], context: SinkContext): Promise | void; -} diff --git a/packages/contrail/src/core/types.ts b/packages/contrail/src/core/types.ts index 4819b81..0be2255 100644 --- a/packages/contrail/src/core/types.ts +++ b/packages/contrail/src/core/types.ts @@ -261,11 +261,6 @@ export interface ContrailConfig { /** Expose the notifyOfUpdate HTTP endpoint. Off by default. * 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 record-ingest commit - * on both live and backfill paths, with failures isolated so a throwing - * sink never blocks ingestion. */ - sinks?: import("./sinks/types").Sink[]; /** Labels module configuration. When set, contrail subscribes to the * configured labelers, indexes their labels into a single `labels` table, * and hydrates `record.labels` onto `listRecords` / `getRecord` / profile diff --git a/packages/contrail/src/core/validation.ts b/packages/contrail/src/core/validation.ts index e38d464..bdc95d8 100644 --- a/packages/contrail/src/core/validation.ts +++ b/packages/contrail/src/core/validation.ts @@ -116,7 +116,7 @@ export function prepareRecordValidation( return context; } -/** Validate one parsed create/update before filters, projections, or sinks. */ +/** Validate one parsed create/update before filters or projection. */ export async function validateCanonicalRecord( config: ContrailConfig, event: IngestEvent, diff --git a/packages/contrail/src/index.ts b/packages/contrail/src/index.ts index c83d456..6ce88b9 100644 --- a/packages/contrail/src/index.ts +++ b/packages/contrail/src/index.ts @@ -14,7 +14,6 @@ export { validateExternalUrl, } from "./core/client"; export type { ResolvedIdentity } from "./core/client"; -export * from "./core/sinks/types"; // Ingestion and maintenance. export * from "./core/ingest"; diff --git a/packages/contrail/tests/constellation-ingest.test.ts b/packages/contrail/tests/constellation-ingest.test.ts index 134eeef..ef13dbf 100644 --- a/packages/contrail/tests/constellation-ingest.test.ts +++ b/packages/contrail/tests/constellation-ingest.test.ts @@ -9,7 +9,6 @@ import { initSchema, queryRecords, resolveConfig, - type RecordEvent, } from "../src/index"; import { createSqliteDatabase } from "../src/adapters/sqlite"; @@ -60,7 +59,7 @@ async function cidFor(record: unknown) { return toString(await create(CODEC_DCBOR, encode(record))); } -async function setup(sinkBatches: RecordEvent[][]) { +async function setup() { const config = resolveConfig({ namespace: "com.example", profiles: [], @@ -81,13 +80,6 @@ async function setup(sinkBatches: RecordEvent[][]) { url: "https://constellation.example.com", userAgent: "contrail-test", }, - sinks: [ - { - async onRecords(records) { - sinkBatches.push(records); - }, - }, - ], }); const db = createSqliteDatabase(":memory:"); await initSchema(db, config); @@ -123,7 +115,6 @@ async function setup(sinkBatches: RecordEvent[][]) { ], config, ); - sinkBatches.length = 0; return { config, db }; } @@ -132,9 +123,8 @@ afterEach(() => { }); describe("Constellation enrichment through shared ingestion", () => { - it("uses normal validation/feed/sink projection and is idempotent", async () => { - const sinkBatches: RecordEvent[][] = []; - const { config, db } = await setup(sinkBatches); + it("uses normal validation, feed projection, and source ordering idempotently", async () => { + const { config, db } = await setup(); const fetchSpy = vi.fn(async () => new Response( JSON.stringify({ @@ -180,11 +170,5 @@ describe("Constellation enrichment through shared ingestion", () => { .bind(FOLLOWER, EVENT_URI) .first(), ).toEqual({ actor: FOLLOWER, uri: EVENT_URI }); - expect(sinkBatches).toHaveLength(1); - expect(sinkBatches[0]).toHaveLength(1); - expect(sinkBatches[0][0]).toMatchObject({ - kind: "created", - collection: FOLLOW, - }); }); }); diff --git a/packages/contrail/tests/profiles-ingest.test.ts b/packages/contrail/tests/profiles-ingest.test.ts index 6033e4a..a082998 100644 --- a/packages/contrail/tests/profiles-ingest.test.ts +++ b/packages/contrail/tests/profiles-ingest.test.ts @@ -9,7 +9,6 @@ import { resolveConfig, resolveProfiles, type Database, - type RecordEvent, } from "../src/index"; import { createSqliteDatabase } from "../src/adapters/sqlite"; @@ -54,7 +53,7 @@ async function cidFor(record: unknown) { return toString(await create(CODEC_DCBOR, encode(record))); } -async function setup(options?: { rejectB?: boolean; sink?: RecordEvent[][] }) { +async function setup(options?: { rejectB?: boolean }) { const config = resolveConfig({ namespace: "com.example", logger, @@ -75,15 +74,6 @@ async function setup(options?: { rejectB?: boolean; sink?: RecordEvent[][] }) { validation: { lexicons: [profileLexicon(PROFILE_A), profileLexicon(PROFILE_B)], }, - sinks: options?.sink - ? [ - { - async onRecords(records) { - options.sink!.push(records); - }, - }, - ] - : undefined, }); const db = createSqliteDatabase(":memory:"); await initSchema(db, config); @@ -127,11 +117,9 @@ afterEach(() => { }); describe("profile enrichment through shared ingestion", () => { - it("checks each profile collection and applies validation, optional FTS, sinks, and source metadata", async () => { - const sinkBatches: RecordEvent[][] = []; - const { config, db } = await setup({ sink: sinkBatches }); + it("checks each profile collection and applies validation, optional FTS, and source metadata", async () => { + const { config, db } = await setup(); await seedProfileA(db, config); - sinkBatches.length = 0; const profileB = { $type: PROFILE_B, displayName: "Fetched B" }; const fetchSpy = vi.fn(async (input: RequestInfo | URL) => { const url = new URL(input instanceof Request ? input.url : String(input)); @@ -154,8 +142,6 @@ describe("profile enrichment through shared ingestion", () => { PROFILE_A, PROFILE_B, ]); - expect(sinkBatches).toHaveLength(1); - expect(sinkBatches[0][0].collection).toBe(PROFILE_B); expect( await db .prepare( @@ -178,10 +164,8 @@ describe("profile enrichment through shared ingestion", () => { }); it("does not return or persist a fetched profile rejected by normal admission", async () => { - const sinkBatches: RecordEvent[][] = []; - const { config, db } = await setup({ rejectB: true, sink: sinkBatches }); + const { config, db } = await setup({ rejectB: true }); await seedProfileA(db, config); - sinkBatches.length = 0; const rejected = { $type: PROFILE_B, displayName: "Rejected" }; vi.stubGlobal( "fetch", @@ -203,8 +187,5 @@ describe("profile enrichment through shared ingestion", () => { expect( (await queryRecords(db, config, { collection: "profileB" })).records, ).toHaveLength(0); - // The exclusion is ordering state, not a visible create/delete, so it does - // not reach extensions when no prior profile existed. - expect(sinkBatches).toHaveLength(0); }); }); diff --git a/packages/contrail/tests/sinks.test.ts b/packages/contrail/tests/sinks.test.ts deleted file mode 100644 index 9b57e63..0000000 --- a/packages/contrail/tests/sinks.test.ts +++ /dev/null @@ -1,166 +0,0 @@ -/** 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 `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 { ingestRecords, queryRecords } from "../src/index"; -import { resolveConfig } from "../src/index"; -import type { - ContrailConfig, - IngestEvent, - RecordEvent, - Sink, - SinkContext, -} from "../src/index"; - -const ALICE = "did:plc:alice"; -const EVENT_NSID = "community.lexicon.calendar.event"; -const EVENT_URI = `at://${ALICE}/${EVENT_NSID}/abc`; - -/** Captures every onRecords call for assertions. */ -class RecordingSink implements Sink { - calls: { events: RecordEvent[]; ctx: SinkContext }[] = []; - async onRecords(events: RecordEvent[], ctx: SinkContext): Promise { - this.calls.push({ events, ctx }); - } -} - -function configWithSinks(sinks: Sink[]): ContrailConfig { - return resolveConfig({ - namespace: "test.sinks", - collections: { event: { collection: EVENT_NSID } }, - sinks, - }); -} - -function createEvent(overrides: Partial = {}): IngestEvent { - return { - uri: EVENT_URI, - did: ALICE, - collection: EVENT_NSID, - rkey: "abc", - operation: "create", - cid: "bafytest", - record: JSON.stringify({ name: "Launch party" }), - time_us: 1_700_000_000_000_000, - indexed_at: 1_700_000_000_000, - ...overrides, - }; -} - -async function freshDb(config: ContrailConfig) { - const db = createSqliteDatabase(":memory:"); - await initSchema(db, config); - return db; -} - -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 ingestRecords(db, [createEvent()], config); - - expect(sink.calls).toHaveLength(1); - expect(sink.calls[0].ctx.phase).toBe("live"); - // One event per record — not the collection:/actor: pair the realtime - // pubsub emits. - expect(sink.calls[0].events).toHaveLength(1); - expect(sink.calls[0].events[0]).toMatchObject({ - kind: "created", - uri: EVENT_URI, - did: ALICE, - collection: EVENT_NSID, - rkey: "abc", - cid: "bafytest", - record: { name: "Launch party" }, - }); - }); - - it("delivers deleted events carrying identity only", async () => { - const sink = new RecordingSink(); - const config = configWithSinks([sink]); - const db = await freshDb(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({ - kind: "deleted", - uri: EVENT_URI, - did: ALICE, - collection: EVENT_NSID, - rkey: "abc", - }); - }); - - it("forwards phase=backfill", async () => { - const sink = new RecordingSink(); - const config = configWithSinks([sink]); - const db = await freshDb(config); - - await ingestRecords(db, [createEvent()], config, { phase: "backfill" }); - - expect(sink.calls[0].ctx.phase).toBe("backfill"); - }); - - it("fans out to multiple sinks without nesting", async () => { - const a = new RecordingSink(); - const b = new RecordingSink(); - const config = configWithSinks([a, b]); - const db = await freshDb(config); - - await ingestRecords(db, [createEvent()], config); - - expect(a.calls).toHaveLength(1); - expect(b.calls).toHaveLength(1); - }); - - it("isolates a throwing sink: the record still commits and other sinks still fire", async () => { - const errors: unknown[][] = []; - const throwing: Sink = { - onRecords() { - throw new Error("boom"); - }, - }; - const after = new RecordingSink(); - const config = configWithSinks([throwing, after]); - config.logger = { - log() {}, - warn() {}, - error: (...args: unknown[]) => { - errors.push(args); - }, - }; - const db = await freshDb(config); - - await ingestRecords(db, [createEvent()], config); - - // The throw blocked neither the later sink... - expect(after.calls).toHaveLength(1); - // ...nor the DB commit... - const result = await queryRecords(db, config, { collection: EVENT_NSID }); - expect(result.records).toHaveLength(1); - // ...and the failure was logged, not propagated. - expect(errors).toHaveLength(1); - }); - - it("does not fire when there are no events", async () => { - const sink = new RecordingSink(); - const config = configWithSinks([sink]); - const db = await freshDb(config); - - await ingestRecords(db, [], config); - - expect(sink.calls).toHaveLength(0); - }); -}); diff --git a/packages/contrail/tests/source-ordering.test.ts b/packages/contrail/tests/source-ordering.test.ts index 1632916..576cba1 100644 --- a/packages/contrail/tests/source-ordering.test.ts +++ b/packages/contrail/tests/source-ordering.test.ts @@ -74,17 +74,8 @@ async function visibleName(db: Database, resolved = config()) { } describe("durable source ordering", () => { - it("rejects stale and duplicate updates before projections and sinks", async () => { - const sinkBatches: string[][] = []; - const resolved = config({ - sinks: [ - { - async onRecords(records) { - sinkBatches.push(records.map((record) => record.uri)); - }, - }, - ], - }); + it("rejects stale and duplicate updates before projection", async () => { + const resolved = config(); const db = await setup(resolved); const newest = mutation({ operation: "update", @@ -126,7 +117,6 @@ describe("durable source ordering", () => { expect(await visibleName(db, resolved)).toBe("new"); expect(stale.accepted).toHaveLength(0); expect(stale.dropped.superseded).toBe(2); - expect(sinkBatches).toHaveLength(1); }); it("does not let a stale delete remove a newer record", async () => { -- 2.51.2 From 746b31c31935e5c8103e5dd3e142fb0640456066 Mon Sep 17 00:00:00 2001 From: Florian <45694132+flo-bit@users.noreply.github.com> Date: Thu, 6 Aug 2026 02:43:10 +0200 Subject: [PATCH 2/2] Release sink removal as patch --- .changeset/remove-sinks.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.changeset/remove-sinks.md b/.changeset/remove-sinks.md index 6592144..d6aadb5 100644 --- a/.changeset/remove-sinks.md +++ b/.changeset/remove-sinks.md @@ -1,5 +1,5 @@ --- -"@atmo-dev/contrail": minor +"@atmo-dev/contrail": patch --- Remove the best-effort post-commit sink API. Contrail now limits core ingestion to transactional SQL projection instead of invoking external callbacks after live or historical commits.