Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
99 commits
Select commit Hold shift + click to select a range
a9ed834
fix(db): unify ordered window semantics
KyleAMathews Aug 25, 2026
530d4ec
Merge remote-tracking branch 'origin/main' into codex/loadsubset-tota…
KyleAMathews Aug 25, 2026
c476cb3
Merge branch 'codex/loadsubset-outcome-plumbing' into codex/loadsubse…
KyleAMathews Aug 25, 2026
b4dca01
Merge branch 'codex/loadsubset-coverage-registry' into codex/loadsubs…
KyleAMathews Aug 25, 2026
1fbcec7
fix(db): prove ordered window boundaries
KyleAMathews Aug 25, 2026
29244d5
fix(db): reset ordered replay coverage
KyleAMathews Aug 25, 2026
5e09952
Merge branch 'codex/loadsubset-coverage-registry' into codex/loadsubs…
KyleAMathews Aug 25, 2026
9c701e0
fix(db): track ordered refinement fallback
KyleAMathews Aug 25, 2026
3158ac0
test(db): harden ordered pagination oracle
KyleAMathews Aug 25, 2026
d63b6c4
test(db): scale replay oracle timeout
KyleAMathews Aug 25, 2026
c224b6f
fix(db): keep ordered replay refinement atomic
KyleAMathews Aug 25, 2026
5e9e7d9
Merge branch 'codex/loadsubset-coverage-registry' into codex/loadsubs…
KyleAMathews Aug 25, 2026
3b1ba84
fix(db): preserve legacy ordered subset settlement
KyleAMathews Aug 25, 2026
42e754b
fix(db): refresh ordered coverage after SSE batches
KyleAMathews Aug 25, 2026
f43b857
fix(db): retain SSE changes during boundary refinement
KyleAMathews Aug 25, 2026
7f77e02
fix(db): revalidate changed refinement prefixes
KyleAMathews Aug 25, 2026
6877a0b
fix(db): scope ordered coverage to window revisions
KyleAMathews Aug 25, 2026
4fdb207
Merge branch 'codex/loadsubset-coverage-registry' into codex/loadsubs…
KyleAMathews Aug 25, 2026
d248e01
Merge branch 'codex/loadsubset-coverage-registry' into codex/loadsubs…
KyleAMathews Aug 25, 2026
34e578f
refactor(db): simplify ordered window state
KyleAMathews Aug 25, 2026
b5a1772
fix(db): preserve ordered window evidence
KyleAMathews Aug 26, 2026
f38795a
Merge commit '58402160' into codex/loadsubset-total-order
KyleAMathews Aug 26, 2026
03f7a58
fix(db): preserve optimized reverse ordering
KyleAMathews Aug 26, 2026
10f389c
Merge commit '049e0ce2' into codex/loadsubset-total-order
KyleAMathews Aug 26, 2026
8243c46
fix(db): retry refined ordered continuations
KyleAMathews Aug 26, 2026
32b1e1d
chore: add ordered window changeset
KyleAMathews Aug 26, 2026
7b232e4
ci: apply automated fixes
autofix-ci[bot] Aug 26, 2026
b18a663
chore: format ordered window changes
KyleAMathews Aug 26, 2026
73c408c
Merge remote-tracking branch 'origin/codex/loadsubset-total-order' in…
KyleAMathews Aug 26, 2026
337e6e8
fix(db): separate request shape from coverage
KyleAMathews Aug 26, 2026
4ec5e94
Merge branch 'codex/loadsubset-coverage-registry' into codex/loadsubs…
KyleAMathews Aug 26, 2026
b24457c
Merge branch 'codex/loadsubset-coverage-registry' into codex/loadsubs…
KyleAMathews Aug 26, 2026
cb84d47
Merge branch 'codex/loadsubset-coverage-registry' into codex/loadsubs…
KyleAMathews Aug 26, 2026
c9cc4d1
fix(db): revoke stale ordered continuations
KyleAMathews Aug 26, 2026
55a576c
test(db): cover ordered evidence boundaries
KyleAMathews Aug 26, 2026
f944398
docs(db): name outcome-free subset completion
KyleAMathews Aug 26, 2026
6d7c48d
test(electric): align subset ownership expectations
KyleAMathews Aug 26, 2026
03ee53b
Merge coverage registry into total order
KyleAMathews Aug 27, 2026
3189f00
fix(db): preserve ordered replay boundaries
KyleAMathews Aug 27, 2026
dc6ca62
test(db): close ordered boundary oracle gaps
KyleAMathews Aug 27, 2026
5ba3524
fix(db): publish ordered replays atomically
KyleAMathews Aug 27, 2026
e4a8a63
fix(db): retain incomplete replay publications
KyleAMathews Aug 27, 2026
f2ac578
test(db): close replay lifecycle oracle gaps
KyleAMathews Aug 27, 2026
2d7194e
docs(db): preserve replay publication laws
KyleAMathews Aug 27, 2026
3135310
test(db): pin replay cancellation laws
KyleAMathews Aug 27, 2026
ca61a6c
fix(db): make ordered continuation progress honest
KyleAMathews Aug 27, 2026
18ccb43
test(db): close ordered progress oracle gaps
KyleAMathews Aug 27, 2026
26c0e6b
fix(db): advance past excluded continuation rows
KyleAMathews Aug 27, 2026
c407fbd
fix(db): align effect continuation progress
KyleAMathews Aug 27, 2026
6c36f84
fix(db): complete effect truncate replay
KyleAMathews Aug 27, 2026
bebff2b
fix(db): unify ordered publication restoration
KyleAMathews Aug 27, 2026
7dbc3eb
fix(db): retain empty ordered publications
KyleAMathews Aug 28, 2026
f83dd91
fix(db): isolate active replacement cursors
KyleAMathews Aug 28, 2026
cb5f9f7
fix(db): scope replay continuation state
KyleAMathews Aug 28, 2026
889e2db
fix(db): rebuild ordered replay acquisitions
KyleAMathews Aug 28, 2026
075f7df
fix(db): scope ordered replay evidence
KyleAMathews Aug 28, 2026
7480aed
fix(db): gate synchronous ordered evidence
KyleAMathews Aug 28, 2026
e6e3a14
fix(db): retain failed ordered publications
KyleAMathews Aug 28, 2026
b9734b8
fix(db): reconcile failed publication transitions
KyleAMathews Aug 28, 2026
c286111
fix(db): preserve failed publication provenance
KyleAMathews Aug 28, 2026
752b99f
fix(db): retain subset acquisition provenance
KyleAMathews Aug 28, 2026
7b30c26
fix(db): preserve exact sync request provenance
KyleAMathews Aug 28, 2026
174c46e
fix(db): preserve deduplicated request provenance
KyleAMathews Aug 28, 2026
45c8da3
fix(db): preserve same-version request authority
KyleAMathews Aug 28, 2026
1f0b3c1
fix(db): keep metadata out of row authority
KyleAMathews Aug 28, 2026
1fdd20e
fix(db): isolate released request authority
KyleAMathews Aug 28, 2026
e32524b
fix(db): separate publication authority from cleanup debt
KyleAMathews Aug 28, 2026
3f19330
fix(db): retire released subset authority
KyleAMathews Aug 28, 2026
9459144
fix(db): unify subset release transitions
KyleAMathews Aug 28, 2026
d063242
fix(db): guard subset acquisition lifetimes
KyleAMathews Aug 28, 2026
f70df19
fix(db): reject obsolete subset authority
KyleAMathews Aug 28, 2026
97574d4
fix(db): settle replay before callbacks
KyleAMathews Aug 28, 2026
308b849
fix(db): settle replay errors by attempt
KyleAMathews Aug 28, 2026
2db26a1
fix(db): retain released replay barriers
KyleAMathews Aug 28, 2026
9830a7c
fix(db): contain replay callback failures
KyleAMathews Aug 28, 2026
d0c7bac
fix(db): bind replay callbacks to attempts
KyleAMathews Aug 28, 2026
48bbf9e
fix(db): retain replay callbacks before publish
KyleAMathews Aug 28, 2026
86fc3a4
fix(db): dedupe replay failure attribution
KyleAMathews Aug 28, 2026
8ae224b
fix(db): track replay failure occurrences
KyleAMathews Aug 28, 2026
5a66bbb
fix(db): preserve replay failure occurrences
KyleAMathews Aug 28, 2026
0fecd4d
fix(db): preserve cleanup failure multiplicity
KyleAMathews Aug 28, 2026
5207c0c
fix(db): preserve automatic cleanup failures
KyleAMathews Aug 28, 2026
a1afbaf
fix(db): preserve nested cleanup provenance
KyleAMathews Aug 28, 2026
b76c368
fix(db): preserve recursive cleanup provenance
KyleAMathews Aug 28, 2026
14b25b4
fix(db): distinguish cleanup failure occurrences
KyleAMathews Aug 28, 2026
63b3239
fix(db): preserve callback cleanup provenance
KyleAMathews Aug 28, 2026
474c6d1
fix(db): preserve nested callback failures
KyleAMathews Aug 28, 2026
6d56871
fix(db): retain replay failures through teardown
KyleAMathews Aug 28, 2026
a0ca2dc
fix(db): preserve queued replay failures on teardown
KyleAMathews Aug 28, 2026
b70ebdd
fix(db): preserve replay error delivery
KyleAMathews Aug 28, 2026
ebd55c2
fix(db): preserve recursive subset failure identity
KyleAMathews Aug 28, 2026
74e972a
fix(db): preserve async subset failure identity
KyleAMathews Aug 28, 2026
b0731e5
fix(db): contain propagated subset failures
KyleAMathews Aug 28, 2026
7a47c97
fix(db): scope subset failure propagation
KyleAMathews Aug 28, 2026
0754c83
fix(db): preserve ordinary recursive failure identity
KyleAMathews Aug 28, 2026
5b23796
fix(db): preserve replay entry failure identity
KyleAMathews Aug 28, 2026
5dc1fc2
fix(db): retain replay failures through teardown
KyleAMathews Aug 28, 2026
f11e874
fix(db): retain adapter-caught cleanup failures
KyleAMathews Aug 28, 2026
fed6992
fix(db): retain nested replay failure order
KyleAMathews Aug 28, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .changeset/unify-ordered-window-semantics.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
---
'@tanstack/db': patch
'@tanstack/db-ivm': patch
---

