From bc2f4c1b9eb6bb15867efdf6e8ac5256c912fa52 Mon Sep 17 00:00:00 2001 From: Florian <45694132+flo-bit@users.noreply.github.com> Date: Fri, 21 Aug 2026 09:45:48 +0200 Subject: [PATCH] Add transactional projection change log foundation --- .changeset/durable-projection-log.md | 7 + packages/contrail/README.md | 26 + packages/contrail/src/adapters/postgres.ts | 15 + packages/contrail/src/adapters/sqlite.ts | 1 + packages/contrail/src/core/backfill.ts | 1 + packages/contrail/src/core/bootstrap.ts | 1 + packages/contrail/src/core/change-log.ts | 520 ++++++++++++++++++ packages/contrail/src/core/constellation.ts | 1 + packages/contrail/src/core/db/records.ts | 157 +++++- packages/contrail/src/core/db/schema.ts | 62 ++- packages/contrail/src/core/ingest.ts | 198 ++++--- packages/contrail/src/core/jetstream.ts | 1 + packages/contrail/src/core/persistent.ts | 1 + packages/contrail/src/core/router/notify.ts | 2 +- packages/contrail/src/core/router/profiles.ts | 4 +- packages/contrail/src/core/types.ts | 135 +++++ packages/contrail/src/index.ts | 6 + packages/contrail/tests/change-log.test.ts | 493 +++++++++++++++++ packages/contrail/tests/ingest.test.ts | 4 +- .../tests/postgres-concurrent-init.test.ts | 2 +- packages/contrail/tests/postgres-e2e.test.ts | 79 ++- packages/contrail/tests/postgres.test.ts | 2 +- packages/contrail/tests/schema.test.ts | 53 +- turbo.json | 2 +- 24 files changed, 1655 insertions(+), 118 deletions(-) create mode 100644 .changeset/durable-projection-log.md create mode 100644 packages/contrail/src/core/change-log.ts create mode 100644 packages/contrail/tests/change-log.test.ts diff --git a/.changeset/durable-projection-log.md b/.changeset/durable-projection-log.md new file mode 100644 index 0000000..6a3e1e5 --- /dev/null +++ b/.changeset/durable-projection-log.md @@ -0,0 +1,7 @@ +--- +"@atmo-dev/contrail": minor +--- + +Add the first transactional projection change-log milestone. Optional static consumer definitions now create a fresh-generation log, durable registrations, and collection/phase coverage ledger. Winning logical URI changes append compact references atomically with canonical records, derived projections, tombstones, and source checkpoints; disabled configurations create no log tables or append writes. + +Harden all projection writers with transaction-time predecessor guards and bounded conflict retries so overlapping cron, persistent, notify, and backfill work cannot commit stale canonical or derived state. Enabling or changing log coverage on a populated generation fails closed pending explicit quiet-boundary migration tooling. diff --git a/packages/contrail/README.md b/packages/contrail/README.md index 8a1e40d..99614f1 100644 --- a/packages/contrail/README.md +++ b/packages/contrail/README.md @@ -72,6 +72,32 @@ After a write to a user's PDS, `contrail.notify(uri)` can fetch the authoritativ Contrail stores source event time, repository revision, source cursor, CID, and local index time separately from record/application time. Durable tombstones reject stale resurrection, and live Jetstream projection commits its exact yielded cursor in the same transaction. A successful PDS `listRecords` page is a current authoritative observation, so it supersedes older durable state without a redundant version read; its version writes and page cursor still commit atomically. Tombstones are retained indefinitely; authoritative rebuild/retention tooling is planned separately. +Projection winner selection is guarded again inside the write transaction. Overlapping cron, persistent, notify, or backfill writers cannot commit a stale canonical row, derived projection, tombstone, or source checkpoint; a changed predecessor rolls the complete attempt back and retries from fresh durable state. + +## Transactional change log (experimental) + +Fresh empty generations may opt into a compact transactional projection change log: + +```ts +const config = { + // ... + changes: { + consumers: { + search: { + collections: ["community.lexicon.calendar.event"], + phases: ["historical", "live"], + initial: "current", + requiredForActivation: true, + }, + }, + }, +} satisfies ContrailConfig; +``` + +Static definitions contain no handlers, URLs, clients, credentials, or secrets. Contrail registers them with a random database-generation ID and collection/phase coverage ledger. A winning logical put/delete appends one compact URI/version reference in the same transaction as canonical and derived state plus the source checkpoint. Duplicate, stale, same-CID, absent-delete, rejected, and rolled-back mutations append nothing. Record bodies are hydrated from current state by the later delivery layer rather than copied into the log. + +This milestone intentionally exposes only the atomic log foundation; leased claim/hydrate/ack delivery and current-state bootstrap APIs follow separately. Enabling logging on a populated database, changing coverage, or disabling an existing log fails closed until explicit quiet-boundary migration tooling lands. With no configured consumers, no change-log tables or append writes exist. + ## Local development A project containing only `contrail.config.ts` can start a complete local service: diff --git a/packages/contrail/src/adapters/postgres.ts b/packages/contrail/src/adapters/postgres.ts index 310be5d..32586e2 100644 --- a/packages/contrail/src/adapters/postgres.ts +++ b/packages/contrail/src/adapters/postgres.ts @@ -16,6 +16,21 @@ const BIGINT_COLUMNS = new Set([ "last_seen_at", "resolved_at", "total", + "revision", + "head_position", + "retained_floor_position", + "position", + "acknowledged_position", + "bootstrap_anchor_position", + "bootstrap_target_position", + "from_position", + "through_position", + "lease_expires_at", + "next_attempt_at", + "last_success_at", + "last_error_at", + "created_at", + "updated_at", ]); function normalizeRow(row: any): any { diff --git a/packages/contrail/src/adapters/sqlite.ts b/packages/contrail/src/adapters/sqlite.ts index 34b2511..57dee1e 100644 --- a/packages/contrail/src/adapters/sqlite.ts +++ b/packages/contrail/src/adapters/sqlite.ts @@ -9,6 +9,7 @@ interface SqliteStatement extends Statement { export function createSqliteDatabase(path: string): Database { const raw = new DatabaseSync(path); raw.exec("PRAGMA journal_mode = WAL"); + raw.exec("PRAGMA busy_timeout = 5000"); function wrapStatement( sql: string, diff --git a/packages/contrail/src/core/backfill.ts b/packages/contrail/src/core/backfill.ts index f4d799f..e84ad58 100644 --- a/packages/contrail/src/core/backfill.ts +++ b/packages/contrail/src/core/backfill.ts @@ -429,6 +429,7 @@ async function backfillUserAttempt( ) .bind(nextCursor ?? null, pageDone ? 1 : 0, now, did, collection); const result = await ingestRecords(db, events, config, { + phase: "historical", skipReplayDetection: options?.skipReplayDetection, skipFeedFanout: true, knownDids: options?.knownDids, diff --git a/packages/contrail/src/core/bootstrap.ts b/packages/contrail/src/core/bootstrap.ts index 7071cbe..e4e5757 100644 --- a/packages/contrail/src/core/bootstrap.ts +++ b/packages/contrail/src/core/bootstrap.ts @@ -496,6 +496,7 @@ export class DatabaseBootstrapTarget implements BootstrapTarget { ): Promise { const knownDids = await this.getKnownDids(); const result = await ingestRecords(this.db, events, this.config, { + phase: "historical", knownDids, skipDerivedProjections: this.options.deferDerivedProjections === true, authoritativeSourceObservation, diff --git a/packages/contrail/src/core/change-log.ts b/packages/contrail/src/core/change-log.ts new file mode 100644 index 0000000..6505ce9 --- /dev/null +++ b/packages/contrail/src/core/change-log.ts @@ -0,0 +1,520 @@ +import type { SqlDialect } from "./dialect"; +import type { + ContrailConfig, + Database, + IngestEvent, + ProjectionPhase, + Statement, +} from "./types"; +import { + canonicalChangeDefinitions, + changeConsumerPhases, + changeLogCoverage, + changesEnabled, +} from "./types"; + +export const MAX_CHANGE_BATCH_CHANGES = 500; +export const MAX_CHANGE_BATCH_BYTES = 512_000; + +export interface RecordChange { + id: string; + kind: "record"; + operation: "put" | "delete"; + uri: string; + did: string; + collection: string; + rkey: string; + cid: string | null; + version: { + sourceId: string; + sourceEpoch: string | null; + sourceRevision: string | null; + sourceTimeUs: number; + sourceCursor: string | null; + }; +} + +type StoredRecordChange = Omit; + +export interface ChangeLogState { + generation: string; + head: string; + retainedFloor: string; + createdAt: number; +} + +interface ChangeLogStateRow { + generation_id: string; + head_position: number | string; + retained_floor_position: number | string; + definitions_json: string; + created_at: number | string; +} + +interface ChangeConsumerRow { + consumer_id: string; + generation_id: string; + acknowledged_position: number | string; + configured_collections_json: string; + configured_phases_json: string; + initial_mode: string; + required_for_activation: number | string; + bootstrap_state: string; + bootstrap_anchor_position: number | string | null; +} + +interface ChangeCoverageRow { + generation_id: string; + collection: string; + phase: ProjectionPhase; + from_position: number | string; + through_position: number | string | null; +} + +export interface ChangeLogSchemaProbe { + exists: boolean; + state: ChangeLogStateRow | null; +} + +/** Optional physical schema. It is absent when no consumers are configured. */ +export function buildChangeLogSchema( + config: ContrailConfig, + dialect: SqlDialect, +): string[] { + if (!changesEnabled(config)) return []; + const bigint = dialect.bigintType; + return [ + `CREATE TABLE IF NOT EXISTS change_log_state ( + id INTEGER PRIMARY KEY CHECK (id = 1), + generation_id TEXT NOT NULL, + head_position ${bigint} NOT NULL, + retained_floor_position ${bigint} NOT NULL, + definitions_json TEXT NOT NULL, + created_at ${bigint} NOT NULL + )`, + `CREATE TABLE IF NOT EXISTS change_batches ( + generation_id TEXT NOT NULL, + position ${bigint} NOT NULL, + projection_transaction_id TEXT NOT NULL, + source_id TEXT NOT NULL, + source_epoch TEXT, + source_cursor TEXT, + phase TEXT NOT NULL CHECK (phase IN ('historical', 'live')), + changes_json TEXT NOT NULL, + change_count INTEGER NOT NULL, + created_at ${bigint} NOT NULL, + PRIMARY KEY (generation_id, position) + )`, + "CREATE UNIQUE INDEX IF NOT EXISTS idx_change_batches_transaction ON change_batches(generation_id, projection_transaction_id)", + "CREATE INDEX IF NOT EXISTS idx_change_batches_created ON change_batches(generation_id, created_at)", + `CREATE TABLE IF NOT EXISTS change_consumers ( + consumer_id TEXT PRIMARY KEY, + generation_id TEXT NOT NULL, + acknowledged_position ${bigint} NOT NULL, + configured_collections_json TEXT NOT NULL, + configured_phases_json TEXT NOT NULL, + initial_mode TEXT NOT NULL CHECK (initial_mode IN ('current', 'future', 'history')), + required_for_activation INTEGER NOT NULL DEFAULT 0, + bootstrap_state TEXT NOT NULL CHECK (bootstrap_state IN ('pending', 'scanning', 'catching-up', 'activating', 'ready', 'error', 'reset-required')), + bootstrap_anchor_position ${bigint}, + bootstrap_scan_collection TEXT, + bootstrap_scan_cursor TEXT, + bootstrap_target_position ${bigint}, + bootstrap_token TEXT, + lease_owner TEXT, + lease_expires_at ${bigint}, + attempts INTEGER NOT NULL DEFAULT 0, + next_attempt_at ${bigint}, + last_success_at ${bigint}, + last_error_code TEXT, + last_error_at ${bigint}, + updated_at ${bigint} NOT NULL + )`, + `CREATE TABLE IF NOT EXISTS change_log_coverage ( + generation_id TEXT NOT NULL, + collection TEXT NOT NULL, + phase TEXT NOT NULL CHECK (phase IN ('historical', 'live')), + from_position ${bigint} NOT NULL, + through_position ${bigint}, + PRIMARY KEY (generation_id, collection, phase, from_position) + )`, + ]; +} + +function missingTable(error: unknown): boolean { + if (!error || typeof error !== "object") return false; + if ((error as { code?: unknown }).code === "42P01") return true; + return /no such table|does not exist/i.test( + String((error as { message?: unknown }).message ?? ""), + ); +} + +export async function probeChangeLogSchema( + db: Database, +): Promise { + try { + const state = await db + .prepare( + `SELECT generation_id, head_position, retained_floor_position, + definitions_json, created_at + FROM change_log_state WHERE id = 1`, + ) + .first(); + return { exists: true, state }; + } catch (error) { + if (missingTable(error)) return { exists: false, state: null }; + throw error; + } +} + +/** Enabling the first log without an old-writer quiet boundary is supported + * only on an empty projection generation in this milestone. */ +export async function assertFreshChangeLogGeneration( + db: Database, +): Promise { + const row = await db + .prepare( + `SELECT CASE WHEN + EXISTS (SELECT 1 FROM record_versions LIMIT 1) OR + EXISTS (SELECT 1 FROM cursor LIMIT 1) OR + EXISTS (SELECT 1 FROM source_position LIMIT 1) OR + EXISTS (SELECT 1 FROM bootstrap_state LIMIT 1) OR + EXISTS (SELECT 1 FROM backfills LIMIT 1) OR + EXISTS (SELECT 1 FROM discovery LIMIT 1) OR + EXISTS (SELECT 1 FROM identities LIMIT 1) + THEN 1 ELSE 0 END AS active`, + ) + .first<{ active: number | string }>(); + if (Number(row?.active ?? 0) !== 0) { + throw new Error( + "Transactional change logging can currently be enabled only on a fresh empty generation; build a fresh generation instead of racing existing projection writers", + ); + } +} + +function canonicalCollections(collections: string[]): string { + return JSON.stringify([...collections].sort()); +} + +function canonicalPhases(phases: ProjectionPhase[]): string { + return JSON.stringify([...phases].sort()); +} + +/** Initialize one immutable milestone-1 logging definition. Later milestones + * add explicit quiet-boundary operations for changing this durable definition. */ +export async function initializeChangeLog( + db: Database, + config: ContrailConfig, +): Promise { + if (!changesEnabled(config)) return; + + const definitions = canonicalChangeDefinitions(config); + const now = Date.now(); + const candidateGeneration = crypto.randomUUID(); + const statements: Statement[] = [ + db + .prepare( + `INSERT INTO change_log_state + (id, generation_id, head_position, retained_floor_position, + definitions_json, created_at) + VALUES (1, ?, 0, 0, ?, ?) + ON CONFLICT(id) DO NOTHING`, + ) + .bind(candidateGeneration, definitions, now), + ]; + + for (const [consumerId, consumer] of Object.entries( + config.changes?.consumers ?? {}, + ).sort(([left], [right]) => left.localeCompare(right))) { + const initialReady = consumer.initial === "current" ? "pending" : "ready"; + statements.push( + db + .prepare( + `INSERT INTO change_consumers + (consumer_id, generation_id, acknowledged_position, + configured_collections_json, configured_phases_json, initial_mode, + required_for_activation, bootstrap_state, + bootstrap_anchor_position, attempts, updated_at) + SELECT ?, generation_id, head_position, ?, ?, ?, ?, ?, + CASE WHEN ? = 'current' THEN head_position ELSE NULL END, + 0, ? + FROM change_log_state + WHERE id = 1 AND definitions_json = ? + ON CONFLICT(consumer_id) DO NOTHING`, + ) + .bind( + consumerId, + canonicalCollections(consumer.collections), + canonicalPhases(changeConsumerPhases(consumer)), + consumer.initial, + consumer.requiredForActivation === true ? 1 : 0, + initialReady, + consumer.initial, + now, + definitions, + ), + ); + } + + for (const item of changeLogCoverage(config)) { + statements.push( + db + .prepare( + `INSERT INTO change_log_coverage + (generation_id, collection, phase, from_position, through_position) + SELECT generation_id, ?, ?, head_position, NULL + FROM change_log_state + WHERE id = 1 AND definitions_json = ? + ON CONFLICT(generation_id, collection, phase, from_position) + DO NOTHING`, + ) + .bind(item.collection, item.phase, definitions), + ); + } + + // Keep initialization under conservative D1 statement limits. The durable + // definitions_json winner makes these resumable chunks safe under concurrent + // initialization; init does not return until the complete set verifies. + await db.batch(statements.slice(0, 1)); + for (let index = 1; index < statements.length; index += 50) { + await db.batch(statements.slice(index, index + 50)); + } + await assertChangeLogDefinition(db, config); +} + +async function assertChangeLogDefinition( + db: Database, + config: ContrailConfig, +): Promise { + const state = (await probeChangeLogSchema(db)).state; + const definitions = canonicalChangeDefinitions(config); + if (!state || state.definitions_json !== definitions) { + throw new Error( + "Durable change consumer definitions differ from configuration; changing consumers or coverage requires a fresh generation in this milestone", + ); + } + + const consumers = await db + .prepare( + `SELECT consumer_id, generation_id, acknowledged_position, + configured_collections_json, configured_phases_json, initial_mode, + required_for_activation, bootstrap_state, + bootstrap_anchor_position + FROM change_consumers ORDER BY consumer_id`, + ) + .all(); + const expectedConsumers = Object.entries(config.changes?.consumers ?? {}).sort( + ([left], [right]) => left.localeCompare(right), + ); + if (consumers.results.length !== expectedConsumers.length) { + throw new Error("Durable change consumer registration is incomplete"); + } + for (let index = 0; index < expectedConsumers.length; index++) { + const [id, expected] = expectedConsumers[index]!; + const actual = consumers.results[index]!; + const bootstrapState = expected.initial === "current" ? "pending" : "ready"; + if ( + actual.consumer_id !== id || + actual.generation_id !== state.generation_id || + actual.configured_collections_json !== canonicalCollections(expected.collections) || + actual.configured_phases_json !== canonicalPhases(changeConsumerPhases(expected)) || + actual.initial_mode !== expected.initial || + Number(actual.required_for_activation) !== + (expected.requiredForActivation === true ? 1 : 0) || + actual.bootstrap_state !== bootstrapState || + Number(actual.acknowledged_position) !== 0 || + (expected.initial === "current" + ? Number(actual.bootstrap_anchor_position) !== 0 + : actual.bootstrap_anchor_position !== null) + ) { + throw new Error(`Durable change consumer ${id} is incompatible`); + } + } + + const coverage = await db + .prepare( + `SELECT generation_id, collection, phase, from_position, through_position + FROM change_log_coverage + ORDER BY collection, phase, from_position`, + ) + .all(); + const expectedCoverage = changeLogCoverage(config); + if (coverage.results.length !== expectedCoverage.length) { + throw new Error("Durable change-log coverage is incomplete"); + } + for (let index = 0; index < expectedCoverage.length; index++) { + const actual = coverage.results[index]!; + const expected = expectedCoverage[index]!; + if ( + actual.generation_id !== state.generation_id || + actual.collection !== expected.collection || + actual.phase !== expected.phase || + Number(actual.from_position) !== 0 || + actual.through_position !== null + ) { + throw new Error("Durable change-log coverage is incompatible"); + } + } +} + +export async function getChangeLogState( + db: Database, +): Promise { + const probe = await probeChangeLogSchema(db); + if (!probe.state) return null; + return { + generation: probe.state.generation_id, + head: String(probe.state.head_position), + retainedFloor: String(probe.state.retained_floor_position), + createdAt: Number(probe.state.created_at), + }; +} + +function bounded(value: string | null, label: string, maximum: number): void { + if (value !== null && value.length > maximum) { + throw new Error(`${label} exceeds ${maximum} characters`); + } +} + +function logicalChanges( + events: IngestEvent[], + existing: ReadonlyMap, + config: ContrailConfig, + phase: ProjectionPhase, +): StoredRecordChange[] { + const covered = new Set( + changeLogCoverage(config) + .filter((item) => item.phase === phase) + .map((item) => item.collection), + ); + if (covered.size === 0) return []; + + // projectEvents already selects one source winner per URI. Keep this final + // reduction defensive for direct internal callers. + const final = new Map(); + for (const event of events) final.set(event.uri, event); + + const changes: StoredRecordChange[] = []; + for (const event of final.values()) { + if (!covered.has(event.collection)) continue; + const prior = existing.get(event.uri); + const deleted = event.operation === "delete"; + const visibleChange = deleted + ? prior !== undefined + : prior === undefined || + prior.cid !== event.cid || + (event.cid === null && prior.record !== event.record); + if (!visibleChange) continue; + + const source = event.source; + const change: StoredRecordChange = { + kind: "record", + operation: deleted ? "delete" : "put", + uri: event.uri, + did: event.did, + collection: event.collection, + rkey: event.rkey, + cid: deleted ? null : event.cid, + version: { + sourceId: source?.id ?? "legacy-caller", + sourceEpoch: source?.epoch ?? null, + sourceRevision: source?.revision ?? null, + sourceTimeUs: source?.time_us ?? event.time_us, + sourceCursor: source?.cursor ?? null, + }, + }; + bounded(change.uri, "change URI", 2_048); + bounded(change.did, "change DID", 2_048); + bounded(change.collection, "change collection", 512); + bounded(change.rkey, "change rkey", 512); + bounded(change.cid, "change CID", 512); + bounded(change.version.sourceId, "change source ID", 128); + bounded(change.version.sourceEpoch, "change source epoch", 256); + bounded(change.version.sourceRevision, "change source revision", 2_048); + bounded(change.version.sourceCursor, "change source cursor", 2_048); + if ( + !Number.isSafeInteger(change.version.sourceTimeUs) || + change.version.sourceTimeUs < 0 + ) { + throw new Error("change source time must be a non-negative safe integer"); + } + changes.push(change); + } + return changes; +} + +function commonValue( + values: Array, +): string | null { + if (values.length === 0) return null; + const first = values[0]!; + return values.every((value) => value === first) ? first : null; +} + +/** Statements appended after canonical/derived projection and before source + * checkpoints in the same database transaction. */ +export function appendChangeLogStatements( + db: Database, + events: IngestEvent[], + existing: ReadonlyMap, + config: ContrailConfig, + phase: ProjectionPhase, +): Statement[] { + if (!changesEnabled(config)) return []; + const changes = logicalChanges(events, existing, config, phase); + if (changes.length === 0) return []; + if (changes.length > MAX_CHANGE_BATCH_CHANGES) { + throw new Error( + `Projection change batch contains ${changes.length} changes; maximum is ${MAX_CHANGE_BATCH_CHANGES}`, + ); + } + const serialized = JSON.stringify(changes); + const bytes = new TextEncoder().encode(serialized).byteLength; + if (bytes > MAX_CHANGE_BATCH_BYTES) { + throw new Error( + `Projection change batch contains ${bytes} encoded bytes; maximum is ${MAX_CHANGE_BATCH_BYTES}`, + ); + } + + const sourceIds = changes.map((change) => change.version.sourceId); + const sourceId = commonValue(sourceIds) ?? "mixed"; + const sourceEpoch = commonValue( + changes.map((change) => change.version.sourceEpoch), + ); + const sourceCursor = commonValue( + changes.map((change) => change.version.sourceCursor), + ); + const transactionId = crypto.randomUUID(); + const now = Date.now(); + + return [ + db + .prepare( + `UPDATE change_log_state + SET head_position = head_position + 1 + WHERE id = 1`, + ), + db + .prepare( + `INSERT INTO change_batches + (generation_id, position, projection_transaction_id, source_id, + source_epoch, source_cursor, phase, changes_json, change_count, + created_at) + VALUES ( + (SELECT generation_id FROM change_log_state WHERE id = 1), + (SELECT head_position FROM change_log_state WHERE id = 1), + ?, ?, ?, ?, ?, ?, ?, ? + )`, + ) + .bind( + transactionId, + sourceId, + sourceEpoch, + sourceCursor, + phase, + serialized, + changes.length, + now, + ), + ]; +} diff --git a/packages/contrail/src/core/constellation.ts b/packages/contrail/src/core/constellation.ts index 221bd88..1c4495d 100644 --- a/packages/contrail/src/core/constellation.ts +++ b/packages/contrail/src/core/constellation.ts @@ -184,6 +184,7 @@ export async function backfillFollowersFromConstellation( if (events.length > 0) { known.add(subjectDid); const ingest = await ingestRecords(db, events, config, { + phase: "historical", knownDids: known, }); inserted += ingest.accepted.length; diff --git a/packages/contrail/src/core/db/records.ts b/packages/contrail/src/core/db/records.ts index e06cac0..e74fb60 100644 --- a/packages/contrail/src/core/db/records.ts +++ b/packages/contrail/src/core/db/records.ts @@ -8,6 +8,7 @@ import type { RecordRow, RecordSource, OrderedSourceConfig, + ProjectionPhase, } from "../types"; import { getNestedValue, @@ -22,6 +23,7 @@ import { normalizeFeedTarget, feedTargetMaxItems, DEFAULT_FOLLOW_SHORT, + changesEnabled, } from "../types"; import { getSearchableFields, @@ -35,6 +37,7 @@ import { sqliteFtsContentExpression, } from "../dialect"; import type { SourcePosition } from "../sources"; +import { appendChangeLogStatements } from "../change-log"; // --- Counts --- @@ -700,9 +703,13 @@ export interface RecordVersionInfo { source_time_us: number; source_cursor: string | null; indexed_at: number; + /** Opaque optimistic-concurrency token; not part of source ordering. */ + projection_token: string; } -function versionForEvent(event: IngestEvent): RecordVersionInfo { +type ComparableRecordVersion = Omit; + +function versionForEvent(event: IngestEvent): ComparableRecordVersion { const source = event.source; return { uri: event.uri, @@ -742,8 +749,8 @@ function operationRank(operation: RecordVersionInfo["operation"]): number { * tie-breakers; notably a delete wins an exact tie so replay cannot resurrect it. */ export function compareRecordVersions( - left: RecordVersionInfo, - right: RecordVersionInfo, + left: ComparableRecordVersion, + right: ComparableRecordVersion, ): number { if ( left.source_revision !== null && @@ -793,7 +800,7 @@ export async function lookupRecordVersions( const placeholders = chunk.map(() => "?").join(","); const rows = await db .prepare( - `SELECT uri, did, collection, rkey, operation, cid, source_id, source_epoch, source_revision, source_time_us, source_cursor, indexed_at FROM record_versions WHERE uri IN (${placeholders})`, + `SELECT uri, did, collection, rkey, operation, cid, source_id, source_epoch, source_revision, source_time_us, source_cursor, indexed_at, projection_token FROM record_versions WHERE uri IN (${placeholders})`, ) .bind(...chunk) .all(); @@ -807,13 +814,13 @@ export interface MutationSelection { superseded: number; } -function selectMutationWinners( +export function selectMutationWinners( events: IngestEvent[], durable: ReadonlyMap, ): MutationSelection { const winners = new Map< string, - { event: IngestEvent; version: RecordVersionInfo; index: number } + { event: IngestEvent; version: ComparableRecordVersion; index: number } >(); let superseded = 0; @@ -918,10 +925,11 @@ const RECORD_UPSERT_BINDINGS = 7; const RECORD_UPSERT_ROWS = Math.floor( MAX_STATEMENT_BINDINGS / RECORD_UPSERT_BINDINGS ); -const RECORD_VERSION_BINDINGS = 12; +const RECORD_VERSION_BINDINGS = 13; const RECORD_VERSION_ROWS = Math.floor( MAX_STATEMENT_BINDINGS / RECORD_VERSION_BINDINGS, ); +const PROJECTION_GUARD_URIS = 40; interface StorageMutation { event: IngestEvent; @@ -932,17 +940,18 @@ function buildRecordVersionStatements( db: Database, events: IngestEvent[], existing: Map, + projectionTokens: ReadonlyMap, ): Statement[] { const statements: Statement[] = []; for (let index = 0; index < events.length; index += RECORD_VERSION_ROWS) { const chunk = events.slice(index, index + RECORD_VERSION_ROWS); const values = chunk - .map(() => "(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)") + .map(() => "(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)") .join(", "); statements.push( db .prepare( - `INSERT INTO record_versions (uri, did, collection, rkey, operation, cid, source_id, source_epoch, source_revision, source_time_us, source_cursor, indexed_at) VALUES ${values} ON CONFLICT(uri) DO UPDATE SET did = excluded.did, collection = excluded.collection, rkey = excluded.rkey, operation = excluded.operation, cid = excluded.cid, source_id = excluded.source_id, source_epoch = excluded.source_epoch, source_revision = excluded.source_revision, source_time_us = excluded.source_time_us, source_cursor = excluded.source_cursor, indexed_at = excluded.indexed_at`, + `INSERT INTO record_versions (uri, did, collection, rkey, operation, cid, source_id, source_epoch, source_revision, source_time_us, source_cursor, indexed_at, projection_token) VALUES ${values} ON CONFLICT(uri) DO UPDATE SET did = excluded.did, collection = excluded.collection, rkey = excluded.rkey, operation = excluded.operation, cid = excluded.cid, source_id = excluded.source_id, source_epoch = excluded.source_epoch, source_revision = excluded.source_revision, source_time_us = excluded.source_time_us, source_cursor = excluded.source_cursor, indexed_at = excluded.indexed_at, projection_token = excluded.projection_token`, ) .bind( ...chunk.flatMap((event) => { @@ -964,6 +973,7 @@ function buildRecordVersionStatements( version.source_time_us, version.source_cursor, version.indexed_at, + projectionTokens.get(event.uri)!, ]; }), ), @@ -1024,6 +1034,71 @@ function buildRecordMutationStatements( return statements; } +function buildProjectionGuardStatements( + db: Database, + events: IngestEvent[], + predecessors: ReadonlyMap, +): Statement[] { + const uris = [...new Set(events.map((event) => event.uri))]; + const statements: Statement[] = [ + db.prepare( + `INSERT INTO _contrail_projection_state (id, revision, guard) + VALUES (1, 0, 1) ON CONFLICT(id) DO NOTHING`, + ), + // PostgreSQL takes a row lock here. The following statement then receives a + // fresh READ COMMITTED snapshot after any earlier projector commits. SQLite + // and D1 already serialize the containing write batch. + db.prepare( + `UPDATE _contrail_projection_state + SET revision = revision + 1 + WHERE id = 1`, + ), + ]; + for (let index = 0; index < uris.length; index += PROJECTION_GUARD_URIS) { + const chunk = uris.slice(index, index + PROJECTION_GUARD_URIS); + const conditions: string[] = []; + const bindings: string[] = []; + for (const uri of chunk) { + const predecessor = predecessors.get(uri); + if (!predecessor) { + conditions.push( + "NOT EXISTS (SELECT 1 FROM record_versions WHERE uri = ?)", + ); + bindings.push(uri); + continue; + } + if (!predecessor.projection_token) { + throw new Error(`Record version ${uri} has no projection token`); + } + conditions.push( + "EXISTS (SELECT 1 FROM record_versions WHERE uri = ? AND projection_token = ?)", + ); + bindings.push(uri, predecessor.projection_token); + } + statements.push( + db + .prepare( + `UPDATE _contrail_projection_state + SET guard = CASE WHEN ${conditions.join(" AND ")} THEN 1 ELSE 0 END + WHERE id = 1`, + ) + .bind(...bindings), + ); + } + return statements; +} + +/** Adapter-neutral classification for the named optimistic guard constraint. */ +export function isProjectionConflictError(error: unknown): boolean { + if (!error || typeof error !== "object") return false; + const candidate = error as { code?: unknown; constraint?: unknown; message?: unknown }; + return ( + (candidate.code === "23514" && + candidate.constraint === "projection_guard_valid") || + /projection_guard_valid/i.test(String(candidate.message ?? "")) + ); +} + export async function projectEvents( db: Database, events: IngestEvent[], @@ -1033,19 +1108,26 @@ export async function projectEvents( skipFeedFanout?: boolean; /** Skip FTS and relation-count maintenance during canonical bulk loading. */ skipDerivedProjections?: boolean; - /** Pre-fetched existing records — skips the internal lookup when provided */ + /** @deprecated Existing rows are re-read for transaction conflict safety. */ existing?: Map; /** Statements committed after projection in the same database batch. */ trailingStatements?: Statement[]; /** Internal: ingestRecords already checked durable source order. */ sourceOrderingChecked?: boolean; + /** Durable versions observed while selecting source winners. */ + predecessors?: ReadonlyMap; + /** Acquisition phase persisted on an optional change batch. */ + phase?: ProjectionPhase; }, ): Promise { if (events.length === 0) return { applied: [], superseded: 0 }; + const predecessors = + options?.predecessors ?? + (await lookupRecordVersions(db, events.map((event) => event.uri))); const selection = options?.sourceOrderingChecked ? { applied: events, superseded: 0 } - : await selectCurrentMutations(db, events); + : selectMutationWinners(events, predecessors); events = selection.applied; if (events.length === 0) { if (options?.trailingStatements?.length) { @@ -1062,17 +1144,19 @@ export async function projectEvents( (relation) => relation.count !== 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 (!options?.skipReplayDetection) { - existingMap = await lookupExistingRecords(db, events, needRecordContent, config); - } else { - existingMap = new Map(); - } + const needRecordContent = + followCollections.length > 0 || + hasCountingRelations || + changesEnabled(config); + + // Existing state must be read after predecessor selection. A caller-provided + // map can predate that selection and would make derived changes incorrect + // even when the optimistic token guard itself succeeds. + const needExistingState = + !options?.skipReplayDetection || needRecordContent || changesEnabled(config); + const existingMap = needExistingState + ? await lookupExistingRecords(db, events, needRecordContent, config) + : new Map(); const batch: Statement[] = []; // Keep only the final storage mutation for a URI within this atomic batch. @@ -1122,15 +1206,36 @@ export async function projectEvents( } } - // Storage and durable version/tombstone metadata run first so FTS, feeds, - // and count statements in the same atomic batch observe the final records. + const projectionTokens = new Map( + events.map((event) => [event.uri, crypto.randomUUID()] as const), + ); + + // Lock, verify the exact durable predecessors selected by the caller, then + // write storage and version metadata. Any changed token violates the named + // guard constraint and rolls the complete database batch back. batch.unshift( + ...buildProjectionGuardStatements(db, events, predecessors), ...buildRecordMutationStatements(db, storageMutations.values()), - ...buildRecordVersionStatements(db, events, existingMap), + ...buildRecordVersionStatements( + db, + events, + existingMap, + projectionTokens, + ), ); - // Build deduplicated count statements — one UPDATE per unique target. + // Build deduplicated count statements, then append the compact change batch. + // Caller-provided source checkpoints deliberately remain last. batch.push(...buildBatchCountStatements(db, config, countTargets)); + batch.push( + ...appendChangeLogStatements( + db, + events, + existingMap, + config, + options?.phase ?? "live", + ), + ); if (options?.trailingStatements?.length) { batch.push(...options.trailingStatements); } diff --git a/packages/contrail/src/core/db/schema.ts b/packages/contrail/src/core/db/schema.ts index 95fa42d..51dfe2f 100644 --- a/packages/contrail/src/core/db/schema.ts +++ b/packages/contrail/src/core/db/schema.ts @@ -25,9 +25,19 @@ import { getSearchableFields, } from "../search"; import { buildLabelsSchema } from "../labels/schema"; +import { + assertFreshChangeLogGeneration, + buildChangeLogSchema, + initializeChangeLog, + probeChangeLogSchema, +} from "../change-log"; +import { + canonicalChangeDefinitions, + changesEnabled, +} from "../types"; import { getMeta, setMeta } from "./meta"; -export const CONTRAIL_SCHEMA_VERSION = 11; +export const CONTRAIL_SCHEMA_VERSION = 12; const SCHEMA_FINGERPRINT_KEY = "schema_fingerprint"; function getResolved(config: ContrailConfig): ResolvedMaps { @@ -130,7 +140,8 @@ CREATE TABLE IF NOT EXISTS record_versions ( source_revision TEXT, source_time_us ${dialect.bigintType} NOT NULL, source_cursor TEXT, - indexed_at ${dialect.bigintType} NOT NULL + indexed_at ${dialect.bigintType} NOT NULL, + projection_token TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS idx_record_versions_collection ON record_versions(collection); CREATE INDEX IF NOT EXISTS idx_record_versions_did ON record_versions(did); @@ -140,6 +151,12 @@ CREATE TABLE IF NOT EXISTS ingest_diagnostics ( total ${dialect.bigintType} NOT NULL, last_seen_at ${dialect.bigintType} NOT NULL ); +CREATE TABLE IF NOT EXISTS _contrail_projection_state ( + id INTEGER PRIMARY KEY CHECK (id = 1), + revision ${dialect.bigintType} NOT NULL DEFAULT 0, + guard INTEGER NOT NULL DEFAULT 1 + CONSTRAINT projection_guard_valid CHECK (guard = 1) +); `; } @@ -576,6 +593,11 @@ const MIGRATIONS: MigrationOp[] = [ column: "source_epoch", columnDef: "TEXT", }, + { + table: "record_versions", + column: "projection_token", + columnDef: "TEXT", + }, { table: "discovery", column: "retries", @@ -639,8 +661,8 @@ async function seedLegacyRecordVersions( const table = recordsTableName(shortName); await db .prepare( - `INSERT INTO record_versions (uri, did, collection, rkey, operation, cid, source_id, source_epoch, source_revision, source_time_us, source_cursor, indexed_at) - SELECT uri, did, ?, rkey, 'update', cid, 'legacy', NULL, NULL, indexed_at, NULL, indexed_at FROM ${table} + `INSERT INTO record_versions (uri, did, collection, rkey, operation, cid, source_id, source_epoch, source_revision, source_time_us, source_cursor, indexed_at, projection_token) + SELECT uri, did, ?, rkey, 'update', cid, 'legacy', NULL, NULL, indexed_at, NULL, indexed_at, uri FROM ${table} WHERE 1 = 1 ON CONFLICT(uri) DO NOTHING`, ) .bind(collection.collection) @@ -671,6 +693,7 @@ function schemaFingerprint( indexes: string[]; feeds: string[]; fts: string[]; + changes: string[]; }, ): string { return hashStrings([ @@ -682,6 +705,8 @@ function schemaFingerprint( ...ddl.indexes, ...ddl.feeds, ...ddl.fts, + ...ddl.changes, + canonicalChangeDefinitions(config), ...buildCountColumns(config), ...(config.labels ? buildLabelsSchema(dialect) : []), JSON.stringify(MIGRATIONS), @@ -702,12 +727,14 @@ export async function initSchema( const indexes = buildDynamicIndexes(config, dialect); const feeds = buildFeedTables(config, dialect); const fts = buildFtsTables(config, dialect); + const changes = buildChangeLogSchema(config, dialect); const fingerprint = schemaFingerprint(config, dialect, { base, collections, indexes, feeds, fts, + changes, }); if ((await getMeta(db, SCHEMA_FINGERPRINT_KEY)) === fingerprint) { @@ -715,6 +742,13 @@ export async function initSchema( return; } + const priorChangeLog = await probeChangeLogSchema(db); + if (!changesEnabled(config) && priorChangeLog.exists) { + throw new Error( + "Durable change logging cannot be disabled or removed during ordinary initialization", + ); + } + for (const statement of [...base, ...collections, ...indexes, ...feeds]) { await runIdempotentDdl(db, statement); } @@ -729,9 +763,29 @@ export async function initSchema( const hasFeeds = !!(config.feeds && Object.keys(config.feeds).length > 0); await runMigrations(db, hasFeeds); + await db + .prepare( + "UPDATE record_versions SET projection_token = uri WHERE projection_token IS NULL", + ) + .run(); + await db + .prepare( + `INSERT INTO _contrail_projection_state (id, revision, guard) + VALUES (1, 0, 1) ON CONFLICT(id) DO NOTHING`, + ) + .run(); await applyCountColumns(db, config); await seedLegacyRecordVersions(db, config); + if (changesEnabled(config)) { + const concurrentChangeLog = await probeChangeLogSchema(db); + if (!priorChangeLog.state && !concurrentChangeLog.state) { + await assertFreshChangeLogGeneration(db); + } + for (const statement of changes) await runIdempotentDdl(db, statement); + await initializeChangeLog(db, config); + } + for (const apply of options.extraSchemas ?? []) await apply(db); await setMeta(db, SCHEMA_FINGERPRINT_KEY, fingerprint); } diff --git a/packages/contrail/src/core/ingest.ts b/packages/contrail/src/core/ingest.ts index 48527cf..fe1392c 100644 --- a/packages/contrail/src/core/ingest.ts +++ b/packages/contrail/src/core/ingest.ts @@ -4,6 +4,7 @@ import type { Database, IngestEvent, MutationSource, + ProjectionPhase, Statement, } from "./types"; import { @@ -11,9 +12,11 @@ import { resolveCollectionKey, } from "./types"; import { + isProjectionConflictError, + lookupRecordVersions, projectEvents, selectAuthoritativeMutations, - selectCurrentMutations, + selectMutationWinners, type ExistingRecordInfo, } from "./db/records"; import { @@ -114,6 +117,9 @@ export interface IngestRecordsOptions { /** The source response is a current authoritative snapshot, so it supersedes * durable observations without a redundant version lookup. */ authoritativeSourceObservation?: boolean; + /** Acquisition phase for optional durable consumers. Defaults to live for + * backwards-compatible direct ingestRecords() calls. */ + phase?: ProjectionPhase; /** @internal Aggregate private diagnostics for one bulk run. The caller * flushes this bounded object once after concurrent page processing. */ aggregateDiagnostics?: IngestDiagnosticCounts; @@ -277,97 +283,131 @@ export async function ingestRecords( ); } - // Reject duplicate/stale source observations before they can admit dependent - // actors in this batch. The winning versions are persisted with projection. - const ordered = options.authoritativeSourceObservation - ? selectAuthoritativeMutations(accepted) - : await selectCurrentMutations(db, accepted); - dropped.superseded += ordered.superseded; + // Winner reads occur before db.batch on D1, so each projection carries the + // exact predecessor tokens it observed. A concurrent projector changes a + // token, the named transaction guard rolls everything back, and this loop + // repeats selection plus derived-state reads from fresh durable state. + const retryBase = { + unknownActor: dropped.unknownActor, + unknownSubject: dropped.unknownSubject, + superseded: dropped.superseded, + }; + const maximumAttempts = 5; + for (let attempt = 1; attempt <= maximumAttempts; attempt++) { + const attemptDropped: IngestDropCounts = { + ...dropped, + unknownActor: retryBase.unknownActor, + unknownSubject: retryBase.unknownSubject, + superseded: retryBase.superseded, + }; + const predecessors = await lookupRecordVersions( + db, + accepted.map((event) => event.uri), + ); + const ordered = options.authoritativeSourceObservation + ? selectAuthoritativeMutations(accepted) + : selectMutationWinners(accepted, predecessors); + attemptDropped.superseded += ordered.superseded; - const projectionExclusions = ordered.applied.filter((event) => - policyExcluded.has(event), - ); - const admitted = ordered.applied.filter((event) => !policyExcluded.has(event)); + const projectionExclusions = ordered.applied.filter((event) => + policyExcluded.has(event), + ); + const admitted = ordered.applied.filter( + (event) => !policyExcluded.has(event), + ); - const effectiveKnownDids = options.knownDids - ? new Set(options.knownDids) - : undefined; - const discoveredDids: string[] = []; - if (effectiveKnownDids) { + const effectiveKnownDids = options.knownDids + ? new Set(options.knownDids) + : undefined; + const discoveredDids: string[] = []; + if (effectiveKnownDids) { + for (const event of admitted) { + if (event.operation === "delete") continue; + const shortName = resolveCollectionKey(config, event.collection); + const collection = shortName ? config.collections[shortName] : undefined; + if (collection?.discover === false || effectiveKnownDids.has(event.did)) { + continue; + } + effectiveKnownDids.add(event.did); + discoveredDids.push(event.did); + } + } + + const actorFiltered: IngestEvent[] = []; for (const event of admitted) { - if (event.operation === "delete") continue; + if (event.operation === "delete" || !effectiveKnownDids) { + actorFiltered.push(event); + continue; + } const shortName = resolveCollectionKey(config, event.collection); const collection = shortName ? config.collections[shortName] : undefined; - if (collection?.discover === false || effectiveKnownDids.has(event.did)) { + if (collection?.discover !== false || effectiveKnownDids.has(event.did)) { + actorFiltered.push(event); continue; } - effectiveKnownDids.add(event.did); - discoveredDids.push(event.did); + attemptDropped.unknownActor++; } - } - const actorFiltered: IngestEvent[] = []; - for (const event of admitted) { - if (event.operation === "delete" || !effectiveKnownDids) { - actorFiltered.push(event); - continue; - } - const shortName = resolveCollectionKey(config, event.collection); - const collection = shortName ? config.collections[shortName] : undefined; - if (collection?.discover !== false || effectiveKnownDids.has(event.did)) { - actorFiltered.push(event); - continue; - } - dropped.unknownActor++; - } + const subjectFiltered = await filterUnknownSubjects( + db, + config, + actorFiltered, + effectiveKnownDids, + attemptDropped, + ); + const projectionEvents = [...subjectFiltered, ...projectionExclusions]; + const diagnosticCounts: IngestDiagnosticCounts = { + unknown_collection: attemptDropped.unknownCollection, + invalid_json: attemptDropped.invalidRecord, + lexicon_validation: attemptDropped.lexiconValidation, + cid_mismatch: attemptDropped.cidMismatch, + cid_encoding: attemptDropped.cidEncoding, + missing_cid: attemptDropped.missingCid, + record_filter: attemptDropped.recordFilter, + unknown_actor: attemptDropped.unknownActor, + unknown_subject: attemptDropped.unknownSubject, + superseded: attemptDropped.superseded, + }; + const diagnostics = options.aggregateDiagnostics + ? null + : ingestDiagnosticsStatement(db, diagnosticCounts); + const trailingStatements = [ + ...(diagnostics ? [diagnostics] : []), + ...(options.trailingStatements ?? []), + ]; - const subjectFiltered = await filterUnknownSubjects( - db, - config, - actorFiltered, - effectiveKnownDids, - dropped, - ); + try { + if (projectionEvents.length > 0) { + await projectEvents(db, projectionEvents, config, { + ...options, + phase: options.phase ?? "live", + // A pre-fetched visible-row map is not safe after a conflict; the + // projector deliberately reloads it under this predecessor attempt. + existing: undefined, + trailingStatements, + sourceOrderingChecked: true, + predecessors, + }); + } else if (trailingStatements.length > 0) { + await db.batch(trailingStatements); + } + } catch (error) { + if (isProjectionConflictError(error) && attempt < maximumAttempts) { + continue; + } + throw error; + } - const projectionEvents = [ - ...subjectFiltered, - ...projectionExclusions, - ]; - const diagnosticCounts: IngestDiagnosticCounts = { - unknown_collection: dropped.unknownCollection, - invalid_json: dropped.invalidRecord, - lexicon_validation: dropped.lexiconValidation, - cid_mismatch: dropped.cidMismatch, - cid_encoding: dropped.cidEncoding, - missing_cid: dropped.missingCid, - record_filter: dropped.recordFilter, - unknown_actor: dropped.unknownActor, - unknown_subject: dropped.unknownSubject, - superseded: dropped.superseded, - }; - const diagnostics = options.aggregateDiagnostics - ? null - : ingestDiagnosticsStatement(db, diagnosticCounts); - const trailingStatements = [ - ...(diagnostics ? [diagnostics] : []), - ...(options.trailingStatements ?? []), - ]; - if (projectionEvents.length > 0) { - await projectEvents(db, projectionEvents, config, { - ...options, - trailingStatements, - sourceOrderingChecked: true, - }); - } else if (trailingStatements.length > 0) { - await db.batch(trailingStatements); - } - // Aggregate only after the canonical projection/checkpoint transaction - // succeeds, so a rolled-back page cannot inflate private diagnostics. - if (options.aggregateDiagnostics) { - addIngestDiagnosticCounts(options.aggregateDiagnostics, diagnosticCounts); + Object.assign(dropped, attemptDropped); + // Aggregate only after the canonical projection/checkpoint transaction + // succeeds, so a rolled-back attempt cannot inflate private diagnostics. + if (options.aggregateDiagnostics) { + addIngestDiagnosticCounts(options.aggregateDiagnostics, diagnosticCounts); + } + return { accepted: subjectFiltered, dropped, discoveredDids }; } - return { accepted: subjectFiltered, dropped, discoveredDids }; + throw new Error("Projection conflict retry limit exhausted"); } function incrementValidationDrop( diff --git a/packages/contrail/src/core/jetstream.ts b/packages/contrail/src/core/jetstream.ts index b29c8d1..e1285e9 100644 --- a/packages/contrail/src/core/jetstream.ts +++ b/packages/contrail/src/core/jetstream.ts @@ -478,6 +478,7 @@ export async function runIngestCycle( const batch = events.slice(i, i + BATCH_SIZE); const isFinalBatch = i + BATCH_SIZE >= events.length; const result = await ingestRecords(db, batch, config, { + phase: "live", knownDids, // Earlier batches may commit without moving the cursor. A crash replays // them safely; the final batch atomically commits the exact source cursor. diff --git a/packages/contrail/src/core/persistent.ts b/packages/contrail/src/core/persistent.ts index 3b96cc6..877526b 100644 --- a/packages/contrail/src/core/persistent.ts +++ b/packages/contrail/src/core/persistent.ts @@ -161,6 +161,7 @@ async function streamAndFlush( let ingestResult: Awaited>; try { ingestResult = await ingestRecords(db, batch, config, { + phase: "live", knownDids, trailingStatements: [ saveCursorStatement(db, lastTimeUs), diff --git a/packages/contrail/src/core/router/notify.ts b/packages/contrail/src/core/router/notify.ts index 1f11282..3688a27 100644 --- a/packages/contrail/src/core/router/notify.ts +++ b/packages/contrail/src/core/router/notify.ts @@ -255,7 +255,7 @@ export async function processNotifyUris( } const appliedEvents = events.length > 0 - ? (await ingestRecords(db, events, config, { existing })).accepted + ? (await ingestRecords(db, events, config, { existing, phase: "live" })).accepted : []; // The shared ingest path fans these records into feed_items exactly like the cron and diff --git a/packages/contrail/src/core/router/profiles.ts b/packages/contrail/src/core/router/profiles.ts index abaa351..5bdf518 100644 --- a/packages/contrail/src/core/router/profiles.ts +++ b/packages/contrail/src/core/router/profiles.ts @@ -185,7 +185,9 @@ async function fetchMissingProfiles( const events = fetched.filter((event) => event !== null); if (events.length === 0) return {}; - const { accepted } = await ingestRecords(db, events, config); + const { accepted } = await ingestRecords(db, events, config, { + phase: "historical", + }); const result: Record = {}; for (const event of accepted) { if (event.operation === "delete") continue; diff --git a/packages/contrail/src/core/types.ts b/packages/contrail/src/core/types.ts index c472054..bb5860e 100644 --- a/packages/contrail/src/core/types.ts +++ b/packages/contrail/src/core/types.ts @@ -249,6 +249,28 @@ export interface OrderedSourceConfig { epoch: string; } +/** How an accepted mutation entered the logical projection. */ +export type ProjectionPhase = "historical" | "live"; + +export type ChangeConsumerInitialMode = "current" | "future" | "history"; + +/** Static, secret-free definition for one durable change-log consumer. */ +export interface ChangeConsumerConfig { + /** Exact configured collection NSIDs. Short aliases are deliberately rejected. */ + collections: string[]; + /** Projection phases to observe. Defaults to both historical and live. */ + phases?: ProjectionPhase[]; + /** How the consumer establishes its first durable position. */ + initial: ChangeConsumerInitialMode; + /** Whether deployment generation activation may require this consumer. */ + requiredForActivation?: boolean; +} + +export interface ChangeLogConfig { + /** Stable consumer IDs mapped to their static delivery policy. */ + consumers: Record; +} + export type AtprotoServiceAuthMethod = "getFeed" | "notifyOfUpdate"; export interface AtprotoServiceAuthConfig { @@ -284,6 +306,9 @@ export interface ContrailConfig { * cursor is persisted atomically with projected mutations and may be exposed * to clients as a cache invalidation coordinate. */ orderedSource?: OrderedSourceConfig; + /** Optional transactional projection change log. Runtime handlers and + * destination credentials are bound separately and never belong here. */ + changes?: ChangeLogConfig; feeds?: Record; logger?: Logger; /** Expose the notifyOfUpdate HTTP endpoint. Off by default. @@ -655,6 +680,58 @@ function validateShortName(short: string): void { } } +const CHANGE_CONSUMER_ID = /^[a-zA-Z][a-zA-Z0-9_-]{0,63}$/; +const MAX_CHANGE_CONSUMERS = 32; +const MAX_CHANGE_CONSUMER_COLLECTIONS = 64; +const MAX_CHANGE_COVERAGE_PAIRS = 256; +const MAX_CHANGE_DEFINITIONS_BYTES = 64 * 1_024; + +/** Whether this configuration requires the optional transactional change log. */ +export function changesEnabled(config: ContrailConfig): boolean { + return Object.keys(config.changes?.consumers ?? {}).length > 0; +} + +/** Canonical phases for a consumer definition. */ +export function changeConsumerPhases( + consumer: ChangeConsumerConfig, +): ProjectionPhase[] { + return consumer.phases ?? ["historical", "live"]; +} + +/** Canonical collection/phase pairs whose changes must be retained. */ +export function changeLogCoverage( + config: ContrailConfig, +): Array<{ collection: string; phase: ProjectionPhase }> { + const pairs = new Map(); + for (const consumer of Object.values(config.changes?.consumers ?? {})) { + for (const collection of consumer.collections) { + for (const phase of changeConsumerPhases(consumer)) { + pairs.set(`${collection}\0${phase}`, { collection, phase }); + } + } + } + return [...pairs.values()].sort( + (left, right) => + left.collection.localeCompare(right.collection) || + left.phase.localeCompare(right.phase), + ); +} + +/** Stable secret-free representation used by schema/config compatibility checks. */ +export function canonicalChangeDefinitions(config: ContrailConfig): string { + return JSON.stringify( + Object.entries(config.changes?.consumers ?? {}) + .sort(([left], [right]) => left.localeCompare(right)) + .map(([id, consumer]) => ({ + id, + collections: [...consumer.collections].sort(), + phases: [...changeConsumerPhases(consumer)].sort(), + initial: consumer.initial, + requiredForActivation: consumer.requiredForActivation === true, + })), + ); +} + export function validateConfig(config: ContrailConfig): void { const shortNames = new Set(); for (const [short, colConfig] of Object.entries(config.collections)) { @@ -715,6 +792,64 @@ export function validateConfig(config: ContrailConfig): void { } } + const consumers = Object.entries(config.changes?.consumers ?? {}); + if (consumers.length > MAX_CHANGE_CONSUMERS) { + throw new Error(`changes supports at most ${MAX_CHANGE_CONSUMERS} consumers`); + } + const configuredNsids = new Set(getCollectionNsids(config)); + for (const [id, consumer] of consumers) { + if (!CHANGE_CONSUMER_ID.test(id)) { + throw new Error( + `Invalid change consumer ID "${id}"; use 1-64 letters, digits, underscores, or hyphens`, + ); + } + if ( + !Array.isArray(consumer.collections) || + consumer.collections.length === 0 || + consumer.collections.length > MAX_CHANGE_CONSUMER_COLLECTIONS || + new Set(consumer.collections).size !== consumer.collections.length + ) { + throw new Error( + `Change consumer "${id}" requires 1-${MAX_CHANGE_CONSUMER_COLLECTIONS} unique collection NSIDs`, + ); + } + for (const collection of consumer.collections) { + if (!configuredNsids.has(collection)) { + throw new Error( + `Change consumer "${id}" references unconfigured collection NSID "${collection}"`, + ); + } + } + const phases = changeConsumerPhases(consumer); + if ( + phases.length === 0 || + phases.length > 2 || + new Set(phases).size !== phases.length || + phases.some((phase) => phase !== "historical" && phase !== "live") + ) { + throw new Error( + `Change consumer "${id}" requires unique historical/live phases`, + ); + } + if (!(["current", "future", "history"] as string[]).includes(consumer.initial)) { + throw new Error(`Change consumer "${id}" has an invalid initial mode`); + } + } + const coveragePairs = changeLogCoverage(config).length; + if (coveragePairs > MAX_CHANGE_COVERAGE_PAIRS) { + throw new Error( + `changes requires ${coveragePairs} collection/phase coverage pairs; maximum is ${MAX_CHANGE_COVERAGE_PAIRS}`, + ); + } + const definitionBytes = new TextEncoder().encode( + canonicalChangeDefinitions(config), + ).byteLength; + if (definitionBytes > MAX_CHANGE_DEFINITIONS_BYTES) { + throw new Error( + `changes definitions contain ${definitionBytes} encoded bytes; maximum is ${MAX_CHANGE_DEFINITIONS_BYTES}`, + ); + } + } // Helpers diff --git a/packages/contrail/src/index.ts b/packages/contrail/src/index.ts index bd254c2..0e68f6c 100644 --- a/packages/contrail/src/index.ts +++ b/packages/contrail/src/index.ts @@ -29,6 +29,12 @@ export * from "./core/persistent"; export * from "./core/backfill"; export * from "./core/status"; export * from "./core/diagnostics"; +export { + getChangeLogState, + MAX_CHANGE_BATCH_BYTES, + MAX_CHANGE_BATCH_CHANGES, +} from "./core/change-log"; +export type { ChangeLogState, RecordChange } from "./core/change-log"; export * from "./core/validation"; export * from "./core/search"; export * from "./core/constellation"; diff --git a/packages/contrail/tests/change-log.test.ts b/packages/contrail/tests/change-log.test.ts new file mode 100644 index 0000000..ead13a1 --- /dev/null +++ b/packages/contrail/tests/change-log.test.ts @@ -0,0 +1,493 @@ +import { describe, expect, it } from "vitest"; +import { + Contrail, + createIngestEvent, + getChangeLogState, + ingestRecords, + initSchema, + queryRecords, + resolveConfig, + saveCursorStatement, + type ContrailConfig, + type Database, + type Statement, +} from "../src/index"; +import { createSqliteDatabase } from "../src/adapters/sqlite"; + +const EVENT = "com.example.event"; +const NOTE = "com.example.note"; +const URI = `at://did:plc:alice/${EVENT}/one`; +const logger = { log() {}, warn() {}, error() {} }; + +function config(options: { + changes?: ContrailConfig["changes"]; + includeNote?: boolean; +} = {}) { + return resolveConfig({ + namespace: "com.example", + profiles: [], + logger, + collections: { + event: { collection: EVENT }, + ...(options.includeNote ? { note: { collection: NOTE } } : {}), + }, + changes: options.changes, + }); +} + +function loggedConfig() { + return config({ + changes: { + consumers: { + search: { + collections: [EVENT], + initial: "current", + requiredForActivation: true, + }, + webhooks: { + collections: [EVENT], + phases: ["live"], + initial: "future", + }, + }, + }, + }); +} + +function mutation(options: { + sourceTime: number; + cid?: string | null; + value?: Record; + operation?: "create" | "update" | "delete"; + revision?: string | null; +}) { + const operation = options.operation ?? "update"; + return createIngestEvent({ + uri: URI, + did: "did:plc:alice", + collection: EVENT, + rkey: "one", + operation, + cid: operation === "delete" ? null : (options.cid ?? `cid-${options.sourceTime}`), + value: + operation === "delete" + ? undefined + : (options.value ?? { name: `event-${options.sourceTime}` }), + timeUs: options.sourceTime, + indexedAt: options.sourceTime + 10_000, + source: { + id: "source", + epoch: "epoch", + time_us: options.sourceTime, + revision: options.revision ?? String(options.sourceTime), + cursor: String(options.sourceTime), + }, + }); +} + +async function batches(db: Database) { + return ( + await db + .prepare( + `SELECT generation_id, position, projection_transaction_id, source_id, + source_epoch, source_cursor, phase, changes_json, change_count + FROM change_batches ORDER BY position`, + ) + .all() + ).results; +} + +describe("transactional projection change log", () => { + it("validates bounded static consumer definitions", () => { + const base = { + namespace: "com.example", + profiles: [] as string[], + collections: { event: { collection: EVENT } }, + }; + expect( + () => + new Contrail({ + ...base, + changes: { + consumers: { + "bad id": { collections: [EVENT], initial: "future" }, + }, + }, + }), + ).toThrow("Invalid change consumer ID"); + expect( + () => + new Contrail({ + ...base, + changes: { + consumers: { + search: { collections: ["event"], initial: "future" }, + }, + }, + }), + ).toThrow("unconfigured collection NSID"); + expect( + () => + new Contrail({ + ...base, + changes: { + consumers: { + search: { + collections: [EVENT], + phases: ["live", "live"], + initial: "future", + }, + }, + }, + }), + ).toThrow("unique historical/live phases"); + }); + + it("has no change-log schema or writes when disabled", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = config(); + await initSchema(db, resolved); + + const tables = await db + .prepare( + "SELECT name FROM sqlite_master WHERE type = 'table' AND name LIKE 'change_%' ORDER BY name", + ) + .all<{ name: string }>(); + expect(tables.results).toEqual([]); + expect( + await db + .prepare( + "SELECT name FROM sqlite_master WHERE type = 'table' AND name = '_contrail_projection_state'", + ) + .first(), + ).not.toBeNull(); + + await ingestRecords(db, [mutation({ sourceTime: 1 })], resolved); + expect(await getChangeLogState(db)).toBeNull(); + }); + + it("initializes a fresh generation, registrations, and coverage ledger", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = loggedConfig(); + await initSchema(db, resolved); + + const state = await getChangeLogState(db); + expect(state).toMatchObject({ head: "0", retainedFloor: "0" }); + expect(state?.generation).toMatch(/^[0-9a-f-]{36}$/); + + const consumers = await db + .prepare( + `SELECT consumer_id, configured_collections_json, + configured_phases_json, initial_mode, + required_for_activation, bootstrap_state, + acknowledged_position, bootstrap_anchor_position + FROM change_consumers ORDER BY consumer_id`, + ) + .all(); + expect(consumers.results).toEqual([ + { + consumer_id: "search", + configured_collections_json: JSON.stringify([EVENT]), + configured_phases_json: JSON.stringify(["historical", "live"]), + initial_mode: "current", + required_for_activation: 1, + bootstrap_state: "pending", + acknowledged_position: 0, + bootstrap_anchor_position: 0, + }, + { + consumer_id: "webhooks", + configured_collections_json: JSON.stringify([EVENT]), + configured_phases_json: JSON.stringify(["live"]), + initial_mode: "future", + required_for_activation: 0, + bootstrap_state: "ready", + acknowledged_position: 0, + bootstrap_anchor_position: null, + }, + ]); + + const coverage = await db + .prepare( + `SELECT collection, phase, from_position, through_position + FROM change_log_coverage ORDER BY collection, phase`, + ) + .all(); + expect(coverage.results).toEqual([ + { + collection: EVENT, + phase: "historical", + from_position: 0, + through_position: null, + }, + { + collection: EVENT, + phase: "live", + from_position: 0, + through_position: null, + }, + ]); + + // Initialization is idempotent and retains one random database generation. + await initSchema(db, resolved); + expect((await getChangeLogState(db))?.generation).toBe(state?.generation); + }); + + it("appends only committed logical current-state changes", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = loggedConfig(); + await initSchema(db, resolved); + + const original = mutation({ + sourceTime: 100, + cid: "cid-one", + value: { name: "one" }, + }); + await ingestRecords(db, [original], resolved, { phase: "historical" }); + expect((await getChangeLogState(db))?.head).toBe("1"); + + // Exact replay and a newer source observation of the same immutable state + // may update source metadata but do not wake current-state consumers. + await ingestRecords(db, [original], resolved, { phase: "historical" }); + await ingestRecords( + db, + [ + mutation({ + sourceTime: 110, + cid: "cid-one", + value: { name: "one" }, + }), + ], + resolved, + { phase: "live" }, + ); + expect((await getChangeLogState(db))?.head).toBe("1"); + + await ingestRecords( + db, + [ + mutation({ + sourceTime: 200, + cid: "cid-two", + value: { name: "two" }, + }), + ], + resolved, + { phase: "live" }, + ); + await ingestRecords( + db, + [mutation({ sourceTime: 300, operation: "delete" })], + resolved, + { phase: "live" }, + ); + await ingestRecords( + db, + [mutation({ sourceTime: 400, operation: "delete" })], + resolved, + { phase: "live" }, + ); + + const rows = await batches(db); + expect(rows).toHaveLength(3); + expect(rows.map((row) => [row.position, row.phase, row.change_count])).toEqual([ + [1, "historical", 1], + [2, "live", 1], + [3, "live", 1], + ]); + const changes = rows.map((row) => JSON.parse(row.changes_json)[0]); + expect(changes[0]).toMatchObject({ + kind: "record", + operation: "put", + uri: URI, + cid: "cid-one", + version: { sourceId: "source", sourceTimeUs: 100 }, + }); + expect(changes[0]).not.toHaveProperty("id"); + expect(changes[0]).not.toHaveProperty("record"); + expect(changes[2]).toMatchObject({ + operation: "delete", + uri: URI, + cid: null, + }); + expect((await getChangeLogState(db))?.head).toBe("3"); + }); + + it("reduces multiple mutations for one URI to the final state", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = loggedConfig(); + await initSchema(db, resolved); + + await ingestRecords( + db, + [ + mutation({ sourceTime: 1, cid: "cid-one" }), + mutation({ sourceTime: 2, cid: "cid-two" }), + mutation({ sourceTime: 3, cid: "cid-three" }), + ], + resolved, + ); + + const rows = await batches(db); + expect(rows).toHaveLength(1); + expect(rows[0].change_count).toBe(1); + expect(JSON.parse(rows[0].changes_json)[0].cid).toBe("cid-three"); + }); + + it("does not log a phase outside every consumer's coverage", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = config({ + changes: { + consumers: { + webhook: { + collections: [EVENT], + phases: ["live"], + initial: "future", + }, + }, + }, + }); + await initSchema(db, resolved); + + await ingestRecords(db, [mutation({ sourceTime: 1 })], resolved, { + phase: "historical", + }); + expect((await getChangeLogState(db))?.head).toBe("0"); + expect(await batches(db)).toEqual([]); + }); + + it("rolls projection, log head, batch, and checkpoint back together", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = loggedConfig(); + await initSchema(db, resolved); + await db.prepare("CREATE TABLE change_failure (value TEXT UNIQUE)").run(); + await db.prepare("INSERT INTO change_failure VALUES ('duplicate')").run(); + + await expect( + ingestRecords(db, [mutation({ sourceTime: 1 })], resolved, { + trailingStatements: [ + saveCursorStatement(db, 1), + db.prepare("INSERT INTO change_failure VALUES ('duplicate')"), + ], + }), + ).rejects.toThrow(); + + expect((await getChangeLogState(db))?.head).toBe("0"); + expect(await batches(db)).toEqual([]); + expect( + await db.prepare("SELECT uri FROM records_event").first(), + ).toBeNull(); + expect(await db.prepare("SELECT time_us FROM cursor").first()).toBeNull(); + }); + + it("gives each fresh database a distinct generation", async () => { + const resolved = loggedConfig(); + const first = createSqliteDatabase(":memory:"); + const second = createSqliteDatabase(":memory:"); + await initSchema(first, resolved); + await initSchema(second, resolved); + expect((await getChangeLogState(first))?.generation).not.toBe( + (await getChangeLogState(second))?.generation, + ); + }); + + it("fails closed for unsafe enable, disable, and definition changes", async () => { + const populated = createSqliteDatabase(":memory:"); + const disabled = config(); + await initSchema(populated, disabled); + await ingestRecords(populated, [mutation({ sourceTime: 1 })], disabled); + await expect(initSchema(populated, loggedConfig())).rejects.toThrow( + "fresh empty generation", + ); + expect( + await populated + .prepare( + "SELECT name FROM sqlite_master WHERE type='table' AND name='change_log_state'", + ) + .first(), + ).toBeNull(); + + const initialized = createSqliteDatabase(":memory:"); + await initSchema(initialized, loggedConfig()); + await expect(initSchema(initialized, config())).rejects.toThrow( + "cannot be disabled", + ); + const changed = config({ + changes: { + consumers: { + replacement: { collections: [EVENT], initial: "future" }, + }, + }, + }); + await expect(initSchema(initialized, changed)).rejects.toThrow( + "definitions differ", + ); + }); + + it("retries a losing overlapping writer from fresh durable state", async () => { + const real = createSqliteDatabase(":memory:"); + const resolved = loggedConfig(); + await initSchema(real, resolved); + + let arrivals = 0; + let releaseBoth!: () => void; + const both = new Promise((resolve) => { + releaseBoth = resolve; + }); + let releaseNewer!: () => void; + const newerDone = new Promise((resolve) => { + releaseNewer = resolve; + }); + + function overlapping(role: "older" | "newer"): Database { + let projectionBatches = 0; + return { + prepare(sql: string): Statement { + return real.prepare(sql); + }, + async batch(statements: Statement[]) { + projectionBatches++; + if (projectionBatches > 1) return real.batch(statements); + arrivals++; + if (arrivals === 2) releaseBoth(); + await both; + if (role === "older") await newerDone; + try { + return await real.batch(statements); + } finally { + if (role === "newer") releaseNewer(); + } + }, + dialect: real.dialect, + }; + } + + const olderDb = overlapping("older"); + const newerDb = overlapping("newer"); + const [older] = await Promise.all([ + ingestRecords( + olderDb, + [mutation({ sourceTime: 100, cid: "cid-old", value: { name: "old" } })], + resolved, + { trailingStatements: [saveCursorStatement(real, 100)] }, + ), + ingestRecords( + newerDb, + [mutation({ sourceTime: 200, cid: "cid-new", value: { name: "new" } })], + resolved, + { trailingStatements: [saveCursorStatement(real, 200)] }, + ), + ]); + + expect(older.dropped.superseded).toBe(1); + const visible = await queryRecords(real, resolved, { collection: "event" }); + expect(JSON.parse(visible.records[0]!.record!).name).toBe("new"); + expect((await getChangeLogState(real))?.head).toBe("1"); + expect(JSON.parse((await batches(real))[0].changes_json)[0].cid).toBe( + "cid-new", + ); + expect( + await real.prepare("SELECT time_us FROM cursor WHERE id = 1").first(), + ).toEqual({ time_us: 200 }); + }); +}); diff --git a/packages/contrail/tests/ingest.test.ts b/packages/contrail/tests/ingest.test.ts index eefa59b..c70cd74 100644 --- a/packages/contrail/tests/ingest.test.ts +++ b/packages/contrail/tests/ingest.test.ts @@ -106,7 +106,9 @@ describe("ingestRecords", () => { const versionInserts = sql.filter((statement) => statement.startsWith("INSERT INTO record_versions"), ); - expect(versionInserts).toHaveLength(4); + // Thirteen bindings include the optimistic projection token, so seven + // versions fit under D1's 100-binding statement ceiling. + expect(versionInserts).toHaveLength(5); expect( versionInserts.every( (statement) => (statement.match(/\?/g) ?? []).length <= 100, diff --git a/packages/contrail/tests/postgres-concurrent-init.test.ts b/packages/contrail/tests/postgres-concurrent-init.test.ts index 954bbb3..dd6fdae 100644 --- a/packages/contrail/tests/postgres-concurrent-init.test.ts +++ b/packages/contrail/tests/postgres-concurrent-init.test.ts @@ -71,7 +71,7 @@ if (!PG_URL) { const tables = await pool.query( `SELECT tablename FROM pg_tables WHERE schemaname = 'public' AND (tablename LIKE 'records_%' OR tablename LIKE 'fts_%' - OR tablename IN ('_contrail_meta', 'backfills', 'backfill_state', 'discovery', 'cursor', 'identities', 'record_versions', 'ingest_diagnostics', 'feed_items', 'feed_prune_cursor', 'feed_backfills'))` + OR tablename IN ('_contrail_meta', '_contrail_projection_state', 'backfills', 'backfill_state', 'discovery', 'cursor', 'source_position', 'bootstrap_state', 'bootstrap_snapshot_progress', 'identities', 'record_versions', 'ingest_diagnostics', 'feed_items', 'feed_prune_cursor', 'feed_backfills', 'change_log_state', 'change_batches', 'change_consumers', 'change_log_coverage'))` ); for (const { tablename } of tables.rows) { await pool.query(`DROP TABLE IF EXISTS ${tablename} CASCADE`); diff --git a/packages/contrail/tests/postgres-e2e.test.ts b/packages/contrail/tests/postgres-e2e.test.ts index e5032ea..0e5d499 100644 --- a/packages/contrail/tests/postgres-e2e.test.ts +++ b/packages/contrail/tests/postgres-e2e.test.ts @@ -18,7 +18,7 @@ import { saveCursor, } from "../src/index"; import { resolveConfig } from "../src/index"; -import type { Database } from "../src/index"; +import type { Database, Statement } from "../src/index"; import { resolveHydrates, resolveReferences } from "../src/index"; import { makeEvent } from "./helpers"; @@ -78,7 +78,7 @@ if (!PG_URL) { const tables = await pool.query( `SELECT tablename FROM pg_tables WHERE schemaname = 'public' AND (tablename LIKE 'records_%' OR tablename LIKE 'fts_%' - OR tablename IN ('_contrail_meta', 'backfills', 'backfill_state', 'discovery', 'cursor', 'identities', 'record_versions', 'ingest_diagnostics', 'feed_items', 'feed_prune_cursor', 'feed_backfills'))` + OR tablename IN ('_contrail_meta', '_contrail_projection_state', 'backfills', 'backfill_state', 'discovery', 'cursor', 'source_position', 'bootstrap_state', 'bootstrap_snapshot_progress', 'identities', 'record_versions', 'ingest_diagnostics', 'feed_items', 'feed_prune_cursor', 'feed_backfills', 'change_log_state', 'change_batches', 'change_consumers', 'change_log_coverage'))` ); for (const { tablename } of tables.rows) { await pool.query(`DROP TABLE IF EXISTS ${tablename} CASCADE`); @@ -171,6 +171,81 @@ if (!PG_URL) { expect(result.records).toHaveLength(0); }); + it("retries an overlapping stale projector after the transaction lock", async () => { + const uri = "at://did:plc:test/community.lexicon.calendar.event/concurrent"; + let arrivals = 0; + let releaseBoth!: () => void; + const both = new Promise((resolve) => { + releaseBoth = resolve; + }); + let releaseNewer!: () => void; + const newerDone = new Promise((resolve) => { + releaseNewer = resolve; + }); + + const overlap = (role: "older" | "newer"): Database => { + let writes = 0; + return { + prepare(sql: string): Statement { + return db.prepare(sql); + }, + async batch(statements: Statement[]) { + writes++; + if (writes > 1) return db.batch(statements); + arrivals++; + if (arrivals === 2) releaseBoth(); + await both; + if (role === "older") await newerDone; + try { + return await db.batch(statements); + } finally { + if (role === "newer") releaseNewer(); + } + }, + dialect: db.dialect, + }; + }; + + const [older] = await Promise.all([ + ingestRecords( + overlap("older"), + [ + makeEvent({ + uri, + rkey: "concurrent", + cid: "cid-old", + record: { name: "old", mode: "online" }, + time_us: 100, + indexed_at: 100, + }), + ], + TEST_CONFIG, + ), + ingestRecords( + overlap("newer"), + [ + makeEvent({ + uri, + rkey: "concurrent", + cid: "cid-new", + record: { name: "new", mode: "online" }, + time_us: 200, + indexed_at: 200, + }), + ], + TEST_CONFIG, + ), + ]); + + expect(older.dropped.superseded).toBe(1); + const row = await pool.query( + "SELECT record, cid FROM records_community_lexicon_calendar_event WHERE uri = $1", + [uri], + ); + expect(row.rows[0]).toMatchObject({ cid: "cid-new" }); + expect(row.rows[0].record.name).toBe("new"); + }); + it("does nothing for empty events", async () => { await ingestRecords(db, []); const cursor = await getLastCursor(db); diff --git a/packages/contrail/tests/postgres.test.ts b/packages/contrail/tests/postgres.test.ts index 2738dc6..0b7b530 100644 --- a/packages/contrail/tests/postgres.test.ts +++ b/packages/contrail/tests/postgres.test.ts @@ -57,7 +57,7 @@ if (!PG_URL) { const tables = await pool.query( `SELECT tablename FROM pg_tables WHERE schemaname = 'public' AND (tablename LIKE 'records_%' OR tablename LIKE 'fts_%' - OR tablename IN ('_contrail_meta', 'backfills', 'backfill_state', 'discovery', 'cursor', 'identities', 'record_versions', 'ingest_diagnostics', 'feed_items', 'feed_prune_cursor', 'feed_backfills'))` + OR tablename IN ('_contrail_meta', '_contrail_projection_state', 'backfills', 'backfill_state', 'discovery', 'cursor', 'source_position', 'bootstrap_state', 'bootstrap_snapshot_progress', 'identities', 'record_versions', 'ingest_diagnostics', 'feed_items', 'feed_prune_cursor', 'feed_backfills', 'change_log_state', 'change_batches', 'change_consumers', 'change_log_coverage'))` ); for (const { tablename } of tables.rows) { await pool.query(`DROP TABLE IF EXISTS ${tablename} CASCADE`); diff --git a/packages/contrail/tests/schema.test.ts b/packages/contrail/tests/schema.test.ts index 15643dc..3a95149 100644 --- a/packages/contrail/tests/schema.test.ts +++ b/packages/contrail/tests/schema.test.ts @@ -21,6 +21,8 @@ describe("initSchema", () => { expect(names).toContain("identities"); expect(names).toContain("record_versions"); expect(names).toContain("ingest_diagnostics"); + expect(names).toContain("_contrail_projection_state"); + expect(names).not.toContain("change_log_state"); const backfillColumns = await db .prepare("PRAGMA table_info(backfills)") @@ -115,6 +117,52 @@ describe("initSchema", () => { ); }); + it("migrates optimistic tokens onto existing version rows", async () => { + const db = createTestDb(); + await db + .prepare( + `CREATE TABLE record_versions ( + uri TEXT PRIMARY KEY, + did TEXT NOT NULL, + collection TEXT NOT NULL, + rkey TEXT NOT NULL, + operation TEXT NOT NULL, + cid TEXT, + source_id TEXT NOT NULL, + source_revision TEXT, + source_time_us BIGINT NOT NULL, + source_cursor TEXT, + indexed_at BIGINT NOT NULL + )`, + ) + .run(); + await db + .prepare( + `INSERT INTO record_versions + (uri, did, collection, rkey, operation, cid, source_id, + source_revision, source_time_us, source_cursor, indexed_at) + VALUES (?, ?, ?, ?, 'update', ?, 'legacy', NULL, 1, NULL, 1)`, + ) + .bind( + "at://did:plc:legacy/community.lexicon.calendar.event/one", + "did:plc:legacy", + "community.lexicon.calendar.event", + "one", + "cid-legacy", + ) + .run(); + + await initSchema(db, TEST_CONFIG); + expect( + await db + .prepare("SELECT projection_token FROM record_versions") + .first(), + ).toEqual({ + projection_token: + "at://did:plc:legacy/community.lexicon.calendar.event/one", + }); + }); + it("seeds ordering metadata for visible rows from older schemas", async () => { const db = createTestDb(); await db @@ -149,7 +197,7 @@ describe("initSchema", () => { const version = await db .prepare( - "SELECT operation, source_id, source_time_us, cid FROM record_versions WHERE uri = ?", + "SELECT operation, source_id, source_time_us, cid, projection_token FROM record_versions WHERE uri = ?", ) .bind("at://did:plc:legacy/community.lexicon.calendar.event/one") .first<{ @@ -157,12 +205,15 @@ describe("initSchema", () => { source_id: string; source_time_us: number; cid: string; + projection_token: string; }>(); expect(version).toEqual({ operation: "update", source_id: "legacy", source_time_us: 200, cid: "cid-legacy", + projection_token: + "at://did:plc:legacy/community.lexicon.calendar.event/one", }); }); diff --git a/turbo.json b/turbo.json index 94b9292..2279a64 100644 --- a/turbo.json +++ b/turbo.json @@ -10,7 +10,7 @@ "dependsOn": ["^build"] }, "test": { - "dependsOn": ["^build"], + "dependsOn": ["^build", "build"], "outputs": [] }, "dev": { -- 2.51.2