diff --git a/.changeset/durable-projection-log.md b/.changeset/durable-projection-log.md index 08d7e18..7318ebc 100644 --- a/.changeset/durable-projection-log.md +++ b/.changeset/durable-projection-log.md @@ -4,4 +4,4 @@ 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. Add independent bounded consumer leases, filtered/coalesced claims, set-oriented current-state hydration, CAS acknowledgement, failure backoff, lease renewal, private status, and manual retry APIs. Add crash-safe current-state snapshot/tail/activation bootstrap, safe additive consumers over existing coverage, required-consumer readiness gates, consumer-aware bounded pruning, and audited explicit skip operations. Private `contrail changes` commands cover status, retry, prune, and skip. Add fair bounded Worker delivery after ingestion/retries, best-effort immediate notify wakes, runtime handler validation, deadline cancellation, isolated retry scheduling, and a persistent delivery supervisor. Include an app-owned atmo.rsvp Meilisearch reference consumer with task-success acknowledgement, hidden/delete convergence, candidate-index snapshot/tail bootstrap, and idempotent generation-marker activation. Enabling or expanding log coverage on a populated generation fails closed pending explicit quiet-boundary migration tooling. +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. Add independent bounded consumer leases, filtered/coalesced claims, set-oriented current-state hydration, CAS acknowledgement, failure backoff, lease renewal, private status, and manual retry APIs. Add crash-safe current-state snapshot/tail/activation bootstrap, safe additive consumers over existing coverage, required-consumer readiness gates, consumer-aware bounded pruning, and audited explicit skip operations. Current-state consumers require both projection phases, and candidate destination tokens are scoped strictly to bootstrap deliveries. Private `contrail changes` commands cover status, retry, prune, and skip. Add fair bounded Worker delivery after ingestion/retries, best-effort immediate notify wakes, runtime handler validation, deadline cancellation, isolated retry scheduling, and a persistent delivery supervisor. Include an app-owned atmo.rsvp Meilisearch reference consumer with task-success acknowledgement, hidden/delete convergence, candidate-index snapshot/tail bootstrap, and idempotent generation-marker activation. Enabling or expanding 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 22ac887..947a24d 100644 --- a/packages/contrail/README.md +++ b/packages/contrail/README.md @@ -149,7 +149,7 @@ if (claim) { Claims coalesce repeated URIs, hydrate in set-oriented collection queries, and resolve delete/recreate races from newest canonical state. Consumers lease and progress independently; irrelevant position ranges advance without invoking a handler. Delivery is intentionally at least once—a destination success followed by an acknowledgement crash causes duplicate delivery. Handlers must be idempotent by stable record/document key. -`initial: "current"` uses a durable snapshot-plus-tail coordinator. Repeatedly claim and idempotently acknowledge `contrail.changes.claimSnapshotPage()`, then drain `claimBootstrapChanges()` through its fixed target using ordinary hydrate/ack. Finally claim the stable generation-scoped activation token with `claimActivation()`, perform an idempotent destination swap, and call `completeActivation()`. A crash replays the same URI page, tail range, or activation token. Records updated or deleted while the keyset scan races are corrected by the anchored tail. +`initial: "current"` uses a durable snapshot-plus-tail coordinator and must observe both `historical` and `live` phases so every mutation racing the keyset scan is retained for reconciliation. Repeatedly claim and idempotently acknowledge `contrail.changes.claimSnapshotPage()`, then drain `claimBootstrapChanges()` through its fixed target using ordinary hydrate/ack. Snapshot and bootstrap-tail deliveries carry the candidate destination token; ordinary deliveries after activation do not. Finally claim the stable generation-scoped activation token with `claimActivation()`, perform an idempotent destination swap, and call `completeActivation()`. A crash replays the same URI page, tail range, or activation token. Records updated or deleted while the keyset scan races are corrected by the anchored tail. A consumer can be added to a populated log when all of its collection/phase pairs were already covered; its current/future anchor is the atomic current head, while history starts at the retained floor. Expanding coverage still fails closed without a fresh generation or explicit old-writer quiet boundary. `contrail changes status`, `retry`, `prune`, and explicitly confirmed `skip` expose private operations for SQLite or Wrangler D1 deployments. Status includes a conservative projection/head/batch/ack write plan. Pruning is bounded by the slowest durable consumer/bootstrap anchor. Skip records a bounded private audit reason and never occurs implicitly. Disabling or removing an existing log remains fail-closed. With no configured consumers, no change-log tables or append writes exist. diff --git a/packages/contrail/src/core/change-log.ts b/packages/contrail/src/core/change-log.ts index 7db6713..1220a3a 100644 --- a/packages/contrail/src/core/change-log.ts +++ b/packages/contrail/src/core/change-log.ts @@ -363,7 +363,7 @@ export async function initializeChangeLog( const statements: Statement[] = []; for (const [consumerId, consumer] of Object.entries( config.changes?.consumers ?? {}, - ).sort(([left], [right]) => left.localeCompare(right))) { + ).sort(([left], [right]) => (left < right ? -1 : left > right ? 1 : 0))) { const initialReady = consumer.initial === "current" ? "pending" : "ready"; statements.push( db @@ -444,17 +444,17 @@ async function assertChangeLogDefinition( FROM change_consumers ORDER BY consumer_id`, ) .all(); - const expectedConsumers = Object.entries(config.changes?.consumers ?? {}).sort( - ([left], [right]) => left.localeCompare(right), - ); + const expectedConsumers = Object.entries(config.changes?.consumers ?? {}); 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 consumersById = new Map( + consumers.results.map((consumer) => [consumer.consumer_id, consumer]), + ); + for (const [id, expected] of expectedConsumers) { + const actual = consumersById.get(id); if ( - actual.consumer_id !== id || + !actual || actual.generation_id !== state.generation_id || actual.configured_collections_json !== canonicalCollections(expected.collections) || actual.configured_phases_json !== canonicalPhases(changeConsumerPhases(expected)) || @@ -477,13 +477,14 @@ async function assertChangeLogDefinition( 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]!; + const coverageByPair = new Map( + coverage.results.map((item) => [`${item.collection}\0${item.phase}`, item]), + ); + for (const expected of expectedCoverage) { + const actual = coverageByPair.get(`${expected.collection}\0${expected.phase}`); if ( + !actual || actual.generation_id !== state.generation_id || - actual.collection !== expected.collection || - actual.phase !== expected.phase || Number(actual.from_position) !== 0 || actual.through_position !== null ) { diff --git a/packages/contrail/src/core/changes.ts b/packages/contrail/src/core/changes.ts index 636684f..3627147 100644 --- a/packages/contrail/src/core/changes.ts +++ b/packages/contrail/src/core/changes.ts @@ -631,9 +631,9 @@ async function claimChangeRange( changes: [...coalesced.values()], attempt: Number(leased.attempts) + 1, leaseExpiresAt, - ...(leased.bootstrap_token === null - ? {} - : { bootstrapToken: leased.bootstrap_token }), + ...(bootstrap && leased.bootstrap_token !== null + ? { bootstrapToken: leased.bootstrap_token } + : {}), ...(bootstrapTarget === null ? {} : { bootstrapTarget }), leaseOwner: owner, }; diff --git a/packages/contrail/src/core/types.ts b/packages/contrail/src/core/types.ts index bb5860e..0544dde 100644 --- a/packages/contrail/src/core/types.ts +++ b/packages/contrail/src/core/types.ts @@ -258,7 +258,7 @@ export type ChangeConsumerInitialMode = "current" | "future" | "history"; export interface ChangeConsumerConfig { /** Exact configured collection NSIDs. Short aliases are deliberately rejected. */ collections: string[]; - /** Projection phases to observe. Defaults to both historical and live. */ + /** Projection phases to observe. Defaults to both; `initial: "current"` requires both. */ phases?: ProjectionPhase[]; /** How the consumer establishes its first durable position. */ initial: ChangeConsumerInitialMode; @@ -686,6 +686,11 @@ const MAX_CHANGE_CONSUMER_COLLECTIONS = 64; const MAX_CHANGE_COVERAGE_PAIRS = 256; const MAX_CHANGE_DEFINITIONS_BYTES = 64 * 1_024; +/** Locale-independent order for durable definitions compared across runtimes. */ +function compareCanonicalText(left: string, right: string): number { + return left < right ? -1 : left > right ? 1 : 0; +} + /** Whether this configuration requires the optional transactional change log. */ export function changesEnabled(config: ContrailConfig): boolean { return Object.keys(config.changes?.consumers ?? {}).length > 0; @@ -712,8 +717,8 @@ export function changeLogCoverage( } return [...pairs.values()].sort( (left, right) => - left.collection.localeCompare(right.collection) || - left.phase.localeCompare(right.phase), + compareCanonicalText(left.collection, right.collection) || + compareCanonicalText(left.phase, right.phase), ); } @@ -721,7 +726,7 @@ export function changeLogCoverage( export function canonicalChangeDefinitions(config: ContrailConfig): string { return JSON.stringify( Object.entries(config.changes?.consumers ?? {}) - .sort(([left], [right]) => left.localeCompare(right)) + .sort(([left], [right]) => compareCanonicalText(left, right)) .map(([id, consumer]) => ({ id, collections: [...consumer.collections].sort(), @@ -831,6 +836,14 @@ export function validateConfig(config: ContrailConfig): void { `Change consumer "${id}" requires unique historical/live phases`, ); } + if ( + consumer.initial === "current" && + !(phases.includes("historical") && phases.includes("live")) + ) { + throw new Error( + `Current-state change consumer "${id}" must observe both historical and live phases`, + ); + } if (!(["current", "future", "history"] as string[]).includes(consumer.initial)) { throw new Error(`Change consumer "${id}" has an invalid initial mode`); } diff --git a/packages/contrail/tests/change-bootstrap.test.ts b/packages/contrail/tests/change-bootstrap.test.ts index 435ede2..8a23e1b 100644 --- a/packages/contrail/tests/change-bootstrap.test.ts +++ b/packages/contrail/tests/change-bootstrap.test.ts @@ -206,6 +206,9 @@ describe("current-state change consumer bootstrap", () => { await apply(db, withSearch, event({ rkey: "d", time: 6 })); const ordinary = await claimChanges(db, "search", { now: 400 }); expect(ordinary).toMatchObject({ from: "2", through: "3" }); + expect(ordinary).not.toHaveProperty("bootstrapToken"); + const ordinaryDelivery = await hydrateChanges(db, withSearch, ordinary!); + expect(ordinaryDelivery).not.toHaveProperty("destinationToken"); }); it("persists snapshot failure backoff and resumes the same page", async () => { @@ -243,7 +246,7 @@ describe("current-state change consumer bootstrap", () => { await initSchema(db, eventOnly); const expanded = config({ keeper: { collections: [EVENT], phases: ["live"], initial: "history" }, - notes: { collections: [NOTE], phases: ["live"], initial: "current" }, + notes: { collections: [NOTE], initial: "current" }, }); await expect(initSchema(db, expanded)).rejects.toThrow( "expands collection/phase coverage", diff --git a/packages/contrail/tests/change-log.test.ts b/packages/contrail/tests/change-log.test.ts index bb5a3c3..17f3303 100644 --- a/packages/contrail/tests/change-log.test.ts +++ b/packages/contrail/tests/change-log.test.ts @@ -142,6 +142,21 @@ describe("transactional projection change log", () => { }, }), ).toThrow("unique historical/live phases"); + expect( + () => + new Contrail({ + ...base, + changes: { + consumers: { + search: { + collections: [EVENT], + phases: ["live"], + initial: "current", + }, + }, + }, + }), + ).toThrow("must observe both historical and live phases"); }); it("has no change-log schema or writes when disabled", async () => { @@ -412,6 +427,25 @@ describe("transactional projection change log", () => { ); }); + it("matches durable consumer IDs independently of locale sort order", async () => { + const db = createSqliteDatabase(":memory:"); + const resolved = config({ + changes: { + consumers: { + a_: { collections: [EVENT], initial: "future" }, + "a-": { collections: [EVENT], initial: "future" }, + }, + }, + }); + + await initSchema(db, resolved); + await expect(initSchema(db, resolved)).resolves.toBeUndefined(); + const rows = await db + .prepare("SELECT consumer_id FROM change_consumers ORDER BY consumer_id") + .all<{ consumer_id: string }>(); + expect(rows.results.map((row) => row.consumer_id)).toEqual(["a-", "a_"]); + }); + it("fails closed for unsafe enable, disable, and definition changes", async () => { const populated = createSqliteDatabase(":memory:"); const disabled = config();