diff --git a/.changeset/fix-shared-sqlite-hydration-fairness.md b/.changeset/fix-shared-sqlite-hydration-fairness.md
new file mode 100644
index 0000000000..cd5b79d342
--- /dev/null
+++ b/.changeset/fix-shared-sqlite-hydration-fairness.md
@@ -0,0 +1,7 @@
+---
+'@tanstack/db-sqlite-persistence-core': patch
+'@tanstack/browser-db-sqlite-persistence': patch
+'@tanstack/electron-db-sqlite-persistence': patch
+---
+
+Schedule complete SQLite hydration units fairly without holding coordinator work inside the local hydration scope. Preserve per-Collection leader adapter routing, mutation results across transport retries, terminal coordinator disposal, real-adapter restart order, and promise-discovered shared scheduling.
diff --git a/.github/workflows/e2e-tests.yml b/.github/workflows/e2e-tests.yml
index a340502b46..6b3756ff74 100644
--- a/.github/workflows/e2e-tests.yml
+++ b/.github/workflows/e2e-tests.yml
@@ -74,6 +74,11 @@ jobs:
cd examples/react/start-ssr-e2e
pnpm exec playwright install --with-deps chromium
+ - name: Run Browser SQLite OPFS fairness E2E tests
+ run: |
+ cd packages/browser-db-sqlite-persistence
+ pnpm test:opfs-fairness
+
- name: Run React Start SSR E2E tests
run: |
cd examples/react/start-ssr-e2e
diff --git a/docs/contributing/oracle-coverage.md b/docs/contributing/oracle-coverage.md
index 06eb175c2a..039f032032 100644
--- a/docs/contributing/oracle-coverage.md
+++ b/docs/contributing/oracle-coverage.md
@@ -110,7 +110,7 @@ comment and the current API/architecture contract before extending its model.
| Opaque backend pagination | [window oracle](https://github.com/TanStack/db/blob/main/packages/query-db-collection/tests/cursor-pagination.oracle.test.ts), [cache histories](https://github.com/TanStack/db/blob/main/packages/query-db-collection/tests/cursor-pagination.cache-oracle.test.ts), [cache publication](https://github.com/TanStack/db/blob/main/packages/query-db-collection/tests/cursor-pagination.publication-oracle.test.ts), [browser acquisition boundaries](https://github.com/TanStack/db/blob/main/packages/query-db-collection/tests/cursor-pagination.boundary-oracle.test.ts), [QueryCollection integration](https://github.com/TanStack/db/blob/main/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](https://github.com/TanStack/db/blob/main/packages/electric-db-collection/tests/electric-oracle.property.test.ts), [PostgreSQL semantics](https://github.com/TanStack/db/blob/main/packages/electric-db-collection/e2e/sql-predicate-semantics.e2e.test.ts), [TrailBase contract](https://github.com/TanStack/db/blob/main/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. |
| PowerSync | [tests](https://github.com/TanStack/db/tree/main/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](https://github.com/TanStack/db/blob/main/packages/db-sqlite-persistence-core/tests/persisted.test.ts), [driver contracts](https://github.com/TanStack/db/blob/main/packages/db-sqlite-persistence-core/tests/contracts/sqlite-driver-contract.ts), [browser OPFS lifecycle](https://github.com/TanStack/db/blob/main/packages/browser-db-sqlite-persistence/tests/opfs-page-lifecycle-oracle.test.ts), [worker diagnostics](https://github.com/TanStack/db/blob/main/packages/browser-db-sqlite-persistence/tests/opfs-worker-diagnostics-oracle.test.ts), [113-law manifest](https://github.com/TanStack/db/blob/main/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](https://github.com/TanStack/db/blob/main/packages/db-sqlite-persistence-core/tests/persisted.test.ts), [driver contracts](https://github.com/TanStack/db/blob/main/packages/db-sqlite-persistence-core/tests/contracts/sqlite-driver-contract.ts), [shared-driver fairness](https://github.com/TanStack/db/blob/main/packages/browser-db-sqlite-persistence/tests/shared-driver-fairness-oracle.test.ts), [browser OPFS lifecycle](https://github.com/TanStack/db/blob/main/packages/browser-db-sqlite-persistence/tests/opfs-page-lifecycle-oracle.test.ts), [worker diagnostics](https://github.com/TanStack/db/blob/main/packages/browser-db-sqlite-persistence/tests/opfs-worker-diagnostics-oracle.test.ts), [113-law manifest](https://github.com/TanStack/db/blob/main/packages/db-collection-e2e/src/fixtures/persisted-conformance-manifest.ts) | Cache/remote rejection/peer/reopen histories, exact driver results, K=1 complete-logical-hydrate fairness over one shared Browser SQLite driver, controlled page/worker ownership, and diagnostic-cause retention. The fairness owner runs one identical generated property in retained fixed-seed and seedless-random lanes, or only a requested checked seed+shrink-path replay. Its grammar controls reconstruct the known witness, exercise bounded marginals, reject an invalid storm boundary, and ablate tail order; a named persist-first FIFO wrong answer calibrates the checker. It does not establish elapsed-time latency, unbounded eventuality, multi-process coordination, or a browser matrix; Chromium OPFS separately refines the provider boundary. 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](https://github.com/TanStack/db/blob/main/packages/offline-transactions/tests/KeyScheduler.property.test.ts), [leadership](https://github.com/TanStack/db/blob/main/packages/offline-transactions/tests/leadership-replay.property.test.ts), [settlement](https://github.com/TanStack/db/blob/main/packages/offline-transactions/tests/transaction-settlement.property.test.ts), [serialization](https://github.com/TanStack/db/blob/main/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](https://github.com/TanStack/db/blob/main/packages/react-db/tests/conformance.test.tsx), [React pagination](https://github.com/TanStack/db/blob/main/packages/react-db/tests/infinite-query-conformance.test.tsx), [shared suites](https://github.com/TanStack/db/tree/main/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. |
| Structural values and ordered primitives | [hash values](https://github.com/TanStack/db/blob/main/packages/db-ivm/tests/hash.property.test.ts), [hash graphs](https://github.com/TanStack/db/blob/main/packages/db-ivm/tests/hash-graph.property.test.ts), [mixed hash graphs](https://github.com/TanStack/db/blob/main/packages/db-ivm/tests/hash-mixed-graph.property.test.ts), [hash retry](https://github.com/TanStack/db/blob/main/packages/db-ivm/tests/hash-failure-retry.property.test.ts), [comparison](https://github.com/TanStack/db/blob/main/packages/db/tests/comparison.property.test.ts), [deep equality](https://github.com/TanStack/db/blob/main/packages/db/tests/utils.property.test.ts), [cursor](https://github.com/TanStack/db/blob/main/packages/db/tests/cursor.property.test.ts), [indexes](https://github.com/TanStack/db/blob/main/packages/db/tests/index-update.property.test.ts), [query identity](https://github.com/TanStack/db/blob/main/packages/db/tests/query/identity-output-shape-oracle.test.ts) | Independent flat values, graph topology, algebraic laws, Map/group/sort recomputation, expression denotation, and compiled output bags. Hash collision freedom is not promised. Unsupported composite cursors reject. |
diff --git a/packages/browser-db-sqlite-persistence/e2e/shared-driver-fairness.opfs.html b/packages/browser-db-sqlite-persistence/e2e/shared-driver-fairness.opfs.html
new file mode 100644
index 0000000000..215878cfd3
--- /dev/null
+++ b/packages/browser-db-sqlite-persistence/e2e/shared-driver-fairness.opfs.html
@@ -0,0 +1,13 @@
+
+
+
+
+
+
+ Shared driver OPFS fairness oracle
+
+
+
+
+
+
diff --git a/packages/browser-db-sqlite-persistence/e2e/shared-driver-fairness.opfs.spec.ts b/packages/browser-db-sqlite-persistence/e2e/shared-driver-fairness.opfs.spec.ts
new file mode 100644
index 0000000000..849de99ccc
--- /dev/null
+++ b/packages/browser-db-sqlite-persistence/e2e/shared-driver-fairness.opfs.spec.ts
@@ -0,0 +1,102 @@
+/**
+ * Browser checkpoint assertions for the OPFS refinement. Expected public rows
+ * are rebuilt from the scenario IDs rather than from production output. The
+ * neutral case proves cold-query reach; the storm case requires no K=1
+ * violation. Failure-before-checkpoint, semantic mismatch, driver cleanup, and
+ * OPFS cleanup remain distinct outcomes so setup or teardown cannot satisfy the
+ * scheduling law.
+ */
+import { expect, test } from '@playwright/test'
+import type { Page } from '@playwright/test'
+import type { OPFSOracleResult } from './shared-driver-fairness.opfs'
+
+async function readOracleResult(
+ page: Page,
+ mode: `neutral` | `storm`,
+): Promise {
+ let rejectPageError!: (error: Error) => void
+ const pageError = new Promise((_resolve, reject) => {
+ rejectPageError = reject
+ })
+ void pageError.catch(() => undefined)
+ const onPageError = (error: Error) => {
+ rejectPageError(
+ new Error(
+ `OPFS fairness page failed before publishing a result: ${error.message}`,
+ ),
+ )
+ }
+ page.on(`pageerror`, onPageError)
+
+ try {
+ await page.goto(`/e2e/shared-driver-fairness.opfs.html?mode=${mode}`)
+ await Promise.race([
+ page.waitForFunction(
+ () => window.__tanstackDriverFairnessOracle !== undefined,
+ ),
+ pageError,
+ ])
+ return page.evaluate(() => window.__tanstackDriverFairnessOracle!)
+ } finally {
+ page.off(`pageerror`, onPageError)
+ }
+}
+
+function expectedHydratedCollections(scenarioId: string, count: number) {
+ return Array.from({ length: count }, (_, index) => ({
+ collectionId: `${scenarioId}-hydrate-${index}`,
+ rows: [
+ { id: `row-${index}-0`, value: index * 10 },
+ { id: `row-${index}-1`, value: index * 10 + 1 },
+ ],
+ }))
+}
+
+test(`real Chromium OPFS fixture reaches and cleans up cold hydration`, async ({
+ page,
+}) => {
+ const result = await readOracleResult(page, `neutral`)
+
+ if (result.status !== `complete`) throw new Error(result.primaryFailure)
+ expect(result.provider).toBe(`Chromium OPFSCoopSyncVFS worker`)
+ expect(result.observation.admittedHydrateIds).toHaveLength(2)
+ // These are actual public Collection rows captured after preload, compared
+ // with seed values built independently by this browser assertion.
+ expect(result.observation.hydratedCollections).toEqual(
+ expectedHydratedCollections(`opfs-neutral-reach`, 2),
+ )
+ expect(
+ result.observation.rawDequeues.some((entry) =>
+ entry.sql.startsWith(`SELECT key, value, metadata, row_version FROM`),
+ ),
+ ).toBe(true)
+ expect(result.observation.cleanupFailures).toEqual([])
+ expect(result.opfsCleanupFailures).toEqual([])
+})
+
+test(`real Chromium OPFS fixture bounds pending cold hydration behind persists`, async ({
+ page,
+}) => {
+ const result = await readOracleResult(page, `storm`)
+
+ if (result.status !== `complete`) throw new Error(result.primaryFailure)
+ expect(result.provider).toBe(`Chromium OPFSCoopSyncVFS worker`)
+ expect(result.observation.admittedHydrateIds).toHaveLength(4)
+ expect(result.observation.hydratedCollections).toEqual(
+ expectedHydratedCollections(`opfs-fixed-persist-storm`, 4),
+ )
+ // This is the semantic RED checkpoint. Setup, wall time, and cleanup are
+ // reported independently and cannot satisfy this assertion.
+ if (result.violation !== undefined) {
+ throw new Error(
+ `real OPFS fairness mismatch: ${JSON.stringify(result.violation)}; ` +
+ `logical completion order: ${JSON.stringify(result.observation.logicalCompletionOrder)}; ` +
+ `driver admissions: ${result.observation.driverAdmissions.length}; ` +
+ `raw dequeues: ${result.observation.rawDequeues.length}; ` +
+ `driver cleanup diagnostics: ${JSON.stringify(result.observation.cleanupFailures)}; ` +
+ `OPFS cleanup diagnostics: ${JSON.stringify(result.opfsCleanupFailures)}`,
+ )
+ }
+ expect(result.observation.cleanupFailures).toEqual([])
+ expect(result.opfsCleanupFailures).toEqual([])
+})
diff --git a/packages/browser-db-sqlite-persistence/e2e/shared-driver-fairness.opfs.ts b/packages/browser-db-sqlite-persistence/e2e/shared-driver-fairness.opfs.ts
new file mode 100644
index 0000000000..a443f680f2
--- /dev/null
+++ b/packages/browser-db-sqlite-persistence/e2e/shared-driver-fairness.opfs.ts
@@ -0,0 +1,160 @@
+/**
+ * Real-provider refinement of the shared-driver fairness oracle. The page runs
+ * the same legal neutral and fixed-storm histories through public `preload()`,
+ * the browser/core adapter boundary, and Chromium's OPFSCoopSyncVFS worker. It
+ * freezes logical completions, public rows, raw dequeue reach, and the K=1
+ * violation result before cleanup, then reports provider cleanup separately.
+ * This fixture adds real OPFS/worker evidence; it does not claim multi-tab,
+ * multi-process, non-Chromium, latency, or unbounded-eventuality coverage.
+ */
+import { openBrowserWASQLiteOPFSDatabase } from '../src/index'
+import {
+ findSharedDriverFairnessViolation,
+ observeSharedDriverFairness,
+} from '../tests/shared-driver-fairness-oracle'
+import type {
+ SharedDriverFairnessObservation,
+ SharedDriverFairnessScenario,
+ SharedDriverFairnessViolation,
+} from '../tests/shared-driver-fairness-oracle'
+
+export type OPFSOracleResult =
+ | {
+ status: `complete`
+ provider: `Chromium OPFSCoopSyncVFS worker`
+ observation: SharedDriverFairnessObservation
+ violation: SharedDriverFairnessViolation | undefined
+ opfsCleanupFailures: ReadonlyArray
+ }
+ | {
+ status: `failed-before-checkpoint`
+ provider: `Chromium OPFSCoopSyncVFS worker`
+ primaryFailure: string
+ opfsCleanupFailures: ReadonlyArray
+ }
+
+declare global {
+ interface Window {
+ __tanstackDriverFairnessOracle?: OPFSOracleResult
+ }
+}
+
+async function removeOPFSArtifacts(
+ databaseName: string,
+): Promise> {
+ const failures: Array = []
+ const root = await navigator.storage.getDirectory()
+ for (const suffix of [``, `-journal`, `-wal`]) {
+ try {
+ await root.removeEntry(`${databaseName}${suffix}`)
+ } catch (error) {
+ if (!(error instanceof DOMException && error.name === `NotFoundError`)) {
+ failures.push(
+ `${databaseName}${suffix}: ${error instanceof Error ? error.message : String(error)}`,
+ )
+ }
+ }
+ }
+
+ // OPFSCoopSyncVFS creates one private temporary access-handle directory per
+ // worker. This page owns its isolated origin for the fixture run.
+ const iterableRoot = root as FileSystemDirectoryHandle & {
+ entries: () => AsyncIterableIterator<[string, FileSystemHandle]>
+ }
+ for await (const [name, handle] of iterableRoot.entries()) {
+ if (handle.kind !== `directory` || !name.startsWith(`.ahp-`)) continue
+ try {
+ await root.removeEntry(name, { recursive: true })
+ } catch (error) {
+ failures.push(
+ `${name}: ${error instanceof Error ? error.message : String(error)}`,
+ )
+ }
+ }
+ return failures
+}
+
+function scenarioFromLocation(): SharedDriverFairnessScenario {
+ const mode = new URL(location.href).searchParams.get(`mode`)
+ if (mode === `neutral`) {
+ return {
+ id: `opfs-neutral-reach`,
+ work: [0, 1].map((index) => ({
+ kind: `hydrate` as const,
+ id: `hydrate-${index}`,
+ seededRows: [
+ { id: `row-${index}-0`, value: index * 10 },
+ { id: `row-${index}-1`, value: index * 10 + 1 },
+ ],
+ })),
+ }
+ }
+ return {
+ id: `opfs-fixed-persist-storm`,
+ work: [
+ ...Array.from({ length: 5 }, (_, index) => ({
+ kind: `persist` as const,
+ id: `persist-${index}`,
+ mutationsPerPersist: 2,
+ })),
+ ...Array.from({ length: 4 }, (_, index) => ({
+ kind: `hydrate` as const,
+ id: `hydrate-${index}`,
+ seededRows: [
+ { id: `row-${index}-0`, value: index * 10 },
+ { id: `row-${index}-1`, value: index * 10 + 1 },
+ ],
+ })),
+ ],
+ }
+}
+
+async function run(): Promise {
+ const status = document.querySelector(`#oracle-status`)
+ const databaseName = `ws5b-${crypto.randomUUID()}.sqlite`
+ let observation: SharedDriverFairnessObservation | undefined
+ let violation: SharedDriverFairnessViolation | undefined
+ let primaryFailure: string | undefined
+ let opfsCleanupFailures: ReadonlyArray = []
+ try {
+ observation = await observeSharedDriverFairness(
+ () => openBrowserWASQLiteOPFSDatabase({ databaseName }),
+ scenarioFromLocation(),
+ )
+ // Freeze the reached semantic checkpoint before attempting OPFS cleanup.
+ violation = findSharedDriverFairnessViolation(observation)
+ } catch (error) {
+ primaryFailure = error instanceof Error ? error.message : String(error)
+ }
+
+ try {
+ opfsCleanupFailures = await removeOPFSArtifacts(databaseName)
+ } catch (cleanupError) {
+ opfsCleanupFailures = [
+ cleanupError instanceof Error
+ ? cleanupError.message
+ : String(cleanupError),
+ ]
+ }
+
+ if (observation) {
+ window.__tanstackDriverFairnessOracle = {
+ status: `complete`,
+ provider: `Chromium OPFSCoopSyncVFS worker`,
+ observation,
+ violation,
+ opfsCleanupFailures,
+ }
+ if (status) status.value = `complete`
+ } else {
+ window.__tanstackDriverFairnessOracle = {
+ status: `failed-before-checkpoint`,
+ provider: `Chromium OPFSCoopSyncVFS worker`,
+ primaryFailure: primaryFailure ?? `unknown failure before checkpoint`,
+ opfsCleanupFailures,
+ }
+ if (status) status.value = `failed-before-checkpoint`
+ }
+}
+
+void run()
diff --git a/packages/browser-db-sqlite-persistence/package.json b/packages/browser-db-sqlite-persistence/package.json
index 9bae7bb83c..26b2b70fef 100644
--- a/packages/browser-db-sqlite-persistence/package.json
+++ b/packages/browser-db-sqlite-persistence/package.json
@@ -23,6 +23,8 @@
"dev": "vite build --watch",
"lint": "eslint . --fix",
"test": "vitest --run",
+ "test:oracles": "vitest --run tests/shared-driver-fairness-oracle.test.ts",
+ "test:opfs-fairness": "playwright test --config playwright.opfs.config.ts",
"test:e2e": "pnpm --filter @tanstack/db-ivm build && pnpm --filter @tanstack/db build && pnpm --filter @tanstack/db-sqlite-persistence-core build && pnpm --filter @tanstack/browser-db-sqlite-persistence build && vitest --config vitest.e2e.config.ts --run"
},
"type": "module",
@@ -56,6 +58,7 @@
},
"devDependencies": {
"@journeyapps/wa-sqlite": "^1.4.1",
+ "@playwright/test": "^1.60.0",
"@types/better-sqlite3": "^7.6.13",
"@vitest/coverage-istanbul": "^3.2.4",
"better-sqlite3": "^12.6.2"
diff --git a/packages/browser-db-sqlite-persistence/playwright.opfs.config.ts b/packages/browser-db-sqlite-persistence/playwright.opfs.config.ts
new file mode 100644
index 0000000000..ecfa8af0a6
--- /dev/null
+++ b/packages/browser-db-sqlite-persistence/playwright.opfs.config.ts
@@ -0,0 +1,25 @@
+import { defineConfig } from '@playwright/test'
+
+const baseURL = `http://127.0.0.1:4185`
+const browserChannel =
+ process.env.PLAYWRIGHT_CHANNEL ?? (process.env.CI ? undefined : `chrome`)
+
+export default defineConfig({
+ testDir: `./e2e`,
+ testMatch: `shared-driver-fairness.opfs.spec.ts`,
+ timeout: 60_000,
+ fullyParallel: false,
+ workers: 1,
+ use: {
+ baseURL,
+ ...(browserChannel ? { channel: browserChannel } : {}),
+ headless: true,
+ trace: `retain-on-failure`,
+ },
+ webServer: {
+ command: `vite --config vite.opfs.config.ts --host 127.0.0.1 --port 4185`,
+ reuseExistingServer: false,
+ timeout: 120_000,
+ url: `${baseURL}/e2e/shared-driver-fairness.opfs.html`,
+ },
+})
diff --git a/packages/browser-db-sqlite-persistence/src/browser-coordinator.ts b/packages/browser-db-sqlite-persistence/src/browser-coordinator.ts
index 1babddc5a7..edf19b77af 100644
--- a/packages/browser-db-sqlite-persistence/src/browser-coordinator.ts
+++ b/packages/browser-db-sqlite-persistence/src/browser-coordinator.ts
@@ -1,6 +1,7 @@
import { safeRandomUUID } from '@tanstack/db-sqlite-persistence-core'
import type {
ApplyLocalMutationsResponse,
+ HydrationPersistenceAdapter,
PersistedCollectionCoordinator,
PersistedIndexSpec,
PersistedMutationEnvelope,
@@ -20,6 +21,7 @@ const RPC_RETRY_ATTEMPTS = 2
const RPC_RETRY_DELAY_MS = 200
const WRITER_LOCK_BUSY_RETRY_MS = 50
const WRITER_LOCK_MAX_RETRIES = 20
+const APPLIED_ENVELOPE_RETENTION_MS = 60_000
// ---------------------------------------------------------------------------
// Internal types
@@ -81,6 +83,12 @@ type CollectionState = {
subscribers: Set<(message: ProtocolEnvelope) => void>
}
+type AppliedEnvelope = {
+ collectionId: string
+ appliedAt: number
+ response: ApplyLocalMutationsResponse
+}
+
// Adapter with pullSince support
type AdapterWithPullSince = PersistenceAdapter & {
pullSince?: (
@@ -122,10 +130,18 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
private readonly nodeId = safeRandomUUID()
private readonly dbName: string
private adapter: AdapterWithPullSince | null
+ private readonly collectionAdapters = new Map()
private readonly channel: BroadcastChannel
private readonly collections = new Map()
private readonly pendingRPCs = new Map()
- private readonly appliedEnvelopeIds = new Map()
+ private readonly appliedEnvelopeIds = new Map()
+ private readonly applyingEnvelopeIds = new Map<
+ string,
+ Promise
+ >()
+ private appliedEnvelopePruneTimer: ReturnType | null = null
+ private readonly disposedPromise: Promise
+ private rejectDisposed: ((error: Error) => void) | null = null
private disposed = false
/** Method indirection to prevent TypeScript from narrowing `disposed` across awaits */
@@ -133,18 +149,27 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
return this.disposed
}
- private requireAdapter(): AdapterWithPullSince {
- if (!this.adapter) {
+ private requireAdapter(collectionId: string): AdapterWithPullSince {
+ let adapter = this.collectionAdapters.get(collectionId)
+ if (!adapter && this.adapter) {
+ adapter = this.adapter
+ this.collectionAdapters.set(collectionId, adapter)
+ }
+ if (!adapter) {
throw new Error(
`BrowserCollectionCoordinator: adapter not set. Call setAdapter() before using leader-side operations.`,
)
}
- return this.adapter
+ return adapter
}
constructor(options: BrowserCollectionCoordinatorOptions) {
this.dbName = options.dbName
this.adapter = options.adapter ?? null
+ this.disposedPromise = new Promise((_resolve, reject) => {
+ this.rejectDisposed = reject
+ })
+ void this.disposedPromise.catch(() => undefined)
this.channel = new BroadcastChannel(`tsdb:coord:${this.dbName}`)
this.channel.onmessage = (event: MessageEvent) => {
this.onChannelMessage(event.data)
@@ -160,6 +185,14 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
this.adapter = adapter
}
+ /** Register the adapter that owns leader-side work for one collection. */
+ setAdapterForCollection(
+ collectionId: string,
+ adapter: AdapterWithPullSince,
+ ): void {
+ this.collectionAdapters.set(collectionId, adapter)
+ }
+
// -----------------------------------------------------------------------
// PersistedCollectionCoordinator interface
// -----------------------------------------------------------------------
@@ -176,6 +209,7 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
state.subscribers.add(onMessage)
return () => {
state.subscribers.delete(onMessage)
+ this.releaseCollectionIfUnused(collectionId, state)
}
}
@@ -221,9 +255,16 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
collectionId: string,
signature: string,
spec: PersistedIndexSpec,
+ scopedAdapter?: HydrationPersistenceAdapter,
+ localEnsureCompleted = false,
): Promise {
if (this.isLeader(collectionId)) {
- await this.requireAdapter().ensureIndex(collectionId, signature, spec)
+ if (localEnsureCompleted) return
+ await (scopedAdapter ?? this.requireAdapter(collectionId)).ensureIndex(
+ collectionId,
+ signature,
+ spec,
+ )
return
}
@@ -270,13 +311,19 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
async pullSince(
collectionId: string,
fromRowVersion: number,
+ scopedAdapter?: HydrationPersistenceAdapter,
): Promise {
if (this.isLeader(collectionId)) {
- return this.handlePullSince(collectionId, {
- type: `rpc:pullSince:req`,
- rpcId: safeRandomUUID(),
- fromRowVersion,
- })
+ // A scoped adapter is a leader-local capability and never crosses RPC.
+ return this.handlePullSince(
+ collectionId,
+ {
+ type: `rpc:pullSince:req`,
+ rpcId: safeRandomUUID(),
+ fromRowVersion,
+ },
+ scopedAdapter,
+ )
}
return this.sendRPC(collectionId, {
@@ -291,7 +338,11 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
// -----------------------------------------------------------------------
dispose(): void {
+ if (this.disposed) return
this.disposed = true
+ const disposedError = new Error(`coordinator disposed`)
+ this.rejectDisposed?.(disposedError)
+ this.rejectDisposed = null
for (const [collectionId, state] of this.collections) {
this.releaseLeadership(collectionId, state)
@@ -299,12 +350,20 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
for (const [, pending] of this.pendingRPCs) {
clearTimeout(pending.timer)
- pending.reject(new Error(`coordinator disposed`))
+ pending.reject(disposedError)
}
this.pendingRPCs.clear()
+ if (this.appliedEnvelopePruneTimer !== null) {
+ clearTimeout(this.appliedEnvelopePruneTimer)
+ this.appliedEnvelopePruneTimer = null
+ }
+ this.appliedEnvelopeIds.clear()
+ this.applyingEnvelopeIds.clear()
+
this.channel.close()
this.collections.clear()
+ this.collectionAdapters.clear()
}
// -----------------------------------------------------------------------
@@ -344,11 +403,16 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
lockName,
{ signal: abortController.signal },
async () => {
- if (this.isDisposed()) return
+ if (
+ this.isDisposed() ||
+ this.collections.get(collectionId) !== state
+ ) {
+ return
+ }
try {
// Restore stream position from DB before claiming leadership
- const adapter = this.requireAdapter()
+ const adapter = this.requireAdapter(collectionId)
if (adapter.getStreamPosition) {
const pos = await adapter.getStreamPosition(collectionId)
state.latestTerm = pos.latestTerm
@@ -356,6 +420,13 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
state.latestRowVersion = pos.latestRowVersion
}
+ if (
+ this.isDisposed() ||
+ this.collections.get(collectionId) !== state
+ ) {
+ return
+ }
+
state.latestTerm++
state.isLeader = true
@@ -393,7 +464,7 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
}
// Re-acquire if not disposed (leadership was released by another means)
- if (!this.isDisposed()) {
+ if (!this.isDisposed() && this.collections.get(collectionId) === state) {
void this.acquireLeadership(collectionId, state)
}
}
@@ -413,6 +484,23 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
state.isLeader = false
}
+ private releaseCollectionIfUnused(
+ collectionId: string,
+ state: CollectionState,
+ ): void {
+ if (
+ state.subscribers.size > 0 ||
+ this.collections.get(collectionId) !== state
+ ) {
+ return
+ }
+
+ this.releaseLeadership(collectionId, state)
+ this.collections.delete(collectionId)
+ this.collectionAdapters.delete(collectionId)
+ this.deleteAppliedEnvelopeIdsForCollection(collectionId)
+ }
+
private emitHeartbeat(collectionId: string, state: CollectionState): void {
const envelope: ProtocolEnvelope = {
v: 1,
@@ -486,16 +574,22 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
collectionId: string,
request: RPCRequest,
): Promise {
+ if (this.isDisposed()) throw new Error(`coordinator disposed`)
let lastError: Error | undefined
for (let attempt = 0; attempt <= RPC_RETRY_ATTEMPTS; attempt++) {
if (attempt > 0) {
- await sleep(RPC_RETRY_DELAY_MS * attempt)
+ await Promise.race([
+ sleep(RPC_RETRY_DELAY_MS * attempt),
+ this.disposedPromise,
+ ])
}
+ if (this.isDisposed()) throw new Error(`coordinator disposed`)
try {
return await this.sendRPCOnce(collectionId, request)
} catch (error) {
+ if (this.isDisposed()) throw error
lastError = error instanceof Error ? error : new Error(String(error))
}
}
@@ -610,7 +704,7 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
},
): Promise {
await this.withWriterLock(() =>
- this.requireAdapter().ensureIndex(
+ this.requireAdapter(collectionId).ensureIndex(
collectionId,
request.signature,
request.spec,
@@ -632,17 +726,43 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
mutations: Array
},
): Promise {
- // Dedupe by envelopeId
- if (this.appliedEnvelopeIds.has(request.envelopeId)) {
- return {
- type: `rpc:applyLocalMutations:res`,
- rpcId: request.rpcId,
- ok: false,
- code: `CONFLICT`,
- error: `envelope ${request.envelopeId} already applied`,
+ const envelopeKey = JSON.stringify([collectionId, request.envelopeId])
+ this.pruneAppliedEnvelopeIds()
+ const appliedEnvelope = this.appliedEnvelopeIds.get(envelopeKey)
+ if (appliedEnvelope) {
+ return { ...appliedEnvelope.response, rpcId: request.rpcId }
+ }
+
+ const applyingEnvelope = this.applyingEnvelopeIds.get(envelopeKey)
+ if (applyingEnvelope) {
+ return { ...(await applyingEnvelope), rpcId: request.rpcId }
+ }
+
+ const application = this.applyLocalMutationsOnce(
+ collectionId,
+ request,
+ envelopeKey,
+ )
+ this.applyingEnvelopeIds.set(envelopeKey, application)
+ try {
+ return await application
+ } finally {
+ if (this.applyingEnvelopeIds.get(envelopeKey) === application) {
+ this.applyingEnvelopeIds.delete(envelopeKey)
}
}
+ }
+ private async applyLocalMutationsOnce(
+ collectionId: string,
+ request: {
+ type: `rpc:applyLocalMutations:req`
+ rpcId: string
+ envelopeId: string
+ mutations: Array
+ },
+ envelopeKey: string,
+ ): Promise {
const state = this.collections.get(collectionId)
if (!state || !state.isLeader) {
return {
@@ -676,12 +796,28 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
}
await this.withWriterLock(() =>
- this.requireAdapter().applyCommittedTx(collectionId, tx),
+ this.requireAdapter(collectionId).applyCommittedTx(collectionId, tx),
)
- // Track envelope for dedup
- this.appliedEnvelopeIds.set(request.envelopeId, Date.now())
+ const response: ApplyLocalMutationsResponse = {
+ type: `rpc:applyLocalMutations:res`,
+ rpcId: request.rpcId,
+ ok: true,
+ term,
+ seq,
+ latestRowVersion: rowVersion,
+ acceptedMutationIds: request.mutations.map((m) => m.mutationId),
+ }
+
+ // Retain the completed result so transport retries observe the same
+ // durable outcome without applying the envelope again.
+ this.appliedEnvelopeIds.set(envelopeKey, {
+ collectionId,
+ appliedAt: Date.now(),
+ response,
+ })
this.pruneAppliedEnvelopeIds()
+ this.scheduleAppliedEnvelopePrune()
// Broadcast tx:committed to all tabs
const changedRows = request.mutations
@@ -715,15 +851,7 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
subscriber(txCommitted)
}
- return {
- type: `rpc:applyLocalMutations:res`,
- rpcId: request.rpcId,
- ok: true,
- term,
- seq,
- latestRowVersion: rowVersion,
- acceptedMutationIds: request.mutations.map((m) => m.mutationId),
- }
+ return response
}
private async handlePullSince(
@@ -733,10 +861,11 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
rpcId: string
fromRowVersion: number
},
+ scopedAdapter?: HydrationPersistenceAdapter,
): Promise {
const state = this.collections.get(collectionId)
- const adapter = this.requireAdapter()
+ const adapter = scopedAdapter ?? this.requireAdapter(collectionId)
if (!adapter.pullSince) {
return {
type: `rpc:pullSince:res`,
@@ -809,13 +938,43 @@ export class BrowserCollectionCoordinator implements PersistedCollectionCoordina
// -----------------------------------------------------------------------
private pruneAppliedEnvelopeIds(): void {
- // Keep envelopes for 60 seconds for dedup
- const cutoff = Date.now() - 60_000
- for (const [id, ts] of this.appliedEnvelopeIds) {
- if (ts < cutoff) {
+ const cutoff = Date.now() - APPLIED_ENVELOPE_RETENTION_MS
+ for (const [id, applied] of this.appliedEnvelopeIds) {
+ if (applied.appliedAt <= cutoff) {
+ this.appliedEnvelopeIds.delete(id)
+ }
+ }
+ }
+
+ private deleteAppliedEnvelopeIdsForCollection(collectionId: string): void {
+ for (const [id, applied] of this.appliedEnvelopeIds) {
+ if (applied.collectionId === collectionId) {
this.appliedEnvelopeIds.delete(id)
}
}
+ this.scheduleAppliedEnvelopePrune()
+ }
+
+ private scheduleAppliedEnvelopePrune(): void {
+ if (this.appliedEnvelopePruneTimer !== null) {
+ clearTimeout(this.appliedEnvelopePruneTimer)
+ this.appliedEnvelopePruneTimer = null
+ }
+ if (this.disposed || this.appliedEnvelopeIds.size === 0) return
+
+ let earliestAppliedAt = Number.POSITIVE_INFINITY
+ for (const applied of this.appliedEnvelopeIds.values()) {
+ earliestAppliedAt = Math.min(earliestAppliedAt, applied.appliedAt)
+ }
+ const delay = Math.max(
+ 0,
+ earliestAppliedAt + APPLIED_ENVELOPE_RETENTION_MS - Date.now(),
+ )
+ this.appliedEnvelopePruneTimer = setTimeout(() => {
+ this.appliedEnvelopePruneTimer = null
+ this.pruneAppliedEnvelopeIds()
+ this.scheduleAppliedEnvelopePrune()
+ }, delay)
}
}
diff --git a/packages/browser-db-sqlite-persistence/src/browser-persistence.ts b/packages/browser-db-sqlite-persistence/src/browser-persistence.ts
index b36e564900..c9ac081fbf 100644
--- a/packages/browser-db-sqlite-persistence/src/browser-persistence.ts
+++ b/packages/browser-db-sqlite-persistence/src/browser-persistence.ts
@@ -131,32 +131,36 @@ export function createBrowserWASQLitePersistence(
})
adapterCache.set(cacheKey, adapter)
- // Wire the adapter into the multi-tab coordinator so it can handle
- // leader-side RPCs (applyCommittedTx, pullSince, ensureIndex, etc.)
- if (resolvedCoordinator instanceof BrowserCollectionCoordinator) {
- resolvedCoordinator.setAdapter(adapter)
- }
-
return adapter
}
const createCollectionPersistence = (
mode: PersistedCollectionMode,
schemaVersion: number | undefined,
- ): PersistedCollectionPersistence => ({
- adapter: getAdapterForCollection(mode, schemaVersion),
- coordinator: resolvedCoordinator,
- })
+ collectionId?: string,
+ ): PersistedCollectionPersistence => {
+ const adapter = getAdapterForCollection(mode, schemaVersion)
+ if (
+ collectionId !== undefined &&
+ resolvedCoordinator instanceof BrowserCollectionCoordinator
+ ) {
+ resolvedCoordinator.setAdapterForCollection(collectionId, adapter)
+ }
+ return { adapter, coordinator: resolvedCoordinator }
+ }
const defaultPersistence = createCollectionPersistence(
`sync-absent`,
undefined,
)
+ if (resolvedCoordinator instanceof BrowserCollectionCoordinator) {
+ resolvedCoordinator.setAdapter(defaultPersistence.adapter)
+ }
return {
...defaultPersistence,
- resolvePersistenceForCollection: ({ mode, schemaVersion }) =>
- createCollectionPersistence(mode, schemaVersion),
+ resolvePersistenceForCollection: ({ collectionId, mode, schemaVersion }) =>
+ createCollectionPersistence(mode, schemaVersion, collectionId),
resolvePersistenceForMode: (mode) =>
createCollectionPersistence(mode, undefined),
}
diff --git a/packages/browser-db-sqlite-persistence/src/wa-sqlite-driver.ts b/packages/browser-db-sqlite-persistence/src/wa-sqlite-driver.ts
index c1ae91b9e0..3a6be4c901 100644
--- a/packages/browser-db-sqlite-persistence/src/wa-sqlite-driver.ts
+++ b/packages/browser-db-sqlite-persistence/src/wa-sqlite-driver.ts
@@ -1,4 +1,7 @@
-import { InvalidPersistedCollectionConfigError } from '@tanstack/db-sqlite-persistence-core'
+import {
+ InvalidPersistedCollectionConfigError,
+ SQLITE_DRIVER_SHARED_LOGICAL_SCHEDULING_KEY,
+} from '@tanstack/db-sqlite-persistence-core'
import type { SQLiteDriver } from '@tanstack/db-sqlite-persistence-core'
export type BrowserWASQLiteDatabase = {
@@ -36,6 +39,7 @@ function assertDatabaseShape(
}
export class BrowserWASQLiteDriver implements SQLiteDriver {
+ readonly [SQLITE_DRIVER_SHARED_LOGICAL_SCHEDULING_KEY] = {}
private readonly database: BrowserWASQLiteDatabase
private queue: Promise = Promise.resolve()
private nextSavepointId = 1
@@ -46,29 +50,33 @@ export class BrowserWASQLiteDriver implements SQLiteDriver {
this.database = options.database
}
- async exec(sql: string): Promise {
- await this.enqueue(async () => {
+ exec(sql: string): Promise {
+ return this.enqueue(async () => {
await this.database.execute(sql)
})
}
- async query(
+ query(
sql: string,
params: ReadonlyArray = [],
): Promise> {
return this.enqueue(() => this.database.execute(sql, params))
}
- async run(sql: string, params: ReadonlyArray = []): Promise {
- await this.enqueue(async () => {
+ run(sql: string, params: ReadonlyArray = []): Promise {
+ return this.enqueue(async () => {
await this.database.execute(sql, params)
})
}
- async transaction(
+ transaction(
fn: (transactionDriver: SQLiteDriver) => Promise,
): Promise {
- assertTransactionCallbackHasDriverArg(fn)
+ try {
+ assertTransactionCallbackHasDriverArg(fn)
+ } catch (error) {
+ return Promise.reject(error)
+ }
return this.enqueue(async () => {
await this.database.execute(`BEGIN IMMEDIATE`)
@@ -87,7 +95,7 @@ export class BrowserWASQLiteDriver implements SQLiteDriver {
})
}
- async transactionWithDriver(
+ transactionWithDriver(
fn: (transactionDriver: SQLiteDriver) => Promise,
): Promise {
return this.transaction(fn)
@@ -144,6 +152,14 @@ export class BrowserWASQLiteDriver implements SQLiteDriver {
private enqueue(operation: () => Promise | T): Promise {
const queuedOperation = this.queue.then(operation, operation)
+ // Brand the exact Promise returned to callers. Transparent wrappers may
+ // preserve this identity for late discovery; wrappers that create a new
+ // Promise must forward the driver key before adapter construction.
+ Object.defineProperty(
+ queuedOperation,
+ SQLITE_DRIVER_SHARED_LOGICAL_SCHEDULING_KEY,
+ { value: this[SQLITE_DRIVER_SHARED_LOGICAL_SCHEDULING_KEY] },
+ )
this.queue = queuedOperation.then(
() => undefined,
() => undefined,
diff --git a/packages/browser-db-sqlite-persistence/tests/browser-coordinator.test.ts b/packages/browser-db-sqlite-persistence/tests/browser-coordinator.test.ts
index 5f554006e2..a73b76cb62 100644
--- a/packages/browser-db-sqlite-persistence/tests/browser-coordinator.test.ts
+++ b/packages/browser-db-sqlite-persistence/tests/browser-coordinator.test.ts
@@ -1,7 +1,31 @@
+/**
+ * # Which cross-window coordinator facts survive retries and disposal?
+ *
+ * Contract: leader work uses the adapter registered for that Collection. A
+ * retry after response loss observes the first durable mutation result without
+ * applying it again. Disposal is terminal for current and future RPC attempts.
+ *
+ * Model and history grammar: an independent collection-to-adapter map selects
+ * one call recorder; a completed envelope retains one result; a disposed flag
+ * permits no request posts. Generated histories vary Collection registration
+ * order, RPC kind, mutation count, response delivery, and disposal phase.
+ *
+ * Production path and observations: paired real coordinators communicate
+ * through controlled BroadcastChannel and Web Locks fixtures. The laws record
+ * exact adapter calls, durable transaction count, returned result, promise
+ * settlement, and request posts.
+ *
+ * Known omissions: these laws do not establish native browser scheduling or
+ * concurrent delivery of one in-flight envelope. Electron keeps focused
+ * parity witnesses for the same protocol boundaries.
+ */
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
+import fc from 'fast-check'
+import { BasicIndex, createCollection } from '@tanstack/db'
+import { persistedCollectionOptions } from '@tanstack/db-sqlite-persistence-core'
import { BrowserCollectionCoordinator } from '../src/browser-coordinator'
-import type { BrowserCollectionCoordinatorOptions } from '../src/browser-coordinator'
import type { PersistenceAdapter } from '@tanstack/db-sqlite-persistence-core'
+import type { BrowserCollectionCoordinatorOptions } from '../src/browser-coordinator'
// ---------------------------------------------------------------------------
// BroadcastChannel mock
@@ -12,6 +36,8 @@ const channels: Map<
string,
Set<{ onmessage: MessageHandler | null }>
> = new Map()
+let dropNextMessageWhen: ((data: unknown) => boolean) | undefined
+let observePostedMessage: ((data: unknown) => void) | undefined
class MockBroadcastChannel {
readonly name: string
@@ -26,6 +52,11 @@ class MockBroadcastChannel {
}
postMessage(data: unknown): void {
+ observePostedMessage?.(data)
+ if (dropNextMessageWhen?.(data)) {
+ dropNextMessageWhen = undefined
+ return
+ }
const peers = channels.get(this.name)
if (!peers) return
// Deliver to all other instances on same channel (simulating cross-tab)
@@ -162,6 +193,8 @@ function installGlobals(): void {
}
function cleanupGlobals(): void {
+ dropNextMessageWhen = undefined
+ observePostedMessage = undefined
channels.clear()
heldLocks.clear()
lockQueues.clear()
@@ -228,6 +261,31 @@ async function flush(ms: number = 10): Promise {
await new Promise((resolve) => setTimeout(resolve, ms))
}
+type Deferred = {
+ promise: Promise
+ resolve: () => void
+}
+
+function createDeferred(): Deferred {
+ let resolve!: () => void
+ const promise = new Promise((settle) => {
+ resolve = settle
+ })
+ return { promise, resolve }
+}
+
+type CoordinatorInspection = {
+ collectionAdapters: Map
+ collections: Map
+ appliedEnvelopeIds: Map
+}
+
+function inspectCoordinator(
+ coordinator: BrowserCollectionCoordinator,
+): CoordinatorInspection {
+ return coordinator as unknown as CoordinatorInspection
+}
+
// ---------------------------------------------------------------------------
// Tests
// ---------------------------------------------------------------------------
@@ -460,6 +518,319 @@ describe(`BrowserCollectionCoordinator`, () => {
coord.dispose()
})
+
+ it(`scopes envelope deduplication by collection`, async () => {
+ const alphaAdapter = createStubAdapter()
+ const betaAdapter = createStubAdapter()
+ const coordinator = createCoordinator(alphaAdapter)
+ coordinator.setAdapterForCollection(`alpha`, alphaAdapter)
+ coordinator.setAdapterForCollection(`beta`, betaAdapter)
+ coordinator.subscribe(`alpha`, () => {})
+ coordinator.subscribe(`beta`, () => {})
+ await flush(50)
+
+ const internals = coordinator as unknown as {
+ handleApplyLocalMutations: (
+ collectionId: string,
+ request: {
+ type: `rpc:applyLocalMutations:req`
+ rpcId: string
+ envelopeId: string
+ mutations: Array<{
+ mutationId: string
+ type: `insert`
+ key: string
+ value: { id: string }
+ }>
+ },
+ ) => Promise<{ ok: boolean; rpcId: string }>
+ }
+
+ try {
+ const alpha = await internals.handleApplyLocalMutations(`alpha`, {
+ type: `rpc:applyLocalMutations:req`,
+ rpcId: `alpha-rpc`,
+ envelopeId: `shared-envelope`,
+ mutations: [
+ {
+ mutationId: `alpha-mutation`,
+ type: `insert`,
+ key: `alpha`,
+ value: { id: `alpha` },
+ },
+ ],
+ })
+ const beta = await internals.handleApplyLocalMutations(`beta`, {
+ type: `rpc:applyLocalMutations:req`,
+ rpcId: `beta-rpc`,
+ envelopeId: `shared-envelope`,
+ mutations: [
+ {
+ mutationId: `beta-mutation`,
+ type: `insert`,
+ key: `beta`,
+ value: { id: `beta` },
+ },
+ ],
+ })
+
+ expect({
+ alpha,
+ alphaApplies: alphaAdapter.appliedTxs,
+ beta,
+ betaApplies: betaAdapter.appliedTxs,
+ }).toMatchObject({
+ alpha: { ok: true, rpcId: `alpha-rpc` },
+ alphaApplies: [{ collectionId: `alpha` }],
+ beta: { ok: true, rpcId: `beta-rpc` },
+ betaApplies: [{ collectionId: `beta` }],
+ })
+ } finally {
+ coordinator.dispose()
+ }
+ })
+
+ it(`coalesces an envelope retry while its first write is in flight`, async () => {
+ const adapter = createStubAdapter()
+ const firstApplyEntered = createDeferred()
+ const releaseFirstApply = createDeferred()
+ let applyCalls = 0
+ adapter.applyCommittedTx = async (collectionId, tx) => {
+ applyCalls++
+ if (applyCalls === 1) {
+ firstApplyEntered.resolve()
+ await releaseFirstApply.promise
+ }
+ adapter.appliedTxs.push({ collectionId, txId: tx.txId })
+ }
+ const coordinator = createCoordinator(adapter)
+ coordinator.subscribe(`todos`, () => {})
+ await flush(50)
+
+ const internals = coordinator as unknown as {
+ handleApplyLocalMutations: (
+ collectionId: string,
+ request: {
+ type: `rpc:applyLocalMutations:req`
+ rpcId: string
+ envelopeId: string
+ mutations: Array<{
+ mutationId: string
+ type: `insert`
+ key: string
+ value: { id: string }
+ }>
+ },
+ ) => Promise<{
+ ok: boolean
+ rpcId: string
+ term?: number
+ seq?: number
+ latestRowVersion?: number
+ }>
+ }
+ const request = {
+ type: `rpc:applyLocalMutations:req` as const,
+ envelopeId: `in-flight-envelope`,
+ mutations: [
+ {
+ mutationId: `mutation`,
+ type: `insert` as const,
+ key: `row`,
+ value: { id: `row` },
+ },
+ ],
+ }
+
+ try {
+ const first = internals.handleApplyLocalMutations(`todos`, {
+ ...request,
+ rpcId: `first-rpc`,
+ })
+ await firstApplyEntered.promise
+ const retry = internals.handleApplyLocalMutations(`todos`, {
+ ...request,
+ rpcId: `retry-rpc`,
+ })
+ await Promise.resolve()
+ releaseFirstApply.resolve()
+
+ const [firstResponse, retryResponse] = await Promise.all([first, retry])
+ expect({
+ applyCalls,
+ first: { ...firstResponse, rpcId: undefined },
+ retry: { ...retryResponse, rpcId: undefined },
+ }).toEqual({
+ applyCalls: 1,
+ first: { ...retryResponse, rpcId: undefined },
+ retry: { ...firstResponse, rpcId: undefined },
+ })
+ } finally {
+ releaseFirstApply.resolve()
+ coordinator.dispose()
+ }
+ })
+
+ it(`replays the successful mutation result when its first response is lost`, async () => {
+ vi.useFakeTimers()
+ const adapter = createStubAdapter()
+ let leadershipRead!: () => void
+ const leadershipReadPromise = new Promise((resolve) => {
+ leadershipRead = resolve
+ })
+ let firstSuccessDropped!: () => void
+ const firstSuccessDroppedPromise = new Promise((resolve) => {
+ firstSuccessDropped = resolve
+ })
+ adapter.getStreamPosition = async () => {
+ leadershipRead()
+ return { latestTerm: 0, latestSeq: 0, latestRowVersion: 0 }
+ }
+
+ const leader = createCoordinator(adapter)
+ const follower = createCoordinator(adapter)
+
+ try {
+ leader.subscribe(`todos`, () => {})
+ await leadershipReadPromise
+ await Promise.resolve()
+ follower.subscribe(`todos`, () => {})
+
+ dropNextMessageWhen = (data) => {
+ const payload = (
+ data as { payload?: { type?: string; ok?: boolean } }
+ ).payload
+ if (
+ payload?.type === `rpc:applyLocalMutations:res` &&
+ payload.ok === true
+ ) {
+ firstSuccessDropped()
+ return true
+ }
+ return false
+ }
+
+ const responsePromise = follower.requestApplyLocalMutations(`todos`, [
+ {
+ mutationId: `mut-lost-response`,
+ type: `insert`,
+ key: `lost-response`,
+ value: { id: `lost-response`, title: `Persisted once` },
+ },
+ ])
+
+ await firstSuccessDroppedPromise
+ expect(adapter.appliedTxs).toHaveLength(1)
+
+ await vi.advanceTimersByTimeAsync(10_200)
+ const response = await responsePromise
+ expect(response).toMatchObject({
+ ok: true,
+ acceptedMutationIds: [`mut-lost-response`],
+ })
+ expect(adapter.appliedTxs).toHaveLength(1)
+ } finally {
+ follower.dispose()
+ leader.dispose()
+ vi.useRealTimers()
+ }
+ })
+
+ it(`preserves one durable mutation result across generated response-delivery histories`, async () => {
+ vi.useFakeTimers()
+ let run = 0
+ try {
+ await fc.assert(
+ fc.asyncProperty(
+ fc.record({
+ mutationCount: fc.integer({ min: 1, max: 4 }),
+ dropFirstSuccess: fc.boolean(),
+ }),
+ async ({ mutationCount, dropFirstSuccess }) => {
+ run++
+ const collectionId = `delivery-history-${run}`
+ const adapter = createStubAdapter()
+ let leadershipRead!: () => void
+ const leadershipReadPromise = new Promise((resolve) => {
+ leadershipRead = resolve
+ })
+ let firstSuccessDropped!: () => void
+ const firstSuccessDroppedPromise = new Promise(
+ (resolve) => {
+ firstSuccessDropped = resolve
+ },
+ )
+ adapter.getStreamPosition = async () => {
+ leadershipRead()
+ return { latestTerm: 0, latestSeq: 0, latestRowVersion: 0 }
+ }
+
+ const leader = createCoordinator(adapter)
+ const follower = createCoordinator(adapter)
+ try {
+ leader.subscribe(collectionId, () => {})
+ await leadershipReadPromise
+ await Promise.resolve()
+ follower.subscribe(collectionId, () => {})
+
+ if (dropFirstSuccess) {
+ dropNextMessageWhen = (data) => {
+ const payload = (
+ data as { payload?: { type?: string; ok?: boolean } }
+ ).payload
+ if (
+ payload?.type === `rpc:applyLocalMutations:res` &&
+ payload.ok === true
+ ) {
+ firstSuccessDropped()
+ return true
+ }
+ return false
+ }
+ }
+
+ const mutations = Array.from(
+ { length: mutationCount },
+ (_, index) => ({
+ mutationId: `mut-${run}-${index}`,
+ type: `insert` as const,
+ key: `${index}`,
+ value: { id: `${index}`, title: `row ${index}` },
+ }),
+ )
+ const responsePromise = follower.requestApplyLocalMutations(
+ collectionId,
+ mutations,
+ )
+ if (dropFirstSuccess) {
+ await firstSuccessDroppedPromise
+ await vi.advanceTimersByTimeAsync(10_200)
+ }
+
+ const response = await responsePromise
+ expect(response).toMatchObject({
+ ok: true,
+ acceptedMutationIds: mutations.map(
+ (mutation) => mutation.mutationId,
+ ),
+ })
+ expect(adapter.appliedTxs).toHaveLength(1)
+ } finally {
+ dropNextMessageWhen = undefined
+ follower.dispose()
+ leader.dispose()
+ for (let microtask = 0; microtask < 4; microtask++) {
+ await Promise.resolve()
+ }
+ }
+ },
+ ),
+ { seed: 1868, numRuns: 12, endOnFailure: true },
+ )
+ } finally {
+ vi.useRealTimers()
+ }
+ })
})
describe(`RPC - pullSince`, () => {
@@ -528,6 +899,243 @@ describe(`BrowserCollectionCoordinator`, () => {
coord.dispose()
})
+ it(`does not repeat successful leader-local index creation`, async () => {
+ const adapter = createStubAdapter()
+ let leadershipRead!: () => void
+ const leadershipReadPromise = new Promise((resolve) => {
+ leadershipRead = resolve
+ })
+ adapter.getStreamPosition = async () => {
+ leadershipRead()
+ return { latestTerm: 0, latestSeq: 0, latestRowVersion: 0 }
+ }
+ adapter.ensureIndex = vi.fn().mockResolvedValue(undefined)
+
+ const coord = createCoordinator(adapter)
+ try {
+ coord.subscribe(`todos`, () => {})
+ await leadershipReadPromise
+ await Promise.resolve()
+
+ const spec = { expressionSql: [`title`] }
+ await adapter.ensureIndex(`todos`, `idx-once`, spec)
+ await coord.requestEnsurePersistedIndex(
+ `todos`,
+ `idx-once`,
+ spec,
+ adapter,
+ true,
+ )
+
+ expect(adapter.ensureIndex).toHaveBeenCalledOnce()
+ } finally {
+ coord.dispose()
+ }
+ })
+
+ it(`keeps production collection bootstrap index work exact-once on the leader`, async () => {
+ const adapter = createStubAdapter()
+ adapter.ensureIndex = vi.fn().mockResolvedValue(undefined)
+ const coordinator = createCoordinator(adapter)
+ const leaderReady = createDeferred()
+ observePostedMessage = (data) => {
+ const envelope = data as {
+ collectionId?: string
+ payload?: { type?: string }
+ }
+ if (
+ envelope.collectionId === `leader-bootstrap-index` &&
+ envelope.payload?.type === `leader:heartbeat`
+ ) {
+ leaderReady.resolve()
+ }
+ }
+ const releaseLeader = coordinator.subscribe(
+ `leader-bootstrap-index`,
+ () => {},
+ )
+ await leaderReady.promise
+ observePostedMessage = undefined
+
+ const collection = createCollection(
+ persistedCollectionOptions<{ id: string; title: string }, string>({
+ id: `leader-bootstrap-index`,
+ getKey: (row) => row.id,
+ defaultIndexType: BasicIndex,
+ sync: {
+ sync: ({ markReady }) => {
+ markReady()
+ },
+ },
+ persistence: { adapter, coordinator },
+ }),
+ )
+ collection.createIndex((row) => row.title, { name: `title` })
+
+ try {
+ await collection.preload()
+ expect(adapter.ensureIndex).toHaveBeenCalledOnce()
+ } finally {
+ await collection.cleanup()
+ releaseLeader()
+ coordinator.dispose()
+ }
+ })
+
+ it(`reaches crossed production RPC leadership only after both hydration scopes exit`, async () => {
+ const bothLocalIndexesEntered = createDeferred()
+ const bothLeaderRPCsEntered = createDeferred()
+ const releaseSchedulerCycle = createDeferred()
+ let localIndexEntries = 0
+ let leaderRPCEntries = 0
+
+ const createTab = () => {
+ const adapter = createStubAdapter()
+ const baseEnsureIndex = adapter.ensureIndex.bind(adapter)
+ let hydrationScopeActive = false
+ const leaderRPCScopeObservations: Array = []
+ const scopedAdapter: PersistenceAdapter = {
+ ...adapter,
+ ensureIndex: async (...args) => {
+ localIndexEntries++
+ if (localIndexEntries === 2) bothLocalIndexesEntered.resolve()
+ await bothLocalIndexesEntered.promise
+ await baseEnsureIndex(...args)
+ },
+ }
+ adapter.runInHydrationScope = async (task) => {
+ hydrationScopeActive = true
+ try {
+ return await task(scopedAdapter)
+ } finally {
+ hydrationScopeActive = false
+ }
+ }
+ adapter.ensureIndex = async (...args) => {
+ leaderRPCScopeObservations.push(hydrationScopeActive)
+ leaderRPCEntries++
+ if (leaderRPCEntries === 2) bothLeaderRPCsEntered.resolve()
+ if (hydrationScopeActive) await releaseSchedulerCycle.promise
+ await baseEnsureIndex(...args)
+ }
+ return { adapter, leaderRPCScopeObservations }
+ }
+
+ const tab1 = createTab()
+ const tab2 = createTab()
+ const coordinator1 = createCoordinator(tab1.adapter)
+ const coordinator2 = createCoordinator(tab2.adapter)
+ const leaderCollections = new Set()
+ const bothLeadersReady = createDeferred()
+ observePostedMessage = (data) => {
+ const envelope = data as {
+ collectionId?: string
+ payload?: { type?: string }
+ }
+ if (
+ envelope.payload?.type === `leader:heartbeat` &&
+ (envelope.collectionId === `crossed-a` ||
+ envelope.collectionId === `crossed-b`)
+ ) {
+ leaderCollections.add(envelope.collectionId)
+ if (leaderCollections.size === 2) bothLeadersReady.resolve()
+ }
+ }
+ coordinator1.setAdapterForCollection(`crossed-b`, tab1.adapter)
+ coordinator2.setAdapterForCollection(`crossed-a`, tab2.adapter)
+ const releaseLeaderB = coordinator1.subscribe(`crossed-b`, () => {})
+ const releaseLeaderA = coordinator2.subscribe(`crossed-a`, () => {})
+ await bothLeadersReady.promise
+ observePostedMessage = undefined
+
+ const createFollowerCollection = (
+ id: string,
+ adapter: PersistenceAdapter,
+ coordinator: BrowserCollectionCoordinator,
+ ) => {
+ const collection = createCollection(
+ persistedCollectionOptions<{ id: string; title: string }, string>({
+ id,
+ getKey: (row) => row.id,
+ defaultIndexType: BasicIndex,
+ sync: {
+ sync: ({ markReady }) => {
+ markReady()
+ },
+ },
+ persistence: { adapter, coordinator },
+ }),
+ )
+ collection.createIndex((row) => row.title, { name: `${id}-title` })
+ return collection
+ }
+ const followerA = createFollowerCollection(
+ `crossed-a`,
+ tab1.adapter,
+ coordinator1,
+ )
+ const followerB = createFollowerCollection(
+ `crossed-b`,
+ tab2.adapter,
+ coordinator2,
+ )
+ const preloadA = Promise.resolve(followerA.preload())
+ const preloadB = Promise.resolve(followerB.preload())
+ void preloadA.catch(() => undefined)
+ void preloadB.catch(() => undefined)
+
+ try {
+ await bothLeaderRPCsEntered.promise
+ expect({
+ tab1: tab1.leaderRPCScopeObservations,
+ tab2: tab2.leaderRPCScopeObservations,
+ }).toEqual({ tab1: [false], tab2: [false] })
+ await Promise.all([preloadA, preloadB])
+ } finally {
+ releaseSchedulerCycle.resolve()
+ await Promise.all([preloadA, preloadB]).catch(() => undefined)
+ await Promise.all([followerA.cleanup(), followerB.cleanup()])
+ releaseLeaderA()
+ releaseLeaderB()
+ coordinator2.dispose()
+ coordinator1.dispose()
+ }
+ })
+
+ it(`uses a supplied leader-local adapter when local work is not complete`, async () => {
+ const registeredAdapter = createStubAdapter()
+ const scopedAdapter = createStubAdapter()
+ let leadershipRead!: () => void
+ const leadershipReadPromise = new Promise((resolve) => {
+ leadershipRead = resolve
+ })
+ registeredAdapter.getStreamPosition = async () => {
+ leadershipRead()
+ return { latestTerm: 0, latestSeq: 0, latestRowVersion: 0 }
+ }
+ registeredAdapter.ensureIndex = vi.fn().mockResolvedValue(undefined)
+ scopedAdapter.ensureIndex = vi.fn().mockResolvedValue(undefined)
+
+ const coord = createCoordinator(registeredAdapter)
+ try {
+ coord.subscribe(`todos`, () => {})
+ await leadershipReadPromise
+ await Promise.resolve()
+
+ await coord.requestEnsurePersistedIndex(
+ `todos`,
+ `idx-scoped`,
+ { expressionSql: [`title`] },
+ scopedAdapter,
+ )
+
+ expect(scopedAdapter.ensureIndex).toHaveBeenCalledOnce()
+ expect(registeredAdapter.ensureIndex).not.toHaveBeenCalled()
+ } finally {
+ coord.dispose()
+ }
+ })
+
it(`follower routes ensurePersistedIndex to leader`, async () => {
const adapter = createStubAdapter()
adapter.ensureIndex = vi.fn().mockResolvedValue(undefined)
@@ -548,6 +1156,351 @@ describe(`BrowserCollectionCoordinator`, () => {
leader.dispose()
follower.dispose()
})
+
+ it(`keeps leader RPC work on the adapter registered for its collection`, async () => {
+ const todosAdapter = createStubAdapter()
+ const notesAdapter = createStubAdapter()
+ let leadershipRead!: () => void
+ const leadershipReadPromise = new Promise((resolve) => {
+ leadershipRead = resolve
+ })
+ todosAdapter.ensureIndex = vi.fn().mockResolvedValue(undefined)
+ notesAdapter.ensureIndex = vi.fn().mockResolvedValue(undefined)
+ todosAdapter.getStreamPosition = async () => {
+ leadershipRead()
+ return { latestTerm: 0, latestSeq: 0, latestRowVersion: 0 }
+ }
+
+ const leader = createCoordinator(todosAdapter)
+ const follower = createCoordinator(notesAdapter)
+
+ try {
+ leader.subscribe(`todos`, () => {})
+ await leadershipReadPromise
+ await Promise.resolve()
+ expect(leader.isLeader(`todos`)).toBe(true)
+ follower.subscribe(`todos`, () => {})
+
+ // Resolving a later collection variant must not replace the adapter
+ // already owning leader-side work for `todos`.
+ leader.setAdapter(notesAdapter)
+ await follower.requestEnsurePersistedIndex(`todos`, `idx-todos`, {
+ expressionSql: [`title`],
+ })
+
+ expect(todosAdapter.ensureIndex).toHaveBeenCalledOnce()
+ expect(notesAdapter.ensureIndex).not.toHaveBeenCalled()
+ } finally {
+ follower.dispose()
+ leader.dispose()
+ }
+ })
+
+ it(`replaces a cached default adapter with its collection registration`, async () => {
+ const defaultAdapter = createStubAdapter()
+ const replacementAdapter = createStubAdapter()
+ let leadershipRead!: () => void
+ const leadershipReadPromise = new Promise((resolve) => {
+ leadershipRead = resolve
+ })
+ defaultAdapter.getStreamPosition = async () => {
+ leadershipRead()
+ return { latestTerm: 0, latestSeq: 0, latestRowVersion: 0 }
+ }
+ defaultAdapter.ensureIndex = vi.fn().mockResolvedValue(undefined)
+ replacementAdapter.ensureIndex = vi.fn().mockResolvedValue(undefined)
+
+ const coordinator = createCoordinator(defaultAdapter)
+ try {
+ coordinator.subscribe(`todos`, () => {})
+ await leadershipReadPromise
+ await Promise.resolve()
+
+ await coordinator.requestEnsurePersistedIndex(`todos`, `idx-default`, {
+ expressionSql: [`title`],
+ })
+ coordinator.setAdapterForCollection(`todos`, replacementAdapter)
+ await coordinator.requestEnsurePersistedIndex(
+ `todos`,
+ `idx-replacement`,
+ { expressionSql: [`title`] },
+ )
+
+ expect(defaultAdapter.ensureIndex).toHaveBeenCalledOnce()
+ expect(replacementAdapter.ensureIndex).toHaveBeenCalledOnce()
+ expect(replacementAdapter.ensureIndex).toHaveBeenCalledWith(
+ `todos`,
+ `idx-replacement`,
+ { expressionSql: [`title`] },
+ )
+ } finally {
+ coordinator.dispose()
+ }
+ })
+
+ it(`routes generated collection and RPC histories through their owning adapter`, async () => {
+ let run = 0
+ await fc.assert(
+ fc.asyncProperty(
+ fc.integer({ min: 2, max: 4 }).chain((collectionCount) =>
+ fc.record({
+ collectionCount: fc.constant(collectionCount),
+ registrationOrder: fc.shuffledSubarray(
+ Array.from({ length: collectionCount }, (_, index) => index),
+ { minLength: collectionCount, maxLength: collectionCount },
+ ),
+ targetIndex: fc.integer({ min: 0, max: collectionCount - 1 }),
+ rpcKind: fc.constantFrom(`index`, `pull`, `mutation`),
+ }),
+ ),
+ async ({
+ collectionCount,
+ registrationOrder,
+ targetIndex,
+ rpcKind,
+ }) => {
+ run++
+ const collectionIds = Array.from(
+ { length: collectionCount },
+ (_, index) => `routing-${run}-${index}`,
+ )
+ const adapters = collectionIds.map(() => createStubAdapter())
+ const ensureCalls = adapters.map(() => vi.fn())
+ const pullCalls = adapters.map(() => vi.fn())
+ adapters.forEach((adapter, index) => {
+ adapter.ensureIndex =
+ ensureCalls[index]!.mockResolvedValue(undefined)
+ adapter.pullSince = pullCalls[index]!.mockResolvedValue({
+ latestRowVersion: index,
+ requiresFullReload: false,
+ changedKeys: [],
+ deletedKeys: [],
+ })
+ })
+
+ let leadershipRead!: () => void
+ const leadershipReadPromise = new Promise((resolve) => {
+ leadershipRead = resolve
+ })
+ adapters[targetIndex]!.getStreamPosition = async () => {
+ leadershipRead()
+ return { latestTerm: 0, latestSeq: 0, latestRowVersion: 0 }
+ }
+
+ const leader = createCoordinator(adapters[0])
+ const follower = createCoordinator(adapters.at(-1))
+ const targetCollectionId = collectionIds[targetIndex]!
+ try {
+ for (const index of registrationOrder) {
+ leader.setAdapterForCollection(
+ collectionIds[index]!,
+ adapters[index]!,
+ )
+ }
+ leader.subscribe(targetCollectionId, () => {})
+ await leadershipReadPromise
+ await Promise.resolve()
+ follower.subscribe(targetCollectionId, () => {})
+
+ if (rpcKind === `index`) {
+ await follower.requestEnsurePersistedIndex(
+ targetCollectionId,
+ `idx-${run}`,
+ { expressionSql: [`title`] },
+ )
+ } else if (rpcKind === `pull`) {
+ await follower.pullSince(targetCollectionId, 0)
+ } else {
+ const response = await follower.requestApplyLocalMutations(
+ targetCollectionId,
+ [
+ {
+ mutationId: `mut-${run}`,
+ type: `insert`,
+ key: `${run}`,
+ value: { id: `${run}`, title: `row ${run}` },
+ },
+ ],
+ )
+ expect(response.ok).toBe(true)
+ }
+
+ adapters.forEach((adapter, index) => {
+ const callCount =
+ rpcKind === `index`
+ ? ensureCalls[index]!.mock.calls.length
+ : rpcKind === `pull`
+ ? pullCalls[index]!.mock.calls.length
+ : adapter.appliedTxs.length
+ expect(callCount).toBe(index === targetIndex ? 1 : 0)
+ })
+ } finally {
+ follower.dispose()
+ leader.dispose()
+ for (let microtask = 0; microtask < 4; microtask++) {
+ await Promise.resolve()
+ }
+ }
+ },
+ ),
+ { seed: 1868, numRuns: 18, endOnFailure: true },
+ )
+ })
+ })
+
+ describe(`collection and retry-result retention`, () => {
+ it(`releases generated collection-owned state after the last subscriber`, async () => {
+ let run = 0
+ await fc.assert(
+ fc.asyncProperty(
+ fc.integer({ min: 1, max: 4 }).chain((collectionCount) =>
+ fc.record({
+ collectionCount: fc.constant(collectionCount),
+ releaseOrder: fc.shuffledSubarray(
+ Array.from({ length: collectionCount }, (_, index) => index),
+ { minLength: collectionCount, maxLength: collectionCount },
+ ),
+ }),
+ ),
+ async ({ collectionCount, releaseOrder }) => {
+ run++
+ const coordinator = createCoordinator()
+ const collectionIds = Array.from(
+ { length: collectionCount },
+ (_, index) => `lifecycle-${run}-${index}`,
+ )
+ const readyCollections = new Set()
+ const allReady = createDeferred()
+ observePostedMessage = (data) => {
+ const envelope = data as {
+ collectionId?: string
+ payload?: { type?: string }
+ }
+ if (
+ envelope.payload?.type === `leader:heartbeat` &&
+ envelope.collectionId &&
+ collectionIds.includes(envelope.collectionId)
+ ) {
+ readyCollections.add(envelope.collectionId)
+ if (readyCollections.size === collectionCount) {
+ allReady.resolve()
+ }
+ }
+ }
+ for (const collectionId of collectionIds) {
+ coordinator.setAdapterForCollection(
+ collectionId,
+ createStubAdapter(),
+ )
+ }
+ const releases = collectionIds.map((collectionId) =>
+ coordinator.subscribe(collectionId, () => {}),
+ )
+
+ try {
+ await allReady.promise
+ observePostedMessage = undefined
+ expect(
+ inspectCoordinator(coordinator).collectionAdapters.size,
+ ).toBe(collectionCount)
+
+ let remaining = collectionCount
+ for (const releasedIndex of releaseOrder) {
+ releases[releasedIndex]!()
+ remaining--
+ expect({
+ adapters:
+ inspectCoordinator(coordinator).collectionAdapters.size,
+ collections: inspectCoordinator(coordinator).collections.size,
+ }).toEqual({
+ adapters: remaining,
+ collections: remaining,
+ })
+ }
+ } finally {
+ observePostedMessage = undefined
+ for (const release of releases) release()
+ coordinator.dispose()
+ }
+ },
+ ),
+ { seed: 1868, numRuns: 12, endOnFailure: true },
+ )
+ })
+
+ it(`bounds generated retry-result histories by retention and collection lifecycle`, async () => {
+ vi.useFakeTimers()
+ let run = 0
+ try {
+ await fc.assert(
+ fc.asyncProperty(
+ fc.record({
+ mutationCount: fc.integer({ min: 1, max: 4 }),
+ terminal: fc.constantFrom(`retention`, `unsubscribe`, `dispose`),
+ }),
+ async ({ mutationCount, terminal }) => {
+ run++
+ vi.setSystemTime(0)
+ const adapter = createStubAdapter()
+ const coordinator = createCoordinator(adapter)
+ const collectionId = `envelope-retention-${run}`
+ const ready = createDeferred()
+ observePostedMessage = (data) => {
+ const envelope = data as {
+ collectionId?: string
+ payload?: { type?: string }
+ }
+ if (
+ envelope.collectionId === collectionId &&
+ envelope.payload?.type === `leader:heartbeat`
+ ) {
+ ready.resolve()
+ }
+ }
+ const release = coordinator.subscribe(collectionId, () => {})
+
+ try {
+ await ready.promise
+ observePostedMessage = undefined
+ for (let index = 0; index < mutationCount; index++) {
+ await coordinator.requestApplyLocalMutations(collectionId, [
+ {
+ mutationId: `mutation-${run}-${index}`,
+ type: `insert`,
+ key: `${index}`,
+ value: { id: `${index}`, title: `row ${index}` },
+ },
+ ])
+ }
+ expect(
+ inspectCoordinator(coordinator).appliedEnvelopeIds.size,
+ ).toBe(mutationCount)
+
+ if (terminal === `retention`) {
+ await vi.advanceTimersByTimeAsync(60_000)
+ } else if (terminal === `unsubscribe`) {
+ release()
+ } else {
+ coordinator.dispose()
+ }
+
+ expect(
+ inspectCoordinator(coordinator).appliedEnvelopeIds.size,
+ ).toBe(0)
+ } finally {
+ observePostedMessage = undefined
+ release()
+ coordinator.dispose()
+ vi.clearAllTimers()
+ }
+ },
+ ),
+ { seed: 1868, numRuns: 12, endOnFailure: true },
+ )
+ } finally {
+ vi.useRealTimers()
+ }
+ })
})
describe(`dispose`, () => {
@@ -561,5 +1514,192 @@ describe(`BrowserCollectionCoordinator`, () => {
// Should not throw after disposal
expect(coord.isLeader(`todos`)).toBe(false)
})
+
+ it(`settles an in-flight RPC at dispose without posting retries`, async () => {
+ vi.useFakeTimers()
+ const adapter = createStubAdapter()
+ let leadershipRead!: () => void
+ const leadershipReadPromise = new Promise((resolve) => {
+ leadershipRead = resolve
+ })
+ let firstRequestPosted!: () => void
+ const firstRequestPostedPromise = new Promise((resolve) => {
+ firstRequestPosted = resolve
+ })
+ adapter.getStreamPosition = async () => {
+ leadershipRead()
+ return { latestTerm: 0, latestSeq: 0, latestRowVersion: 0 }
+ }
+
+ const leader = createCoordinator(adapter)
+ const follower = createCoordinator(adapter)
+ let requestPosts = 0
+ let settled = false
+
+ try {
+ leader.subscribe(`todos`, () => {})
+ await leadershipReadPromise
+ await Promise.resolve()
+ follower.subscribe(`todos`, () => {})
+
+ observePostedMessage = (data) => {
+ const type = (data as { payload?: { type?: string } }).payload?.type
+ if (type === `rpc:ensurePersistedIndex:req`) {
+ requestPosts++
+ if (requestPosts === 1) firstRequestPosted()
+ }
+ }
+ dropNextMessageWhen = (data) =>
+ (data as { payload?: { type?: string } }).payload?.type ===
+ `rpc:ensurePersistedIndex:req`
+
+ const request = follower
+ .requestEnsurePersistedIndex(`todos`, `idx-dispose`, {
+ expressionSql: [`title`],
+ })
+ .then(
+ () => {
+ settled = true
+ },
+ () => {
+ settled = true
+ },
+ )
+
+ await firstRequestPostedPromise
+ follower.dispose()
+ for (let microtask = 0; microtask < 8; microtask++) {
+ await Promise.resolve()
+ }
+ const settledAtDispose = settled
+
+ await vi.advanceTimersByTimeAsync(21_000)
+ await request
+
+ expect({ settledAtDispose, requestPosts }).toEqual({
+ settledAtDispose: true,
+ requestPosts: 1,
+ })
+ } finally {
+ follower.dispose()
+ leader.dispose()
+ vi.useRealTimers()
+ }
+ })
+
+ it(`makes disposal terminal across generated RPC kinds and lifecycle phases`, async () => {
+ vi.useFakeTimers()
+ let run = 0
+ try {
+ await fc.assert(
+ fc.asyncProperty(
+ fc.record({
+ rpcKind: fc.constantFrom(`index`, `pull`, `mutation`),
+ disposePhase: fc.constantFrom(
+ `before-request`,
+ `pending`,
+ `retry-delay`,
+ ),
+ }),
+ async ({ rpcKind, disposePhase }) => {
+ run++
+ const collectionId = `dispose-history-${run}`
+ const adapter = createStubAdapter()
+ let leadershipRead!: () => void
+ const leadershipReadPromise = new Promise((resolve) => {
+ leadershipRead = resolve
+ })
+ let firstRequestPosted!: () => void
+ const firstRequestPostedPromise = new Promise((resolve) => {
+ firstRequestPosted = resolve
+ })
+ adapter.getStreamPosition = async () => {
+ leadershipRead()
+ return { latestTerm: 0, latestSeq: 0, latestRowVersion: 0 }
+ }
+
+ const leader = createCoordinator(adapter)
+ const follower = createCoordinator(adapter)
+ let requestPosts = 0
+ const requestType =
+ rpcKind === `index`
+ ? `rpc:ensurePersistedIndex:req`
+ : rpcKind === `pull`
+ ? `rpc:pullSince:req`
+ : `rpc:applyLocalMutations:req`
+
+ try {
+ leader.subscribe(collectionId, () => {})
+ await leadershipReadPromise
+ await Promise.resolve()
+ follower.subscribe(collectionId, () => {})
+
+ observePostedMessage = (data) => {
+ if (
+ (data as { payload?: { type?: string } }).payload?.type ===
+ requestType
+ ) {
+ requestPosts++
+ if (requestPosts === 1) firstRequestPosted()
+ }
+ }
+ dropNextMessageWhen = (data) =>
+ (data as { payload?: { type?: string } }).payload?.type ===
+ requestType
+
+ if (disposePhase === `before-request`) {
+ follower.dispose()
+ }
+ const request =
+ rpcKind === `index`
+ ? follower.requestEnsurePersistedIndex(
+ collectionId,
+ `idx-${run}`,
+ { expressionSql: [`title`] },
+ )
+ : rpcKind === `pull`
+ ? follower.pullSince(collectionId, 0)
+ : follower.requestApplyLocalMutations(collectionId, [
+ {
+ mutationId: `mut-${run}`,
+ type: `insert`,
+ key: `${run}`,
+ value: { id: `${run}`, title: `row ${run}` },
+ },
+ ])
+ const settled = request.then(
+ () => `resolved` as const,
+ () => `rejected` as const,
+ )
+
+ if (disposePhase !== `before-request`) {
+ await firstRequestPostedPromise
+ if (disposePhase === `retry-delay`) {
+ await vi.advanceTimersByTimeAsync(10_000)
+ }
+ follower.dispose()
+ }
+
+ await expect(settled).resolves.toBe(`rejected`)
+ expect(requestPosts).toBe(
+ disposePhase === `before-request` ? 0 : 1,
+ )
+ } finally {
+ observePostedMessage = undefined
+ dropNextMessageWhen = undefined
+ follower.dispose()
+ leader.dispose()
+ for (let microtask = 0; microtask < 4; microtask++) {
+ await Promise.resolve()
+ }
+ }
+ },
+ ),
+ { seed: 1868, numRuns: 18, endOnFailure: true },
+ )
+ } finally {
+ vi.useRealTimers()
+ }
+ })
})
})
diff --git a/packages/browser-db-sqlite-persistence/tests/shared-driver-fairness-oracle.test.ts b/packages/browser-db-sqlite-persistence/tests/shared-driver-fairness-oracle.test.ts
new file mode 100644
index 0000000000..297c59e576
--- /dev/null
+++ b/packages/browser-db-sqlite-persistence/tests/shared-driver-fairness-oracle.test.ts
@@ -0,0 +1,555 @@
+/**
+ * Node campaign for the shared-driver K=1 contract documented in
+ * `shared-driver-fairness-oracle.ts`. Fixed histories prove neutral reach and a
+ * persist storm; generated legal histories vary both lane sizes, mutation
+ * width, and tail order. Public rows and logical completion checkpoints are
+ * checked independently of production scheduling, and an executable
+ * named persist-first FIFO wrong answer proves the checker rejects the original
+ * fault. Grammar controls reconstruct that witness, exercise the bounded
+ * marginals, reject a nearby invalid storm, and ablate tail-order variation.
+ * The identical generated property runs in retained fixed-seed and seedless
+ * random lanes. Supplying both TANSTACK_DB_DRIVER_FAIRNESS_SEED and
+ * TANSTACK_DB_DRIVER_FAIRNESS_PATH replaces those lanes with one checked
+ * seed+shrink-path replay. An axis-removal calibration proves that removing
+ * tail-order variation loses hydrate-before-later-persist histories.
+ */
+import { mkdtempSync, rmSync } from 'node:fs'
+import { tmpdir } from 'node:os'
+import { join } from 'node:path'
+import fc from 'fast-check'
+import { describe, expect, it } from 'vitest'
+import {
+ SHARED_DRIVER_FAIRNESS_BOUND,
+ createPersistFirstFaultObservation,
+ findSharedDriverFairnessViolation,
+ observeSharedDriverFairness,
+} from './shared-driver-fairness-oracle'
+import { createWASQLiteTestDatabase } from './helpers/wa-sqlite-test-db'
+import type {
+ SharedDriverFairnessObservation,
+ SharedDriverFairnessOptions,
+ SharedDriverFairnessScenario,
+ SharedDriverFairnessWork,
+} from './shared-driver-fairness-oracle'
+
+const FIXED_SEED = 165_905
+const GENERATED_RUNS = 12
+const HYDRATE_COUNT_RANGE = { min: 2, max: 7 } as const
+const PERSIST_COUNT_RANGE = { min: 3, max: 7 } as const
+const MUTATIONS_PER_PERSIST_RANGE = { min: 1, max: 3 } as const
+
+type FairnessReplayEnvironment = Record
+
+type FairnessReplayConfig = {
+ seed: number
+ path: string
+}
+
+type FairnessPropertyMode = {
+ label: `fixed-seed` | `seedless-random` | `checked-replay`
+ seed?: number
+ path?: string
+}
+
+function readFairnessReplayConfig(
+ environment: FairnessReplayEnvironment = process.env,
+): FairnessReplayConfig | undefined {
+ const seedText = environment.TANSTACK_DB_DRIVER_FAIRNESS_SEED
+ const path = environment.TANSTACK_DB_DRIVER_FAIRNESS_PATH
+ if (seedText === undefined && path === undefined) return undefined
+ if (seedText === undefined || path === undefined) {
+ throw new Error(
+ `TANSTACK_DB_DRIVER_FAIRNESS_SEED and TANSTACK_DB_DRIVER_FAIRNESS_PATH must be supplied together`,
+ )
+ }
+
+ const seed = Number(seedText)
+ if (seedText.trim() === `` || !Number.isSafeInteger(seed)) {
+ throw new Error(`TANSTACK_DB_DRIVER_FAIRNESS_SEED must be an integer`)
+ }
+ if (!/^\d+(?::\d+)*$/.test(path)) {
+ throw new Error(
+ `TANSTACK_DB_DRIVER_FAIRNESS_PATH must contain colon-separated nonnegative integers`,
+ )
+ }
+ return { seed, path }
+}
+
+function createFairnessPropertyModes(
+ replay: FairnessReplayConfig | undefined,
+): ReadonlyArray {
+ if (replay) {
+ return [{ label: `checked-replay`, ...replay }]
+ }
+ return [
+ { label: `fixed-seed`, seed: FIXED_SEED },
+ { label: `seedless-random` },
+ ]
+}
+
+function readFairnessCalibration(
+ environment: FairnessReplayEnvironment = process.env,
+): SharedDriverFairnessOptions {
+ const fault = environment.TANSTACK_DB_DRIVER_FAIRNESS_CALIBRATION
+ if (fault === undefined) return {}
+ if (fault !== `persist-first-fifo`) {
+ throw new Error(
+ `TANSTACK_DB_DRIVER_FAIRNESS_CALIBRATION must be persist-first-fifo`,
+ )
+ }
+ return { schedulingFault: fault }
+}
+
+const fairnessPropertyModes = createFairnessPropertyModes(
+ readFairnessReplayConfig(),
+)
+const fairnessCalibration = readFairnessCalibration()
+
+type WorkKind = SharedDriverFairnessWork[`kind`]
+
+type GeneratedFairnessHistory = {
+ hydrateCount: number
+ persistCount: number
+ mutationsPerPersist: number
+ orderedKinds: ReadonlyArray
+ replayHistory: string
+}
+
+type GeneratedTailEntry = {
+ kind: WorkKind
+ token: string
+}
+
+function createScenario(
+ id: string,
+ orderedKinds: ReadonlyArray,
+ mutationsPerPersist = 1,
+): SharedDriverFairnessScenario {
+ let hydrateIndex = 0
+ let persistIndex = 0
+ return {
+ id,
+ work: orderedKinds.map((kind) => {
+ if (kind === `persist`) {
+ const index = persistIndex++
+ return {
+ kind,
+ id: `persist-${index}`,
+ mutationsPerPersist,
+ }
+ }
+ const index = hydrateIndex++
+ return {
+ kind,
+ id: `hydrate-${index}`,
+ seededRows: [
+ { id: `row-${index}-0`, value: index * 10 },
+ { id: `row-${index}-1`, value: index * 10 + 1 },
+ ],
+ }
+ }),
+ }
+}
+
+function expectedHydratedCollections(scenario: SharedDriverFairnessScenario) {
+ return scenario.work.flatMap((work) =>
+ work.kind === `hydrate`
+ ? [
+ {
+ collectionId: `${scenario.id}-${work.id}`,
+ rows: work.seededRows.map((row) => ({ ...row })),
+ },
+ ]
+ : [],
+ )
+}
+
+async function withNodeScenario(
+ scenario: SharedDriverFairnessScenario,
+ assertion: (
+ observation: SharedDriverFairnessObservation,
+ ) => void | Promise,
+ options: SharedDriverFairnessOptions = {},
+): Promise {
+ const directory = mkdtempSync(join(tmpdir(), `db-driver-fairness-`))
+ let primaryFailure: unknown
+ try {
+ const observation = await observeSharedDriverFairness(
+ () =>
+ createWASQLiteTestDatabase({
+ filename: join(directory, `state.sqlite`),
+ }),
+ scenario,
+ options,
+ )
+ await assertion(observation)
+ } catch (error) {
+ primaryFailure = error
+ }
+
+ let nodeCleanupFailure: unknown
+ try {
+ rmSync(directory, { recursive: true, force: true })
+ } catch (error) {
+ nodeCleanupFailure = error
+ }
+
+ if (primaryFailure !== undefined) {
+ if (nodeCleanupFailure !== undefined) {
+ const primaryMessage =
+ primaryFailure instanceof Error
+ ? primaryFailure.message
+ : String(primaryFailure)
+ const cleanupMessage =
+ nodeCleanupFailure instanceof Error
+ ? nodeCleanupFailure.message
+ : String(nodeCleanupFailure)
+ throw new Error(
+ `${primaryMessage}; node cleanup diagnostics: ${cleanupMessage}`,
+ )
+ }
+ throw primaryFailure
+ }
+ if (nodeCleanupFailure !== undefined) throw nodeCleanupFailure
+}
+
+function expectObservationReach(
+ observation: SharedDriverFairnessObservation,
+): void {
+ const expectedHydrates = expectedHydratedCollections(observation.scenario)
+ expect(observation.admittedHydrateIds).toEqual(
+ expectedHydrates.map(({ collectionId }) => collectionId),
+ )
+ expect(observation.hydrationCompletions).toHaveLength(expectedHydrates.length)
+ // Actual rows come from the public Collection after preload. The expected
+ // rows above are derived independently from the generated seed history.
+ expect(observation.hydratedCollections).toEqual(expectedHydrates)
+}
+
+function expectFairObservation(
+ observation: SharedDriverFairnessObservation,
+): void {
+ expectObservationReach(observation)
+ const violation = findSharedDriverFairnessViolation(
+ observation,
+ SHARED_DRIVER_FAIRNESS_BOUND,
+ )
+ if (violation) {
+ throw new Error(
+ `shared-driver fairness mismatch at ${violation.checkpoint.collectionId}: ` +
+ `${violation.checkpoint.completedPersistIds.length} persists completed ` +
+ `(maximum ${violation.expectedMaximumCompletedPersists}), ` +
+ `${violation.checkpoint.pendingPersistCount} remained pending ` +
+ `(minimum ${violation.expectedMinimumPendingPersists}); ` +
+ `permitted completed ids: ${JSON.stringify(violation.permittedCompletedPersistIds)}; ` +
+ `cleanup diagnostics: ${JSON.stringify(observation.cleanupFailures)}`,
+ )
+ }
+ expect(observation.cleanupFailures).toEqual([])
+}
+
+function buildGeneratedHistory(
+ hydrateCount: number,
+ persistCount: number,
+ mutationsPerPersist: number,
+ generatedTail: ReadonlyArray,
+): GeneratedFairnessHistory {
+ return {
+ hydrateCount,
+ persistCount,
+ mutationsPerPersist,
+ orderedKinds: [
+ `persist` as const,
+ ...generatedTail.map(({ kind }) => kind),
+ ],
+ replayHistory: `p0,${generatedTail.map(({ token }) => token).join(`,`)}`,
+ }
+}
+
+function createGeneratedHistoryArbitrary(
+ tailOrderAxis: `varied` | `removed`,
+): fc.Arbitrary {
+ return fc
+ .record({
+ hydrateCount: fc.integer(HYDRATE_COUNT_RANGE),
+ persistCount: fc.integer(PERSIST_COUNT_RANGE),
+ mutationsPerPersist: fc.integer(MUTATIONS_PER_PERSIST_RANGE),
+ })
+ .chain(({ hydrateCount, persistCount, mutationsPerPersist }) => {
+ const tail = [
+ ...Array.from({ length: persistCount - 1 }, (_, index) => ({
+ kind: `persist` as const,
+ token: `p${index + 1}`,
+ })),
+ ...Array.from({ length: hydrateCount }, (_, index) => ({
+ kind: `hydrate` as const,
+ token: `h${index}`,
+ })),
+ ]
+ const orderedTail =
+ tailOrderAxis === `varied`
+ ? fc.shuffledSubarray(tail, {
+ minLength: tail.length,
+ maxLength: tail.length,
+ })
+ : fc.constant(tail)
+ return orderedTail.map((generatedTail) =>
+ buildGeneratedHistory(
+ hydrateCount,
+ persistCount,
+ mutationsPerPersist,
+ generatedTail,
+ ),
+ )
+ })
+}
+
+function reachesHydrateBeforeLaterPersist(
+ history: GeneratedFairnessHistory,
+): boolean {
+ let hydrateSeen = false
+ for (const kind of history.orderedKinds.slice(1)) {
+ if (kind === `hydrate`) hydrateSeen = true
+ if (kind === `persist` && hydrateSeen) return true
+ }
+ return false
+}
+
+const generatedHistoryArbitrary = createGeneratedHistoryArbitrary(`varied`)
+const tailOrderAblatedArbitrary = createGeneratedHistoryArbitrary(`removed`)
+let generatedPropertyExecutions = 0
+
+const generatedFairnessProperty = fc.asyncProperty(
+ generatedHistoryArbitrary,
+ async ({
+ hydrateCount,
+ persistCount,
+ mutationsPerPersist,
+ orderedKinds,
+ replayHistory,
+ }) => {
+ generatedPropertyExecutions++
+ const observationId =
+ `generated-h${hydrateCount}-p${persistCount}-m${mutationsPerPersist}-` +
+ replayHistory.replaceAll(`,`, `-`)
+ await withNodeScenario(
+ createScenario(observationId, orderedKinds, mutationsPerPersist),
+ expectFairObservation,
+ fairnessCalibration,
+ )
+ },
+)
+
+describe(`shared BrowserWASQLiteDriver fairness oracle`, () => {
+ it(`selects fixed and seedless lanes by default and only replay when requested`, () => {
+ expect(createFairnessPropertyModes(readFairnessReplayConfig({}))).toEqual([
+ { label: `fixed-seed`, seed: FIXED_SEED },
+ { label: `seedless-random` },
+ ])
+ expect(
+ createFairnessPropertyModes(
+ readFairnessReplayConfig({
+ TANSTACK_DB_DRIVER_FAIRNESS_SEED: `42`,
+ TANSTACK_DB_DRIVER_FAIRNESS_PATH: `0:1`,
+ }),
+ ),
+ ).toEqual([{ label: `checked-replay`, seed: 42, path: `0:1` }])
+ expect(readFairnessCalibration({})).toEqual({})
+ expect(
+ readFairnessCalibration({
+ TANSTACK_DB_DRIVER_FAIRNESS_CALIBRATION: `persist-first-fifo`,
+ }),
+ ).toEqual({ schedulingFault: `persist-first-fifo` })
+ })
+
+ it.each([
+ {
+ name: `seed without path`,
+ environment: { TANSTACK_DB_DRIVER_FAIRNESS_SEED: `42` },
+ message: `must be supplied together`,
+ },
+ {
+ name: `path without seed`,
+ environment: { TANSTACK_DB_DRIVER_FAIRNESS_PATH: `0` },
+ message: `must be supplied together`,
+ },
+ {
+ name: `non-integer seed`,
+ environment: {
+ TANSTACK_DB_DRIVER_FAIRNESS_SEED: `4.2`,
+ TANSTACK_DB_DRIVER_FAIRNESS_PATH: `0`,
+ },
+ message: `must be an integer`,
+ },
+ {
+ name: `invalid shrink path`,
+ environment: {
+ TANSTACK_DB_DRIVER_FAIRNESS_SEED: `42`,
+ TANSTACK_DB_DRIVER_FAIRNESS_PATH: `0:-1`,
+ },
+ message: `colon-separated nonnegative integers`,
+ },
+ ])(`rejects incomplete or invalid replay coordinates: $name`, (probe) => {
+ expect(() => readFairnessReplayConfig(probe.environment)).toThrow(
+ probe.message,
+ )
+ })
+
+ it(`reconstructs the known persist-first witness, covers range marginals, and rejects a nearby invalid storm`, async () => {
+ const knownWitness = buildGeneratedHistory(2, 4, 1, [
+ { kind: `persist`, token: `p1` },
+ { kind: `persist`, token: `p2` },
+ { kind: `persist`, token: `p3` },
+ { kind: `hydrate`, token: `h0` },
+ { kind: `hydrate`, token: `h1` },
+ ])
+ expect(knownWitness).toEqual({
+ hydrateCount: 2,
+ persistCount: 4,
+ mutationsPerPersist: 1,
+ orderedKinds: [
+ `persist`,
+ `persist`,
+ `persist`,
+ `persist`,
+ `hydrate`,
+ `hydrate`,
+ ],
+ replayHistory: `p0,p1,p2,p3,h0,h1`,
+ })
+
+ const marginalHistories = [
+ buildGeneratedHistory(2, 3, 1, [
+ { kind: `persist`, token: `p1` },
+ { kind: `persist`, token: `p2` },
+ { kind: `hydrate`, token: `h0` },
+ { kind: `hydrate`, token: `h1` },
+ ]),
+ buildGeneratedHistory(7, 7, 3, [
+ ...Array.from({ length: 6 }, (_, index) => ({
+ kind: `persist` as const,
+ token: `p${index + 1}`,
+ })),
+ ...Array.from({ length: 7 }, (_, index) => ({
+ kind: `hydrate` as const,
+ token: `h${index}`,
+ })),
+ ]),
+ ]
+ expect(
+ marginalHistories.map((history) => [
+ history.hydrateCount,
+ history.persistCount,
+ history.mutationsPerPersist,
+ ]),
+ ).toEqual([
+ [2, 3, 1],
+ [7, 7, 3],
+ ])
+
+ await expect(
+ withNodeScenario(
+ createScenario(`invalid-storm-boundary`, [
+ `hydrate`,
+ `persist`,
+ `hydrate`,
+ ]),
+ expectFairObservation,
+ ),
+ ).rejects.toThrow(`must begin with its already-running persist`)
+ })
+
+ it(`reaches cold hydration and preserves independently seeded public rows without queued persists`, async () => {
+ await withNodeScenario(
+ createScenario(`neutral-reach`, [`hydrate`, `hydrate`, `hydrate`]),
+ (observation) => {
+ expectFairObservation(observation)
+ expect(
+ observation.rawDequeues.some((entry) =>
+ entry.sql.startsWith(
+ `SELECT key, value, metadata, row_version FROM`,
+ ),
+ ),
+ ).toBe(true)
+ },
+ )
+ })
+
+ it(`completes a pending cold hydrate before an unrelated persist backlog drains`, async () => {
+ await withNodeScenario(
+ createScenario(
+ `fixed-persist-storm`,
+ [`persist`, `persist`, `persist`, `hydrate`, `hydrate`],
+ 2,
+ ),
+ expectFairObservation,
+ )
+ })
+
+ it.each(fairnessPropertyModes)(
+ `bounds persist completions for generated ordered cold-hydrate/persist histories in $label mode`,
+ async (mode) => {
+ const executionsBefore = generatedPropertyExecutions
+ await fc.assert(generatedFairnessProperty, {
+ numRuns: GENERATED_RUNS,
+ verbose: 2,
+ ...(mode.seed === undefined ? {} : { seed: mode.seed }),
+ ...(mode.path === undefined ? {} : { path: mode.path }),
+ })
+ expect(generatedPropertyExecutions).toBeGreaterThan(executionsBefore)
+ },
+ )
+
+ it(`proves tail-order ablation removes a hydrate-before-later-persist history`, () => {
+ const sampleOptions = { seed: FIXED_SEED, numRuns: GENERATED_RUNS }
+ const completeGrammar = fc.sample(generatedHistoryArbitrary, sampleOptions)
+ const tailOrderAblatedGrammar = fc.sample(
+ tailOrderAblatedArbitrary,
+ sampleOptions,
+ )
+
+ expect(completeGrammar.some(reachesHydrateBeforeLaterPersist)).toBe(true)
+ expect(tailOrderAblatedGrammar.some(reachesHydrateBeforeLaterPersist)).toBe(
+ false,
+ )
+ })
+
+ it(`rejects the named persist-first FIFO wrong answer through the real package path`, async () => {
+ await withNodeScenario(
+ createScenario(
+ `persist-first-fixture-fault`,
+ [`persist`, `persist`, `persist`, `persist`, `hydrate`, `hydrate`],
+ 1,
+ ),
+ (observation) => {
+ expectObservationReach(observation)
+ const violation = findSharedDriverFairnessViolation(observation)
+ if (!violation) {
+ throw new Error(
+ `persist-first scheduling mutant escaped the fairness checker; ` +
+ `cleanup diagnostics: ${JSON.stringify(observation.cleanupFailures)}`,
+ )
+ }
+ expect(violation.checkpoint.completedPersistIds).toHaveLength(4)
+ expect(violation.checkpoint.pendingPersistCount).toBe(0)
+ expect(observation.cleanupFailures).toEqual([])
+ },
+ { schedulingFault: `persist-first-fifo` },
+ )
+ })
+
+ it(`calibrates the checker against a synthetic persist-first observation`, () => {
+ const scenario = createScenario(
+ `persist-first-checker-calibration`,
+ [`persist`, `persist`, `persist`, `persist`, `hydrate`, `hydrate`],
+ 1,
+ )
+ const fault = createPersistFirstFaultObservation(scenario)
+
+ expect(findSharedDriverFairnessViolation(fault)).toMatchObject({
+ checkpoint: fault.hydrationCompletions[0],
+ expectedMaximumCompletedPersists: 1,
+ expectedMinimumPendingPersists: 3,
+ })
+ })
+})
diff --git a/packages/browser-db-sqlite-persistence/tests/shared-driver-fairness-oracle.ts b/packages/browser-db-sqlite-persistence/tests/shared-driver-fairness-oracle.ts
new file mode 100644
index 0000000000..01666bb94d
--- /dev/null
+++ b/packages/browser-db-sqlite-persistence/tests/shared-driver-fairness-oracle.ts
@@ -0,0 +1,753 @@
+/**
+ * # When does a cold hydrate get a turn on a shared SQLite driver?
+ *
+ * Contract and source: RFC #1659 accepts K=1 complete-logical-cold-hydrate
+ * scheduling. The persist already executing when the storm begins is
+ * non-preemptible. After that, at most one additional persist may complete
+ * between consecutive hydrate completions, with FIFO identity preserved
+ * inside the hydrate and persist lanes.
+ *
+ * History grammar and domain: a legal storm has unique work IDs, starts with
+ * that already-running persist when any persist exists, and then permutes
+ * complete hydrate and persist requests. Hydrates contain nonempty unique
+ * seeded rows; persists contain one or more mutations. The generated campaign
+ * varies 2..7 hydrates, 3..7 persists, 1..3 mutations per persist, and sampled
+ * tail permutations. The minimum grammar still has two hydrate checkpoints and
+ * a persist backlog beyond K=1. Hydrate count supplies repeated checkpoints;
+ * persist count varies the backlog; mutation width proves K counts logical
+ * persists rather than their rows; tail order creates overlapping lane work.
+ * Neutral histories contain only hydrates.
+ *
+ * Independent model and production boundary: `createFairnessReference`
+ * computes permitted completed persist IDs from the ordered history and K; it
+ * does not import or simulate the production scheduler. The driver exercises
+ * public `Collection.preload()` through persisted collection options, the core
+ * adapter, and one real `BrowserWASQLiteDriver`. Each preload completion is a
+ * checkpoint after the complete logical hydrate, not after an individual SQL
+ * statement.
+ *
+ * Observed public facts: admitted and completed logical IDs, completed and
+ * pending persists at every hydrate checkpoint, independently seeded public
+ * collection rows, raw SQL dequeue reach, and cleanup diagnostics. Known
+ * omissions: the oracle does not establish elapsed-time latency, unbounded
+ * eventuality, multi-process coordination, or a browser matrix; the Chromium
+ * OPFS fixture separately refines the provider boundary.
+ *
+ * Challenge and replay: the named persist-first FIFO wrong answer must violate
+ * the same K=1 checker while using the public/core/driver path. The generated
+ * property has fixed-seed and seedless-random lanes. Supplying both
+ * TANSTACK_DB_DRIVER_FAIRNESS_SEED and TANSTACK_DB_DRIVER_FAIRNESS_PATH runs
+ * only their checked oracle replay. Tail-order ablation must lose the legal
+ * hydrate-before-later-persist cell. The grammar also reconstructs that known
+ * wrong-answer witness, covers its bounded marginals, and rejects a storm that
+ * omits the already-running persist boundary. Cleanup preserves the primary
+ * failure and reports secondary resource-release diagnostics separately.
+ */
+import { createCollection } from '../../db/src/index'
+import { persistedCollectionOptions } from '../src/index'
+import { BrowserWASQLiteDriver } from '../src/wa-sqlite-driver'
+import {
+ SingleProcessCoordinator,
+ createSQLiteCorePersistenceAdapter,
+} from '../../db-sqlite-persistence-core/src/index'
+import type { Collection } from '../../db/src/index'
+import type {
+ PersistedCollectionPersistence,
+ PersistedTx,
+ SQLiteDriver,
+} from '../../db-sqlite-persistence-core/src/index'
+import type { BrowserWASQLiteDatabase } from '../src/index'
+
+export const SHARED_DRIVER_FAIRNESS_BOUND = 1
+
+export type FairnessRow = {
+ id: string
+ value: number
+}
+
+export type SharedDriverFairnessWork =
+ | {
+ kind: `hydrate`
+ id: string
+ seededRows: ReadonlyArray
+ }
+ | {
+ kind: `persist`
+ id: string
+ mutationsPerPersist: number
+ }
+
+export type SharedDriverFairnessScenario = {
+ id: string
+ /**
+ * Logical admission history. A storm begins with the one persist that is
+ * already non-preemptibly running; the remaining hydrate and persist work
+ * may be interleaved in any order.
+ */
+ work: ReadonlyArray
+}
+
+export type RawDriverDequeue = {
+ ordinal: number
+ sql: string
+ params: ReadonlyArray
+}
+
+export type RawDriverAdmission = RawDriverDequeue & {
+ kind: `exec` | `query` | `run` | `transaction`
+}
+
+export type HydrationCompletionCheckpoint = {
+ collectionId: string
+ completionOrdinal: number
+ completedPersistIds: ReadonlyArray
+ pendingPersistCount: number
+ rawDequeueCount: number
+}
+
+export type HydratedCollectionRows = {
+ collectionId: string
+ rows: ReadonlyArray
+}
+
+export type SharedDriverFairnessObservation = {
+ scenario: SharedDriverFairnessScenario
+ admittedHydrateIds: ReadonlyArray
+ logicalCompletionOrder: ReadonlyArray
+ hydrationCompletions: ReadonlyArray
+ driverAdmissions: ReadonlyArray
+ rawDequeues: ReadonlyArray
+ hydratedCollections: ReadonlyArray
+ cleanupFailures: ReadonlyArray
+}
+
+export type SharedDriverFairnessViolation = {
+ checkpoint: HydrationCompletionCheckpoint
+ expectedCollectionId: string
+ expectedMaximumCompletedPersists: number
+ expectedMinimumPendingPersists: number
+ permittedCompletedPersistIds: ReadonlyArray
+ unexpectedCompletedPersistIds: ReadonlyArray
+}
+
+export type SharedDriverFairnessOptions = {
+ /** Test-only hostile control that preserves the current global FIFO fault. */
+ schedulingFault?: `persist-first-fifo`
+}
+
+type Deferred = {
+ promise: Promise
+ resolve: () => void
+}
+
+type CloseableSQLiteDriver = SQLiteDriver & {
+ transactionWithDriver: (
+ fn: (transactionDriver: SQLiteDriver) => Promise,
+ ) => Promise
+ close: () => Promise
+}
+
+type FairnessReferenceCheckpoint = {
+ collectionId: string
+ expectedMaximumCompletedPersists: number
+ expectedMinimumPendingPersists: number
+ permittedCompletedPersistIds: ReadonlyArray
+}
+
+export type BrowserWASQLiteDatabaseFactory = () =>
+ | BrowserWASQLiteDatabase
+ | Promise
+
+function createDeferred(): Deferred {
+ let resolve!: () => void
+ const promise = new Promise((settle) => {
+ resolve = settle
+ })
+ return { promise, resolve }
+}
+
+function normalizeSql(sql: string): string {
+ return sql.replace(/\s+/g, ` `).trim()
+}
+
+function failureMessage(error: unknown): string {
+ return error instanceof Error ? error.message : String(error)
+}
+
+function collectionIdFor(
+ scenario: SharedDriverFairnessScenario,
+ work: SharedDriverFairnessWork,
+): string {
+ return `${scenario.id}-${work.id}`
+}
+
+function validateScenario(scenario: SharedDriverFairnessScenario): void {
+ const hydrateWork = scenario.work.filter((work) => work.kind === `hydrate`)
+ const persistWork = scenario.work.filter((work) => work.kind === `persist`)
+ if (hydrateWork.length < 1) {
+ throw new Error(`work must contain at least one hydrate`)
+ }
+ if (persistWork.length > 0 && scenario.work[0]?.kind !== `persist`) {
+ throw new Error(
+ `a storm history must begin with its already-running persist`,
+ )
+ }
+ if (
+ new Set(scenario.work.map((work) => work.id)).size !== scenario.work.length
+ ) {
+ throw new Error(`work ids must be unique`)
+ }
+ for (const work of scenario.work) {
+ if (work.kind === `persist` && work.mutationsPerPersist < 1) {
+ throw new Error(`mutationsPerPersist must be at least one`)
+ }
+ if (work.kind === `hydrate`) {
+ if (work.seededRows.length < 1) {
+ throw new Error(`each hydrate must contain at least one seeded row`)
+ }
+ if (
+ new Set(work.seededRows.map((row) => row.id)).size !==
+ work.seededRows.length
+ ) {
+ throw new Error(`seeded row ids must be unique within each hydrate`)
+ }
+ }
+ }
+}
+
+class ObservedDatabase implements BrowserWASQLiteDatabase {
+ readonly rawDequeues: Array = []
+ private heldBegin: Deferred | undefined
+ private beginEntered: Deferred | undefined
+
+ constructor(private readonly database: BrowserWASQLiteDatabase) {}
+
+ holdNextTransactionBegin(): { entered: Promise; release: () => void } {
+ if (this.heldBegin) {
+ throw new Error(`a transaction begin is already held`)
+ }
+ this.heldBegin = createDeferred()
+ this.beginEntered = createDeferred()
+ return {
+ entered: this.beginEntered.promise,
+ release: () => this.releaseHeldBegin(),
+ }
+ }
+
+ clearTrace(): void {
+ this.rawDequeues.length = 0
+ }
+
+ async execute(
+ sql: string,
+ params: ReadonlyArray = [],
+ ): Promise> {
+ const normalizedSql = normalizeSql(sql)
+ this.rawDequeues.push({
+ ordinal: this.rawDequeues.length,
+ sql: normalizedSql,
+ params: [...params],
+ })
+
+ if (normalizedSql === `BEGIN IMMEDIATE` && this.heldBegin) {
+ const heldBegin = this.heldBegin
+ this.beginEntered?.resolve()
+ await heldBegin.promise
+ }
+
+ return this.database.execute(sql, params)
+ }
+
+ async close(): Promise {
+ this.releaseHeldBegin()
+ await this.database.close?.()
+ }
+
+ private releaseHeldBegin(): void {
+ this.heldBegin?.resolve()
+ this.heldBegin = undefined
+ this.beginEntered = undefined
+ }
+}
+
+/**
+ * A test-only hostile scheduler. It admits calls through the same public/core/
+ * BrowserWASQLiteDriver path but serializes them in global admission order, so
+ * a future fair production driver cannot accidentally make this control pass.
+ */
+class PersistFirstFIFOFaultDriver implements SQLiteDriver {
+ private tail = Promise.resolve()
+
+ constructor(private readonly driver: CloseableSQLiteDriver) {}
+
+ exec(sql: string): Promise {
+ return this.enqueue(() => this.driver.exec(sql))
+ }
+
+ query(
+ sql: string,
+ params: ReadonlyArray = [],
+ ): Promise> {
+ return this.enqueue(() => this.driver.query(sql, params))
+ }
+
+ run(sql: string, params: ReadonlyArray = []): Promise {
+ return this.enqueue(() => this.driver.run(sql, params))
+ }
+
+ transaction(
+ fn: (transactionDriver: SQLiteDriver) => Promise,
+ ): Promise {
+ return this.enqueue(() => this.driver.transaction(fn))
+ }
+
+ transactionWithDriver(
+ fn: (transactionDriver: SQLiteDriver) => Promise,
+ ): Promise {
+ return this.enqueue(() => this.driver.transactionWithDriver(fn))
+ }
+
+ close(): Promise {
+ return this.enqueue(() => this.driver.close())
+ }
+
+ private enqueue(operation: () => Promise): Promise {
+ const result = this.tail.then(operation)
+ this.tail = result.then(
+ () => undefined,
+ () => undefined,
+ )
+ return result
+ }
+}
+
+class AdmissionObservedDriver implements SQLiteDriver {
+ readonly admissions: Array = []
+
+ constructor(
+ private readonly driver: CloseableSQLiteDriver,
+ private readonly onAdmission?: (entry: RawDriverAdmission) => void,
+ ) {}
+
+ exec(sql: string): Promise {
+ this.record(`exec`, sql, [])
+ return this.driver.exec(sql)
+ }
+
+ query(
+ sql: string,
+ params: ReadonlyArray = [],
+ ): Promise> {
+ this.record(`query`, sql, params)
+ return this.driver.query(sql, params)
+ }
+
+ run(sql: string, params: ReadonlyArray = []): Promise {
+ this.record(`run`, sql, params)
+ return this.driver.run(sql, params)
+ }
+
+ transaction(
+ fn: (transactionDriver: SQLiteDriver) => Promise,
+ ): Promise {
+ this.record(`transaction`, `BEGIN IMMEDIATE`, [])
+ return this.driver.transaction(fn)
+ }
+
+ transactionWithDriver(
+ fn: (transactionDriver: SQLiteDriver) => Promise,
+ ): Promise {
+ this.record(`transaction`, `BEGIN IMMEDIATE`, [])
+ return this.driver.transactionWithDriver(fn)
+ }
+
+ close(): Promise {
+ return this.driver.close()
+ }
+
+ clearAdmissions(): void {
+ this.admissions.length = 0
+ }
+
+ private record(
+ kind: RawDriverAdmission[`kind`],
+ sql: string,
+ params: ReadonlyArray,
+ ): void {
+ const entry = {
+ ordinal: this.admissions.length,
+ kind,
+ sql: normalizeSql(sql),
+ params: [...params],
+ }
+ this.admissions.push(entry)
+ this.onAdmission?.(entry)
+ }
+}
+
+function createSharedPersistence(
+ database: BrowserWASQLiteDatabase,
+ onAdmission?: (entry: RawDriverAdmission) => void,
+ schedulingFault?: SharedDriverFairnessOptions[`schedulingFault`],
+): {
+ driver: AdmissionObservedDriver
+ persistence: PersistedCollectionPersistence
+} {
+ const productionDriver = new BrowserWASQLiteDriver({ database })
+ const scheduledDriver =
+ schedulingFault === `persist-first-fifo`
+ ? new PersistFirstFIFOFaultDriver(productionDriver)
+ : productionDriver
+ const driver = new AdmissionObservedDriver(scheduledDriver, onAdmission)
+ const adapter = createSQLiteCorePersistenceAdapter({
+ driver,
+ schemaMismatchPolicy: `sync-absent-error`,
+ appliedTxPruneMaxRows: 0,
+ appliedTxPruneMaxAgeSeconds: 0,
+ })
+ return {
+ driver,
+ persistence: {
+ adapter,
+ coordinator: new SingleProcessCoordinator(),
+ },
+ }
+}
+
+function createPersistedTx(
+ collectionId: string,
+ sequence: number,
+ mutationsPerPersist: number,
+): PersistedTx {
+ return {
+ txId: `${collectionId}-tx-${sequence}`,
+ term: 1,
+ seq: sequence,
+ rowVersion: sequence,
+ mutations: Array.from({ length: mutationsPerPersist }, (_, index) => ({
+ type: `insert` as const,
+ key: `${sequence}-${index}`,
+ value: {
+ id: `${sequence}-${index}`,
+ value: sequence * 100 + index,
+ },
+ })),
+ }
+}
+
+function createFairnessReference(
+ scenario: SharedDriverFairnessScenario,
+ fairnessBound: number,
+): ReadonlyArray {
+ const persistIds = scenario.work
+ .filter((work) => work.kind === `persist`)
+ .map((work) => collectionIdFor(scenario, work))
+ const hydrateIds = scenario.work
+ .filter((work) => work.kind === `hydrate`)
+ .map((work) => collectionIdFor(scenario, work))
+
+ return hydrateIds.map((collectionId, hydrateIndex) => {
+ // The first persist is the one already executing when hydrate work reaches
+ // the queue. The ordered reference then permits at most K additional
+ // persists between successive hydrate completions, preserving FIFO order
+ // within each logical lane.
+ const expectedMaximumCompletedPersists = Math.min(
+ persistIds.length,
+ persistIds.length === 0 ? 0 : 1 + hydrateIndex * fairnessBound,
+ )
+ return {
+ collectionId,
+ expectedMaximumCompletedPersists,
+ expectedMinimumPendingPersists:
+ persistIds.length - expectedMaximumCompletedPersists,
+ permittedCompletedPersistIds: persistIds.slice(
+ 0,
+ expectedMaximumCompletedPersists,
+ ),
+ }
+ })
+}
+
+export function findSharedDriverFairnessViolation(
+ observation: SharedDriverFairnessObservation,
+ fairnessBound = SHARED_DRIVER_FAIRNESS_BOUND,
+): SharedDriverFairnessViolation | undefined {
+ const reference = createFairnessReference(observation.scenario, fairnessBound)
+ for (const [
+ hydrateIndex,
+ checkpoint,
+ ] of observation.hydrationCompletions.entries()) {
+ const expected = reference[hydrateIndex]
+ if (!expected) continue
+ const permitted = new Set(expected.permittedCompletedPersistIds)
+ const unexpectedCompletedPersistIds = checkpoint.completedPersistIds.filter(
+ (id) => !permitted.has(id),
+ )
+ if (
+ checkpoint.collectionId !== expected.collectionId ||
+ checkpoint.completedPersistIds.length >
+ expected.expectedMaximumCompletedPersists ||
+ checkpoint.pendingPersistCount <
+ expected.expectedMinimumPendingPersists ||
+ unexpectedCompletedPersistIds.length > 0
+ ) {
+ return {
+ checkpoint,
+ expectedCollectionId: expected.collectionId,
+ expectedMaximumCompletedPersists:
+ expected.expectedMaximumCompletedPersists,
+ expectedMinimumPendingPersists: expected.expectedMinimumPendingPersists,
+ permittedCompletedPersistIds: expected.permittedCompletedPersistIds,
+ unexpectedCompletedPersistIds,
+ }
+ }
+ }
+ return undefined
+}
+
+/** Checker calibration only; the executable hostile control uses the fixture. */
+export function createPersistFirstFaultObservation(
+ scenario: SharedDriverFairnessScenario,
+): SharedDriverFairnessObservation {
+ const persistIds = scenario.work
+ .filter((work) => work.kind === `persist`)
+ .map((work) => collectionIdFor(scenario, work))
+ const firstHydrate = scenario.work.find((work) => work.kind === `hydrate`)
+ if (!firstHydrate) throw new Error(`fault observation requires a hydrate`)
+ const firstHydrateId = collectionIdFor(scenario, firstHydrate)
+ return {
+ scenario,
+ admittedHydrateIds: [firstHydrateId],
+ logicalCompletionOrder: [
+ ...persistIds.map((id) => `persist:${id}`),
+ `hydrate:${firstHydrateId}`,
+ ],
+ hydrationCompletions: [
+ {
+ collectionId: firstHydrateId,
+ completionOrdinal: persistIds.length,
+ completedPersistIds: persistIds,
+ pendingPersistCount: 0,
+ rawDequeueCount: persistIds.length + 1,
+ },
+ ],
+ rawDequeues: [],
+ driverAdmissions: [],
+ hydratedCollections: [],
+ cleanupFailures: [],
+ }
+}
+
+/**
+ * Runs the public persisted-collection startup path over one shared
+ * BrowserWASQLiteDriver. The expected scheduler is deliberately not imported:
+ * the oracle observes only logical admissions/completions and raw SQL dequeue.
+ */
+export async function observeSharedDriverFairness(
+ openDatabase: BrowserWASQLiteDatabaseFactory,
+ scenario: SharedDriverFairnessScenario,
+ options: SharedDriverFairnessOptions = {},
+): Promise {
+ validateScenario(scenario)
+
+ const hydrateWork = scenario.work.filter((work) => work.kind === `hydrate`)
+ const persistWork = scenario.work.filter((work) => work.kind === `persist`)
+ const hydrateIds = hydrateWork.map((work) => collectionIdFor(scenario, work))
+ const persistIds = persistWork.map((work) => collectionIdFor(scenario, work))
+ const admittedHydrateIds: Array = []
+ const allHydratesRequested = createDeferred()
+
+ // A separate connection and adapter own seed/setup. Closing and reopening
+ // makes every measured hydration cold at the adapter and driver layers.
+ const seedDatabase = await openDatabase()
+ const seed = createSharedPersistence(seedDatabase)
+ let seedPrimaryFailure: unknown
+ try {
+ for (const work of hydrateWork) {
+ const collectionId = collectionIdFor(scenario, work)
+ await seed.persistence.adapter.applyCommittedTx(collectionId, {
+ txId: `${collectionId}-seed`,
+ term: 1,
+ seq: 1,
+ rowVersion: 1,
+ mutations: work.seededRows.map((row) => ({
+ type: `insert` as const,
+ key: row.id,
+ value: { ...row },
+ })),
+ })
+ }
+ for (const collectionId of persistIds) {
+ await seed.persistence.adapter.loadSubset(collectionId, {})
+ }
+ } catch (error) {
+ seedPrimaryFailure = error
+ }
+ let seedCleanupFailure: unknown
+ try {
+ await seed.driver.close()
+ } catch (error) {
+ seedCleanupFailure = error
+ }
+ if (seedPrimaryFailure !== undefined) {
+ if (seedCleanupFailure !== undefined) {
+ throw new Error(
+ `${failureMessage(seedPrimaryFailure)}; seed cleanup diagnostics: ${failureMessage(seedCleanupFailure)}`,
+ )
+ }
+ throw seedPrimaryFailure
+ }
+ if (seedCleanupFailure !== undefined) throw seedCleanupFailure
+
+ const observedDatabase = new ObservedDatabase(await openDatabase())
+ const { driver, persistence } = createSharedPersistence(
+ observedDatabase,
+ undefined,
+ options.schedulingFault,
+ )
+
+ const collections: Array> = []
+ const persistPromises: Array> = []
+ const preloadPromises: Array> = []
+ const logicalCompletionOrder: Array = []
+ const completedPersistIds: Array = []
+ const hydrationCompletions: Array = []
+ const observedRows = new Map>()
+ const cleanupFailures: Array = []
+ let releaseHeldBegin: (() => void) | undefined
+ let beginEntered: Promise | undefined
+ let persistSequence = 0
+ let primaryFailure: unknown
+ let observation: SharedDriverFairnessObservation | undefined
+
+ try {
+ // Cache only the unrelated persist tables. Hydrate tables remain cold.
+ for (const collectionId of persistIds) {
+ await persistence.adapter.loadSubset(collectionId, {})
+ }
+
+ observedDatabase.clearTrace()
+ driver.clearAdmissions()
+
+ if (persistIds.length > 0) {
+ const heldBegin = observedDatabase.holdNextTransactionBegin()
+ releaseHeldBegin = heldBegin.release
+ beginEntered = heldBegin.entered
+ }
+
+ // Admit the explicit legal history while the first persist is held. The
+ // first item is the non-preemptible transaction; every tail permutation
+ // is therefore observable at the same deterministic release checkpoint.
+ for (const work of scenario.work) {
+ const collectionId = collectionIdFor(scenario, work)
+ if (work.kind === `persist`) {
+ persistSequence += 1
+ const sequence = persistSequence
+ const persist = persistence.adapter
+ .applyCommittedTx(
+ collectionId,
+ createPersistedTx(collectionId, sequence, work.mutationsPerPersist),
+ )
+ .then(() => {
+ completedPersistIds.push(collectionId)
+ logicalCompletionOrder.push(`persist:${collectionId}`)
+ })
+ persistPromises.push(persist)
+ continue
+ }
+
+ const collection = createCollection(
+ persistedCollectionOptions({
+ id: collectionId,
+ getKey: (row) => row.id,
+ persistence,
+ }),
+ )
+ collections.push(collection)
+ const preloadRequest = collection.preload()
+ admittedHydrateIds.push(collectionId)
+ if (admittedHydrateIds.length === hydrateIds.length) {
+ allHydratesRequested.resolve()
+ }
+ const preload = preloadRequest.then(() => {
+ logicalCompletionOrder.push(`hydrate:${collectionId}`)
+ observedRows.set(
+ collectionId,
+ collection.toArray.map((row) => ({ id: row.id, value: row.value })),
+ )
+ hydrationCompletions.push({
+ collectionId,
+ completionOrdinal: logicalCompletionOrder.length - 1,
+ completedPersistIds: [...completedPersistIds],
+ pendingPersistCount: persistIds.length - completedPersistIds.length,
+ rawDequeueCount: observedDatabase.rawDequeues.length,
+ })
+ })
+ preloadPromises.push(preload)
+ }
+
+ await beginEntered
+ await allHydratesRequested.promise
+ // Every public preload request is now pending. Releasing the one
+ // non-preemptible persist here leaves production responsible for admitting
+ // and completing each full logical hydrate under the approved scheduler.
+ releaseHeldBegin?.()
+ releaseHeldBegin = undefined
+
+ await Promise.all([...persistPromises, ...preloadPromises])
+
+ observation = {
+ scenario,
+ admittedHydrateIds,
+ logicalCompletionOrder,
+ hydrationCompletions,
+ driverAdmissions: driver.admissions.map((entry) => ({
+ ...entry,
+ params: [...entry.params],
+ })),
+ rawDequeues: observedDatabase.rawDequeues.map((entry) => ({
+ ...entry,
+ params: [...entry.params],
+ })),
+ hydratedCollections: hydrateWork.map((work) => {
+ const collectionId = collectionIdFor(scenario, work)
+ return {
+ collectionId,
+ rows:
+ observedRows.get(collectionId)?.map((row) => ({ ...row })) ?? [],
+ }
+ }),
+ cleanupFailures,
+ }
+ } catch (error) {
+ primaryFailure = error
+ } finally {
+ releaseHeldBegin?.()
+ await Promise.allSettled([...persistPromises, ...preloadPromises])
+ for (const collection of collections) {
+ try {
+ await collection.cleanup()
+ } catch (error) {
+ cleanupFailures.push(failureMessage(error))
+ }
+ }
+ try {
+ await driver.close()
+ } catch (error) {
+ cleanupFailures.push(failureMessage(error))
+ }
+ }
+
+ if (primaryFailure !== undefined) {
+ if (cleanupFailures.length > 0) {
+ throw new Error(
+ `${failureMessage(primaryFailure)}; active cleanup diagnostics: ${JSON.stringify(cleanupFailures)}`,
+ )
+ }
+ throw primaryFailure
+ }
+ if (!observation) {
+ throw new Error(`shared-driver observation ended without a result`)
+ }
+ return observation
+}
diff --git a/packages/browser-db-sqlite-persistence/tsconfig.json b/packages/browser-db-sqlite-persistence/tsconfig.json
index 5b14f299c7..b8bc1c02dc 100644
--- a/packages/browser-db-sqlite-persistence/tsconfig.json
+++ b/packages/browser-db-sqlite-persistence/tsconfig.json
@@ -19,6 +19,14 @@
]
}
},
- "include": ["src", "tests", "e2e", "vite.config.ts", "vitest.e2e.config.ts"],
+ "include": [
+ "src",
+ "tests",
+ "e2e",
+ "vite.config.ts",
+ "vite.opfs.config.ts",
+ "vitest.e2e.config.ts",
+ "playwright.opfs.config.ts"
+ ],
"exclude": ["node_modules", "dist"]
}
diff --git a/packages/browser-db-sqlite-persistence/vite.opfs.config.ts b/packages/browser-db-sqlite-persistence/vite.opfs.config.ts
new file mode 100644
index 0000000000..c66d394740
--- /dev/null
+++ b/packages/browser-db-sqlite-persistence/vite.opfs.config.ts
@@ -0,0 +1,30 @@
+import { dirname, resolve } from 'node:path'
+import { fileURLToPath } from 'node:url'
+import { defineConfig } from 'vite'
+
+const packageDirectory = dirname(fileURLToPath(import.meta.url))
+
+export default defineConfig({
+ base: `./`,
+ // wa-sqlite locates its sibling WASM file through import.meta.url. Keeping
+ // the module out of Vite's dependency prebundle preserves that relationship
+ // for this real-browser fixture.
+ optimizeDeps: {
+ exclude: [`@journeyapps/wa-sqlite`],
+ },
+ resolve: {
+ alias: {
+ '@tanstack/db': resolve(packageDirectory, `../db/src`),
+ '@tanstack/db-ivm': resolve(packageDirectory, `../db-ivm/src`),
+ '@tanstack/db-sqlite-persistence-core': resolve(
+ packageDirectory,
+ `../db-sqlite-persistence-core/src`,
+ ),
+ },
+ },
+ server: {
+ fs: {
+ allow: [resolve(packageDirectory, `../..`)],
+ },
+ },
+})
diff --git a/packages/db-sqlite-persistence-core/src/persisted.ts b/packages/db-sqlite-persistence-core/src/persisted.ts
index 9ee0c810f7..c1322c4f6b 100644
--- a/packages/db-sqlite-persistence-core/src/persisted.ts
+++ b/packages/db-sqlite-persistence-core/src/persisted.ts
@@ -220,6 +220,21 @@ export type PersistedRowScanOptions = {
metadataOnly?: boolean
}
+export type PersistencePullSinceResult =
+ | {
+ latestRowVersion: number
+ requiresFullReload: true
+ }
+ | {
+ latestRowVersion: number
+ requiresFullReload: false
+ changedKeys: Array
+ deletedKeys: Array
+ deltas?: Array<
+ ReplayableTxDelta, string | number>
+ >
+ }
+
export type PersistedTx<
T extends object = Record,
TKey extends string | number = string | number,
@@ -250,6 +265,21 @@ export type PersistedTx<
collectionMetadataMutations?: Array
}
+/**
+ * Opaque identity shared by every adapter over the same physical SQLite driver.
+ * The core adapter uses it to serialize non-preemptible logical operations
+ * while alternating one regular operation between queued hydrations.
+ *
+ * Delegating drivers must forward this property before they are passed to a
+ * SQLite persistence adapter. A driver may also brand each returned Promise
+ * with the same key for late capability discovery, but a wrapper must return
+ * that exact Promise: discovering the key after a call cannot retroactively
+ * schedule the wrapper's first logical operation.
+ */
+export const SQLITE_DRIVER_SHARED_LOGICAL_SCHEDULING_KEY = Symbol.for(
+ `tanstack-db.sqlite-driver-supports-shared-logical-scheduling`,
+)
+
export interface PersistenceAdapter {
loadSubset: (
collectionId: string,
@@ -281,9 +311,27 @@ export interface PersistenceAdapter {
latestSeq: number
latestRowVersion: number
}>
+ /**
+ * Runs one complete logical hydrate as a non-preemptible scheduler unit.
+ * The callback must use the supplied unscheduled adapter for all nested
+ * persistence work and must not retain it after the callback settles.
+ */
+ runInHydrationScope?: (
+ task: (adapter: HydrationPersistenceAdapter) => Promise,
+ ) => Promise
+ /** Whether hydration scopes currently enter a shared driver scheduler. */
+ isHydrationScopeScheduled?: () => boolean
+}
+
+export type HydrationPersistenceAdapter = PersistenceAdapter & {
+ pullSince?: (
+ collectionId: string,
+ fromRowVersion: number,
+ ) => Promise
}
export interface SQLiteDriver {
+ readonly [SQLITE_DRIVER_SHARED_LOGICAL_SCHEDULING_KEY]?: object
exec: (sql: string) => Promise
query: (
sql: string,
@@ -298,6 +346,24 @@ export interface SQLiteDriver {
) => Promise
}
+/**
+ * Forwards a driver's shared logical scheduling identity to a transparent
+ * delegating driver before the wrapper is used by a persistence adapter.
+ */
+export function forwardSQLiteDriverSharedLogicalScheduling<
+ TDriver extends SQLiteDriver,
+>(source: SQLiteDriver, wrapper: TDriver): TDriver {
+ const key = source[SQLITE_DRIVER_SHARED_LOGICAL_SCHEDULING_KEY]
+ if (key) {
+ Object.defineProperty(
+ wrapper,
+ SQLITE_DRIVER_SHARED_LOGICAL_SCHEDULING_KEY,
+ { value: key },
+ )
+ }
+ return wrapper
+}
+
export interface PersistedCollectionCoordinator {
getNodeId: () => string
subscribe: (
@@ -311,18 +377,27 @@ export interface PersistedCollectionCoordinator {
collectionId: string,
options: LoadSubsetOptions,
) => Promise
+ /**
+ * Requests leader-side index creation. The scoped adapter is leader-local
+ * and is never serialized to a follower. `localEnsureCompleted` lets the
+ * built-in leader avoid repeating successful local work.
+ */
requestEnsurePersistedIndex: (
collectionId: string,
signature: string,
spec: PersistedIndexSpec,
+ scopedAdapter?: HydrationPersistenceAdapter,
+ localEnsureCompleted?: boolean,
) => Promise
requestApplyLocalMutations?: (
collectionId: string,
mutations: Array,
) => Promise
+ /** The scoped adapter is leader-local and is never serialized to a follower. */
pullSince?: (
collectionId: string,
fromRowVersion: number,
+ scopedAdapter?: HydrationPersistenceAdapter,
) => Promise
}
@@ -803,9 +878,9 @@ class PersistedCollectionRuntime<
truncate: null,
metadata: null,
}
- private started = false
private startupMetadataPromise: Promise | null = null
private startPromise: Promise | null = null
+ private hasAttemptedStartup = false
private resumeBaselinePromise: Promise | null = null
private lifecycleGeneration = 0
private internalApplyDepth = 0
@@ -900,13 +975,106 @@ class PersistedCollectionRuntime<
}
}
+ private runInHydrationScope(
+ task: (adapter: HydrationPersistenceAdapter) => Promise,
+ adapter: HydrationPersistenceAdapter = this.persistence.adapter,
+ ): Promise {
+ if (adapter.runInHydrationScope) {
+ return adapter.runInHydrationScope(task)
+ }
+ return Promise.resolve().then(() => task(adapter))
+ }
+
async ensureStarted(): Promise {
if (this.startPromise) {
return this.startPromise
}
const lifecycleGeneration = this.lifecycleGeneration
- this.startPromise = this.startInternal(lifecycleGeneration)
+ const isRestart = this.hasAttemptedStartup
+ this.hasAttemptedStartup = true
+ let resolveStartupMetadata!: () => void
+ let rejectStartupMetadata!: (error: unknown) => void
+ this.startupMetadataPromise = new Promise((resolve, reject) => {
+ resolveStartupMetadata = resolve
+ rejectStartupMetadata = reject
+ })
+ void this.startupMetadataPromise.catch(() => undefined)
+
+ this.startPromise = (async () => {
+ const loadStartupMetadata = async (
+ adapter: HydrationPersistenceAdapter,
+ ) => {
+ if (lifecycleGeneration !== this.lifecycleGeneration) {
+ resolveStartupMetadata()
+ return false
+ }
+
+ try {
+ await this.loadStartupMetadataInternal(lifecycleGeneration, adapter)
+ resolveStartupMetadata()
+ return lifecycleGeneration === this.lifecycleGeneration
+ } catch (error) {
+ rejectStartupMetadata(error)
+ throw error
+ }
+ }
+
+ let startup:
+ | {
+ appliedCursor: number | undefined
+ indexBootstrapSnapshot: Array
+ completedLocalIndexSignatures: Set
+ }
+ | undefined
+ const scheduleStartupAsOneHydrate =
+ this.persistence.adapter.runInHydrationScope !== undefined &&
+ (!isRestart ||
+ (this.persistence.adapter.isHydrationScopeScheduled?.() ?? true))
+ if (scheduleStartupAsOneHydrate) {
+ startup = await this.applyMutex.run(async () => {
+ const result = await this.runInHydrationScope(async (adapter) => {
+ if (!(await loadStartupMetadata(adapter))) return undefined
+ return this.startInternal(lifecycleGeneration, adapter)
+ })
+ await this.flushQueuedTxCommittedUnsafe()
+ return result
+ })
+ } else {
+ // Preserve the existing unscheduled-adapter lifecycle contract: a
+ // replacement upstream may start while stale hydration is settling.
+ if (await loadStartupMetadata(this.persistence.adapter)) {
+ startup = await this.applyMutex.run(async () => {
+ const result =
+ this.persistence.adapter.isHydrationScopeScheduled?.()
+ ? await this.runInHydrationScope((adapter) =>
+ this.startInternal(lifecycleGeneration, adapter),
+ )
+ : await this.startInternal(
+ lifecycleGeneration,
+ this.persistence.adapter,
+ )
+ await this.flushQueuedTxCommittedUnsafe()
+ return result
+ })
+ }
+ }
+ if (
+ startup !== undefined &&
+ lifecycleGeneration === this.lifecycleGeneration
+ ) {
+ await this.requestCoordinatorPersistedIndexes(
+ startup.indexBootstrapSnapshot,
+ startup.completedLocalIndexSignatures,
+ )
+ if (
+ startup.appliedCursor !== undefined &&
+ lifecycleGeneration === this.lifecycleGeneration
+ ) {
+ await this.waitForAppliedReceiptsAfter(startup.appliedCursor)
+ }
+ }
+ })()
return this.startPromise
}
@@ -921,26 +1089,43 @@ class PersistedCollectionRuntime<
if (lifecycleGeneration !== this.lifecycleGeneration) return
if (this.syncMode !== `on-demand`) return
- await this.hydrateBaseline(lifecycleGeneration)
+ const appliedCursor = await this.applyMutex.run(async () => {
+ const result = await this.runInHydrationScope((adapter) =>
+ this.hydrateBaseline(lifecycleGeneration, adapter),
+ )
+ await this.flushQueuedTxCommittedUnsafe()
+ return result
+ })
+ if (
+ appliedCursor !== undefined &&
+ lifecycleGeneration === this.lifecycleGeneration
+ ) {
+ await this.waitForAppliedReceiptsAfter(appliedCursor)
+ }
})()
return this.resumeBaselinePromise
}
- private async hydrateBaseline(lifecycleGeneration: number): Promise {
- if (lifecycleGeneration !== this.lifecycleGeneration) return
+ private async hydrateBaseline(
+ lifecycleGeneration: number,
+ adapter: HydrationPersistenceAdapter,
+ ): Promise {
+ if (lifecycleGeneration !== this.lifecycleGeneration) return undefined
const baseline = {}
this.activeSubsets.set(this.getSubsetKey(baseline), baseline)
const appliedCursor = this.appliedReceiptSequence
- await this.applyMutex.run(async () => {
- if (lifecycleGeneration !== this.lifecycleGeneration) return
- await this.hydrateSubsetUnsafe(baseline, {
+ await this.hydrateSubsetUnsafe(
+ baseline,
+ {
requestRemoteEnsure: false,
lifecycleGeneration,
- })
- })
- if (lifecycleGeneration !== this.lifecycleGeneration) return
- await this.waitForAppliedReceiptsAfter(appliedCursor)
+ },
+ adapter,
+ )
+ return lifecycleGeneration === this.lifecycleGeneration
+ ? appliedCursor
+ : undefined
}
async ensureStartupMetadataLoaded(): Promise {
@@ -948,41 +1133,52 @@ class PersistedCollectionRuntime<
return this.startupMetadataPromise
}
- const lifecycleGeneration = this.lifecycleGeneration
- this.startupMetadataPromise =
- this.loadStartupMetadataInternal(lifecycleGeneration)
- return this.startupMetadataPromise
+ void this.ensureStarted()
+ return this.startupMetadataPromise!
}
- private async startInternal(lifecycleGeneration: number): Promise {
- if (this.started) {
- return
- }
-
- this.started = true
-
- await this.ensureStartupMetadataLoaded()
- if (lifecycleGeneration !== this.lifecycleGeneration) return
+ private async startInternal(
+ lifecycleGeneration: number,
+ adapter: HydrationPersistenceAdapter,
+ ): Promise<
+ | {
+ appliedCursor: number | undefined
+ indexBootstrapSnapshot: Array
+ completedLocalIndexSignatures: Set
+ }
+ | undefined
+ > {
+ if (lifecycleGeneration !== this.lifecycleGeneration) return undefined
const indexBootstrapSnapshot = this.collection?.getIndexMetadata() ?? []
this.attachIndexLifecycleListeners()
- await this.bootstrapPersistedIndexes(indexBootstrapSnapshot)
- if (lifecycleGeneration !== this.lifecycleGeneration) return
-
- if (this.syncMode !== `on-demand`) {
- await this.hydrateBaseline(lifecycleGeneration)
- }
+ const completedLocalIndexSignatures = await this.bootstrapPersistedIndexes(
+ indexBootstrapSnapshot,
+ adapter,
+ )
+ if (lifecycleGeneration !== this.lifecycleGeneration) return undefined
+
+ const appliedCursor =
+ this.syncMode !== `on-demand`
+ ? await this.hydrateBaseline(lifecycleGeneration, adapter)
+ : undefined
+ return lifecycleGeneration === this.lifecycleGeneration
+ ? {
+ appliedCursor,
+ indexBootstrapSnapshot,
+ completedLocalIndexSignatures,
+ }
+ : undefined
}
private async loadStartupMetadataInternal(
lifecycleGeneration: number,
+ adapter: HydrationPersistenceAdapter,
): Promise {
// Restore stream position from the database so that new mutations
// don't collide with previously applied transactions.
- if (this.persistence.adapter.getStreamPosition) {
- const position = await this.persistence.adapter.getStreamPosition(
- this.collectionId,
- )
+ if (adapter.getStreamPosition) {
+ const position = await adapter.getStreamPosition(this.collectionId)
if (lifecycleGeneration !== this.lifecycleGeneration) return
this.observeStreamPosition(
position.latestTerm,
@@ -991,19 +1187,20 @@ class PersistedCollectionRuntime<
)
}
- const collectionMetadata = await this.loadCollectionMetadataSnapshot()
+ const collectionMetadata =
+ await this.loadCollectionMetadataSnapshot(adapter)
if (lifecycleGeneration !== this.lifecycleGeneration) return
this.replaceCollectionMetadataSnapshot(collectionMetadata)
}
- private async loadCollectionMetadataSnapshot(): Promise<
- Array<{ key: string; value: unknown }>
- > {
- if (!this.persistence.adapter.loadCollectionMetadata) {
+ private async loadCollectionMetadataSnapshot(
+ adapter: HydrationPersistenceAdapter,
+ ): Promise> {
+ if (!adapter.loadCollectionMetadata) {
return []
}
- return this.persistence.adapter.loadCollectionMetadata(this.collectionId)
+ return adapter.loadCollectionMetadata(this.collectionId)
}
private replaceCollectionMetadataSnapshot(
@@ -1047,14 +1244,21 @@ class PersistedCollectionRuntime<
): Promise {
const lifecycleGeneration = this.lifecycleGeneration
this.activeSubsets.set(this.getSubsetKey(options), options)
-
const appliedCursor = this.appliedReceiptSequence
- await this.applyMutex.run(() =>
- this.hydrateSubsetUnsafe(options, {
- requestRemoteEnsure: this.mode === `sync-present`,
- lifecycleGeneration,
- }),
- )
+
+ await this.applyMutex.run(async () => {
+ await this.runInHydrationScope((adapter) =>
+ this.hydrateSubsetUnsafe(
+ options,
+ {
+ requestRemoteEnsure: this.mode === `sync-present`,
+ lifecycleGeneration,
+ },
+ adapter,
+ ),
+ )
+ await this.flushQueuedTxCommittedUnsafe()
+ })
if (lifecycleGeneration !== this.lifecycleGeneration) return
await this.waitForAppliedReceiptsAfter(appliedCursor)
@@ -1092,12 +1296,19 @@ class PersistedCollectionRuntime<
async forceReloadSubset(options: LoadSubsetOptions): Promise {
const lifecycleGeneration = this.lifecycleGeneration
// A one-shot refresh does not acquire an enduring subscription lease.
- await this.applyMutex.run(() =>
- this.hydrateSubsetUnsafe(options, {
- requestRemoteEnsure: false,
- lifecycleGeneration,
- }),
- )
+ await this.applyMutex.run(async () => {
+ await this.runInHydrationScope((adapter) =>
+ this.hydrateSubsetUnsafe(
+ options,
+ {
+ requestRemoteEnsure: false,
+ lifecycleGeneration,
+ },
+ adapter,
+ ),
+ )
+ await this.flushQueuedTxCommittedUnsafe()
+ })
}
queueHydrationBufferedTransaction(
@@ -1244,7 +1455,6 @@ class PersistedCollectionRuntime<
private advanceLifecycle(): void {
this.lifecycleGeneration++
- this.started = false
this.startupMetadataPromise = null
this.startPromise = null
this.resumeBaselinePromise = null
@@ -1271,8 +1481,9 @@ class PersistedCollectionRuntime<
private loadSubsetRowsUnsafe(
options: LoadSubsetOptions,
+ adapter: HydrationPersistenceAdapter,
): Promise> {
- return this.persistence.adapter.loadSubset(this.collectionId, options, {
+ return adapter.loadSubset(this.collectionId, options, {
requiredIndexSignatures: this.getRequiredIndexSignatures(),
}) as Promise>
}
@@ -1302,10 +1513,11 @@ class PersistedCollectionRuntime<
requestRemoteEnsure: boolean
lifecycleGeneration: number
},
+ adapter: HydrationPersistenceAdapter,
): Promise {
this.hydratingGeneration = config.lifecycleGeneration
try {
- const rows = await this.loadSubsetRowsUnsafe(options)
+ const rows = await this.loadSubsetRowsUnsafe(options, adapter)
if (config.lifecycleGeneration !== this.lifecycleGeneration) return
this.applyRowsToCollection(rows)
@@ -1315,8 +1527,7 @@ class PersistedCollectionRuntime<
}
}
- await this.flushQueuedHydrationTransactionsUnsafe()
- await this.flushQueuedTxCommittedUnsafe()
+ await this.flushQueuedHydrationTransactionsUnsafe(adapter)
if (config.requestRemoteEnsure) {
this.queueRemoteSubsetEnsure(options)
@@ -1398,14 +1609,16 @@ class PersistedCollectionRuntime<
})
}
- private async flushQueuedHydrationTransactionsUnsafe(): Promise {
+ private async flushQueuedHydrationTransactionsUnsafe(
+ adapter: HydrationPersistenceAdapter,
+ ): Promise {
while (this.queuedHydrationTransactions.length > 0) {
const transaction = this.queuedHydrationTransactions.shift()
if (!transaction) {
continue
}
try {
- await this.applyBufferedSyncTransactionUnsafe(transaction)
+ await this.applyBufferedSyncTransactionUnsafe(transaction, adapter)
} catch (error) {
transaction.rejectApplied?.(error)
for (const abandoned of this.queuedHydrationTransactions) {
@@ -1419,6 +1632,7 @@ class PersistedCollectionRuntime<
private async applyBufferedSyncTransactionUnsafe(
transaction: BufferedSyncTransaction,
+ adapter: HydrationPersistenceAdapter,
): Promise {
if (transaction.signal?.aborted) {
transaction.rejectApplied?.(new SyncTransactionAbortedError())
@@ -1432,7 +1646,10 @@ class PersistedCollectionRuntime<
}
const applyToCollection = (): SyncAppliedReceipt => {
- begin()
+ // Buffered source replay is part of persistence hydration. Apply it
+ // immediately so it cannot wait for a persisting mutation whose
+ // persistence is queued behind this hydrate's apply mutex.
+ begin({ immediate: true })
if (transaction.truncate) {
truncate?.()
@@ -1481,7 +1698,10 @@ class PersistedCollectionRuntime<
}
if (!transaction.internal) {
- await this.persistAndBroadcastExternalSyncTransactionUnsafe(transaction)
+ await this.persistAndBroadcastExternalSyncTransactionUnsafe(
+ transaction,
+ adapter,
+ )
}
transaction.resolveApplied?.()
} catch (error) {
@@ -1492,6 +1712,7 @@ class PersistedCollectionRuntime<
private async persistAndBroadcastExternalSyncTransactionUnsafe(
transaction: BufferedSyncTransaction,
+ adapter: HydrationPersistenceAdapter = this.persistence.adapter,
): Promise {
if (transaction.internal) {
return
@@ -1521,7 +1742,7 @@ class PersistedCollectionRuntime<
const tx = this.createPersistedTxFromOperations(transaction, streamPosition)
- await this.persistence.adapter.applyCommittedTx(this.collectionId, tx)
+ await adapter.applyCommittedTx(this.collectionId, tx)
this.publishTxCommittedEvent(
this.createTxCommittedPayload({
term: tx.term,
@@ -1976,7 +2197,12 @@ class PersistedCollectionRuntime<
if (isCollectionResetPayload(payload)) {
void this.applyMutex
- .run(() => this.truncateAndReloadUnsafe())
+ .run(async () => {
+ await this.runInHydrationScope((adapter) =>
+ this.truncateAndReloadUnsafe(adapter),
+ )
+ await this.flushQueuedTxCommittedUnsafe()
+ })
.catch((error) => {
console.warn(`Failed to process collection reset message:`, error)
})
@@ -2031,7 +2257,11 @@ class PersistedCollectionRuntime<
txCommitted.latestRowVersion,
)
- await this.invalidateFromCommittedTxUnsafe(txCommitted)
+ await this.invalidateFromCommittedTxUnsafe(
+ txCommitted,
+ this.persistence.adapter,
+ )
+ await this.flushQueuedTxCommittedUnsafe()
}
private async recoverFromSeqGapUnsafe(): Promise {
@@ -2048,25 +2278,40 @@ class PersistedCollectionRuntime<
pullResponse.latestSeq,
pullResponse.latestRowVersion,
)
- if (pullResponse.requiresFullReload || !pullResponse.deltas) {
- await this.reloadActiveSubsetsUnsafe()
+ if (pullResponse.requiresFullReload) {
+ await this.runInHydrationScope((adapter) =>
+ this.reloadActiveSubsetsUnsafe(adapter),
+ )
return
}
-
- for (const delta of pullResponse.deltas) {
- await this.invalidateFromCommittedTxUnsafe({
- type: `tx:committed`,
- term: pullResponse.latestTerm,
- seq: pullResponse.latestSeq,
- txId: delta.txId,
- latestRowVersion: delta.latestRowVersion,
- requiresFullReload: false,
- changedRows: delta.changedRows,
- deletedKeys: delta.deletedKeys,
- rowMetadataMutations: delta.rowMetadataMutations,
- collectionMetadataMutations: delta.collectionMetadataMutations,
- })
+ const deltas = pullResponse.deltas
+ if (!deltas) {
+ await this.runInHydrationScope((adapter) =>
+ this.reloadActiveSubsetsUnsafe(adapter),
+ )
+ return
}
+
+ await this.runInHydrationScope(async (adapter) => {
+ for (const delta of deltas) {
+ await this.invalidateFromCommittedTxUnsafe(
+ {
+ type: `tx:committed`,
+ term: pullResponse.latestTerm,
+ seq: pullResponse.latestSeq,
+ txId: delta.txId,
+ latestRowVersion: delta.latestRowVersion,
+ requiresFullReload: false,
+ changedRows: delta.changedRows,
+ deletedKeys: delta.deletedKeys,
+ rowMetadataMutations: delta.rowMetadataMutations,
+ collectionMetadataMutations:
+ delta.collectionMetadataMutations,
+ },
+ adapter,
+ )
+ }
+ })
return
}
} catch (error) {
@@ -2074,7 +2319,9 @@ class PersistedCollectionRuntime<
}
}
- await this.truncateAndReloadUnsafe()
+ await this.runInHydrationScope((adapter) =>
+ this.truncateAndReloadUnsafe(adapter),
+ )
if (this.mode === `sync-present`) {
for (const options of this.activeSubsets.values()) {
@@ -2083,7 +2330,9 @@ class PersistedCollectionRuntime<
}
}
- private async truncateAndReloadUnsafe(): Promise {
+ private async truncateAndReloadUnsafe(
+ adapter: HydrationPersistenceAdapter,
+ ): Promise {
if (this.syncControls.begin && this.syncControls.commit) {
this.withInternalApply(() => {
this.syncControls.begin?.({ immediate: true })
@@ -2092,21 +2341,28 @@ class PersistedCollectionRuntime<
})
}
- await this.reloadActiveSubsetsUnsafe()
+ await this.reloadActiveSubsetsUnsafe(adapter)
}
private async invalidateFromCommittedTxUnsafe(
txCommitted: TxCommitted,
+ adapter: HydrationPersistenceAdapter,
): Promise {
+ const reloadActiveSubsets = () =>
+ this.runInHydrationScope(
+ (scopedAdapter) => this.reloadActiveSubsetsUnsafe(scopedAdapter),
+ adapter,
+ )
+
if (txCommitted.requiresFullReload) {
- await this.reloadActiveSubsetsUnsafe()
+ await reloadActiveSubsets()
return
}
const changedKeyCount =
txCommitted.changedRows.length + txCommitted.deletedKeys.length
if (changedKeyCount > TARGETED_INVALIDATION_KEY_LIMIT) {
- await this.reloadActiveSubsetsUnsafe()
+ await reloadActiveSubsets()
return
}
@@ -2121,7 +2377,7 @@ class PersistedCollectionRuntime<
// Has paginated subsets — fall back to full reload.
// Targeted invalidation for paginated subsets is deferred to a future iteration.
- await this.reloadActiveSubsetsUnsafe()
+ await reloadActiveSubsets()
}
private async applyTargetedInvalidationUnsafe(
@@ -2180,7 +2436,9 @@ class PersistedCollectionRuntime<
})
}
- private async reloadActiveSubsetsUnsafe(): Promise {
+ private async reloadActiveSubsetsUnsafe(
+ adapter: HydrationPersistenceAdapter,
+ ): Promise {
const lifecycleGeneration = this.lifecycleGeneration
const activeSubsetOptions =
this.activeSubsets.size > 0
@@ -2190,10 +2448,11 @@ class PersistedCollectionRuntime<
this.hydratingGeneration = lifecycleGeneration
try {
const mergedRows = new Map()
- const collectionMetadata = await this.loadCollectionMetadataSnapshot()
+ const collectionMetadata =
+ await this.loadCollectionMetadataSnapshot(adapter)
if (lifecycleGeneration !== this.lifecycleGeneration) return
for (const options of activeSubsetOptions) {
- const subsetRows = await this.loadSubsetRowsUnsafe(options)
+ const subsetRows = await this.loadSubsetRowsUnsafe(options, adapter)
if (lifecycleGeneration !== this.lifecycleGeneration) return
for (const row of subsetRows) {
mergedRows.set(row.key, {
@@ -2217,8 +2476,7 @@ class PersistedCollectionRuntime<
}
}
- await this.flushQueuedHydrationTransactionsUnsafe()
- await this.flushQueuedTxCommittedUnsafe()
+ await this.flushQueuedHydrationTransactionsUnsafe(adapter)
}
private attachIndexLifecycleListeners(): void {
@@ -2243,16 +2501,33 @@ class PersistedCollectionRuntime<
private async bootstrapPersistedIndexes(
indexMetadataSnapshot?: Array,
- ): Promise {
+ adapter: HydrationPersistenceAdapter = this.persistence.adapter,
+ ): Promise> {
const collection = this.collection
if (!collection && !indexMetadataSnapshot) {
- return
+ return new Set()
}
const indexMetadata =
indexMetadataSnapshot ?? collection?.getIndexMetadata() ?? []
+ const completedLocalIndexSignatures = new Set()
+ for (const metadata of indexMetadata) {
+ if (await this.ensureLocalPersistedIndex(metadata, adapter)) {
+ completedLocalIndexSignatures.add(metadata.signature)
+ }
+ }
+ return completedLocalIndexSignatures
+ }
+
+ private async requestCoordinatorPersistedIndexes(
+ indexMetadata: Array,
+ completedLocalIndexSignatures: ReadonlySet,
+ ): Promise {
for (const metadata of indexMetadata) {
- await this.ensurePersistedIndex(metadata)
+ await this.requestCoordinatorPersistedIndex(
+ metadata,
+ completedLocalIndexSignatures.has(metadata.signature),
+ )
}
}
@@ -2271,24 +2546,47 @@ class PersistedCollectionRuntime<
private async ensurePersistedIndex(
indexMetadata: CollectionIndexMetadata,
+ adapter: HydrationPersistenceAdapter = this.persistence.adapter,
): Promise {
+ const completedLocally = await this.ensureLocalPersistedIndex(
+ indexMetadata,
+ adapter,
+ )
+ await this.requestCoordinatorPersistedIndex(indexMetadata, completedLocally)
+ }
+
+ private async ensureLocalPersistedIndex(
+ indexMetadata: CollectionIndexMetadata,
+ adapter: HydrationPersistenceAdapter,
+ ): Promise {
const spec = this.buildPersistedIndexSpec(indexMetadata)
try {
- await this.persistence.adapter.ensureIndex(
+ await adapter.ensureIndex(
this.collectionId,
indexMetadata.signature,
spec,
)
+ return true
} catch (error) {
console.warn(`Failed to ensure persisted index in adapter:`, error)
+ return false
}
+ }
+
+ private async requestCoordinatorPersistedIndex(
+ indexMetadata: CollectionIndexMetadata,
+ completedLocally: boolean,
+ ): Promise {
+ const spec = this.buildPersistedIndexSpec(indexMetadata)
try {
await this.persistence.coordinator.requestEnsurePersistedIndex(
this.collectionId,
indexMetadata.signature,
spec,
+ completedLocally ? this.persistence.adapter : undefined,
+ completedLocally,
)
} catch (error) {
console.warn(
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..075651111a 100644
--- a/packages/db-sqlite-persistence-core/src/sqlite-core-adapter.ts
+++ b/packages/db-sqlite-persistence-core/src/sqlite-core-adapter.ts
@@ -8,12 +8,14 @@ import {
InvalidPersistedStorageKeyEncodingError,
} from './errors'
import {
+ SQLITE_DRIVER_SHARED_LOGICAL_SCHEDULING_KEY,
createPersistedTableName,
decodePersistedStorageKey,
encodePersistedStorageKey,
} from './persisted'
import type { LoadSubsetOptions } from '@tanstack/db'
import type {
+ HydrationPersistenceAdapter,
PersistedIndexSpec,
PersistedRowScanOptions,
PersistedScannedRow,
@@ -71,6 +73,159 @@ export type SQLitePullSinceResult =
deltas: Array, TKey>>
}
+type ScheduledOperationKind = `regular` | `hydrate`
+
+type ScheduledOperation = {
+ kind: ScheduledOperationKind
+ task: () => Promise
+ resolve: (value: T) => void
+ reject: (error: unknown) => void
+}
+
+class SharedPersistenceScheduler {
+ private readonly regularQueue: Array> = []
+ private readonly hydrateQueue: Array> = []
+ private running = false
+ private lastCompletedKind: ScheduledOperationKind | undefined
+
+ runRegular(task: () => Promise): Promise {
+ return this.enqueue(`regular`, task)
+ }
+
+ runHydrate(task: () => Promise): Promise {
+ return this.enqueue(`hydrate`, task)
+ }
+
+ adoptRunningHydrate(completion: Promise): void {
+ if (this.running) return
+ this.running = true
+ const finish = () => {
+ this.lastCompletedKind = `hydrate`
+ this.running = false
+ this.drain()
+ }
+ void completion.then(finish, finish)
+ }
+
+ private enqueue(
+ kind: ScheduledOperationKind,
+ task: () => Promise,
+ ): Promise {
+ const result = new Promise((resolve, reject) => {
+ const operation: ScheduledOperation = {
+ kind,
+ task,
+ resolve,
+ reject,
+ }
+ const queue = kind === `hydrate` ? this.hydrateQueue : this.regularQueue
+ queue.push(operation as ScheduledOperation)
+ })
+ this.drain()
+ return result
+ }
+
+ private drain(): void {
+ if (this.running) return
+
+ const operation = this.takeNext()
+ if (!operation) return
+
+ this.running = true
+ void this.execute(operation)
+ }
+
+ private async execute(operation: ScheduledOperation): Promise {
+ try {
+ operation.resolve(await operation.task())
+ } catch (error) {
+ operation.reject(error)
+ } finally {
+ this.lastCompletedKind = operation.kind
+ this.running = false
+ this.drain()
+ }
+ }
+
+ private takeNext(): ScheduledOperation | undefined {
+ // Hydrates get priority after the currently running non-preemptible unit.
+ // While both lanes remain queued, alternate one regular operation after
+ // each hydrate (K=1), preserving FIFO order within each lane.
+ if (this.hydrateQueue.length > 0) {
+ if (
+ this.regularQueue.length > 0 &&
+ this.lastCompletedKind === `hydrate`
+ ) {
+ return this.regularQueue.shift()
+ }
+ return this.hydrateQueue.shift()
+ }
+ return this.regularQueue.shift()
+ }
+}
+
+const sharedPersistenceSchedulers = new WeakMap<
+ object,
+ SharedPersistenceScheduler
+>()
+const observedDriverSchedulingKeys = new WeakMap