Use one total order and applied row provenance for lazy ordered windows so cursor pagination, live updates, and window replay preserve exact results.
8 changes: 8 additions & 0 deletions packages/db-ivm/src/utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,14 @@ function range(start: number, end: number): Array<number> {
export function compareKeys(a: string | number, b: string | number): number {
// Same type: compare directly
if (typeof a === typeof b) {
if (typeof a === `number` && typeof b === `number`) {
const aIsNaN = Number.isNaN(a)
const bIsNaN = Number.isNaN(b)
if (aIsNaN || bIsNaN) {
if (aIsNaN && bIsNaN) return 0
return aIsNaN ? 1 : -1
}
}
if (a < b) return -1
if (a > b) return 1
return 0
Expand Down
10 changes: 9 additions & 1 deletion packages/db-ivm/tests/utils.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { describe, expect, it } from 'vitest'
import { Temporal } from 'temporal-polyfill'
import { DefaultMap, serializeValue } from '../src/utils.js'
import { DefaultMap, compareKeys, serializeValue } from '../src/utils.js'
import { hash } from '../src/hashing/index.js'

describe(`DefaultMap`, () => {
Expand Down Expand Up @@ -30,6 +30,14 @@ describe(`DefaultMap`, () => {
})
})

describe(`compareKeys`, () => {
it(`orders finite numeric keys before NaN`, () => {
expect(compareKeys(1, Number.NaN)).toBeLessThan(0)
expect(compareKeys(Number.NaN, 1)).toBeGreaterThan(0)
expect(compareKeys(Number.NaN, Number.NaN)).toBe(0)
})
})

describe(`serializeValue`, () => {
it(`preserves the established JSON form for ordinary keys`, () => {
expect(serializeValue(`user1`)).toBe(`"user1"`)
Expand Down
2 changes: 1 addition & 1 deletion packages/db/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
"dev": "vite build --watch",
"lint": "eslint . --fix",
"test": "vitest --run",
"test:oracles": "vitest --run tests/collection-sync-reentrancy.test.ts tests/collection-subscription-replay-oracle.property.test.ts tests/query/coverage-registry-oracle.property.test.ts tests/query/load-subset-full-flow-oracle.property.test.ts tests/query/load-subset-projection-oracle.property.test.ts tests/query/includes-oracle.property.test.ts tests/query/includes-collection-oracle.property.test.ts tests/query/includes-cross-formulation-oracle.property.test.ts tests/query/includes-temporal-oracle.test.ts tests/query/includes-optimistic-oracle.property.test.ts tests/query/includes-publication-oracle.test.ts tests/query/includes-query-shape-oracle.test.ts tests/query/includes-work-counter-oracle.test.ts tests/query/includes-context-transport-oracle.test.ts"
"test:oracles": "vitest --run tests/collection-sync-reentrancy.test.ts tests/collection-subscription-replay-oracle.property.test.ts tests/query/coverage-registry-oracle.property.test.ts tests/query/load-subset-full-flow-oracle.property.test.ts tests/query/load-subset-projection-oracle.property.test.ts tests/query/includes-oracle.property.test.ts tests/query/includes-collection-oracle.property.test.ts tests/query/includes-cross-formulation-oracle.property.test.ts tests/query/includes-temporal-oracle.test.ts tests/query/includes-optimistic-oracle.property.test.ts tests/query/includes-publication-oracle.test.ts tests/query/includes-query-shape-oracle.test.ts tests/query/includes-work-counter-oracle.test.ts tests/query/includes-context-transport-oracle.test.ts tests/query/pagination-oracle.property.test.ts"
},
"type": "module",
"main": "dist/cjs/index.cjs",
Expand Down
70 changes: 30 additions & 40 deletions packages/db/src/collection/change-events.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,9 @@ import {
optimizeExpressionWithIndexes,
} from '../utils/index-optimization.js'
import { ensureIndexForField } from '../indexes/auto-index.js'
import { makeComparator } from '../utils/comparison.js'
import { ReverseIndex } from '../indexes/reverse-index.js'
import { buildCompareOptions } from '../query/compiler/order-by'
import { TotalOrder } from '../query/total-order.js'
import type {
ChangeMessage,
CollectionLike,
Expand Down Expand Up @@ -358,7 +359,30 @@ function getOrderedKeys<T extends object, TKey extends string | number>(
// Take the keys that match the filter and limit
// if no limit is provided `index.keyCount` is used,
// i.e. we will take all keys that match the filter
return index.takeFromStart(limit ?? index.keyCount, filterFn)
if (!(index instanceof ReverseIndex)) {
return index.takeFromStart(limit ?? index.keyCount, filterFn)
}

// Reversing a value index also reverses keys inside an equal-value
// bucket, but query TotalOrder keeps its public-key tie-break ascending.
// Refine all matching indexed rows locally so a limit cannot cut the
// wrong side of a tied boundary.
const totalOrder = new TotalOrder(orderBy, collection)
const indexedEntries = index
.takeFromStart(index.keyCount, filterFn)
.flatMap((key) => {
const value = collection.get(key)
return value === undefined ? [] : [{ key, value }]
})
indexedEntries.sort((left, right) =>
totalOrder.compareEntries(
[left.key, left.value],
[right.key, right.value],
),
)
return indexedEntries
.slice(0, limit ?? indexedEntries.length)
.map(({ key }) => key)
}
}
}
Expand All @@ -375,24 +399,10 @@ function getOrderedKeys<T extends object, TKey extends string | number>(
}
}

// Sort using makeComparator
const compare = (a: { key: TKey; value: T }, b: { key: TKey; value: T }) => {
for (const clause of orderBy) {
const compareFn = makeComparator(clause.compareOptions)

// Extract values for comparison
const aValue = extractValueFromItem(a.value, clause.expression)
const bValue = extractValueFromItem(b.value, clause.expression)

const result = compareFn(aValue, bValue)
if (result !== 0) {
return result
}
}
return 0
}

allItems.sort(compare)
const totalOrder = new TotalOrder(orderBy, collection)
allItems.sort((left, right) =>
totalOrder.compareEntries([left.key, left.value], [right.key, right.value]),
)
const sortedKeys = allItems.map((item) => item.key)

// Apply limit if provided
Expand All @@ -403,23 +413,3 @@ function getOrderedKeys<T extends object, TKey extends string | number>(
// if no limit is provided, we will return all keys
return sortedKeys
}

/**
* Helper function to extract a value from an item based on an expression
*/
function extractValueFromItem(item: any, expression: BasicExpression): any {
if (expression.type === `ref`) {
const propRef = expression
let value = item
for (const pathPart of propRef.path) {
value = value?.[pathPart]
}
return value
} else if (expression.type === `val`) {
return expression.value
} else {
// It must be a function
const evaluator = compileSingleRowExpression(expression)
return evaluator(item as Record<string, unknown>)
}
}
80 changes: 75 additions & 5 deletions packages/db/src/collection/state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,11 @@ import { SortedMap } from '../SortedMap'
import { enrichRowWithVirtualProps } from '../virtual-props.js'
import { SyncTransactionAbortedError } from '../errors.js'
import { createDeferred } from '../deferred'
import {
copySyncRequestProvenance,
getSyncRequestProvenance,
setSyncRequestProvenance,
} from '../load-subset-request-provenance.js'
import { DIRECT_TRANSACTION_METADATA_KEY } from './transaction-metadata.js'
import type {
VirtualOrigin,
Expand Down Expand Up @@ -42,6 +47,8 @@ interface PendingSyncedTransaction<
deletes: Set<TKey>
}
preserveHydrationSeedKeys?: boolean
/** Exact subset request whose commit produced this transaction. */
requestSignal?: AbortSignal
/**
* When true, this transaction should be processed immediately even if there
* are persisting user transactions. Used by manual write operations (writeInsert,
Expand Down Expand Up @@ -339,13 +346,15 @@ export class CollectionStateManager<
: this.enrichWithVirtualProps(change.previousValue, change.key)
: undefined

return {
const enriched = {
key: change.key,
type: change.type,
value: enrichedValue,
previousValue: enrichedPreviousValue,
metadata: change.metadata,
} as ChangeMessage<WithVirtualProps<TOutput, TKey>, TKey>
copySyncRequestProvenance(change, enriched)
return enriched
}

/**
Expand Down Expand Up @@ -940,13 +949,50 @@ export class CollectionStateManager<
const changedKeys = new Set<TKey>()
for (const transaction of committedSyncedTransactions) {
for (const operation of transaction.operations) {
changedKeys.add(operation.key as TKey)
const key = operation.key as TKey
changedKeys.add(key)
}
for (const [key] of transaction.rowMetadataWrites) {
changedKeys.add(key)
}
}

type AppliedRequestProvenance = {
version: {
value: TOutput | undefined
origin: VirtualOrigin | undefined
}
hasOrdinarySource: boolean
requestSignals: Set<AbortSignal>
}
const requestProvenanceByKey = new Map<TKey, AppliedRequestProvenance>()
const provenanceForSignal = (signal: AbortSignal | undefined) => ({
hasOrdinarySource: signal === undefined,
requestSignals:
signal === undefined
? new Set<AbortSignal>()
: new Set<AbortSignal>([signal]),
})
const recordRequestProvenance = (
key: TKey,
signal: AbortSignal | undefined,
) => {
const version = {
value: this.syncedData.get(key),
origin: this.rowOrigins.get(key),
}
const previous = requestProvenanceByKey.get(key)
if (previous !== undefined && deepEquals(previous.version, version)) {
if (signal === undefined) previous.hasOrdinarySource = true
else previous.requestSignals.add(signal)
return
}
requestProvenanceByKey.set(key, {
version,
...provenanceForSignal(signal),
})
}

const virtualSnapshotKeys = new Set(changedKeys)
for (const key of this.pendingOptimisticDirectUpserts) {
virtualSnapshotKeys.add(key)
Expand Down Expand Up @@ -1013,7 +1059,16 @@ export class CollectionStateManager<
truncateOptimisticSnapshot?.upserts.get(key) ||
this.syncedData.get(key)
if (previousValue !== undefined) {
events.push({ type: `delete`, key, value: previousValue })
const event: ChangeMessage<TOutput, TKey> = {
type: `delete`,
key,
value: previousValue,
}
setSyncRequestProvenance(
event,
provenanceForSignal(transaction.requestSignal),
)
events.push(event)
}
}

Expand Down Expand Up @@ -1105,6 +1160,7 @@ export class CollectionStateManager<
this.pendingOptimisticDirectDeletes.delete(key)
break
}
recordRequestProvenance(key, transaction.requestSignal)
if (!transaction.preserveHydrationSeedKeys) {
this.hydrationSeedKeys.delete(key)
this.hydratedKeys.delete(key)
Expand All @@ -1114,9 +1170,9 @@ export class CollectionStateManager<
for (const [key, metadataWrite] of transaction.rowMetadataWrites) {
if (metadataWrite.type === `delete`) {
this.syncedMetadata.delete(key)
continue
} else {
this.syncedMetadata.set(key, metadataWrite.value)
}
this.syncedMetadata.set(key, metadataWrite.value)
}

for (const [
Expand Down Expand Up @@ -1226,6 +1282,7 @@ export class CollectionStateManager<
this.isThisCollection(mutation.collection) &&
mutation.optimistic
) {
requestProvenanceByKey.delete(mutation.key)
switch (mutation.type) {
case `insert`:
case `update`:
Expand Down Expand Up @@ -1261,13 +1318,15 @@ export class CollectionStateManager<
this.pendingOptimisticUpserts.delete(key)
this.pendingLocalOrigins.delete(key)
}
requestProvenanceByKey.delete(key)
}
for (const key of this.pendingOptimisticDirectDeletes) {
if (!changedKeys.has(key)) {
changedKeys.add(key)
}
this.pendingOptimisticDeletes.delete(key)
this.pendingLocalOrigins.delete(key)
requestProvenanceByKey.delete(key)
}
this.pendingOptimisticDirectUpserts.clear()
this.pendingOptimisticDirectDeletes.clear()
Expand Down Expand Up @@ -1379,6 +1438,17 @@ export class CollectionStateManager<
}
}

for (const event of events) {
if (getSyncRequestProvenance(event) !== undefined) continue
const provenance = requestProvenanceByKey.get(event.key)
if (provenance !== undefined) {
setSyncRequestProvenance(event, {
hasOrdinarySource: provenance.hasOrdinarySource,
requestSignals: new Set(provenance.requestSignals),
})
}
}

// Update cached size after synced data changes
this.size = this.calculateSize()

Expand Down
Loading
Loading