diff --git a/.changeset/preserve-resume-baseline-integrity.md b/.changeset/preserve-resume-baseline-integrity.md new file mode 100644 index 0000000000..2cb5517d15 --- /dev/null +++ b/.changeset/preserve-resume-baseline-integrity.md @@ -0,0 +1,6 @@ +--- +'@tanstack/db-sqlite-persistence-core': patch +'@tanstack/electric-db-collection': patch +--- + +Preserve persisted resume integrity with atomic SQLite baseline evidence and stale-writer rejection, and refresh uncertified Electric baselines before publishing resumed data. diff --git a/docs/contributing/oracle-coverage.md b/docs/contributing/oracle-coverage.md index 763060ab54..09dc733978 100644 --- a/docs/contributing/oracle-coverage.md +++ b/docs/contributing/oracle-coverage.md @@ -52,9 +52,9 @@ comment and the current API/architecture contract before extending its model. | Ordered acquisition | [pagination](../../packages/db/tests/query/pagination-oracle.property.test.ts), [ordered work](../../packages/db/tests/query/ordered-work-oracle.property.test.ts), [ordered lifecycle](../../packages/db/tests/query/ordered-lifecycle-oracle.property.test.ts) | Complete finite provider results, inherited collation with exact own-key request options, real lexical/numeric disagreement, pending windows, ties/nulls, ownership and documented repair timing. Request completion is not proof of unrequested source extent. | | Join equality and cold acquisition | `packages/db/tests/query/cold-join-reconciliation-oracle.test.ts` | Independent recomputation for cold acquisition plus direct join/predicate equivalence across established equality domains. Binary/string and nullish classes, replacement histories, raw on-demand values, and both scan/auto-index paths are explicit; compound join syntax is not claimed. | | Opaque backend pagination | [window oracle](../../packages/query-db-collection/tests/cursor-pagination.oracle.test.ts), [cache histories](../../packages/query-db-collection/tests/cursor-pagination.cache-oracle.test.ts), [cache publication](../../packages/query-db-collection/tests/cursor-pagination.publication-oracle.test.ts), [browser acquisition boundaries](../../packages/query-db-collection/tests/cursor-pagination.boundary-oracle.test.ts), [QueryCollection integration](../../packages/query-db-collection/tests/cursor-pagination.integration.test.ts) | Full filter/sort/slice reference, opaque token transport, actual Query cache expiry/invalidation/GC, forced refresh during growth, protocol failure publication/recovery, bounded slice work, nested cancellation/replacement, reader abort, browser retry defaults, manual-write cache isolation, and production window publications. Stable backend sequences; not snapshot guarantees for changing endpoints. Peek-ahead remains enabled. | -| Electric and TrailBase | [Electric histories](../../packages/electric-db-collection/tests/electric-oracle.property.test.ts), [PostgreSQL semantics](../../packages/electric-db-collection/e2e/sql-predicate-semantics.e2e.test.ts), [TrailBase contract](../../packages/trailbase-db-collection/tests/ORACLE.md) | Installed SDK delivery/framing, independent predicates, exact subscription arguments and late errors. SDK fixtures and a real service test earn different credit. | +| Electric and TrailBase | [Electric histories](../../packages/electric-db-collection/tests/electric-oracle.property.test.ts), [recovery histories](../../packages/electric-db-collection/tests/electric-recovery-oracle.test.ts), [held resume snapshots](../../packages/electric-db-collection/tests/electric-resume-snapshot-races.test.ts), [PostgreSQL semantics](../../packages/electric-db-collection/e2e/sql-predicate-semantics.e2e.test.ts), [TrailBase contract](../../packages/trailbase-db-collection/tests/ORACLE.md) | Installed SDK delivery/framing, independent predicates, exact subscription arguments, restart/reset lineage, held certification races, and late errors. The recovery fixtures use a mocked ShapeStream; they do not establish live Electric-service framing or native persistence-host behavior. | | PowerSync | [tests](../../packages/powersync-db-collection/tests), `tests/correctness-oracle.test.ts` | Applied receipt positions crossed with held peers, native SQLite/SDK and cleanup evidence. Run the focused owner with the package's `test:oracles` command. A timeout mutant proves a progress failure, not every value assertion. | -| SQLite persistence and native hosts | [persisted histories](../../packages/db-sqlite-persistence-core/tests/persisted.test.ts), [driver contracts](../../packages/db-sqlite-persistence-core/tests/contracts/sqlite-driver-contract.ts), [browser OPFS lifecycle](../../packages/browser-db-sqlite-persistence/tests/opfs-page-lifecycle-oracle.test.ts), [worker diagnostics](../../packages/browser-db-sqlite-persistence/tests/opfs-worker-diagnostics-oracle.test.ts), [113-law manifest](../../packages/db-collection-e2e/src/fixtures/persisted-conformance-manifest.ts) | Cache/remote rejection/peer/reopen histories, exact driver results, controlled page/worker ownership, and diagnostic-cause retention. Fake workers and synthetic page events do not prove native handle release or real bfcache admission. The manifest excludes progressive and move suites; registration and shim runs are not device execution. | +| SQLite persistence and native hosts | [persisted histories](../../packages/db-sqlite-persistence-core/tests/persisted.test.ts), [reset/resume histories](../../packages/db-sqlite-persistence-core/tests/sqlite-core-adapter.test.ts), [dual-adapter resume snapshots](../../packages/db-sqlite-persistence-core/tests/sqlite-resume-snapshot.test.ts), [driver contracts](../../packages/db-sqlite-persistence-core/tests/contracts/sqlite-driver-contract.ts), [browser OPFS lifecycle](../../packages/browser-db-sqlite-persistence/tests/opfs-page-lifecycle-oracle.test.ts), [worker diagnostics](../../packages/browser-db-sqlite-persistence/tests/opfs-worker-diagnostics-oracle.test.ts), [113-law manifest](../../packages/db-collection-e2e/src/fixtures/persisted-conformance-manifest.ts) | Cache/remote rejection/peer/reopen histories, atomic reset/resume lineage, key-set evidence, dual-adapter races, exact driver results, controlled page/worker ownership, and diagnostic-cause retention. The reset/resume owners use sqlite3 CLI and in-memory node:sqlite seams; they do not prove multi-process WAL, mobile/Tauri, or other native-device execution. Fake workers and synthetic page events do not prove native handle release or real bfcache admission. The manifest excludes progressive and move suites; registration and shim runs are not device execution. | | Offline execution | [scheduler](../../packages/offline-transactions/tests/KeyScheduler.property.test.ts), [leadership](../../packages/offline-transactions/tests/leadership-replay.property.test.ts), [settlement](../../packages/offline-transactions/tests/transaction-settlement.property.test.ts), [serialization](../../packages/offline-transactions/tests/transaction-serializer.property.test.ts) | Declarative FIFO eligibility, per-transaction outcomes, durable state and typed wire trees. Issued work may finish after ownership loss, but new work must not start. Exactly-once network execution is not promised. | | Frameworks | [React conformance](../../packages/react-db/tests/conformance.test.tsx), [React pagination](../../packages/react-db/tests/infinite-query-conformance.test.tsx), [shared suites](../../packages/db-collection-e2e/src/suites) | Exact exposed rows/pages and each framework's own lifecycle cuts. A React witness does not prove Vue/Solid/Angular/Svelte scheduling. Preserve their receiving registrations. | | Small structures and test mechanics | [SortedMap](../../packages/db/tests/SortedMap.test.ts), [cleanup queue](../../packages/db/tests/cleanup-queue.property.test.ts), [guarded replay](../../packages/db/tests/oracle-replay.test.ts) | Map/full-sort and appointment-list models; executed target/seed/path checks. Callback-reentrant scheduling is outside the initial cleanup-queue domain. | diff --git a/docs/guides/collection-options-creator.md b/docs/guides/collection-options-creator.md index 8ad914acc6..9dc4f9d456 100644 --- a/docs/guides/collection-options-creator.md +++ b/docs/guides/collection-options-creator.md @@ -474,7 +474,7 @@ sync: { For complete, production-ready examples, see the collection packages in the TanStack DB repository: - **[@tanstack/query-db-collection](https://github.com/TanStack/db/tree/main/packages/query-db-collection)** - Pattern A: User-provided handlers with full refetch strategy -- **[@tanstack/trailbase-db-collection](https://github.com/TanStack/db/tree/main/packages/trailbase-db-collection)** - Pattern B: Built-in handlers with ID-based tracking +- **[@tanstack/trailbase-db-collection](https://github.com/TanStack/db/tree/main/packages/trailbase-db-collection)** - Pattern B: Built-in handlers with ID-based tracking - **[@tanstack/electric-db-collection](https://github.com/TanStack/db/tree/main/packages/electric-db-collection)** - Pattern A: Transaction ID tracking with complex sync protocols - **[@tanstack/rxdb-db-collection](https://github.com/TanStack/db/tree/main/packages/rxdb-db-collection)** - Pattern B: Built-in handlers that bridge [RxDB](https://rxdb.info) change streams into TanStack DB's sync lifecycle diff --git a/packages/db-sqlite-persistence-core/README.md b/packages/db-sqlite-persistence-core/README.md index 1ed225957f..235df60776 100644 --- a/packages/db-sqlite-persistence-core/README.md +++ b/packages/db-sqlite-persistence-core/README.md @@ -31,6 +31,7 @@ binding. Provide a runtime `SQLiteDriver` implementation from a wrapper package. - `PullSinceResponse` - `CollectionReset` - `PersistedIndexSpec` +- `PersistedKeySetEvidence` - `PersistedTx` - `PersistenceAdapter` - `SQLiteDriver` @@ -64,6 +65,29 @@ and resolves persistence using: This lets runtime wrappers expose one shared persistence instance per database while still handling per-collection schema versions correctly. +### Atomic resume snapshots + +Persistence adapters may implement +`loadResumeSnapshot(collectionId, options)` to let a sync source certify a +persisted resume baseline. One call must read rows, collection metadata, stream +position, reset epoch, and key-set evidence from the same atomic database +snapshot. `includeRows: false` requests the same certification data without +materializing rows; `requiredIndexSignatures` carries the indexes needed by a +row-bearing snapshot. + +`PersistedKeySetEvidence.status` has three states: + +- `consistent`: the persisted rows match the adapter's durable expected-key + ledger. +- `incompatible`: row loss, substitution, or a reset-generation change makes + the saved resume baseline unsafe. +- `unknown`: the adapter has no authoritative pre-migration key set and does + not claim completeness. + +The method is optional so existing adapters remain assignable. Without it, the +wrapper retains the legacy stream-position and metadata path. Adapter methods +are invoked with their receiver and may rely on instance state through `this`. + ### SQLite core adapter APIs - `SQLiteCoreAdapterOptions` diff --git a/packages/db-sqlite-persistence-core/src/persisted.ts b/packages/db-sqlite-persistence-core/src/persisted.ts index 0f1e4112ac..232b2e4e34 100644 --- a/packages/db-sqlite-persistence-core/src/persisted.ts +++ b/packages/db-sqlite-persistence-core/src/persisted.ts @@ -219,6 +219,17 @@ export type PersistedRowScanOptions = { metadataOnly?: boolean } +export type PersistedKeySetEvidence = { + status: `unknown` | `consistent` | `incompatible` +} + +type PersistedResumeGeneration = { + latestTerm: number + latestSeq: number + latestRowVersion: number + resetEpoch: number +} + export type PersistedTx< T extends object = Record, TKey extends string | number = string | number, @@ -261,6 +272,25 @@ export interface PersistenceAdapter { metadata?: unknown }> > + loadResumeSnapshot?: ( + collectionId: string, + ctx?: { + requiredIndexSignatures?: ReadonlyArray + includeRows?: boolean + }, + ) => Promise<{ + rows: Array<{ + key: string | number + value: Record + metadata?: unknown + }> + keySet?: PersistedKeySetEvidence + collectionMetadata: Array<{ key: string; value: unknown }> + latestTerm: number + latestSeq: number + latestRowVersion: number + resetEpoch: number + }> applyCommittedTx: (collectionId: string, tx: PersistedTx) => Promise loadCollectionMetadata?: ( collectionId: string, @@ -279,6 +309,7 @@ export interface PersistenceAdapter { latestTerm: number latestSeq: number latestRowVersion: number + keySet?: PersistedKeySetEvidence }> } @@ -590,6 +621,7 @@ type BufferedSyncTransaction = { > truncate: boolean internal: boolean + expectedResumeGenerationOwner?: symbol signal?: AbortSignal resolveApplied?: () => void rejectApplied?: (error: unknown) => void @@ -806,6 +838,10 @@ class PersistedCollectionRuntime< private startupMetadataPromise: Promise | null = null private startPromise: Promise | null = null private resumeBaselinePromise: Promise | null = null + private resumeCertificationPromise: Promise | null = null + private persistedKeySetEvidence: PersistedKeySetEvidence | undefined + private persistedResumeGeneration: PersistedResumeGeneration | undefined + private resumeGenerationOwner = Symbol(`persisted resume generation owner`) private lifecycleGeneration = 0 private internalApplyDepth = 0 private appliedReceiptSequence = 0 @@ -925,6 +961,40 @@ class PersistedCollectionRuntime< return this.resumeBaselinePromise } + ensureResumeBaselineCertified(): Promise { + if (this.resumeCertificationPromise) { + return this.resumeCertificationPromise + } + + const lifecycleGeneration = this.lifecycleGeneration + this.resumeCertificationPromise = (async () => { + await this.ensureStarted() + if (lifecycleGeneration !== this.lifecycleGeneration) return + + const adapter = this.persistence.adapter + if (!adapter.loadResumeSnapshot) return + const snapshot = await adapter.loadResumeSnapshot(this.collectionId, { + requiredIndexSignatures: this.getRequiredIndexSignatures(), + includeRows: false, + }) + if (lifecycleGeneration !== this.lifecycleGeneration) return + this.bindResumeSnapshotEvidence(snapshot) + })() + return this.resumeCertificationPromise + } + + getPersistedKeySetEvidence(): PersistedKeySetEvidence | undefined { + return this.persistedKeySetEvidence + } + + supportsResumeSnapshot(): boolean { + return this.persistence.adapter.loadResumeSnapshot !== undefined + } + + getResumeGenerationOwner(): symbol { + return this.resumeGenerationOwner + } + private async hydrateBaseline(lifecycleGeneration: number): Promise { if (lifecycleGeneration !== this.lifecycleGeneration) return @@ -936,6 +1006,7 @@ class PersistedCollectionRuntime< await this.hydrateSubsetUnsafe(baseline, { requestRemoteEnsure: false, lifecycleGeneration, + bindKeySetEvidence: true, }) }) if (lifecycleGeneration !== this.lifecycleGeneration) return @@ -976,6 +1047,24 @@ class PersistedCollectionRuntime< private async loadStartupMetadataInternal( lifecycleGeneration: number, ): Promise { + if (this.persistence.adapter.loadResumeSnapshot) { + const snapshot = await this.persistence.adapter.loadResumeSnapshot( + this.collectionId, + { includeRows: false }, + ) + if (lifecycleGeneration !== this.lifecycleGeneration) return + this.persistedResumeGeneration = + this.getResumeSnapshotGeneration(snapshot) + this.persistedKeySetEvidence = snapshot.keySet + this.observeStreamPosition( + snapshot.latestTerm, + snapshot.latestSeq, + snapshot.latestRowVersion, + ) + this.replaceCollectionMetadataSnapshot(snapshot.collectionMetadata) + return + } + // Restore stream position from the database so that new mutations // don't collide with previously applied transactions. if (this.persistence.adapter.getStreamPosition) { @@ -983,6 +1072,10 @@ class PersistedCollectionRuntime< this.collectionId, ) if (lifecycleGeneration !== this.lifecycleGeneration) return + this.persistedKeySetEvidence = + position.keySet?.status === `consistent` + ? { status: `unknown` } + : position.keySet this.observeStreamPosition( position.latestTerm, position.latestSeq, @@ -1247,6 +1340,10 @@ class PersistedCollectionRuntime< this.startupMetadataPromise = null this.startPromise = null this.resumeBaselinePromise = null + this.resumeCertificationPromise = null + this.persistedKeySetEvidence = undefined + this.persistedResumeGeneration = undefined + this.resumeGenerationOwner = Symbol(`persisted resume generation owner`) } private withInternalApply(task: () => TResult): TResult { @@ -1300,14 +1397,38 @@ class PersistedCollectionRuntime< config: { requestRemoteEnsure: boolean lifecycleGeneration: number + bindKeySetEvidence?: boolean }, ): Promise { this.hydratingGeneration = config.lifecycleGeneration try { - const rows = await this.loadSubsetRowsUnsafe(options) + let rows: Array<{ key: TKey; value: T; metadata?: unknown }> + if ( + config.bindKeySetEvidence && + this.persistence.adapter.loadResumeSnapshot + ) { + const snapshot = await this.persistence.adapter.loadResumeSnapshot( + this.collectionId, + { + requiredIndexSignatures: this.getRequiredIndexSignatures(), + includeRows: true, + }, + ) + rows = snapshot.rows as Array<{ + key: TKey + value: T + metadata?: unknown + }> + if (config.lifecycleGeneration !== this.lifecycleGeneration) return + this.bindResumeSnapshotEvidence(snapshot) + } else { + rows = await this.loadSubsetRowsUnsafe(options) + } if (config.lifecycleGeneration !== this.lifecycleGeneration) return - this.applyRowsToCollection(rows) + if (this.persistedKeySetEvidence?.status !== `incompatible`) { + this.applyRowsToCollection(rows) + } } finally { if (this.hydratingGeneration === config.lifecycleGeneration) { this.hydratingGeneration = null @@ -1351,6 +1472,57 @@ class PersistedCollectionRuntime< }) } + private getResumeSnapshotGeneration(snapshot: { + latestTerm: number + latestSeq: number + latestRowVersion: number + resetEpoch: number + }): PersistedResumeGeneration { + return { + latestTerm: snapshot.latestTerm, + latestSeq: snapshot.latestSeq, + latestRowVersion: snapshot.latestRowVersion, + resetEpoch: snapshot.resetEpoch, + } + } + + private isExpectedResumeGeneration( + generation: PersistedResumeGeneration, + ): boolean { + const expected = this.persistedResumeGeneration + return ( + expected !== undefined && + expected.latestTerm === generation.latestTerm && + expected.latestSeq === generation.latestSeq && + expected.latestRowVersion === generation.latestRowVersion && + expected.resetEpoch === generation.resetEpoch + ) + } + + private bindResumeSnapshotEvidence(snapshot: { + keySet?: PersistedKeySetEvidence + latestTerm: number + latestSeq: number + latestRowVersion: number + resetEpoch: number + }): void { + const generation = this.getResumeSnapshotGeneration(snapshot) + // Moving between equally uncertified snapshots cannot make either + // trustworthy; preserve that status so sync performs a fresh replacement. + const remainsUncertified = + this.persistedKeySetEvidence?.status === snapshot.keySet?.status && + snapshot.keySet?.status !== `consistent` + this.observeStreamPosition( + snapshot.latestTerm, + snapshot.latestSeq, + snapshot.latestRowVersion, + ) + this.persistedKeySetEvidence = + this.isExpectedResumeGeneration(generation) || remainsUncertified + ? snapshot.keySet + : { status: `incompatible` } + } + private replaceCollectionSnapshot( rows: Array<{ key: TKey; value: T; metadata?: unknown }>, collectionMetadata: Array<{ key: string; value: unknown }>, @@ -1521,6 +1693,18 @@ class PersistedCollectionRuntime< const tx = this.createPersistedTxFromOperations(transaction, streamPosition) await this.persistence.adapter.applyCommittedTx(this.collectionId, tx) + if ( + transaction.expectedResumeGenerationOwner === + this.resumeGenerationOwner && + this.persistedResumeGeneration !== undefined + ) { + this.persistedResumeGeneration = { + ...this.persistedResumeGeneration, + latestTerm: tx.term, + latestSeq: tx.seq, + latestRowVersion: tx.rowVersion, + } + } this.publishTxCommittedEvent( this.createTxCommittedPayload({ term: tx.term, @@ -2428,6 +2612,12 @@ function createWrappedSyncConfig< startupState.cleanedUp ? Promise.resolve() : runtime.ensureResumeBaselineHydrated(), + certifyPersistedResume: runtime.supportsResumeSnapshot() + ? () => + startupState.cleanedUp + ? Promise.resolve() + : runtime.ensureResumeBaselineCertified() + : undefined, get: (key: TKey) => { if (startupState.cleanedUp) return undefined const openTransaction = getOpenTransaction() @@ -2447,6 +2637,26 @@ function createWrappedSyncConfig< startupState.cleanedUp ? Promise.resolve([]) : runtime.scanPersistedRows(options), + getPersistedKeySetEvidence: runtime.supportsResumeSnapshot() + ? () => + startupState.cleanedUp + ? undefined + : runtime.getPersistedKeySetEvidence() + : undefined, + expectCurrentCommitInResumeSnapshot: + runtime.supportsResumeSnapshot() + ? () => { + if (startupState.cleanedUp) return + const openTransaction = getOpenTransaction() + if (!openTransaction) { + throw new InvalidPersistedCollectionConfigError( + `expectCurrentCommitInResumeSnapshot must be called within an open sync transaction`, + ) + } + openTransaction.expectedResumeGenerationOwner = + runtime.getResumeGenerationOwner() + } + : undefined, set: (key: TKey, value: unknown) => { if (startupState.cleanedUp) return const openTransaction = getOpenTransaction() @@ -2600,6 +2810,8 @@ function createWrappedSyncConfig< openTransaction.collectionMetadataWrites, truncate: openTransaction.truncate, internal: openTransaction.internal, + expectedResumeGenerationOwner: + openTransaction.expectedResumeGenerationOwner, signal, resolveApplied, rejectApplied, @@ -2618,6 +2830,8 @@ function createWrappedSyncConfig< openTransaction.collectionMetadataWrites, truncate: openTransaction.truncate, internal: false, + expectedResumeGenerationOwner: + openTransaction.expectedResumeGenerationOwner, }) } const persisted = persistAfterApplication() diff --git a/packages/db-sqlite-persistence-core/src/sqlite-core-adapter.ts b/packages/db-sqlite-persistence-core/src/sqlite-core-adapter.ts index 69f29fc604..de5eef9b53 100644 --- a/packages/db-sqlite-persistence-core/src/sqlite-core-adapter.ts +++ b/packages/db-sqlite-persistence-core/src/sqlite-core-adapter.ts @@ -15,6 +15,7 @@ import { import type { LoadSubsetOptions } from '@tanstack/db' import type { PersistedIndexSpec, + PersistedKeySetEvidence, PersistedRowScanOptions, PersistedScannedRow, PersistedTx, @@ -1159,36 +1160,150 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { })) } - async applyCommittedTx(collectionId: string, tx: PersistedTx): Promise { + async loadResumeSnapshot( + collectionId: string, + ctx?: { + requiredIndexSignatures?: ReadonlyArray + includeRows?: boolean + }, + ): Promise<{ + rows: Array<{ + key: string | number + value: Record + metadata?: unknown + }> + keySet: PersistedKeySetEvidence + collectionMetadata: Array<{ key: string; value: unknown }> + latestTerm: number + latestSeq: number + latestRowVersion: number + resetEpoch: number + }> { const tableMapping = await this.ensureCollectionReady(collectionId) - const collectionTableSql = quoteIdentifier(tableMapping.tableName) - const tombstoneTableSql = quoteIdentifier(tableMapping.tombstoneTableName) + const includeRows = ctx?.includeRows !== false + if (includeRows) { + await this.touchRequiredIndexes( + collectionId, + ctx?.requiredIndexSignatures, + ) + } - await this.runInTransaction(async (transactionDriver) => { - const alreadyApplied = await transactionDriver.query<{ applied: number }>( - `SELECT 1 AS applied + return this.runInTransaction(async (transactionDriver) => { + const rows = includeRows + ? await this.loadSubsetInternal(tableMapping, {}, transactionDriver) + : [] + const { latestRowVersion, keySet } = await this.readKeySetEvidence( + collectionId, + tableMapping, + transactionDriver, + ) + const collectionMetadataRows = await transactionDriver.query<{ + key: string + value: string + }>( + `SELECT key, value + FROM collection_metadata + WHERE collection_id = ?`, + [collectionId], + ) + const termRows = await transactionDriver.query<{ latest_term: number }>( + `SELECT latest_term + FROM leader_term + WHERE collection_id = ? + LIMIT 1`, + [collectionId], + ) + const seqRows = await transactionDriver.query<{ max_seq: number }>( + `SELECT MAX(seq) AS max_seq FROM applied_tx - WHERE collection_id = ? AND term = ? AND seq = ? + WHERE collection_id = ? AND term = ( + SELECT latest_term FROM leader_term WHERE collection_id = ? LIMIT 1 + )`, + [collectionId, collectionId], + ) + const resetRows = await transactionDriver.query<{ reset_epoch: number }>( + `SELECT reset_epoch + FROM collection_reset_epoch + WHERE collection_id = ? LIMIT 1`, - [collectionId, tx.term, tx.seq], + [collectionId], ) - if (alreadyApplied.length > 0) { - return + return { + rows: rows.map((row) => ({ + key: row.key, + value: row.value, + metadata: row.metadata, + })), + keySet, + collectionMetadata: collectionMetadataRows.map((row) => ({ + key: row.key, + value: deserializePersistedRowValue(row.value), + })), + latestTerm: termRows[0]?.latest_term ?? 0, + latestSeq: seqRows[0]?.max_seq ?? 0, + latestRowVersion, + resetEpoch: resetRows[0]?.reset_epoch ?? 0, } + }) + } + + async applyCommittedTx(collectionId: string, tx: PersistedTx): Promise { + const tableMapping = await this.ensureCollectionReady(collectionId) + const collectionTableSql = quoteIdentifier(tableMapping.tableName) + const tombstoneTableSql = quoteIdentifier(tableMapping.tombstoneTableName) + await this.runInTransaction(async (transactionDriver) => { const versionRows = await transactionDriver.query<{ latest_row_version: number + key_set_evidence_available: number + schema_version: number + already_applied: number }>( - `SELECT latest_row_version + `SELECT + latest_row_version, + key_set_evidence_available, + ( + SELECT schema_version + FROM collection_registry + WHERE collection_id = ? + LIMIT 1 + ) AS schema_version, + EXISTS ( + SELECT 1 + FROM applied_tx + WHERE collection_id = ? AND term = ? AND seq = ? + ) AS already_applied FROM collection_version WHERE collection_id = ? LIMIT 1`, - [collectionId], + [collectionId, collectionId, tx.term, tx.seq, collectionId], ) - const currentRowVersion = versionRows[0]?.latest_row_version ?? 0 + const version = versionRows[0] + + if (!version) { + throw new InvalidPersistedCollectionConfigError( + `Missing persisted version state for collection "${collectionId}"`, + ) + } + if (version.schema_version !== this.schemaVersion) { + throw new InvalidPersistedCollectionConfigError( + `Schema version mismatch for collection "${collectionId}": ` + + `found ${version.schema_version}, expected ${this.schemaVersion}. ` + + `Refusing to apply a committed transaction through a stale cached adapter.`, + ) + } + + if (version.already_applied === 1) { + return + } + + const currentRowVersion = version.latest_row_version const nextRowVersion = Math.max(currentRowVersion + 1, tx.rowVersion) - const replayDelta: ReplayableTxDelta | null = tx.truncate + const replacesPersistedBaseline = tx.truncate === true + const tracksPersistedKeySet = + version.key_set_evidence_available === 1 || replacesPersistedBaseline + const replayDelta: ReplayableTxDelta | null = replacesPersistedBaseline ? null : { txId: tx.txId, @@ -1206,7 +1321,12 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { collectionMetadataMutations: tx.collectionMetadataMutations ?? [], } - if (tx.truncate) { + if (replacesPersistedBaseline) { + await transactionDriver.run( + `DELETE FROM collection_expected_keys + WHERE collection_id = ?`, + [collectionId], + ) await transactionDriver.run(`DELETE FROM ${collectionTableSql}`) await transactionDriver.run(`DELETE FROM ${tombstoneTableSql}`) } @@ -1214,6 +1334,13 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { for (const mutation of tx.mutations) { const encodedKey = encodePersistedStorageKey(mutation.key) if (mutation.type === `delete`) { + if (tracksPersistedKeySet) { + await transactionDriver.run( + `DELETE FROM collection_expected_keys + WHERE collection_id = ? AND key = ?`, + [collectionId, encodedKey], + ) + } await transactionDriver.run( `DELETE FROM ${collectionTableSql} WHERE key = ?`, @@ -1262,6 +1389,14 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { ? mutation.metadata : existingMetadata + if (tracksPersistedKeySet) { + await transactionDriver.run( + `INSERT INTO collection_expected_keys (collection_id, key) + VALUES (?, ?) + ON CONFLICT(collection_id, key) DO NOTHING`, + [collectionId, encodedKey], + ) + } await transactionDriver.run( `INSERT INTO ${collectionTableSql} (key, value, metadata, row_version) VALUES (?, ?, ?, ?) @@ -1333,11 +1468,23 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { } await transactionDriver.run( - `INSERT INTO collection_version (collection_id, latest_row_version) - VALUES (?, ?) - ON CONFLICT(collection_id) DO UPDATE SET - latest_row_version = excluded.latest_row_version`, - [collectionId, nextRowVersion], + `UPDATE collection_version + SET latest_row_version = ?, + key_set_evidence_available = CASE + WHEN ? = 1 THEN 1 + ELSE key_set_evidence_available + END, + key_set_evidence_incompatible = CASE + WHEN ? = 1 THEN 0 + ELSE key_set_evidence_incompatible + END + WHERE collection_id = ?`, + [ + nextRowVersion, + replacesPersistedBaseline ? 1 : 0, + replacesPersistedBaseline ? 1 : 0, + collectionId, + ], ) await transactionDriver.run( @@ -1371,7 +1518,7 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { tx.txId, nextRowVersion, replayDelta ? stableStringify(replayDelta) : null, - tx.truncate ? 1 : 0, + replacesPersistedBaseline ? 1 : 0, ], ) @@ -1512,10 +1659,10 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { latestTerm: number latestSeq: number latestRowVersion: number + keySet?: PersistedKeySetEvidence }> { - await this.ensureCollectionReady(collectionId) - - const [termRows, versionRows, seqRows] = await Promise.all([ + const tableMapping = await this.ensureCollectionReady(collectionId) + const [termRows, version, seqRows] = await Promise.all([ this.driver.query<{ latest_term: number }>( `SELECT latest_term FROM leader_term @@ -1523,13 +1670,7 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { LIMIT 1`, [collectionId], ), - this.driver.query<{ latest_row_version: number }>( - `SELECT latest_row_version - FROM collection_version - WHERE collection_id = ? - LIMIT 1`, - [collectionId], - ), + this.readKeySetEvidence(collectionId, tableMapping, this.driver), this.driver.query<{ max_seq: number }>( `SELECT MAX(seq) AS max_seq FROM applied_tx @@ -1543,7 +1684,66 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { return { latestTerm: termRows[0]?.latest_term ?? 0, latestSeq: seqRows[0]?.max_seq ?? 0, - latestRowVersion: versionRows[0]?.latest_row_version ?? 0, + latestRowVersion: version.latestRowVersion, + keySet: version.keySet, + } + } + + private async readKeySetEvidence( + collectionId: string, + tableMapping: CollectionTableMapping, + driver: SQLiteDriver, + ): Promise<{ + latestRowVersion: number + keySet: PersistedKeySetEvidence + }> { + const collectionTableSql = quoteIdentifier(tableMapping.tableName) + const versionRows = await driver.query<{ + latest_row_version: number + key_set_evidence_available: number + key_set_incompatible: number + }>( + `SELECT + latest_row_version, + key_set_evidence_available, + CASE + WHEN key_set_evidence_available = 0 THEN 0 + WHEN key_set_evidence_incompatible = 1 THEN 1 + WHEN EXISTS ( + SELECT 1 + FROM ${collectionTableSql} AS actual + LEFT JOIN collection_expected_keys AS expected + ON expected.collection_id = ? + AND expected.key = actual.key + WHERE expected.key IS NULL + UNION ALL + SELECT 1 + FROM collection_expected_keys AS expected + LEFT JOIN ${collectionTableSql} AS actual + ON actual.key = expected.key + WHERE expected.collection_id = ? + AND actual.key IS NULL + LIMIT 1 + ) THEN 1 + ELSE 0 + END AS key_set_incompatible + FROM collection_version + WHERE collection_id = ? + LIMIT 1`, + [collectionId, collectionId, collectionId], + ) + const version = versionRows[0] + + return { + latestRowVersion: version?.latest_row_version ?? 0, + keySet: { + status: + version?.key_set_evidence_available !== 1 + ? `unknown` + : version.key_set_incompatible === 1 + ? `incompatible` + : `consistent`, + }, } } @@ -1685,6 +1885,7 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { private async loadSubsetInternal( tableMapping: CollectionTableMapping, options: LoadSubsetOptions, + driver: SQLiteDriver = this.driver, ): Promise>>> { const collectionTableSql = quoteIdentifier(tableMapping.tableName) const whereCompiled = options.where @@ -1705,10 +1906,7 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { queryParams.push(...orderByCompiled.params) } - const storedRows = await this.driver.query( - sql, - queryParams, - ) + const storedRows = await driver.query(sql, queryParams) const parsedRows = decodeStoredSqliteRows(storedRows) const filteredRows = this.applyInMemoryWhere(parsedRows, options.where) @@ -1959,8 +2157,13 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { ON ${tombstoneTableSql} (row_version)`, ) await this.driver.run( - `INSERT INTO collection_version (collection_id, latest_row_version) - VALUES (?, 0) + `INSERT INTO collection_version ( + collection_id, + latest_row_version, + key_set_evidence_available, + key_set_evidence_incompatible + ) + VALUES (?, 0, 1, 0) ON CONFLICT(collection_id) DO NOTHING`, [collectionId], ) @@ -1970,7 +2173,7 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { ON CONFLICT(collection_id) DO NOTHING`, [collectionId], ) - + await this.ensureCollectionKeyEvidenceTriggers(tableName) const mapping = { tableName, tombstoneTableName, @@ -1979,6 +2182,75 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { return mapping } + private async ensureCollectionKeyEvidenceTriggers( + tableName: string, + ): Promise { + const collectionTableSql = quoteIdentifier(tableName) + const tableNameLiteral = toSqliteLiteral(tableName) + const insertTriggerSql = quoteIdentifier(`${tableName}_key_evidence_insert`) + const deleteTriggerSql = quoteIdentifier(`${tableName}_key_evidence_delete`) + const updateTriggerSql = quoteIdentifier(`${tableName}_key_evidence_update`) + + await this.driver.exec( + `CREATE TRIGGER IF NOT EXISTS ${insertTriggerSql} + AFTER INSERT ON ${collectionTableSql} + WHEN NOT EXISTS ( + SELECT 1 + FROM collection_expected_keys + WHERE collection_id = ( + SELECT collection_id + FROM collection_registry + WHERE table_name = ${tableNameLiteral} + ) AND key = NEW.key + ) + BEGIN + UPDATE collection_version + SET key_set_evidence_incompatible = 1 + WHERE collection_id = ( + SELECT collection_id + FROM collection_registry + WHERE table_name = ${tableNameLiteral} + ); + END`, + ) + await this.driver.exec( + `CREATE TRIGGER IF NOT EXISTS ${deleteTriggerSql} + AFTER DELETE ON ${collectionTableSql} + WHEN EXISTS ( + SELECT 1 + FROM collection_expected_keys + WHERE collection_id = ( + SELECT collection_id + FROM collection_registry + WHERE table_name = ${tableNameLiteral} + ) AND key = OLD.key + ) + BEGIN + UPDATE collection_version + SET key_set_evidence_incompatible = 1 + WHERE collection_id = ( + SELECT collection_id + FROM collection_registry + WHERE table_name = ${tableNameLiteral} + ); + END`, + ) + await this.driver.exec( + `CREATE TRIGGER IF NOT EXISTS ${updateTriggerSql} + AFTER UPDATE OF key ON ${collectionTableSql} + WHEN OLD.key <> NEW.key + BEGIN + UPDATE collection_version + SET key_set_evidence_incompatible = 1 + WHERE collection_id = ( + SELECT collection_id + FROM collection_registry + WHERE table_name = ${tableNameLiteral} + ); + END`, + ) + } + private async handleSchemaMismatch( collectionId: string, previousSchemaVersion: number, @@ -1997,6 +2269,27 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { const tombstoneTableSql = quoteIdentifier(tombstoneTableName) await this.runInTransaction(async (transactionDriver) => { + const currentSchemaRows = await transactionDriver.query<{ + schema_version: number + }>( + `SELECT schema_version + FROM collection_registry + WHERE collection_id = ? + LIMIT 1`, + [collectionId], + ) + const currentSchemaVersion = currentSchemaRows[0]?.schema_version + if (currentSchemaVersion === nextSchemaVersion) { + return + } + if (currentSchemaVersion !== previousSchemaVersion) { + throw new InvalidPersistedCollectionConfigError( + `Schema version changed concurrently for collection "${collectionId}": ` + + `found ${currentSchemaVersion ?? `no registry entry`} after observing ${previousSchemaVersion}; ` + + `refusing to reset it to ${nextSchemaVersion}.`, + ) + } + const persistedIndexes = await transactionDriver.query<{ index_name: string }>( @@ -2011,6 +2304,11 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { ) } + await transactionDriver.run( + `DELETE FROM collection_expected_keys + WHERE collection_id = ?`, + [collectionId], + ) await transactionDriver.run(`DELETE FROM ${collectionTableSql}`) await transactionDriver.run(`DELETE FROM ${tombstoneTableSql}`) await transactionDriver.run( @@ -2023,6 +2321,11 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { WHERE collection_id = ?`, [collectionId], ) + await transactionDriver.run( + `DELETE FROM collection_metadata + WHERE collection_id = ?`, + [collectionId], + ) await transactionDriver.run( `UPDATE collection_registry SET schema_version = ?, @@ -2031,10 +2334,17 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { [nextSchemaVersion, collectionId], ) await transactionDriver.run( - `INSERT INTO collection_version (collection_id, latest_row_version) - VALUES (?, 0) + `INSERT INTO collection_version ( + collection_id, + latest_row_version, + key_set_evidence_available, + key_set_evidence_incompatible + ) + VALUES (?, 0, 1, 0) ON CONFLICT(collection_id) DO UPDATE SET - latest_row_version = 0`, + latest_row_version = 0, + key_set_evidence_available = 1, + key_set_evidence_incompatible = 0`, [collectionId], ) await transactionDriver.run( @@ -2110,7 +2420,37 @@ export class SQLiteCorePersistenceAdapter implements PersistenceAdapter { await this.driver.exec( `CREATE TABLE IF NOT EXISTS collection_version ( collection_id TEXT PRIMARY KEY, - latest_row_version INTEGER NOT NULL + latest_row_version INTEGER NOT NULL, + key_set_evidence_available INTEGER NOT NULL DEFAULT 0, + key_set_evidence_incompatible INTEGER NOT NULL DEFAULT 0 + )`, + ) + const collectionVersionColumns = await this.driver.query<{ name: string }>( + `PRAGMA table_info(collection_version)`, + ) + const keyEvidenceColumns = [ + `key_set_evidence_available`, + `key_set_evidence_incompatible`, + ] as const + for (const columnName of keyEvidenceColumns) { + if (collectionVersionColumns.some(({ name }) => name === columnName)) { + continue + } + try { + await this.driver.exec( + `ALTER TABLE collection_version ADD COLUMN ${columnName} INTEGER NOT NULL DEFAULT 0`, + ) + } catch (error) { + if (!isDuplicateColumnAddError(error, columnName)) { + throw error + } + } + } + await this.driver.exec( + `CREATE TABLE IF NOT EXISTS collection_expected_keys ( + collection_id TEXT NOT NULL, + key TEXT NOT NULL, + PRIMARY KEY (collection_id, key) )`, ) await this.driver.exec( diff --git a/packages/db-sqlite-persistence-core/tests/persisted.test-d.ts b/packages/db-sqlite-persistence-core/tests/persisted.test-d.ts index b07f9f787d..336c2ff676 100644 --- a/packages/db-sqlite-persistence-core/tests/persisted.test-d.ts +++ b/packages/db-sqlite-persistence-core/tests/persisted.test-d.ts @@ -1,7 +1,11 @@ import { describe, expectTypeOf, it } from 'vitest' import { createCollection } from '@tanstack/db' import { persistedCollectionOptions } from '../src' -import type { PersistedCollectionUtils, PersistenceAdapter } from '../src' +import type { + PersistedCollectionUtils, + PersistedKeySetEvidence, + PersistenceAdapter, +} from '../src' import type { SyncConfig, UtilsRecord } from '@tanstack/db' type Todo = { @@ -24,6 +28,46 @@ const adapter: PersistenceAdapter = { } describe(`persisted collection types`, () => { + it(`keeps the atomic resume snapshot extension optional and exact`, () => { + const legacyAdapter: PersistenceAdapter = adapter + const snapshotAdapter: PersistenceAdapter = { + ...adapter, + loadResumeSnapshot: (_collectionId, _options) => + Promise.resolve({ + rows: [], + keySet: { status: `consistent` }, + collectionMetadata: [], + latestTerm: 1, + latestSeq: 2, + latestRowVersion: 3, + resetEpoch: 4, + }), + } + type LoadResumeSnapshot = NonNullable< + PersistenceAdapter[`loadResumeSnapshot`] + > + type ResumeSnapshot = Awaited> + + expectTypeOf(legacyAdapter).toMatchTypeOf() + expectTypeOf(snapshotAdapter.loadResumeSnapshot).toMatchTypeOf< + LoadResumeSnapshot | undefined + >() + expectTypeOf[1]>().toEqualTypeOf< + | { + requiredIndexSignatures?: ReadonlyArray + includeRows?: boolean + } + | undefined + >() + expectTypeOf().toEqualTypeOf< + PersistedKeySetEvidence | undefined + >() + + // @ts-expect-error key-set evidence has exactly three supported states + const invalidEvidence: PersistedKeySetEvidence = { status: `verified` } + expectTypeOf(invalidEvidence).toEqualTypeOf() + }) + it(`adds persisted utils in sync-absent mode`, () => { const options = persistedCollectionOptions< Todo, diff --git a/packages/db-sqlite-persistence-core/tests/sqlite-core-adapter.test.ts b/packages/db-sqlite-persistence-core/tests/sqlite-core-adapter.test.ts index de10c93e6c..ffb5935c87 100644 --- a/packages/db-sqlite-persistence-core/tests/sqlite-core-adapter.test.ts +++ b/packages/db-sqlite-persistence-core/tests/sqlite-core-adapter.test.ts @@ -4,6 +4,7 @@ import { copyFileSync, existsSync, mkdtempSync, rmSync } from 'node:fs' import { tmpdir } from 'node:os' import { join } from 'node:path' import { promisify } from 'node:util' +import { fc } from '@fast-check/vitest' import { afterEach, describe, expect, it } from 'vitest' import { IR } from '@tanstack/db' import { SQLiteCorePersistenceAdapter, createPersistedTableName } from '../src' @@ -211,6 +212,421 @@ function createHarness( } } +type ResetResumeHistory = { + fromSchemaVersion: number + rows: Array + resumeKind: `none` | `reset` | `resume` + transition: + | `compatible-reopen` + | `schema-reset` + | `partial-restore` + | `external-row-loss` + reopensBeforeTransition: number + reopensAfterTransition: number + unrelatedMetadataKeys: Array +} + +type ResetResumeObservation = { + checkpoint: `after-persistence-transition-restart` + resetEpoch: number + schemaVersion: number + durableRows: ReadonlyArray + tombstones: ReadonlyArray<{ key: string; rowVersion: number }> + appliedTransactions: ReadonlyArray<{ txId: string; rowVersion: number }> + latestRowVersion: number + metadataKeys: ReadonlyArray + resumeState: unknown +} + +type ResetResumeExpectation = { + metadataKeys: ReadonlyArray + resumeKind: unknown +} + +function destroysPersistedBaseline(history: ResetResumeHistory): boolean { + return ( + history.transition === `schema-reset` || + history.transition === `partial-restore` + ) +} + +function expectedResetResumeMetadata( + history: ResetResumeHistory, +): ResetResumeExpectation { + if (destroysPersistedBaseline(history)) { + return { metadataKeys: [], resumeKind: undefined } + } + + return { + metadataKeys: [ + ...(history.resumeKind === `none` ? [] : [`electric:resume`]), + ...history.unrelatedMetadataKeys.map((key) => `oracle:${key}`), + ].sort(), + resumeKind: history.resumeKind === `none` ? undefined : history.resumeKind, + } +} + +class ResetResumeOracleViolation extends Error { + readonly law: string = `reset-resume.baseline-lineage` + readonly discriminant: string + readonly history: ResetResumeHistory + readonly observation: ResetResumeObservation + readonly cleanupEvidence: string + readonly expected: ResetResumeExpectation + readonly actual: ResetResumeExpectation + + constructor( + history: ResetResumeHistory, + observation: ResetResumeObservation, + cleanupEvidence: string, + cause: unknown, + ) { + super( + `A persisted-baseline transition violated collection-metadata lineage at the ` + + `${observation.checkpoint}. ` + + `history=${JSON.stringify(history)} ` + + `observation=${JSON.stringify(observation)} ` + + `cleanup=${cleanupEvidence}`, + { cause }, + ) + this.name = `ResetResumeOracleViolation` + this.history = structuredClone(history) + this.observation = structuredClone(observation) + this.cleanupEvidence = cleanupEvidence + this.expected = expectedResetResumeMetadata(history) + this.actual = { + metadataKeys: [...observation.metadataKeys], + resumeKind: resumeKindOf(observation.resumeState), + } + this.discriminant = destroysPersistedBaseline(history) + ? `reset-retained-metadata` + : `non-reset-metadata-changed` + } +} + +function hasSameResetResumeFailure( + left: ResetResumeOracleViolation, + right: ResetResumeOracleViolation, +): boolean { + const signature = (failure: ResetResumeOracleViolation): string => + JSON.stringify({ + law: String(failure.law), + discriminant: String(failure.discriminant), + checkpoint: String(failure.observation.checkpoint), + expectedMetadataClass: + failure.expected.metadataKeys.length === 0 ? `empty` : `nonempty`, + actualMetadataClass: + failure.actual.metadataKeys.length === 0 ? `empty` : `nonempty`, + }) + return signature(left) === signature(right) +} + +function resumeKindOf(value: unknown): unknown { + return value && typeof value === `object` + ? (value as Record).kind + : undefined +} + +function expectResetResumeLaw( + history: ResetResumeHistory, + observation: ResetResumeObservation, +): void { + // The core adapter owns the reset transaction: it must clear every metadata + // record coupled to the destroyed baseline. Compatible reopen and raw + // external row loss do not give this generic layer authority to interpret an + // Electric cursor, so they preserve the metadata exactly at this checkpoint. + expect( + { + metadataKeys: observation.metadataKeys, + resumeKind: resumeKindOf(observation.resumeState), + }, + `reset clears all collection metadata; non-reset core transitions preserve it`, + ).toEqual(expectedResetResumeMetadata(history)) +} + +function attachResetResumeCleanupDiagnostics( + primary: unknown, + cleanupEvidence: string, + cleanupFailure: unknown, +): Error { + const error = + primary instanceof Error + ? primary + : new Error(`Reset/resume oracle failed with a non-Error value`, { + cause: primary, + }) + if (!(`cleanupEvidence` in error)) { + Object.defineProperty(error, `cleanupEvidence`, { + value: cleanupEvidence, + enumerable: true, + }) + } + if (cleanupFailure !== undefined) { + Object.defineProperty(error, `cleanupFailure`, { + value: cleanupFailure, + enumerable: true, + }) + } + return error +} + +async function observeResetResumeHistory( + history: ResetResumeHistory, + harnessFactory: SQLiteCoreAdapterHarnessFactory, +): Promise { + let harness: ReturnType | undefined + const collectionId = `reset-resume-oracle` + let observation!: ResetResumeObservation + let primaryFailure: unknown + let cleanupFailure: unknown + let failurePhase: `setup` | `reach` | `law` | undefined + let cleanupEvidence = `not-run` + const expectedRows = structuredClone(history.rows) + const seedRows = structuredClone(history.rows) + const restoreRows = structuredClone(history.rows) + + try { + harness = harnessFactory({ schemaVersion: history.fromSchemaVersion }) + await harness.adapter.applyCommittedTx(collectionId, { + txId: `seed-baseline`, + term: 1, + seq: 1, + rowVersion: 1, + mutations: [ + ...seedRows.map((row) => ({ + type: `insert` as const, + key: row.id, + value: structuredClone(row), + })), + { + type: `delete` as const, + key: `deleted-before-baseline`, + value: { + id: `deleted-before-baseline`, + title: `baseline tombstone`, + createdAt: `2026-01-01T00:00:00.000Z`, + score: -1, + }, + }, + ], + collectionMetadataMutations: [ + ...(history.resumeKind === `none` + ? [] + : [ + { + type: `set` as const, + key: `electric:resume`, + value: + history.resumeKind === `reset` + ? { kind: `reset`, updatedAt: 1 } + : { + kind: `resume`, + offset: `10_0`, + handle: `handle-before-reset`, + shapeId: `shape-before-reset`, + updatedAt: 1, + }, + }, + ]), + ...history.unrelatedMetadataKeys.map((key, index) => ({ + type: `set` as const, + key: `oracle:${key}`, + value: { index }, + })), + ], + }) + + for (let index = 0; index < history.reopensBeforeTransition; index++) { + const reopened = new SQLiteCorePersistenceAdapter({ + driver: harness.driver, + schemaVersion: history.fromSchemaVersion, + }) + await reopened.loadSubset(collectionId, {}) + await reopened.loadCollectionMetadata(collectionId) + } + + const usesSchemaReset = destroysPersistedBaseline(history) + const nextSchemaVersion = usesSchemaReset + ? history.fromSchemaVersion + 1 + : history.fromSchemaVersion + let restarted = new SQLiteCorePersistenceAdapter({ + driver: harness.driver, + schemaVersion: nextSchemaVersion, + schemaMismatchPolicy: `sync-present-reset`, + }) + await restarted.loadSubset(collectionId, {}) + + if (history.transition === `partial-restore`) { + await restarted.applyCommittedTx(collectionId, { + txId: `partial-restore`, + term: 2, + seq: 1, + rowVersion: 2, + mutations: restoreRows.slice(0, -1).map((row) => ({ + type: `insert` as const, + key: row.id, + value: structuredClone(row), + })), + }) + } else if (history.transition === `external-row-loss`) { + const collectionTable = createPersistedTableName(collectionId, `c`) + await harness.driver.run( + `DELETE FROM "${collectionTable}" WHERE json_extract(value, '$.id') = ?`, + [history.rows[0]!.id], + ) + } + + for (let index = 0; index < history.reopensAfterTransition; index++) { + restarted = new SQLiteCorePersistenceAdapter({ + driver: harness.driver, + schemaVersion: nextSchemaVersion, + schemaMismatchPolicy: `sync-present-reset`, + }) + await restarted.loadSubset(collectionId, {}) + } + + const metadata = await restarted.loadCollectionMetadata(collectionId) + const resetEpochRows = await harness.driver.query<{ reset_epoch: number }>( + `SELECT reset_epoch FROM collection_reset_epoch WHERE collection_id = ?`, + [collectionId], + ) + const registryRows = await harness.driver.query<{ schema_version: number }>( + `SELECT schema_version FROM collection_registry WHERE collection_id = ?`, + [collectionId], + ) + const tombstoneTable = createPersistedTableName(collectionId, `t`) + const tombstones = await harness.driver.query<{ + key: string + row_version: number + }>(`SELECT key, row_version FROM "${tombstoneTable}" ORDER BY key`) + const appliedTransactions = await harness.driver.query<{ + tx_id: string + row_version: number + }>( + `SELECT tx_id, row_version FROM applied_tx WHERE collection_id = ? ORDER BY term, seq`, + [collectionId], + ) + const versionRows = await harness.driver.query<{ + latest_row_version: number + }>( + `SELECT latest_row_version FROM collection_version WHERE collection_id = ?`, + [collectionId], + ) + observation = { + checkpoint: `after-persistence-transition-restart`, + resetEpoch: resetEpochRows[0]?.reset_epoch ?? -1, + schemaVersion: registryRows[0]?.schema_version ?? -1, + durableRows: (await restarted.loadSubset(collectionId, {})).sort((a, b) => + String(a.key).localeCompare(String(b.key)), + ), + tombstones: tombstones.map(({ key, row_version }) => ({ + key, + rowVersion: row_version, + })), + appliedTransactions: appliedTransactions.map( + ({ tx_id, row_version }) => ({ + txId: tx_id, + rowVersion: row_version, + }), + ), + latestRowVersion: versionRows[0]?.latest_row_version ?? -1, + metadataKeys: metadata.map(({ key }) => key).sort(), + resumeState: metadata.find(({ key }) => key === `electric:resume`)?.value, + } + + try { + // Positive reach evidence is separate from the semantic accusation. + expect(observation.schemaVersion).toBe(nextSchemaVersion) + expect(observation.resetEpoch).toBe(usesSchemaReset ? 1 : 0) + expect(observation.durableRows).toEqual( + (history.transition === `schema-reset` + ? [] + : history.transition === `partial-restore` + ? expectedRows.slice(0, -1) + : history.transition === `external-row-loss` + ? expectedRows.slice(1) + : expectedRows + ) + .map((value) => ({ key: value.id, value })) + .sort((a, b) => String(a.key).localeCompare(String(b.key))), + ) + expect(observation.tombstones).toEqual( + usesSchemaReset + ? [] + : [{ key: `s:deleted-before-baseline`, rowVersion: 1 }], + ) + expect(observation.appliedTransactions).toEqual( + history.transition === `schema-reset` + ? [] + : history.transition === `partial-restore` + ? [{ txId: `partial-restore`, rowVersion: 2 }] + : [{ txId: `seed-baseline`, rowVersion: 1 }], + ) + expect(observation.latestRowVersion).toBe( + history.transition === `schema-reset` + ? 0 + : history.transition === `partial-restore` + ? 2 + : 1, + ) + } catch (error) { + failurePhase = `reach` + primaryFailure = error + } + + if (primaryFailure === undefined) { + try { + expectResetResumeLaw(history, observation) + } catch (error) { + failurePhase = `law` + primaryFailure = error + } + } + } catch (error) { + if (primaryFailure === undefined) { + failurePhase = `setup` + primaryFailure = error + } + } finally { + if (harness) { + try { + await harness.cleanup() + cleanupEvidence = + `dbPath` in harness && typeof harness.dbPath === `string` + ? existsSync(harness.dbPath) + ? `failed: SQLite file still exists` + : `passed: SQLite file removed after captured checkpoint` + : `passed: registered harness cleanup completed after captured checkpoint` + } catch (cleanupError) { + cleanupEvidence = `failed: ${String(cleanupError)}` + cleanupFailure = cleanupError + } + } + } + + if (primaryFailure !== undefined) { + const failure = + failurePhase === `law` + ? new ResetResumeOracleViolation( + history, + observation, + cleanupEvidence, + primaryFailure, + ) + : primaryFailure + throw attachResetResumeCleanupDiagnostics( + failure, + cleanupEvidence, + cleanupFailure, + ) + } + if (cleanupFailure !== undefined) throw cleanupFailure + if (!cleanupEvidence.startsWith(`passed:`)) { + throw new Error(cleanupEvidence) + } + return observation +} + export type SQLiteCoreAdapterHarnessFactory = ( options?: Omit< ConstructorParameters[0], @@ -990,6 +1406,274 @@ export function runSQLiteCoreAdapterContractSuite( expect(resetRows).toEqual([]) }) + /** + * Reset/resume oracle card + * + * Law and source: a destructive schema reset creates a new persisted + * baseline and clears every collection-metadata record in that reset + * transaction. Same-schema restarts preserve metadata, including for a + * genuinely empty baseline. The production reset policy above and the + * independently reproduced history in + * https://github.com/TanStack/db/issues/1589 establish this narrow law. + * + * Domain and legal histories: zero-to-four committed rows; absent, reset, + * or non-initial resume metadata; same-schema restart, vN -> vN+1 + * sync-present-reset, a partial restore after reset, or out-of-band row + * loss; restarts on either side; unrelated metadata. + * + * Reference: a two-generation lineage relation. Same-schema reopen and raw + * external loss keep the core generation/metadata record; schema reset + * replaces the generation and starts with no metadata. This model does not + * inspect the adapter's SQL branches or infer compatibility from row + * cardinality. + * + * Production path and checkpoint: SQLiteCorePersistenceAdapter commits a + * real SQLite baseline, then a new adapter instance reaches + * ensureCollectionReady/handleSchemaMismatch. At the + * after-persistence-transition-restart checkpoint we inspect registry + * version, reset epoch/generation, complete durable rows, tombstones, + * applied transactions, and collection metadata. + * + * Observed result and known omissions: exact settled rows, exact metadata + * keys, and whether the durable Electric state is absent/reset/resume. The + * external-loss lane proves the generic adapter leaves metadata reachable; + * only the persisted+Electric suite judges whether that cursor is safe to + * consume. The injected loss is not a claim that arbitrary SQL is a + * supported public API. Partial resumed updates and SDK framing retain + * their executable owners in the Electric recovery and framing suites. + * + * Trust: reset_epoch and schema_version prove reach; retained metadata + * after reset is the whole-path fault control, while compatible and + * external-loss histories calibrate preservation. Every generated database + * is removed after evidence capture. The failure reports first and reduced + * histories plus a verified fast-check seed/path replay. + */ + it(`resets collection metadata with its persisted baseline across generated restart histories`, async () => { + const historyArbitrary = fc + .constantFrom< + ResetResumeHistory[`transition`] + >(`compatible-reopen`, `schema-reset`, `partial-restore`, `external-row-loss`) + .chain((transition) => + fc.record({ + fromSchemaVersion: fc.integer({ min: 1, max: 4 }), + rows: fc + .uniqueArray(fc.integer({ min: 0, max: 20 }), { + minLength: + transition === `partial-restore` + ? 2 + : transition === `external-row-loss` + ? 1 + : 0, + maxLength: 4, + }) + .map((ids) => + ids + .sort((left, right) => left - right) + .map((id) => ({ + id: String(id), + title: `row-${id}`, + createdAt: `2026-01-01T00:00:00.000Z`, + score: id, + })), + ), + resumeKind: fc.constantFrom(`none`, `reset`, `resume`), + transition: fc.constant(transition), + reopensBeforeTransition: fc.integer({ min: 0, max: 2 }), + reopensAfterTransition: fc.integer({ min: 0, max: 2 }), + unrelatedMetadataKeys: fc.uniqueArray( + fc.constantFrom(`gc`, `provider`, `custom`), + { maxLength: 3 }, + ), + }), + ) + + const seedText = + process.env.TANSTACK_DB_SQLITE_ORACLE_SEED ?? String(1659) + const seed = Number(seedText) + const path = process.env.TANSTACK_DB_SQLITE_ORACLE_PATH + const runsText = process.env.TANSTACK_DB_SQLITE_ORACLE_RUNS ?? String(24) + const numRuns = Number(runsText) + if (!Number.isSafeInteger(seed)) { + throw new Error(`TANSTACK_DB_SQLITE_ORACLE_SEED must be an integer`) + } + if (!Number.isSafeInteger(numRuns) || numRuns < 1) { + throw new Error(`TANSTACK_DB_SQLITE_ORACLE_RUNS must be positive`) + } + if (path !== undefined && !/^\d+(?::\d+)*$/.test(path)) { + throw new Error( + `TANSTACK_DB_SQLITE_ORACLE_PATH must be a numeric shrink path`, + ) + } + + let originalFailure: ResetResumeOracleViolation | undefined + const property = fc.asyncProperty(historyArbitrary, async (history) => { + try { + await observeResetResumeHistory(history, harnessFactory) + } catch (error) { + if ( + originalFailure === undefined && + error instanceof ResetResumeOracleViolation + ) { + originalFailure = error + } + throw error + } + }) + const failure = await fc.check(property, { + seed, + numRuns, + ...(path === undefined ? {} : { path, endOnFailure: true }), + }) + if (!failure.failed) return + if ( + !(failure.errorInstance instanceof ResetResumeOracleViolation) || + originalFailure === undefined || + failure.counterexamplePath === null + ) { + throw failure.errorInstance + } + if (!hasSameResetResumeFailure(originalFailure, failure.errorInstance)) { + throw new Error( + `Shrinking changed the reset/resume law or observation checkpoint`, + ) + } + + const replay = await fc.check( + fc.asyncProperty(historyArbitrary, async (history) => { + await observeResetResumeHistory(history, harnessFactory) + }), + { + seed: failure.seed, + path: failure.counterexamplePath, + endOnFailure: true, + }, + ) + if ( + !replay.failed || + !(replay.errorInstance instanceof ResetResumeOracleViolation) || + JSON.stringify(replay.counterexample) !== + JSON.stringify(failure.counterexample) || + !hasSameResetResumeFailure(failure.errorInstance, replay.errorInstance) + ) { + throw new Error( + `Reset/resume oracle replay did not reproduce the intended violation`, + ) + } + + const reducedFailure = failure.errorInstance + throw new Error( + `Reset/resume baseline-lineage violation. ` + + `seed=${failure.seed} path=${failure.counterexamplePath} ` + + `law=${reducedFailure.law} ` + + `discriminant=${reducedFailure.discriminant} ` + + `checkpoint=${reducedFailure.observation.checkpoint} ` + + `originalTrace=${JSON.stringify(originalFailure.history)} ` + + `reducedTrace=${JSON.stringify(reducedFailure.history)} ` + + `expected=${JSON.stringify(reducedFailure.expected)} ` + + `actual=${JSON.stringify(reducedFailure.actual)} ` + + `observation=${JSON.stringify(reducedFailure.observation)} ` + + `replay=verified ` + + `cleanup=${reducedFailure.cleanupEvidence}`, + { cause: reducedFailure }, + ) + }, 120_000) + + it(`leaves externally inconsistent metadata reachable for consumer validation`, async () => { + const observation = await observeResetResumeHistory( + { + fromSchemaVersion: 1, + rows: [ + { + id: `1`, + title: `lost externally`, + createdAt: `2026-01-01T00:00:00.000Z`, + score: 1, + }, + { + id: `2`, + title: `survives`, + createdAt: `2026-01-01T00:00:00.000Z`, + score: 2, + }, + ], + resumeKind: `resume`, + transition: `external-row-loss`, + reopensBeforeTransition: 1, + reopensAfterTransition: 1, + unrelatedMetadataKeys: [`provider`], + }, + harnessFactory, + ) + + expect(observation.metadataKeys).toEqual([ + `electric:resume`, + `oracle:provider`, + ]) + expect(resumeKindOf(observation.resumeState)).toBe(`resume`) + }, 30_000) + + it(`requires complete metadata reset while preserving non-reset metadata`, () => { + const compatibleEmpty: ResetResumeHistory = { + fromSchemaVersion: 1, + rows: [], + resumeKind: `resume`, + transition: `compatible-reopen`, + reopensBeforeTransition: 0, + reopensAfterTransition: 1, + unrelatedMetadataKeys: [`provider`], + } + const observation: ResetResumeObservation = { + checkpoint: `after-persistence-transition-restart`, + resetEpoch: 0, + schemaVersion: 1, + durableRows: [], + tombstones: [], + appliedTransactions: [], + latestRowVersion: 1, + metadataKeys: [`electric:resume`, `oracle:provider`], + resumeState: { kind: `resume`, offset: `10_0` }, + } + expect(() => + expectResetResumeLaw(compatibleEmpty, observation), + ).not.toThrow() + expect(() => + expectResetResumeLaw( + { ...compatibleEmpty, transition: `external-row-loss` }, + observation, + ), + ).not.toThrow() + expect(() => + expectResetResumeLaw( + { ...compatibleEmpty, transition: `schema-reset` }, + { ...observation, resetEpoch: 1, schemaVersion: 2 }, + ), + ).toThrowError(expect.objectContaining({ name: `AssertionError` })) + expect(() => + expectResetResumeLaw( + { ...compatibleEmpty, transition: `schema-reset` }, + { + ...observation, + resetEpoch: 1, + schemaVersion: 2, + metadataKeys: [`oracle:provider`], + resumeState: undefined, + }, + ), + ).toThrowError(expect.objectContaining({ name: `AssertionError` })) + expect(() => + expectResetResumeLaw( + { ...compatibleEmpty, transition: `schema-reset` }, + { + ...observation, + resetEpoch: 1, + schemaVersion: 2, + metadataKeys: [], + resumeState: undefined, + }, + ), + ).not.toThrow() + }) + it(`returns pullSince deltas and requiresFullReload when threshold is exceeded`, async () => { const { adapter } = registerContractHarness({ pullSinceReloadThreshold: 1, diff --git a/packages/db-sqlite-persistence-core/tests/sqlite-resume-snapshot.test.ts b/packages/db-sqlite-persistence-core/tests/sqlite-resume-snapshot.test.ts new file mode 100644 index 0000000000..1213b282cc --- /dev/null +++ b/packages/db-sqlite-persistence-core/tests/sqlite-resume-snapshot.test.ts @@ -0,0 +1,765 @@ +import { DatabaseSync } from 'node:sqlite' +import { describe, expect, it } from 'vitest' +import { SQLiteCorePersistenceAdapter, createPersistedTableName } from '../src' +import type { SQLiteDriver } from '../src' + +type CachedSchemaState = { + schemaVersion: number + resetEpoch: number + rows: Array<{ key: string | number; value: Record }> + metadata: Array<{ key: string; value: unknown }> + appliedTransactions: Array<{ + term: number + seq: number + txId: string + rowVersion: number + }> + keyEvidence: { + available: number + incompatible: number + expectedKeys: Array + } +} + +function toBinding(value: unknown): string | number | bigint | null { + if (value === null || value === undefined) return null + if (typeof value === `boolean`) return value ? 1 : 0 + if ( + typeof value === `string` || + typeof value === `number` || + typeof value === `bigint` + ) { + return value + } + return String(value) +} + +function createDriver( + database: DatabaseSync, + failTransactionRun?: (sql: string) => boolean, +): SQLiteDriver { + const driver: SQLiteDriver = { + exec: (sql) => { + database.exec(sql) + return Promise.resolve() + }, + query: (sql, params = []) => + Promise.resolve( + database + .prepare(sql) + .all(...params.map(toBinding)) + .map((row) => ({ ...row })) as Array, + ), + run: (sql, params = []) => { + database.prepare(sql).run(...params.map(toBinding)) + return Promise.resolve() + }, + transaction: async (transaction) => { + database.exec(`BEGIN IMMEDIATE`) + try { + const transactionDriver: SQLiteDriver = { + ...driver, + run: (sql, params) => + failTransactionRun?.(sql) + ? Promise.reject(new Error(`injected transaction failure`)) + : driver.run(sql, params), + } + const result = await transaction(transactionDriver) + database.exec(`COMMIT`) + return result + } catch (error) { + database.exec(`ROLLBACK`) + throw error + } + }, + } + return driver +} + +function deferred() { + let resolve!: () => void + const promise = new Promise((complete) => { + resolve = complete + }) + return { promise, resolve } +} + +async function reachCheckpoint( + promise: Promise, + checkpoint: string, +): Promise { + let timer: ReturnType | undefined + try { + await Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error(`Did not reach checkpoint: ${checkpoint}`)), + 1_000, + ) + }), + ]) + } finally { + if (timer !== undefined) clearTimeout(timer) + } +} + +function closeDatabasePreservingPrimary( + database: DatabaseSync, + primaryFailure: unknown, +): never | void { + let cleanupFailure: unknown + try { + database.close() + } catch (error) { + cleanupFailure = error + } + + if (primaryFailure !== undefined) { + const failure = + primaryFailure instanceof Error + ? primaryFailure + : new Error(`SQLite resume snapshot failed`, { cause: primaryFailure }) + if (cleanupFailure !== undefined) { + Object.defineProperty(failure, `cleanupFailures`, { + value: [cleanupFailure], + enumerable: true, + }) + } + throw failure + } + if (cleanupFailure !== undefined) throw cleanupFailure +} + +async function observeCachedSchemaState( + adapter: SQLiteCorePersistenceAdapter, + driver: SQLiteDriver, + collectionId: string, +): Promise { + const snapshot = await adapter.loadResumeSnapshot(collectionId) + const registryRows = await driver.query<{ schema_version: number }>( + `SELECT schema_version FROM collection_registry WHERE collection_id = ?`, + [collectionId], + ) + const versionRows = await driver.query<{ + key_set_evidence_available: number + key_set_evidence_incompatible: number + }>( + `SELECT key_set_evidence_available, key_set_evidence_incompatible + FROM collection_version + WHERE collection_id = ?`, + [collectionId], + ) + const expectedKeys = await driver.query<{ key: string }>( + `SELECT key + FROM collection_expected_keys + WHERE collection_id = ? + ORDER BY key`, + [collectionId], + ) + const appliedTransactions = await driver.query<{ + term: number + seq: number + tx_id: string + row_version: number + }>( + `SELECT term, seq, tx_id, row_version + FROM applied_tx + WHERE collection_id = ? + ORDER BY term, seq`, + [collectionId], + ) + + return { + schemaVersion: registryRows[0]?.schema_version ?? -1, + resetEpoch: snapshot.resetEpoch, + rows: snapshot.rows + .map(({ key, value }) => ({ key, value })) + .sort((left, right) => String(left.key).localeCompare(String(right.key))), + metadata: snapshot.collectionMetadata.sort((left, right) => + left.key.localeCompare(right.key), + ), + appliedTransactions: appliedTransactions.map( + ({ term, seq, tx_id, row_version }) => ({ + term, + seq, + txId: tx_id, + rowVersion: row_version, + }), + ), + keyEvidence: { + available: versionRows[0]?.key_set_evidence_available ?? -1, + incompatible: versionRows[0]?.key_set_evidence_incompatible ?? -1, + expectedKeys: expectedKeys.map(({ key }) => key), + }, + } +} + +/** + * # Which generation does a SQLite resume snapshot certify? + * + * `loadResumeSnapshot` must return rows, collection metadata, applied position, + * reset epoch, and key-set evidence from one atomic persisted generation. Raw + * key loss makes that evidence incompatible until a full replacement establishes + * a new baseline; concurrent schema migration may advance but never downgrade or + * repeat the observed generation. These laws refine the shared persistence and + * schema-mismatch contracts exercised by sqlite-core-adapter.test.ts. + * + * `CachedSchemaState` is the independent projection: complete rows and metadata, + * transaction position, schema/reset lineage, and expected keys. The history + * grammar crosses external row loss, full replacement, two legacy-schema + * adapters, stale reads, newer-schema observation, and a cached writer racing a + * reset. Expected membership and lineage come from the declared transition, not + * from the adapter's SQL or internal branch structure. + * + * The production driver runs two real `SQLiteCorePersistenceAdapter` instances + * over one node:sqlite database and holds the exact transaction or snapshot + * boundary needed for each interleaving. At the settled snapshot checkpoint it + * compares the entire projected schema state; the held boundary and reset epoch + * are reach witnesses, while compatible reopen and recertifying truncate cases + * prevent an oracle that merely rejects every resume. + * + * This narrow fixture supplies the same-connection concurrency seam that the + * serialized copy-on-commit CLI harness cannot. It does not claim native host + * execution or judge whether a consumer such as Electric may use the certified + * cursor; those remain separate driver-contract and Electric recovery owners. + */ +describe(`SQLite resume snapshots`, () => { + it(`keeps raw key loss sticky until a full replacement recertifies the baseline`, async () => { + const database = new DatabaseSync(`:memory:`) + let primaryFailure: unknown + try { + const collectionId = `resume-ledger` + const tableName = createPersistedTableName(collectionId, `c`) + let rejectCollectionInsert = false + const driver = createDriver( + database, + (sql) => + rejectCollectionInsert && sql.includes(`INSERT INTO "${tableName}"`), + ) + const adapter = new SQLiteCorePersistenceAdapter({ driver }) + await adapter.applyCommittedTx(collectionId, { + txId: `seed`, + term: 1, + seq: 1, + rowVersion: 1, + mutations: [ + { type: `insert`, key: 1, value: { id: 1, name: `one` } }, + { type: `insert`, key: 2, value: { id: 2, name: `two` } }, + ], + collectionMetadataMutations: [ + { type: `set`, key: `cursor`, value: `10_0` }, + ], + }) + + const initial = await adapter.loadResumeSnapshot(collectionId) + expect(initial.keySet).toEqual({ status: `consistent` }) + expect(initial.rows.map(({ key }) => key)).toEqual([1, 2]) + expect(initial.collectionMetadata).toEqual([ + { key: `cursor`, value: `10_0` }, + ]) + + await adapter.applyCommittedTx(collectionId, { + txId: `normal-update`, + term: 1, + seq: 2, + rowVersion: 2, + mutations: [ + { type: `update`, key: 1, value: { id: 1, name: `updated-one` } }, + ], + }) + const updated = await adapter.loadResumeSnapshot(collectionId) + expect(updated.rows.find(({ key }) => key === 1)?.value).toEqual({ + id: 1, + name: `updated-one`, + }) + expect(updated.keySet).toEqual({ status: `consistent` }) + + await adapter.applyCommittedTx(collectionId, { + txId: `normal-delete`, + term: 1, + seq: 3, + rowVersion: 3, + mutations: [{ type: `delete`, key: 2, value: { id: 2, name: `two` } }], + }) + await adapter.applyCommittedTx(collectionId, { + txId: `normal-delete`, + term: 1, + seq: 3, + rowVersion: 3, + mutations: [{ type: `delete`, key: 2, value: { id: 2, name: `two` } }], + }) + const deleted = await adapter.loadResumeSnapshot(collectionId) + expect(deleted.rows.map(({ key }) => key)).toEqual([1]) + expect(deleted.keySet).toEqual({ status: `consistent` }) + + rejectCollectionInsert = true + await expect( + adapter.applyCommittedTx(collectionId, { + txId: `rolled-back-insert`, + term: 1, + seq: 4, + rowVersion: 4, + mutations: [{ type: `insert`, key: 3, value: { id: 3 } }], + }), + ).rejects.toThrow(`injected transaction failure`) + rejectCollectionInsert = false + const rolledBack = await adapter.loadResumeSnapshot(collectionId) + expect(rolledBack.rows.map(({ key }) => key)).toEqual([1]) + expect(rolledBack.keySet).toEqual({ status: `consistent` }) + expect( + await driver.query<{ count: number }>( + `SELECT COUNT(*) AS count FROM collection_expected_keys WHERE collection_id = ?`, + [collectionId], + ), + ).toEqual([{ count: 1 }]) + expect( + await driver.query<{ count: number }>( + `SELECT COUNT(*) AS count FROM applied_tx WHERE collection_id = ? AND tx_id = ?`, + [collectionId, `normal-delete`], + ), + ).toEqual([{ count: 1 }]) + + await adapter.applyCommittedTx(collectionId, { + txId: `restore-second-row`, + term: 1, + seq: 5, + rowVersion: 5, + mutations: [{ type: `insert`, key: 2, value: { id: 2, name: `two` } }], + }) + + await driver.run(`DELETE FROM "${tableName}" WHERE key = ?`, [ + database + .prepare(`SELECT key FROM "${tableName}" ORDER BY key LIMIT 1`) + .get()!.key, + ]) + await adapter.applyCommittedTx(collectionId, { + txId: `metadata-after-loss`, + term: 1, + seq: 6, + rowVersion: 6, + mutations: [], + collectionMetadataMutations: [ + { type: `set`, key: `cursor`, value: `11_0` }, + ], + }) + expect((await adapter.loadResumeSnapshot(collectionId)).keySet).toEqual({ + status: `incompatible`, + }) + + await adapter.applyCommittedTx(collectionId, { + txId: `full-replacement`, + term: 1, + seq: 7, + rowVersion: 7, + truncate: true, + mutations: [ + { type: `insert`, key: 1, value: { id: 1, name: `one` } }, + { type: `insert`, key: 2, value: { id: 2, name: `two` } }, + ], + }) + expect((await adapter.loadResumeSnapshot(collectionId)).keySet).toEqual({ + status: `consistent`, + }) + + await driver.run( + `UPDATE "${tableName}" SET key = key || '-replacement' WHERE rowid = (SELECT MIN(rowid) FROM "${tableName}")`, + ) + const substituted = await adapter.loadResumeSnapshot(collectionId) + expect(substituted.rows).toHaveLength(2) + expect(substituted.keySet).toEqual({ status: `incompatible` }) + } catch (error) { + primaryFailure = error + } finally { + closeDatabasePreservingPrimary(database, primaryFailure) + } + }) + + it(`migrates concurrent legacy schemas and stays unknown until truncate`, async () => { + const database = new DatabaseSync(`:memory:`) + let releasePending = () => {} + let primaryFailure: unknown + try { + const driver = createDriver(database) + const collectionId = `legacy-ledger` + const tableName = createPersistedTableName(collectionId, `c`) + const tombstoneTableName = createPersistedTableName(collectionId, `t`) + await driver.exec( + `CREATE TABLE collection_registry ( + collection_id TEXT PRIMARY KEY, + table_name TEXT NOT NULL UNIQUE, + tombstone_table_name TEXT NOT NULL UNIQUE, + schema_version INTEGER NOT NULL, + updated_at INTEGER NOT NULL + )`, + ) + await driver.run( + `INSERT INTO collection_registry + (collection_id, table_name, tombstone_table_name, schema_version, updated_at) + VALUES (?, ?, ?, 1, 0)`, + [collectionId, tableName, tombstoneTableName], + ) + await driver.exec( + `CREATE TABLE "${tableName}" ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL, + metadata TEXT, + row_version INTEGER NOT NULL + )`, + ) + await driver.exec( + `CREATE TABLE "${tombstoneTableName}" ( + key TEXT PRIMARY KEY, + value TEXT, + row_version INTEGER NOT NULL, + deleted_at TEXT NOT NULL + )`, + ) + await driver.exec( + `CREATE TABLE collection_version ( + collection_id TEXT PRIMARY KEY, + latest_row_version INTEGER NOT NULL + )`, + ) + await driver.run( + `INSERT INTO collection_version (collection_id, latest_row_version) + VALUES (?, 0)`, + [collectionId], + ) + + const bothSawLegacyColumns = deferred() + let legacyColumnReaders = 0 + const migrationRelease = deferred() + releasePending = migrationRelease.resolve + const createMigrationDriver = (): SQLiteDriver => ({ + ...driver, + query: async (sql: string, params?: ReadonlyArray) => { + const rows = await driver.query(sql, params) + if (sql.includes(`PRAGMA table_info(collection_version)`)) { + legacyColumnReaders++ + if (legacyColumnReaders === 2) bothSawLegacyColumns.resolve() + await migrationRelease.promise + } + return rows + }, + }) + const migrated = new SQLiteCorePersistenceAdapter({ + driver: createMigrationDriver(), + }) + const concurrentMigrated = new SQLiteCorePersistenceAdapter({ + driver: createMigrationDriver(), + }) + const initialize = (adapter: SQLiteCorePersistenceAdapter) => + ( + adapter as unknown as { ensureInitialized: () => Promise } + ).ensureInitialized() + const migrations = Promise.all([ + initialize(migrated), + initialize(concurrentMigrated), + ]) + await reachCheckpoint( + bothSawLegacyColumns.promise, + `both adapters observed the legacy collection_version schema`, + ) + migrationRelease.resolve() + await migrations + const migratedColumns = await driver.query<{ name: string }>( + `PRAGMA table_info(collection_version)`, + ) + expect(migratedColumns.map(({ name }) => name)).toEqual( + expect.arrayContaining([ + `key_set_evidence_available`, + `key_set_evidence_incompatible`, + ]), + ) + + expect((await migrated.loadResumeSnapshot(collectionId)).keySet).toEqual({ + status: `unknown`, + }) + await migrated.applyCommittedTx(collectionId, { + txId: `legacy-insert`, + term: 1, + seq: 1, + rowVersion: 1, + mutations: [{ type: `insert`, key: 1, value: { id: 1, n: 1 } }], + }) + expect((await migrated.loadResumeSnapshot(collectionId)).keySet).toEqual({ + status: `unknown`, + }) + + await migrated.applyCommittedTx(collectionId, { + txId: `legacy-replacement`, + term: 1, + seq: 2, + rowVersion: 2, + truncate: true, + mutations: [{ type: `insert`, key: 1, value: { id: 1, n: 2 } }], + }) + expect((await migrated.loadResumeSnapshot(collectionId)).keySet).toEqual({ + status: `consistent`, + }) + } catch (error) { + primaryFailure = error + } finally { + releasePending() + closeDatabasePreservingPrimary(database, primaryFailure) + } + }) + + it(`does not repeat a schema reset observed through a stale adapter read`, async () => { + const database = new DatabaseSync(`:memory:`) + let releasePending = () => {} + let primaryFailure: unknown + try { + const driver = createDriver(database) + const collectionId = `concurrent-schema-reset` + const original = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion: 1, + }) + await original.applyCommittedTx(collectionId, { + txId: `seed`, + term: 1, + seq: 1, + rowVersion: 1, + mutations: [{ type: `insert`, key: 1, value: { id: 1 } }], + collectionMetadataMutations: [ + { type: `set`, key: `cursor`, value: `old` }, + ], + }) + + const staleRead = deferred() + const releaseStaleRead = deferred() + releasePending = releaseStaleRead.resolve + let intercepted = false + const gatedDriver: SQLiteDriver = { + ...driver, + query: async (sql: string, params?: ReadonlyArray) => { + const rows = await driver.query(sql, params) + if (!intercepted && sql.includes(`FROM collection_registry`)) { + intercepted = true + staleRead.resolve() + await releaseStaleRead.promise + } + return rows + }, + } + const staleAdapter = new SQLiteCorePersistenceAdapter({ + driver: gatedDriver, + schemaVersion: 2, + }) + const staleLoad = staleAdapter.loadSubset(collectionId, {}) + await reachCheckpoint( + staleRead.promise, + `stale schema-v1 registry read before competing reset`, + ) + + const winner = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion: 2, + }) + await winner.loadSubset(collectionId, {}) + await winner.applyCommittedTx(collectionId, { + txId: `recovery-write`, + term: 2, + seq: 1, + rowVersion: 1, + mutations: [{ type: `insert`, key: 2, value: { id: 2 } }], + collectionMetadataMutations: [ + { type: `set`, key: `cursor`, value: `new` }, + ], + }) + + releaseStaleRead.resolve() + await staleLoad + const snapshot = await winner.loadResumeSnapshot(collectionId) + expect(snapshot.resetEpoch).toBe(1) + expect(snapshot.rows.map(({ key }) => key)).toEqual([2]) + expect(snapshot.collectionMetadata).toEqual([ + { key: `cursor`, value: `new` }, + ]) + } catch (error) { + primaryFailure = error + } finally { + releasePending() + closeDatabasePreservingPrimary(database, primaryFailure) + } + }) + + it(`refuses to downgrade a newer schema observed after a stale read`, async () => { + const database = new DatabaseSync(`:memory:`) + let releasePending = () => {} + let primaryFailure: unknown + try { + const driver = createDriver(database) + const collectionId = `divergent-schema-reset` + const original = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion: 1, + }) + await original.applyCommittedTx(collectionId, { + txId: `seed`, + term: 1, + seq: 1, + rowVersion: 1, + mutations: [{ type: `insert`, key: 1, value: { id: 1 } }], + }) + + const staleRead = deferred() + const releaseStaleRead = deferred() + releasePending = releaseStaleRead.resolve + let intercepted = false + const staleDriver: SQLiteDriver = { + ...driver, + query: async (sql: string, params?: ReadonlyArray) => { + const rows = await driver.query(sql, params) + if (!intercepted && sql.includes(`FROM collection_registry`)) { + intercepted = true + staleRead.resolve() + await releaseStaleRead.promise + } + return rows + }, + } + const staleV2 = new SQLiteCorePersistenceAdapter({ + driver: staleDriver, + schemaVersion: 2, + }) + const staleLoad = staleV2.loadSubset(collectionId, {}) + await reachCheckpoint( + staleRead.promise, + `stale schema-v1 registry read before schema-v3 reset`, + ) + + const winnerV3 = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion: 3, + }) + await winnerV3.loadSubset(collectionId, {}) + await winnerV3.applyCommittedTx(collectionId, { + txId: `winner-write`, + term: 3, + seq: 1, + rowVersion: 1, + mutations: [{ type: `insert`, key: 3, value: { id: 3 } }], + collectionMetadataMutations: [ + { type: `set`, key: `cursor`, value: `v3` }, + ], + }) + + releaseStaleRead.resolve() + await expect(staleLoad).rejects.toThrow( + `Schema version changed concurrently`, + ) + const snapshot = await winnerV3.loadResumeSnapshot(collectionId) + expect(snapshot.resetEpoch).toBe(1) + expect(snapshot.rows.map(({ key }) => key)).toEqual([3]) + expect(snapshot.collectionMetadata).toEqual([ + { key: `cursor`, value: `v3` }, + ]) + expect( + await driver.query<{ schema_version: number }>( + `SELECT schema_version FROM collection_registry WHERE collection_id = ?`, + [collectionId], + ), + ).toEqual([{ schema_version: 3 }]) + } catch (error) { + primaryFailure = error + } finally { + releasePending() + closeDatabasePreservingPrimary(database, primaryFailure) + } + }) + + it(`rejects a committed transaction from a cached adapter after another adapter resets the schema`, async () => { + const database = new DatabaseSync(`:memory:`) + let primaryFailure: unknown + try { + const driver = createDriver(database) + const collectionId = `cached-schema-write` + const staleV1 = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion: 1, + }) + await staleV1.applyCommittedTx(collectionId, { + txId: `seed-v1`, + term: 1, + seq: 1, + rowVersion: 1, + mutations: [{ type: `insert`, key: 1, value: { id: 1 } }], + collectionMetadataMutations: [ + { type: `set`, key: `cursor`, value: `v1` }, + ], + }) + + const currentV2 = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion: 2, + }) + await currentV2.loadSubset(collectionId, {}) + await currentV2.applyCommittedTx(collectionId, { + txId: `seed-v2`, + term: 2, + seq: 1, + rowVersion: 1, + mutations: [{ type: `insert`, key: 2, value: { id: 2 } }], + collectionMetadataMutations: [ + { type: `set`, key: `cursor`, value: `v2` }, + ], + }) + const before = await observeCachedSchemaState( + currentV2, + driver, + collectionId, + ) + + let lateWriteError: unknown + try { + await staleV1.applyCommittedTx(collectionId, { + txId: `late-v1`, + term: 3, + seq: 1, + rowVersion: 2, + truncate: true, + mutations: [{ type: `insert`, key: 3, value: { id: 3 } }], + collectionMetadataMutations: [ + { type: `set`, key: `cursor`, value: `late-v1` }, + { type: `set`, key: `late`, value: true }, + ], + }) + } catch (error) { + lateWriteError = error + } + const after = await observeCachedSchemaState( + currentV2, + driver, + collectionId, + ) + + expect( + { + error: + lateWriteError instanceof Error + ? { name: lateWriteError.name, message: lateWriteError.message } + : lateWriteError, + before, + after, + }, + `a cached adapter must reject at its transaction boundary without changing any durable state`, + ).toEqual({ + error: { + name: `InvalidPersistedCollectionConfigError`, + message: + `Schema version mismatch for collection "cached-schema-write": ` + + `found 2, expected 1. Refusing to apply a committed transaction through a stale cached adapter.`, + }, + before, + after: before, + }) + } catch (error) { + primaryFailure = error + } finally { + closeDatabasePreservingPrimary(database, primaryFailure) + } + }) +}) diff --git a/packages/electric-db-collection/src/electric.ts b/packages/electric-db-collection/src/electric.ts index dd1d2d5848..94c64ee2fe 100644 --- a/packages/electric-db-collection/src/electric.ts +++ b/packages/electric-db-collection/src/electric.ts @@ -72,6 +72,11 @@ type ElectricSyncMetadataWithHydration = SyncMetadataApi & { whenHydrated?: () => Promise // Capability marker for wrappers predating the hydration barrier. scanPersisted?: unknown + certifyPersistedResume?: () => Promise + getPersistedKeySetEvidence?: () => + | { status: `unknown` | `consistent` | `incompatible` } + | undefined + expectCurrentCommitInResumeSnapshot?: () => void } } @@ -1762,6 +1767,13 @@ function createElectricSync>( | undefined const scanPersisted = persistedMetadata?.row.scanPersisted const whenHydrated = persistedMetadata?.row.whenHydrated + const certifyPersistedResume = + persistedMetadata?.row.certifyPersistedResume + const getPersistedKeySetEvidence = + persistedMetadata?.row.getPersistedKeySetEvidence + const expectCurrentCommitInResumeSnapshot = + persistedMetadata?.row.expectCurrentCommitInResumeSnapshot + const persistedKeySetEvidence = getPersistedKeySetEvidence?.() const persistedResumeState = getNewestElectricResumeState( readPersistedResumeState(), @@ -1780,6 +1792,12 @@ function createElectricSync>( persistedResumeState?.kind === `resume` && scanPersisted !== undefined && whenHydrated === undefined + // A pre-ledger `unknown` baseline cannot justify a non-initial cursor. + // One fresh replacement establishes consistent evidence for later resumes. + const lacksCompletePersistedKeySet = + persistedResumeState?.kind === `resume` && + getPersistedKeySetEvidence !== undefined && + persistedKeySetEvidence?.status !== `consistent` if (hasUnverifiablePersistedResume && !warnedUnverifiableResume) { warnedUnverifiableResume = true console.warn( @@ -1791,6 +1809,7 @@ function createElectricSync>( shapeOptions.handle === undefined && persistedResumeState !== undefined && (persistedResumeState.kind === `reset` || + lacksCompletePersistedKeySet || (!retainsTagState && persistedResumeState.requiresTagState !== false)) const canUsePersistedResume = shapeOptions.offset === undefined && @@ -1799,7 +1818,7 @@ function createElectricSync>( !hasIncompatiblePersistedResume && !hasUnverifiablePersistedResume && // Cached rows do not contain authoritative tag/active-condition state. - // Unknown (older) metadata is conservative; untagged shapes still resume. + // Only a complete adapter ledger can justify a persisted cursor. !needsFullSnapshot const hasExplicitResumeOffset = shapeOptions.offset !== undefined && shapeOptions.offset !== `-1` @@ -1807,6 +1826,8 @@ function createElectricSync>( clearTagTrackingState() } const receivesCompleteRows = shapeOptions.params?.replica === `full` + const requiresKeySetCertification = + canUsePersistedResume && certifyPersistedResume !== undefined // Eager and progressive streams that start after the initial offset can // only apply partial updates when the local materialization is complete. const requiresCompleteResume = @@ -1964,7 +1985,9 @@ function createElectricSync>( metadata?.collection.set(`electric:resume`, resumeState) } - const commitResetResumeMetadataImmediately = () => { + const commitResetResumeMetadataImmediately = ( + expectInResumeSnapshot = false, + ) => { const resetState: ElectricResumeState = { kind: `reset`, updatedAt: Date.now(), @@ -1974,6 +1997,9 @@ function createElectricSync>( if (metadata) { begin({ immediate: true }) metadata.collection.set(`electric:resume`, resetState) + if (expectInResumeSnapshot) { + expectCurrentCommitInResumeSnapshot?.() + } commit() } } @@ -1983,7 +2009,10 @@ function createElectricSync>( hasUnverifiablePersistedResume || (needsFullSnapshot && persistedResumeState.kind === `resume`) ) { - commitResetResumeMetadataImmediately() + // This reset is part of the current runtime's startup decision. The + // persisted wrapper may commit it before loading the atomic baseline, + // so carry ownership of exactly this generation into certification. + commitResetResumeMetadataImmediately(true) } /** @@ -2064,8 +2093,29 @@ function createElectricSync>( const resumeKeysPromise = requiresCompleteResume || freshSnapshotPending - ? whenHydrated?.() - : undefined + ? whenHydrated + ? (async () => { + await whenHydrated() + if ( + persistedKeySetEvidence?.status !== `incompatible` && + getPersistedKeySetEvidence?.()?.status === `incompatible` + ) { + throw new Error( + `Electric persisted resume baseline became incompatible during hydration`, + ) + } + })() + : undefined + : requiresKeySetCertification + ? (async () => { + await certifyPersistedResume() + if (getPersistedKeySetEvidence?.()?.status === `incompatible`) { + throw new Error( + `Electric persisted resume baseline became incompatible during certification`, + ) + } + })() + : undefined let areResumeKeysReady = !resumeKeysPromise const pendingResumeBatches: Array>> = [] let unsubscribeStream: () => void = () => {} diff --git a/packages/electric-db-collection/tests/electric-recovery-oracle.test.ts b/packages/electric-db-collection/tests/electric-recovery-oracle.test.ts index 00b3fc09b4..1f3f77ab9c 100644 --- a/packages/electric-db-collection/tests/electric-recovery-oracle.test.ts +++ b/packages/electric-db-collection/tests/electric-recovery-oracle.test.ts @@ -1,9 +1,14 @@ +import { DatabaseSync } from 'node:sqlite' import { isDeepStrictEqual } from 'node:util' import { fc, test as fcTest } from '@fast-check/vitest' import { beforeEach, describe, expect, it, vi } from 'vitest' import { createCollection } from '@tanstack/db' import { ShapeStream } from '@electric-sql/client' -import { persistedCollectionOptions } from '../../db-sqlite-persistence-core/src' +import { + SQLiteCorePersistenceAdapter, + createPersistedTableName, + persistedCollectionOptions, +} from '../../db-sqlite-persistence-core/src' import { electricCollectionOptions } from '../src/electric' import type { Message, Row } from '@electric-sql/client' import type { @@ -11,13 +16,91 @@ import type { PersistedTx, PersistenceAdapter, ProtocolEnvelope, + SQLiteDriver, } from '../../db-sqlite-persistence-core/src' import type { ElectricCollectionUtils, ElectricSyncMode } from '../src/electric' +/** + * # What remains visible while a persisted Electric replica repairs itself? + * + * Hydrated rows and resume metadata provide the last complete public snapshot. + * A must-refetch starts a private replacement. Until that replacement is fully + * applied, readers may see an earlier permitted snapshot but never a torn mix. + * Failure keeps the old public rows and records repair debt; later success may + * replace them atomically. + * + * A plain persisted row Map and metadata Map form the reference snapshots. The + * driver controls hydration, SDK callbacks, applied receipts, cleanup, restart, + * and eager or progressive mode through the real persistence coordinator and + * Electric adapter. It records every exposed cut, not only the final rows. + * + * Legal histories vary sync mode, hydration and restart timing, reset cause, + * peer publication, stream delta, deletion, and full reload. At each recorded + * publication cut, the complete public rows must refine one permitted snapshot; + * durable rows and resume metadata are compared again after the applied or + * post-restart up-to-date checkpoint. Stale/missing canonical-row controls and + * the repaired-intermediate trace challenge the checker; generated failures + * retain fast-check's seed and shrink path. + * + * The fixture does not establish live HTTP delivery, native SQLite host + * behavior, or callback multiplicity beyond the observations named below. + */ + type Item = Row & { id: number; name: string; stable: string } type Subscriber = (messages: Array>) => void type Exposure = { cut: string; rows: Array } +function persistedResumeKind(value: unknown): unknown { + return value && typeof value === `object` + ? (value as Record).kind + : undefined +} + +function toSqliteBinding(value: unknown): string | number | bigint | null { + if (value === null || value === undefined) return null + if (typeof value === `boolean`) return value ? 1 : 0 + if ( + typeof value === `string` || + typeof value === `number` || + typeof value === `bigint` + ) { + return value + } + return String(value) +} + +function nodeSqliteDriver(database: DatabaseSync): SQLiteDriver { + const driver: SQLiteDriver = { + exec: (sql) => { + database.exec(sql) + return Promise.resolve() + }, + query: (sql, params = []) => + Promise.resolve( + database + .prepare(sql) + .all(...params.map(toSqliteBinding)) + .map((row) => ({ ...row })) as Array, + ), + run: (sql, params = []) => { + database.prepare(sql).run(...params.map(toSqliteBinding)) + return Promise.resolve() + }, + transaction: async (transaction) => { + database.exec(`BEGIN IMMEDIATE`) + try { + const result = await transaction(driver) + database.exec(`COMMIT`) + return result + } catch (error) { + database.exec(`ROLLBACK`) + throw error + } + }, + } + return driver +} + function expectWholeRecoveryTrace( entries: Array, allowed: Array>, @@ -62,6 +145,7 @@ function deferred() { } const oldRow: Item = { id: 1, name: `old`, stable: `stable-1` } +const otherOldRow: Item = { id: 3, name: `other-old`, stable: `stable-3` } const freshRow: Item = { id: 2, name: `fresh`, stable: `stable-2` } const upToDate: Message = { headers: { control: `up-to-date` } } @@ -196,12 +280,514 @@ const scenarios = ([`eager`, `progressive`] as const).flatMap((syncMode) => ), ) +type PersistedRestartScenario = { + rowState: `empty` | `nonempty` + resumeKind: `none` | `reset` | `resume` + transition: + | `compatible-reopen` + | `schema-reset` + | `partial-restore` + | `external-row-loss` + sourceHistory: `up-to-date-only` | `replayed-insert` +} + +const persistedRestartScenarios = ([`empty`, `nonempty`] as const).flatMap( + (rowState) => + ([`none`, `reset`, `resume`] as const).flatMap((resumeKind) => + ( + [ + `compatible-reopen`, + `schema-reset`, + ...(rowState === `nonempty` + ? ([`partial-restore`, `external-row-loss`] as const) + : []), + ] as const + ).flatMap((transition) => + ([`up-to-date-only`, `replayed-insert`] as const).map( + (sourceHistory): PersistedRestartScenario => ({ + rowState, + resumeKind, + transition, + sourceHistory, + }), + ), + ), + ), +) + +type PersistedRestartObservation = { + checkpoint: `post-restart-up-to-date` + requestedOffset: string | undefined + requestedHandle: string | undefined + visibleRows: Array + durableRows: Array + durableResumeKind: unknown + status: string +} + +function expectPersistedRestartLaw( + actual: PersistedRestartObservation, + expected: PersistedRestartObservation, +): void { + expect(actual).toEqual(expected) +} + +function cloneItem(item: Item): Item { + return { id: item.id, name: item.name, stable: item.stable } +} + +function attachPersistedRestartCleanupDiagnostics( + primary: unknown, + cleanupEvidence: string, + cleanupFailures: ReadonlyArray, +): Error { + const error = + primary instanceof Error + ? primary + : new Error(`Persisted restart oracle failed with a non-Error value`, { + cause: primary, + }) + Object.defineProperty(error, `cleanupEvidence`, { + value: cleanupEvidence, + enumerable: true, + }) + if (cleanupFailures.length > 0) { + Object.defineProperty(error, `cleanupFailures`, { + value: [...cleanupFailures], + enumerable: true, + }) + } + return error +} + +async function observePersistedRestart( + scenario: PersistedRestartScenario, +): Promise<{ + observation: PersistedRestartObservation + expected: PersistedRestartObservation + cleanupEvidence: string +}> { + subscribers.length = 0 + vi.clearAllMocks() + let database: DatabaseSync | undefined + let result: + | { + observation: PersistedRestartObservation + expected: PersistedRestartObservation + cleanupEvidence: string + } + | undefined + let deferredFailure: unknown + try { + database = new DatabaseSync(`:memory:`) + const driver = nodeSqliteDriver(database) + const collectionId = `persisted-schema-reset-electric` + const modelInitialRows = + scenario.rowState === `empty` + ? [] + : [cloneItem(oldRow), cloneItem(otherOldRow)] + const modelCanonicalRows = + scenario.sourceHistory === `replayed-insert` + ? [...modelInitialRows, cloneItem(freshRow)] + : modelInitialRows + modelCanonicalRows.sort((left, right) => left.id - right.id) + const productionSeedRows = + scenario.rowState === `empty` + ? [] + : [cloneItem(oldRow), cloneItem(otherOldRow)] + const canResume = + scenario.transition === `compatible-reopen` && + scenario.resumeKind === `resume` + const expected: PersistedRestartObservation = { + checkpoint: `post-restart-up-to-date`, + requestedOffset: canResume ? `10_0` : undefined, + requestedHandle: canResume ? `shape-old` : undefined, + visibleRows: modelCanonicalRows.map(cloneItem), + durableRows: modelCanonicalRows.map(cloneItem), + durableResumeKind: `resume`, + status: `ready`, + } + const resumeState = + scenario.resumeKind === `none` + ? [] + : [ + { + type: `set` as const, + key: `electric:resume`, + value: + scenario.resumeKind === `reset` + ? { kind: `reset`, updatedAt: 1 } + : { + kind: `resume`, + requiresTagState: false, + offset: `10_0`, + handle: `shape-old`, + shapeId: `{"params":{"table":"test_table"},"url":"http://test-url"}`, + updatedAt: 1, + }, + }, + ] + const originalAdapter = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion: 1, + }) + await originalAdapter.applyCommittedTx(collectionId, { + txId: `seed-electric-baseline`, + term: 1, + seq: 1, + rowVersion: 1, + mutations: productionSeedRows.map((row) => ({ + type: `insert` as const, + key: row.id, + value: cloneItem(row), + })), + collectionMetadataMutations: resumeState, + }) + + const usesSchemaReset = + scenario.transition === `schema-reset` || + scenario.transition === `partial-restore` + const schemaVersion = usesSchemaReset ? 2 : 1 + const restartedAdapter = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion, + schemaMismatchPolicy: `sync-present-reset`, + }) + await restartedAdapter.loadSubset(collectionId, {}) + if (scenario.transition === `partial-restore`) { + await restartedAdapter.applyCommittedTx(collectionId, { + txId: `partial-electric-restore`, + term: 2, + seq: 1, + rowVersion: 2, + mutations: [ + { type: `insert`, key: oldRow.id, value: cloneItem(oldRow) }, + ], + }) + } else if (scenario.transition === `external-row-loss`) { + const collectionTable = createPersistedTableName(collectionId, `c`) + await driver.run( + `DELETE FROM "${collectionTable}" WHERE json_extract(value, '$.id') = ?`, + [oldRow.id], + ) + } + + const collection = createCollection( + persistedCollectionOptions< + Item, + string | number, + never, + ElectricCollectionUtils + >({ + ...electricCollectionOptions({ + id: collectionId, + shapeOptions: { + url: `http://test-url`, + params: { table: `test_table` }, + }, + syncMode: `eager`, + getKey: (row) => row.id, + startSync: false, + }), + persistence: { adapter: restartedAdapter }, + }), + ) + + let streamSubscriber: Subscriber | undefined + let subscription: ReturnType | undefined + let observation!: PersistedRestartObservation + let semanticFailure: unknown + let processingFailure: unknown + const cleanupFailures: Array = [] + let cleanupEvidence = `not-run` + let publicationEvents = 0 + try { + collection.startSyncImmediate() + subscription = collection.subscribeChanges( + () => { + publicationEvents++ + }, + { includeInitialState: false }, + ) + await vi.waitFor(() => expect(subscribers).toHaveLength(1)) + streamSubscriber = subscribers[0] + const request = vi.mocked(ShapeStream).mock.calls[0]?.[0] as + | { offset?: string; handle?: string } + | undefined + if (!streamSubscriber || !request) { + throw new Error(`Persisted Electric stream did not reach subscription`) + } + + // The source chooses a legal response from the request production made. + // A resumed request receives only changes since its cursor; a fresh request + // receives the complete current snapshot. Both finish with up-to-date. + const sourceRows = + request.offset === undefined + ? modelCanonicalRows.map(cloneItem) + : scenario.sourceHistory === `replayed-insert` + ? [cloneItem(freshRow)] + : [] + streamSubscriber([ + ...sourceRows.map((row) => change(`insert`, cloneItem(row))), + structuredClone(upToDate), + ]) + await vi.waitFor(() => expect(collection.status).toBe(`ready`)) + await new Promise((resolve) => setTimeout(resolve, 0)) + await new Promise((resolve) => setTimeout(resolve, 0)) + + const durableRows = await restartedAdapter.loadSubset(collectionId, {}) + const durableMetadata = + await restartedAdapter.loadCollectionMetadata(collectionId) + observation = { + checkpoint: `post-restart-up-to-date`, + requestedOffset: request.offset, + requestedHandle: request.handle, + visibleRows: Array.from( + collection.values(), + ({ id, name, stable }) => ({ + id, + name, + stable, + }), + ).sort((left, right) => left.id - right.id), + durableRows: durableRows + .map(({ value }) => value as Item) + .map(({ id, name, stable }) => ({ id, name, stable })) + .sort((left, right) => left.id - right.id), + durableResumeKind: persistedResumeKind( + durableMetadata.find(({ key }) => key === `electric:resume`)?.value, + ), + status: collection.status, + } + try { + expectPersistedRestartLaw(observation, expected) + } catch (error) { + semanticFailure = error + } + } catch (error) { + processingFailure = error + } finally { + let collectionCleanupCompleted = false + try { + subscription?.unsubscribe() + } catch (error) { + cleanupFailures.push(error) + } + try { + await collection.cleanup() + collectionCleanupCompleted = true + } catch (error) { + cleanupFailures.push(error) + } + if (collectionCleanupCompleted && streamSubscriber) { + try { + const beforeLatePublicRows = Array.from( + collection.values(), + ({ id, name, stable }) => ({ id, name, stable }), + ).sort((left, right) => left.id - right.id) + const beforeLateStatus = collection.status + const beforeLateEvents = publicationEvents + const beforeLatePending = + collection._state.pendingSyncedTransactions.map( + ({ committed }) => committed, + ) + const beforeLateRows = await restartedAdapter.loadSubset( + collectionId, + {}, + ) + const beforeLateMetadata = + await restartedAdapter.loadCollectionMetadata(collectionId) + const beforeLateApplied = await driver.query<{ + term: number + seq: number + tx_id: string + row_version: number + }>( + `SELECT term, seq, tx_id, row_version FROM applied_tx WHERE collection_id = ? ORDER BY term, seq`, + [collectionId], + ) + // Deliver both data and the commit boundary. An uncommitted message + // would not challenge late application/publication after retirement. + streamSubscriber([ + change(`insert`, { id: 77, name: `late`, stable: `late` }), + structuredClone(upToDate), + ]) + await new Promise((resolve) => setTimeout(resolve, 0)) + await new Promise((resolve) => setTimeout(resolve, 0)) + const afterLatePublicRows = Array.from( + collection.values(), + ({ id, name, stable }) => ({ id, name, stable }), + ).sort((left, right) => left.id - right.id) + const afterLateRows = await restartedAdapter.loadSubset( + collectionId, + {}, + ) + const afterLateMetadata = + await restartedAdapter.loadCollectionMetadata(collectionId) + const afterLateApplied = await driver.query<{ + term: number + seq: number + tx_id: string + row_version: number + }>( + `SELECT term, seq, tx_id, row_version FROM applied_tx WHERE collection_id = ? ORDER BY term, seq`, + [collectionId], + ) + const afterLatePending = + collection._state.pendingSyncedTransactions.map( + ({ committed }) => committed, + ) + cleanupEvidence = + isDeepStrictEqual(afterLatePublicRows, beforeLatePublicRows) && + collection.status === beforeLateStatus && + publicationEvents === beforeLateEvents && + isDeepStrictEqual(beforeLatePending, []) && + isDeepStrictEqual(afterLatePending, beforeLatePending) && + isDeepStrictEqual(afterLateRows, beforeLateRows) && + isDeepStrictEqual(afterLateMetadata, beforeLateMetadata) && + isDeepStrictEqual(afterLateApplied, beforeLateApplied) + ? `passed: committed retired-stream delivery changed no public rows, events, status, pending transactions, durable rows, metadata, or applied effects` + : `failed: committed retired-stream delivery changed public or pending/durable effects` + } catch (error) { + cleanupFailures.push(error) + } + } else if (collectionCleanupCompleted) { + cleanupEvidence = `passed: collection cleanup completed before a retired stream became available` + } + if (cleanupFailures.length > 0) { + cleanupEvidence = + `failed: ${cleanupFailures.map((error) => String(error)).join(`; `)}; ` + + `retiredProbe=${cleanupEvidence}` + } + } + + if (semanticFailure !== undefined) { + deferredFailure = attachPersistedRestartCleanupDiagnostics( + new Error( + `Persisted Electric reset/resume violation. ` + + `checkpoint=${observation.checkpoint} ` + + `scenario=${JSON.stringify(scenario)} ` + + `expected=${JSON.stringify(expected)} ` + + `actual=${JSON.stringify(observation)} ` + + `cleanup=${cleanupEvidence}`, + { cause: semanticFailure }, + ), + cleanupEvidence, + cleanupFailures, + ) + } else if (processingFailure !== undefined) { + deferredFailure = attachPersistedRestartCleanupDiagnostics( + processingFailure, + cleanupEvidence, + cleanupFailures, + ) + } else if (cleanupFailures.length > 0) { + deferredFailure = new AggregateError(cleanupFailures, cleanupEvidence) + } else if (!cleanupEvidence.startsWith(`passed:`)) { + deferredFailure = new Error(cleanupEvidence) + } else { + result = { + observation, + expected, + cleanupEvidence, + } + } + } catch (error) { + deferredFailure ??= + error instanceof Error + ? error + : new Error(`Persisted restart setup failed with a non-Error value`, { + cause: error, + }) + } finally { + try { + database?.close() + } catch (closeError) { + if (deferredFailure instanceof Error) { + Object.defineProperty(deferredFailure, `databaseCloseFailure`, { + value: closeError, + enumerable: true, + }) + } else { + deferredFailure = closeError + } + } + } + + if (deferredFailure !== undefined) throw deferredFailure + if (!result) throw new Error(`Persisted restart result was not captured`) + return result +} + describe(`persisted Electric recovery laws`, () => { beforeEach(() => { subscribers.length = 0 vi.clearAllMocks() }) + /** + * This matrix is the reached consumer half of the SQLite reset/resume law. + * Its independent model carries only baseline lineage and canonical server + * rows. The real SQLite adapter performs the transition, the real persisted + * wrapper hydrates it, and electricCollectionOptions chooses the ShapeStream + * request. The mock boundary supplies the installed protocol's legal rule: + * resume gets only changes since its cursor; fresh sync gets a full snapshot; + * both end at the post-restart up-to-date checkpoint. + * + * We compare request offset/handle, complete public rows, durable rows, + * durable resume kind, and readiness after schema reset, partial restore, + * and one externally deleted-row fault. Callback multiplicity, real HTTP, + * and arbitrary external edit sequences are omitted. Electric's existing + * recovery properties own partial unseen updates and SDK framing. The + * nonempty + resume + schema-reset + up-to-date-only cell preserves the + * original report at https://github.com/TanStack/db/issues/1589. + */ + it.each(persistedRestartScenarios)( + `couples $resumeKind state to $transition with $rowState rows and $sourceHistory source history`, + async (scenario) => { + const { cleanupEvidence } = await observePersistedRestart(scenario) + expect(cleanupEvidence).toContain(`passed: committed retired-stream`) + }, + ) + + it(`accepts a compatible empty resume and rejects a stale reset request or missing canonical row`, () => { + const compatibleEmpty: PersistedRestartObservation = { + checkpoint: `post-restart-up-to-date`, + requestedOffset: `10_0`, + requestedHandle: `shape-old`, + visibleRows: [], + durableRows: [], + durableResumeKind: `resume`, + status: `ready`, + } + expect(() => + expectPersistedRestartLaw( + { ...compatibleEmpty, visibleRows: [], durableRows: [] }, + compatibleEmpty, + ), + ).not.toThrow() + + const freshNonempty: PersistedRestartObservation = { + ...compatibleEmpty, + requestedOffset: undefined, + requestedHandle: undefined, + visibleRows: [{ ...oldRow }, { ...otherOldRow }], + durableRows: [{ ...oldRow }, { ...otherOldRow }], + } + expect(() => + expectPersistedRestartLaw( + { + ...freshNonempty, + requestedOffset: `10_0`, + requestedHandle: `shape-old`, + visibleRows: [{ ...oldRow }], + durableRows: [{ ...oldRow }], + }, + freshNonempty, + ), + ).toThrowError(expect.objectContaining({ name: `AssertionError` })) + }) + it(`keeps repaired intermediate publications in the persisted recovery record`, async () => { const f = fixture(`eager`) try { diff --git a/packages/electric-db-collection/tests/electric-resume-snapshot-races.test.ts b/packages/electric-db-collection/tests/electric-resume-snapshot-races.test.ts new file mode 100644 index 0000000000..5dc4f378bc --- /dev/null +++ b/packages/electric-db-collection/tests/electric-resume-snapshot-races.test.ts @@ -0,0 +1,766 @@ +import { DatabaseSync } from 'node:sqlite' +import { beforeEach, describe, expect, it, vi } from 'vitest' +import { createCollection } from '@tanstack/db' +import { ShapeStream } from '@electric-sql/client' +import { + SQLiteCorePersistenceAdapter, + createPersistedTableName, + persistedCollectionOptions, +} from '../../db-sqlite-persistence-core/src' +import { electricCollectionOptions } from '../src/electric' +import type { Message, Row } from '@electric-sql/client' +import type { Collection } from '@tanstack/db' +import type { + PersistenceAdapter, + SQLiteDriver, +} from '../../db-sqlite-persistence-core/src' +import type { ElectricCollectionUtils } from '../src/electric' + +type Item = Row & { id: number; name: string } +type Subscriber = (messages: Array>) => void + +const subscribers: Array = [] +let synchronousMessages: Array> | undefined +const mockSubscribe = vi.fn((subscriber: Subscriber) => { + subscribers.push(subscriber) + if (synchronousMessages) subscriber(synchronousMessages) + return vi.fn() +}) + +vi.mock(`@electric-sql/client`, async () => ({ + ...(await vi.importActual(`@electric-sql/client`)), + ShapeStream: vi.fn(() => ({ + subscribe: mockSubscribe, + requestSnapshot: vi.fn().mockResolvedValue(undefined), + fetchSnapshot: vi.fn().mockResolvedValue({ metadata: {}, data: [] }), + forceDisconnectAndRefresh: vi.fn().mockResolvedValue(undefined), + isUpToDate: false, + shapeHandle: `shape-current`, + lastOffset: `20_0`, + })), +})) + +function toBinding(value: unknown): string | number | bigint | null { + if (value === null || value === undefined) return null + if (typeof value === `boolean`) return value ? 1 : 0 + if ( + typeof value === `string` || + typeof value === `number` || + typeof value === `bigint` + ) { + return value + } + return String(value) +} + +function createDriver(database: DatabaseSync): SQLiteDriver { + const driver: SQLiteDriver = { + exec: (sql) => { + database.exec(sql) + return Promise.resolve() + }, + query: (sql, params = []) => + Promise.resolve( + database + .prepare(sql) + .all(...params.map(toBinding)) + .map((row) => ({ ...row })) as Array, + ), + run: (sql, params = []) => { + database.prepare(sql).run(...params.map(toBinding)) + return Promise.resolve() + }, + transaction: async (transaction) => { + database.exec(`BEGIN IMMEDIATE`) + try { + const result = await transaction(driver) + database.exec(`COMMIT`) + return result + } catch (error) { + database.exec(`ROLLBACK`) + throw error + } + }, + } + return driver +} + +function deferred() { + let resolve!: () => void + const promise = new Promise((complete) => { + resolve = complete + }) + return { promise, resolve } +} + +async function reachCheckpoint( + promise: Promise, + checkpoint: string, +): Promise { + let timer: ReturnType | undefined + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error(`Did not reach checkpoint: ${checkpoint}`)), + 1_000, + ) + }), + ]) + } finally { + if (timer !== undefined) clearTimeout(timer) + } +} + +function change( + operation: `insert` | `update` | `delete`, + value: Item, +): Message { + return { key: String(value.id), value, headers: { operation } } +} + +async function runRace( + transition: `none` | `external-row-loss` | `schema-reset` | `committed-write`, + syncMode: `eager` | `on-demand` = `eager`, + legacyUnknown = false, + missingKeySetEvidence = false, + startupReset: `none` | `tag-state` | `shape-identity` = `none`, +): Promise { + const database = new DatabaseSync(`:memory:`) + const driver = createDriver(database) + const collectionId = + `resume-snapshot-${syncMode}-${transition}-` + + `${legacyUnknown}-${missingKeySetEvidence}-${startupReset}` + const laterSnapshotEntered = deferred() + const releaseLaterSnapshot = deferred() + let collection: + | Collection> + | undefined + let unsubscribe: (() => void) | undefined + let primaryFailure: unknown + const cleanupFailures: Array = [] + try { + const seedAdapter = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion: 1, + }) + await seedAdapter.applyCommittedTx(collectionId, { + txId: `seed`, + term: 1, + seq: 1, + rowVersion: 1, + mutations: [ + { type: `insert`, key: 1, value: { id: 1, name: `one` } }, + { type: `insert`, key: 2, value: { id: 2, name: `two` } }, + ], + collectionMetadataMutations: [ + { + type: `set`, + key: `electric:resume`, + value: { + kind: `resume`, + requiresTagState: startupReset === `tag-state`, + offset: `10_0`, + handle: `shape-old`, + shapeId: + startupReset === `shape-identity` + ? `{"params":{"table":"other_table"},"url":"http://test-url"}` + : `{"params":{"table":"test_table"},"url":"http://test-url"}`, + updatedAt: 1, + }, + }, + ], + }) + if (legacyUnknown) { + if (syncMode !== `on-demand`) { + const tableName = createPersistedTableName(collectionId, `c`) + await driver.run( + `DELETE FROM "${tableName}" WHERE json_extract(value, '$.id') = ?`, + [1], + ) + } + await driver.run( + `UPDATE collection_version SET key_set_evidence_available = 0 WHERE collection_id = ?`, + [collectionId], + ) + await driver.run( + `DELETE FROM collection_expected_keys WHERE collection_id = ?`, + [collectionId], + ) + expect( + (await seedAdapter.loadResumeSnapshot(collectionId)).keySet, + ).toEqual({ status: `unknown` }) + } + + const restartedAdapter = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion: 1, + }) + let snapshotCalls = 0 + let laterSnapshotIncludedRows: boolean | undefined + let resumeStateAtLaterSnapshot: unknown + let reportReceiverFailure!: (error: Error) => void + const receiverFailure = new Promise((resolve) => { + reportReceiverFailure = resolve + }) + const gatedAdapter = new Proxy(restartedAdapter, { + get(target, property) { + if (property === `loadResumeSnapshot`) { + return async function ( + this: PersistenceAdapter, + ...args: Parameters + ) { + if (this !== gatedAdapter) { + const error = new Error( + `Persistence adapter lost its receiver during resume certification`, + ) + reportReceiverFailure(error) + throw error + } + snapshotCalls++ + if (snapshotCalls > 1) { + laterSnapshotIncludedRows = args[1]?.includeRows + if (startupReset !== `none`) { + resumeStateAtLaterSnapshot = ( + await target.loadCollectionMetadata(collectionId) + ).find(({ key }) => key === `electric:resume`)?.value + } + laterSnapshotEntered.resolve() + await releaseLaterSnapshot.promise + } + const snapshot = await target.loadResumeSnapshot(...args) + return missingKeySetEvidence + ? { ...snapshot, keySet: undefined } + : snapshot + } + } + const value = Reflect.get(target, property, target) as unknown + return typeof value === `function` ? value.bind(target) : value + }, + }) as unknown as PersistenceAdapter + + collection = createCollection( + persistedCollectionOptions< + Item, + string | number, + never, + ElectricCollectionUtils + >({ + ...electricCollectionOptions({ + id: collectionId, + shapeOptions: { + url: `http://test-url`, + params: { table: `test_table` }, + }, + syncMode, + getKey: (row) => row.id, + startSync: false, + }), + persistence: { adapter: gatedAdapter }, + }), + ) + + let publications = 0 + if (startupReset !== `none`) { + synchronousMessages = [ + change(`insert`, { id: 1, name: `one` }), + change(`insert`, { id: 2, name: `two` }), + { headers: { control: `up-to-date` } }, + ] + } + collection.startSyncImmediate() + const subscription = collection.subscribeChanges( + () => { + publications++ + }, + { includeInitialState: false }, + ) + unsubscribe = () => subscription.unsubscribe() + await reachCheckpoint( + Promise.race([ + laterSnapshotEntered.promise, + receiverFailure.then((error) => { + throw error + }), + ]), + `${syncMode} resume reached its second atomic snapshot`, + ) + await vi.waitFor(() => expect(subscribers).toHaveLength(1)) + const subscriber = subscribers[0]! + const request = vi.mocked(ShapeStream).mock.calls[0]![0] as { + offset?: string + handle?: string + } + const replacesUncertifiedBaseline = + legacyUnknown || missingKeySetEvidence || startupReset !== `none` + if (startupReset !== `none`) { + expect(resumeStateAtLaterSnapshot).toMatchObject({ kind: `reset` }) + } + + if (startupReset === `none`) { + subscriber( + request.offset === undefined + ? [ + change(`insert`, { id: 1, name: `one` }), + change(`insert`, { id: 2, name: `two` }), + { headers: { control: `up-to-date` } }, + ] + : transition === `none` + ? [{ headers: { control: `up-to-date` } }] + : [ + change(`update`, { id: 2, name: `new-two` }), + { headers: { control: `up-to-date` } }, + ], + ) + } + if (transition === `external-row-loss`) { + const tableName = createPersistedTableName(collectionId, `c`) + await driver.run( + `DELETE FROM "${tableName}" WHERE json_extract(value, '$.id') = ?`, + [1], + ) + } else if (transition === `schema-reset`) { + const resettingAdapter = new SQLiteCorePersistenceAdapter({ + driver, + schemaVersion: 2, + }) + await resettingAdapter.loadSubset(collectionId, {}) + } else if (transition === `committed-write`) { + await seedAdapter.applyCommittedTx(collectionId, { + txId: `concurrent-writer`, + term: 1, + seq: 2, + rowVersion: 2, + mutations: [ + { type: `insert`, key: 3, value: { id: 3, name: `three` } }, + ], + }) + } + releaseLaterSnapshot.resolve() + + if (transition === `none`) { + await vi.waitFor(() => expect(collection!.status).toBe(`ready`)) + expect( + Array.from(collection.values(), ({ id, name }) => ({ id, name })), + ).toEqual( + replacesUncertifiedBaseline + ? [ + { id: 1, name: `one` }, + { id: 2, name: `two` }, + ] + : [], + ) + expect( + (await restartedAdapter.loadSubset(collectionId, {})).map( + ({ value }) => value, + ), + ).toEqual([ + { id: 1, name: `one` }, + { id: 2, name: `two` }, + ]) + } else { + await vi.waitFor(() => expect(collection!.status).toBe(`error`)) + await vi.waitFor(async () => { + const metadata = + await restartedAdapter.loadCollectionMetadata(collectionId) + const resumeState = metadata.find( + ({ key }) => key === `electric:resume`, + )?.value + if (transition === `schema-reset`) { + // The atomic SQLite reset already removed the stale cursor. A stale + // adapter must not write another marker after the schema changed. + expect(resumeState).toBeUndefined() + } else { + expect(resumeState).toMatchObject({ kind: `reset` }) + } + }) + expect(Array.from(collection.values())).toEqual([]) + expect(collection.status).not.toBe(`ready`) + expect(publications).toBe(0) + const durableRowsBeforeLateDelivery = await restartedAdapter.loadSubset( + collectionId, + {}, + ) + expect(durableRowsBeforeLateDelivery.map(({ value }) => value)).toEqual( + transition === `external-row-loss` + ? [{ id: 2, name: `two` }] + : transition === `committed-write` + ? [ + { id: 1, name: `one` }, + { id: 2, name: `two` }, + { id: 3, name: `three` }, + ] + : [], + ) + subscriber([ + change(`insert`, { id: 9, name: `late` }), + { headers: { control: `up-to-date` } }, + ]) + await new Promise((resolve) => setTimeout(resolve, 0)) + expect(Array.from(collection.values())).toEqual([]) + expect(await restartedAdapter.loadSubset(collectionId, {})).toEqual( + durableRowsBeforeLateDelivery, + ) + expect(publications).toBe(0) + } + if (replacesUncertifiedBaseline) { + expect(request.offset).toBeUndefined() + expect(request.handle).toBeUndefined() + } else { + expect(request).toMatchObject({ offset: `10_0`, handle: `shape-old` }) + } + expect(laterSnapshotIncludedRows).toBe( + replacesUncertifiedBaseline || syncMode !== `on-demand`, + ) + } catch (error) { + primaryFailure = error + } finally { + releaseLaterSnapshot.resolve() + try { + unsubscribe?.() + } catch (error) { + cleanupFailures.push(error) + } + try { + if (collection) { + await reachCheckpoint( + collection.cleanup(), + `Electric race collection cleanup`, + ) + } + } catch (error) { + cleanupFailures.push(error) + } + try { + database.close() + } catch (error) { + cleanupFailures.push(error) + } + } + + if (primaryFailure !== undefined) { + const failure = + primaryFailure instanceof Error + ? primaryFailure + : new Error(`Resume snapshot race failed`, { cause: primaryFailure }) + if (cleanupFailures.length > 0) { + Object.defineProperty(failure, `cleanupFailures`, { + value: cleanupFailures, + enumerable: true, + }) + } + throw failure + } + if (cleanupFailures.length > 0) { + throw new AggregateError(cleanupFailures, `Resume snapshot cleanup failed`) + } +} + +type LegacyUnknownResumeObservation = { + checkpoint: `post-restart-up-to-date` + migratedKeySet: { status: `unknown` | `consistent` | `incompatible` } + migratedRows: Array + requestedOffset: string | undefined + requestedHandle: string | undefined + sourceDelivery: Array + publicRows: Array + durableRows: Array + status: string +} + +async function observeLegacyUnknownResume(): Promise { + const database = new DatabaseSync(`:memory:`) + const driver = createDriver(database) + const collectionId = `legacy-loss-fixed-witness` + const rowLostBeforeMigration: Item = { + id: 1, + name: `lost-before-ledger`, + } + const survivingRow: Item = { id: 2, name: `survivor` } + const postOffsetRow: Item = { id: 3, name: `post-offset` } + const canonicalSourceSnapshot = [ + rowLostBeforeMigration, + survivingRow, + postOffsetRow, + ] + let collection: + | Collection> + | undefined + let unsubscribe: (() => void) | undefined + let observation: LegacyUnknownResumeObservation | undefined + let primaryFailure: unknown + const cleanupFailures: Array = [] + + try { + // Produce the persisted row/metadata encodings through the real adapter, + // then reduce only the key-evidence schema to its pre-ledger form. + const legacyAdapter = new SQLiteCorePersistenceAdapter({ driver }) + await legacyAdapter.applyCommittedTx(collectionId, { + txId: `legacy-snapshot-at-10`, + term: 1, + seq: 1, + rowVersion: 1, + mutations: [rowLostBeforeMigration, survivingRow].map((row) => ({ + type: `insert` as const, + key: row.id, + value: structuredClone(row), + })), + collectionMetadataMutations: [ + { + type: `set`, + key: `electric:resume`, + value: { + kind: `resume`, + requiresTagState: false, + offset: `10_0`, + handle: `shape-old`, + shapeId: `{"params":{"table":"test_table"},"url":"http://test-url"}`, + updatedAt: 1, + }, + }, + ], + }) + + const collectionTable = createPersistedTableName(collectionId, `c`) + await driver.run( + `DELETE FROM "${collectionTable}" WHERE json_extract(value, '$.id') = ?`, + [rowLostBeforeMigration.id], + ) + await driver.exec( + `DROP TRIGGER "${collectionTable}_key_evidence_insert"; + DROP TRIGGER "${collectionTable}_key_evidence_delete"; + DROP TRIGGER "${collectionTable}_key_evidence_update"`, + ) + await driver.exec(`DROP TABLE collection_expected_keys`) + await driver.exec( + `ALTER TABLE collection_version RENAME TO collection_version_with_ledger`, + ) + await driver.exec( + `CREATE TABLE collection_version ( + collection_id TEXT PRIMARY KEY, + latest_row_version INTEGER NOT NULL + )`, + ) + await driver.exec( + `INSERT INTO collection_version (collection_id, latest_row_version) + SELECT collection_id, latest_row_version + FROM collection_version_with_ledger`, + ) + await driver.exec(`DROP TABLE collection_version_with_ledger`) + + const migratedAdapter = new SQLiteCorePersistenceAdapter({ driver }) + const migratedSnapshot = + await migratedAdapter.loadResumeSnapshot(collectionId) + + collection = createCollection( + persistedCollectionOptions< + Item, + string | number, + never, + ElectricCollectionUtils + >({ + ...electricCollectionOptions({ + id: collectionId, + shapeOptions: { + url: `http://test-url`, + params: { table: `test_table` }, + }, + syncMode: `eager`, + getKey: (row) => row.id, + startSync: false, + }), + persistence: { adapter: migratedAdapter }, + }), + ) + collection.startSyncImmediate() + const subscription = collection.subscribeChanges(() => {}, { + includeInitialState: false, + }) + unsubscribe = () => subscription.unsubscribe() + + await vi.waitFor(() => expect(subscribers).toHaveLength(1)) + const request = vi.mocked(ShapeStream).mock.calls[0]![0] as { + offset?: string + handle?: string + } + // The source obeys the request: a resume receives only changes after its + // cursor; a fresh request receives the independently specified snapshot. + const sourceDelivery = + request.offset === undefined + ? canonicalSourceSnapshot.map((row) => structuredClone(row)) + : [structuredClone(postOffsetRow)] + subscribers[0]!([ + ...sourceDelivery.map((row) => change(`insert`, structuredClone(row))), + { headers: { control: `up-to-date` } }, + ]) + + await vi.waitFor(() => expect(collection!.status).toBe(`ready`)) + await new Promise((resolve) => setTimeout(resolve, 0)) + await new Promise((resolve) => setTimeout(resolve, 0)) + observation = { + checkpoint: `post-restart-up-to-date`, + migratedKeySet: migratedSnapshot.keySet, + migratedRows: migratedSnapshot.rows + .map(({ value }) => value as Item) + .sort((left, right) => left.id - right.id), + requestedOffset: request.offset, + requestedHandle: request.handle, + sourceDelivery, + publicRows: Array.from(collection.values(), ({ id, name }) => ({ + id, + name, + })).sort((left, right) => left.id - right.id), + durableRows: (await migratedAdapter.loadSubset(collectionId, {})) + .map(({ value }) => value as Item) + .sort((left, right) => left.id - right.id), + status: collection.status, + } + } catch (error) { + primaryFailure = error + } finally { + try { + unsubscribe?.() + } catch (error) { + cleanupFailures.push(error) + } + try { + await collection?.cleanup() + } catch (error) { + cleanupFailures.push(error) + } + try { + database.close() + } catch (error) { + cleanupFailures.push(error) + } + } + + if (primaryFailure !== undefined) { + const failure = + primaryFailure instanceof Error + ? primaryFailure + : new Error(`Legacy resume witness failed`, { cause: primaryFailure }) + if (cleanupFailures.length > 0) { + Object.defineProperty(failure, `cleanupFailures`, { + value: cleanupFailures, + enumerable: true, + }) + } + throw failure + } + if (cleanupFailures.length > 0) { + throw new AggregateError( + cleanupFailures, + `Legacy resume witness cleanup failed`, + ) + } + if (!observation) { + throw new Error(`Legacy resume witness did not capture an observation`) + } + return observation +} + +/** + * # Which persisted baseline may Electric resume from during startup races? + * + * The persisted resume law requires the rows, resume metadata, stream position, + * and key-set evidence used for certification to belong to one atomic baseline + * generation. An unverifiable, externally changed, or reset baseline must start + * a fresh source snapshot; a compatible baseline may retain its resume cursor. + * This refines the settled recovery law in electric-recovery-oracle.test.ts and + * the atomic `loadResumeSnapshot` persistence contract. + * + * The reference is the small baseline tuple captured by each case: generation, + * complete key set, resume state, and expected source delivery. Legal histories + * vary eager versus on-demand sync, compatible versus unknown legacy evidence, + * startup reset cause, and row loss, schema reset, or committed write between + * the initial metadata read and certification. No production classifier or SQL + * helper computes the expected public and durable rows. + * + * The production driver uses `SQLiteCorePersistenceAdapter`, the persisted + * Collection wrapper, and `electricCollectionOptions`. It holds the adapter's + * later atomic snapshot, injects the selected transition, then compares the + * ShapeStream offset/handle plus complete public and durable rows at the + * post-restart up-to-date checkpoint. Entering the held snapshot is the reach + * witness; compatible and legacy-unknown controls challenge both resume and + * fresh-snapshot branches. + * + * These deterministic schedules do not model arbitrary external SQL edits, + * native SQLite hosts, or a live Electric service. Those require their separate + * persistence-driver and real-provider owners. + */ +describe(`Electric resume snapshot races`, () => { + beforeEach(() => { + subscribers.length = 0 + synchronousMessages = undefined + vi.clearAllMocks() + }) + + it(`keeps a healthy tagged cache when a fresh reset commits before hydration`, async () => { + await runRace(`none`, `eager`, false, false, `tag-state`) + }) + + it(`keeps a healthy cache when a changed shape commits its reset before hydration`, async () => { + await runRace(`none`, `eager`, false, false, `shape-identity`) + }) + + it(`rejects row loss between resume metadata and baseline hydration`, async () => { + await runRace(`external-row-loss`) + }) + + it(`rejects a schema reset between resume metadata and baseline hydration`, async () => { + await runRace(`schema-reset`) + }) + + it(`conservatively rejects a committed write between startup snapshots`, async () => { + await runRace(`committed-write`) + }) + + it(`freshly replaces an unknown on-demand resume baseline`, async () => { + await runRace(`none`, `on-demand`, true) + }) + + it(`freshly replaces an unknown eager resume baseline`, async () => { + await runRace(`none`, `eager`, true) + }) + + it(`freshly replaces a resume when snapshot evidence is missing`, async () => { + await runRace(`none`, `eager`, false, true) + }) + + it(`rejects row loss during on-demand resume certification`, async () => { + await runRace(`external-row-loss`, `on-demand`) + }) + + it(`rejects a schema reset for an unknown on-demand resume`, async () => { + await runRace(`schema-reset`, `on-demand`, true) + }) + + it(`freshly replaces an unverifiable pre-ledger resume baseline`, async () => { + const observation = await observeLegacyUnknownResume() + expect(observation).toEqual({ + checkpoint: `post-restart-up-to-date`, + migratedKeySet: { status: `unknown` }, + migratedRows: [{ id: 2, name: `survivor` }], + requestedOffset: undefined, + requestedHandle: undefined, + sourceDelivery: [ + { id: 1, name: `lost-before-ledger` }, + { id: 2, name: `survivor` }, + { id: 3, name: `post-offset` }, + ], + publicRows: [ + { id: 1, name: `lost-before-ledger` }, + { id: 2, name: `survivor` }, + { id: 3, name: `post-offset` }, + ], + durableRows: [ + { id: 1, name: `lost-before-ledger` }, + { id: 2, name: `survivor` }, + { id: 3, name: `post-offset` }, + ], + status: `ready`, + }) + }) +})