diff --git a/.changeset/db-ivm-nan-index-prefix.md b/.changeset/db-ivm-nan-index-prefix.md new file mode 100644 index 000000000..683337945 --- /dev/null +++ b/.changeset/db-ivm-nan-index-prefix.md @@ -0,0 +1,5 @@ +--- +'@tanstack/db-ivm': patch +--- + +Fix joins over rows keyed by `NaN`. The index compared row-key prefixes with `===`, which treats `NaN` as unequal to itself, while the `Map` holding the prefixes treats it as one key. A retracted `NaN`-keyed row then never cancelled, so the join published a duplicate row or threw `Mismatching prefixes`. diff --git a/.changeset/joined-result-keys-json.md b/.changeset/joined-result-keys-json.md new file mode 100644 index 000000000..c752bd1bc --- /dev/null +++ b/.changeset/joined-result-keys-json.md @@ -0,0 +1,5 @@ +--- +'@tanstack/db': patch +--- + +Fix joined live queries that dropped or rejected a row when two source key pairs printed alike, such as (`a,b`, `c`) and (`a`, `b,c`), or `1` and `'1'`. Joined result keys are now JSON arrays of the two source keys, with `null` for a missing side: `["a","b"]`, `[4,1]`, `[4,null]`. An infinite numeric key is written as an object, such as `[{"number":"Infinity"},1]`, so it cannot collide with a missing side. Code that looks up joined rows by a hand-built key, such as `collection.get('[1,2]')`, must use the new format. diff --git a/.changeset/perf-many-filtered-live-queries.md b/.changeset/perf-many-filtered-live-queries.md new file mode 100644 index 000000000..d11323746 --- /dev/null +++ b/.changeset/perf-many-filtered-live-queries.md @@ -0,0 +1,5 @@ +--- +'@tanstack/db': patch +--- + +Speed up apps that mount many small filtered live queries. Queries without includes keep their compiled pipeline instead of paying for include materialization, eager subscriptions no longer build an abort error on every unsubscribe, a filtered subscription skips source batches that cannot match its `eq` condition, and unindexed snapshots reject rows by that condition before copying them. With 240 `eq`-filtered live queries, mounting is about 2x faster (indexed) to 2.5x faster (unindexed), and a 50-row update batch is about 2.7x faster. diff --git a/docs/contributing/oracle-coverage.md b/docs/contributing/oracle-coverage.md index 86869f382..361d90871 100644 --- a/docs/contributing/oracle-coverage.md +++ b/docs/contributing/oracle-coverage.md @@ -61,6 +61,9 @@ that test identifiers must copy production's private data structures. | Query DB and observer | Complete | Query-scope row ownership, subset identity and cancellation, failure and recovery, and the per-listener eligibility ledger are literate. | | Ordered acquisition | Complete | Exact demand identity, applied settlement, independent pagination recomputation, request work, lifecycle products, replay authority, source-generation readiness, and transaction-refinement abort boundaries are literate. | | Indexed predicate filtering | Complete for bounded dotted-path collision | An independent row filter checks direct predicates and public callback predicates for a dotted scalar field beside a nested path, with and without an index. Selected-field ordering checks the same distinction through `$selected`. | +| WHERE predicate publication | Complete for bounded predicate and sync-transaction grammar | The contract, Kleene reference evaluator, snapshot and change-history grammar, three subscriber consumers plus direct snapshot, and per-commit key-set refinement check are literate. A fixed same-key update checks live-row and direct-subscriber payloads; a focused descriptor boundary checks stored-row prefilter safety. Generated cleanup and restart histories for filtered subscribers remain open. | +| Joined result keys | Complete for bounded two-source key grammar | The contract, nested-loop pair model, delimiter-, number-like, infinite, and `NaN` key grammar, public join driver, and per-checkpoint pair and key-count check are literate. Joins over subqueries, more than two sources, custom `getKey`, and optimistic mutations remain outside this owner. | +| D2 Index storage | Complete for bounded prefix grammar | The contract, plain-`Map` multiset model, prefixed and unprefixed value grammar, `Index` driver, and per-addition `get`/`has` check are literate. Compaction, presence tracking, and structural payloads remain outside this owner. | | Lazy target path identity | Focused compiler boundary | A same-source union/coalesce witness keeps dotted and nested demand paths distinct during target deduplication. | | Correlated include path identity | Focused public route-context witnesses | One-level and nested includes keep dotted and nested parent paths, including ancestor aliases, distinct. Conditional result paths receive separate routes. Fixed fixtures cover initial reads and selected source updates; other recursive source forms and arbitrary path segments remain outside this witness. | | Join equality and cold acquisition | Complete | Independent cold relational recomputation, acquisition evidence, established equality domains, replacements, and scan/index routes are literate. | @@ -128,6 +131,9 @@ comment and the current API/architecture contract before extending its model. | 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), [Vue synchronous publication](https://github.com/TanStack/db/blob/main/packages/vue-db/tests/useLiveQuery-publication-oracle.test.ts) | Exact exposed rows/pages and each framework's own lifecycle cuts. Vue's synchronous watcher checks insert, update, and nonterminal delete through supplied Collection and identity query inputs; it does not cover an empty final result, multiple changes in one transaction, or query recompilation. A React witness does not prove Vue/Solid/Angular/Svelte scheduling. Preserve their receiving registrations. | | Window-controller overlap | [shared infinite-query conformance](https://github.com/TanStack/db/blob/main/packages/db/tests/conformance/infinite-suite.ts), [React driver](https://github.com/TanStack/db/blob/main/packages/react-db/tests/infinite-query-conformance.test.tsx), [Vue driver](https://github.com/TanStack/db/blob/main/packages/vue-db/tests/infinite-query-conformance.test.ts), [Svelte driver](https://github.com/TanStack/db/blob/main/packages/svelte-db/tests/infinite-query-conformance.svelte.test.ts), [review record](oracle-reviews/dec-03-window-overlap.md) | Four or five ordered source rows, two window-settlement orders, and public snapshots after each settlement, a subscribed observer update, a detached controller read, and explicit recovery. Insertion and removal of the fifth row distinguish both continuation directions while an earlier error remains visible. The detached read follows a source change with no controller subscriber; its preloaded Collection retains the five-row window. The exact overlap uses the exported DB controller in each package realm. Framework hooks expose no controller preload, so this cell does not establish a hook scheduling path; direct hook paging remains in the adjacent shared scenarios. An unsubscribed controller has no notification claim. | | Structural values and ordered primitives | [hash values](https://github.com/TanStack/db/blob/main/packages/db-ivm/tests/hash.property.test.ts), [hash identity](https://github.com/TanStack/db/blob/main/packages/db-ivm/tests/hash-identity-oracle.property.test.ts), [MultiSet consolidation](https://github.com/TanStack/db/blob/main/packages/db-ivm/tests/multiset-consolidate-oracle.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), [BTree Map](https://github.com/TanStack/db/blob/main/packages/db/tests/btree-map-oracle.test.ts), [query identity](https://github.com/TanStack/db/blob/main/packages/db/tests/query/identity-output-shape-oracle.test.ts), [LIKE semantics](https://github.com/TanStack/db/blob/main/packages/db/tests/query/compiler/evaluators.test.ts) | Independent flat values, graph topology, algebraic laws, Map/group/sort recomputation, expression denotation, custom comparator dispatch/reference identity/index compatibility, LIKE wildcard refinement, compiled output bags, declared value-pair agreement between D2 equality keys and IR equality operands, and Collection collation with exact public scan/index order and one-index reuse across repeated reads. The auto-index cells cross BasicIndex and BTreeIndex; scan cells omit indexes. Both paths cover lexical and numeric `en-US` locale defaults plus an explicit ordinary locale override of a numeric default. An unindexed six-row work cell bounds Collection collation resolution to at most two reads per ordered scan: one index-path check and one scan clause. Other locales and ordered histories are outside this owner. The identity pair grammar covers primitives, signed zero, NaN/invalid Date, valid Date/timestamp, binary/Buffer, Temporal, nested references, symbols, and functions; fixed and sampled pairs do not prove every JavaScript value or hash collision freedom. Ref proxies are expressions rather than literal values in this lane. The hash-values owner samples Set/object and Map/object type-marker distinctions in fixed and random campaigns with seed-and-path replay; collision freedom is not promised. The hash-identity owner checks `hash` and `equalHashValues` against one spec-level identity model: primitives with signed zero, NaN, and bigint; reference leaves (functions, registered handles, `File`, binary over 128 bytes); Date, binary, Temporal, and RegExp headers; array length and holes; Map and Set insertion order; and object keys in any order, for acyclic values and for cycles through arrays, Map values, Sets, and plain objects. Map/Set order sensitivity is pinned current behavior, not a promise. Arrays from another realm, RegExps with a `NaN` `lastIndex`, getters, proxies, other cycle carriers, and pairs that differ only in a back-edge target are outside its grammar; pinned cases keep equality's cross-realm length check and strict RegExp field comparison. `hash-work.test.ts` pins how visited values and array/RegExp headers count toward the structural work cap, and that equality copies a shared Map's entries once per distinct pair ([code-weight hash dispatch review](oracle-reviews/code-weight-hash-dispatch.md)). Custom string indexes do not optimize ordinary ranges, and custom cursors remain unsupported. The LIKE owner covers boolean string matching and bounded work, not nullish three-valued logic or a general Unicode collation contract. Unsupported composite cursors reject. The MultiSet consolidation owner checks keyed, single-type, and structural identity, signed zero and NaN in keyed and single-number data, the retained first record, zero-sum removal, and input record contents. It does not assert output order. Its keyed grammar excludes the known keyed identity collisions until #1948 is fixed: cross-type values or keys with the same text, the `\|` delimiter, and symbol or function values ([code-weight consolidation review](oracle-reviews/code-weight-multiset-consolidation.md)). The index owner enforces one invalid-comparator law across both BasicIndex and BTreeIndex: for each comparator, accepted numeric prefix (including the empty prefix), and probe operation (add, update, eq/range lookup, take, build), the operation throws exactly when it receives a `NaN` or non-number result and never when every comparison is valid. A `signed infinity` comparator is the valid control, and a custom collation `compare` supplied through `compareOptions` is checked the same way. Two BTree Map split histories check inserted-key reads, exact whole-range payloads/order, size, and extrema, proving split placement follows the insertion index without a post-mutation comparison. A rejected index add or remove leaves the index refining its accepted rows; update and build are not atomic. At the Collection boundary, both index types and both write paths (optimistic insert, sync commit) crash the collection: status becomes `error` and the next mutation throws, so a usable collection never holds rows its subscribers were not told about. Open: a sync source can still commit writes on a collection in `error` state, which stores unpublished rows; this is existing `markError` behavior and needs a lifecycle-owner witness. Comparator transitivity is outside this evidence. | +| WHERE predicate publication | [WHERE predicate publication oracle](https://github.com/TanStack/db/blob/main/packages/db/tests/query/where-predicate-publication-oracle.property.test.ts) | An independent Kleene evaluator judges `eq`/`not`/`and`/`or` predicates over strings, a normalization-prefixed string, booleans, numbers, `NaN`, a valid Date, `null`, a missing field, and virtual fields. Snapshot histories with pending optimistic inserts compare a live query, direct subscribers with and without initial state, and `currentStateAsChanges`; change histories apply multi-key sync transactions, including reinsertion of a key deleted earlier, and compare every consumer's key set after each commit, on scan and `BasicIndex` paths. Peer subscribers share one field with different literals, plus one on a virtual field, so every published change must reach exactly the subscribers whose predicate it can satisfy; routing mutants that ignored the previous value, ignored the field, delivered a change twice, or routed while stale published rows awaited reconciliation fail here. Fixed witnesses cover a layout-only publication reaching a filtered subscriber as one empty batch, a filtered subscriber's empty Collection-readiness batch, retraction of a vanished row after eager cleanup and restart, and the new row plus exact subscriber update payload when a same-key row stays TRUE. A stale-live-value mutant survives the generated membership checks but fails the payload witness. A focused [property-visibility test](https://github.com/TanStack/db/blob/main/packages/db/tests/query/where-prefilter-property-visibility.test.ts) compares the public unindexed snapshot with the enriched-row contract for inherited, non-enumerable, enumerable-own, and nested getter paths; the pre-fix stored-row shortcut failed three of its four controls. Other property-descriptor and stateful-getter histories remain outside those fixed controls. Hostile mutants for FALSE-for-UNKNOWN `eq`, `or` or number-literal subscription prefilters, a prefilter that ignores the previous value, a skipped readiness batch, a skip while stale published rows await reconciliation, and an unindexed snapshot scan that reads only synced rows while an optimistic insert or delete changes visibility passed the prior `@tanstack/db` suite and fail here. Change routing withholds a dropped row from an `eq` subscription before its sent key could be recorded, so a mutant that records dropped rows as sent is equivalent for routed predicates in this oracle. A focused `collection-subscription.test.ts` witness checks that a dropped insert or update, alone or beside a matching change, cannot advance a limited subscription's page offset under a routed `eq` and an unrouted `or` predicate; the unrouted cases beside a matching change kill that mutant, and the unrouted `alone` case does not. A loose-equality prefilter is an equivalent mutant: it only skips less. Skipping during truncate replay also survives; the replay's baseline diff re-derives the same retraction, so the guard keeps the prior dataflow without its own witness. Generated cleanup and restart histories for filtered subscribers remain open for the lifecycle publication owner. No-op updates, other comparison operators, Temporal and binary operands, joins, ordering, optimistic updates, and truncate are outside this owner. It runs in `@tanstack/db`'s `test:oracles` campaign. The [2026-09-30 review](https://github.com/TanStack/db/blob/main/docs/contributing/oracle-reviews/2026-09-30-where-predicate-and-join-keys.md) records each ORC outcome. | +| Joined result keys | [Joined result key oracle](https://github.com/TanStack/db/blob/main/packages/db/tests/query/join-result-key-oracle.property.test.ts) | A nested-loop model recomputes the pair set of an inner, left, or full join over source keys drawn from plain, comma-bearing, bracket-bearing, and quoted strings, numbers beside the strings that print the same, both infinities, and `NaN`. After preload and each synced group change, the published rows must equal the model's pairs and the result key count must equal the row count. The comma-joined key encoding fails its two pinned histories and both campaigns; plain `JSON.stringify`, which prints infinities as the missing-side `null`, fails the pinned infinity history and both campaigns. A pinned history in which a `NaN`-keyed pair leaves and re-forms fails when the join index compares source-key prefixes with `===`. Joins over subqueries, more than two sources, custom `getKey`, and optimistic mutations are outside this owner. It runs in `@tanstack/db`'s `test:oracles` campaign. The [2026-09-30 review](https://github.com/TanStack/db/blob/main/docs/contributing/oracle-reviews/2026-09-30-where-predicate-and-join-keys.md) records each ORC outcome. | +| D2 Index storage | [Index refinement oracle](https://github.com/TanStack/db/blob/main/packages/db-ivm/tests/index-refinement-oracle.property.test.ts) | A plain-`Map` model sums multiplicities under a type-and-value identity in which `NaN` equals `NaN` and `-0` equals `0`. Histories of up to twelve additions to two keys mix prefixed arrays with `NaN`, `0`, `-0`, `1`, `'1'`, `1n`, and `'a'` prefixes and unprefixed values, with cancellation and reappearance. After each addition, `get` and `has` must equal the model. Comparing prefixes with `===`, which disagrees with the `Map` that holds them, fails three pinned histories and both campaigns; the `-0` control passes under both. Compaction, presence tracking, and structural payloads are outside this owner; the join operator tests and incrementalization law own operator behavior. | | Index suggestions | [collection-size suggestion oracle](https://github.com/TanStack/db/blob/main/packages/db/tests/index-suggestion-oracle.test.ts) | A public filtered live-query Collection checks manual/default, eager, threshold, matching-index, and unrelated-index cases against the documented collection-size suggestion policy. The original manual-silence/eager-noise implementation failed two expected observations. The owner does not judge repeated-query warning volume, slow-query timing, query work, production-build suppression, or every query shape. It runs in `@tanstack/db`'s `test:oracles` campaign. | | Minified DB public API | [consumer bundle check](https://github.com/TanStack/db/blob/main/scripts/test-minified-db.mjs) | CI bundles the public `@tanstack/db` entry with db-ivm inlined through esbuild `minify: true`. It checks every exported error class's current public `name`, the built-in `BasicIndex` resolver name in index metadata and its event, complete query rows (after removing the four documented virtual fields) against an independent array filter/sort/projection, and one live update. The `new.target.name`, public-member mangle, and extra-output-field calibration modes fail at the intended observations. This fixed slice does not replace the unminified generated oracles, exercise framework adapters, or establish stable names for custom index resolvers. | | Boundary refinements | [cleanup/restart](https://github.com/TanStack/db/blob/main/packages/db/tests/collection-cleanup-restart-oracle.test.ts), [issue #1891 async-cleanup review](oracle-reviews/issue-1891-async-cleanup.md), [issue #1891 cleanup-start follow-up](oracle-reviews/issue-1891-cleanup-start-followup.md), [metadata publication](https://github.com/TanStack/db/blob/main/packages/db/tests/collection-metadata-publication-oracle.property.test.ts), [state retention](https://github.com/TanStack/db/blob/main/packages/db/tests/collection-state-retention-oracle.property.test.ts), [acquisition cells](https://github.com/TanStack/db/blob/main/packages/db/tests/collection-subscription-lifecycle-oracle.test.ts), [D2 source reconciliation](https://github.com/TanStack/db/blob/main/packages/db/tests/d2-source-reconciliation-oracle.property.test.ts), [top-K support windows](https://github.com/TanStack/db/blob/main/packages/db-ivm/tests/operators/topk-support-window-oracle.test.ts), [nested Query work](https://github.com/TanStack/db/blob/main/packages/query-db-collection/tests/includes-work-counter-oracle.test.ts) | Explicit lifecycle products, independent source maps and weighted relations, exact publication cuts, support/multiplicity, and value-plus-work observations. These refine the larger subsystem models; they do not replace them. | @@ -180,6 +186,16 @@ remove or detach ownership before invoking adapter unload. These are bounded equivalence arguments, not permission to delete the guards without a separate review. The companion does not establish arbitrary provider callback behavior. +Row canonicalization before the result boundary now runs only for DISTINCT, +which tracks visibility by selected value; removing it fails the functional +projection oracle. Ordering already consolidates each key's batch before top-K +state changes, and materialized relations reduce by public key. Removing the +former parent-route, non-Collection source, join, grouping, ordering, and +includes conditions passed the whole `@tanstack/db` suite and a +`TANSTACK_DB_ORACLE_RUNS_MULTIPLIER=10` oracle campaign. Without the keyed +reduction, `live-query-result-multiplicity.test.ts` witnesses the output +boundary rejecting a flush that changes one key by more than one row. + The Collection lifecycle publication owner checks 32 bounded histories in which cleanup interrupts nested publication deferrals and a restarted sync run later publishes or discards one or two source rows. It compares the next subscriber's diff --git a/docs/contributing/oracle-reviews/2026-09-30-where-predicate-and-join-keys.md b/docs/contributing/oracle-reviews/2026-09-30-where-predicate-and-join-keys.md new file mode 100644 index 000000000..4fac3e047 --- /dev/null +++ b/docs/contributing/oracle-reviews/2026-09-30-where-predicate-and-join-keys.md @@ -0,0 +1,104 @@ +# WHERE predicate publication and joined result key oracle review + +## Reviewed state and claim + +Base: `49dc79d4`, the head of pull request #1956 before this review. This +record reviews the Git tree that contains it. The final commit or pull request +identifies that immutable tree. The review applied the ORC-001 through ORC-014 +requirements of the oracle guide revision that adds ORC-013 and ORC-014. + +The WHERE predicate publication oracle claims that a filtered live-query +Collection, direct subscriptions with and without initial state, peer +subscriptions routed by one `eq` field, and `currentStateAsChanges` publish +exactly the rows whose predicate is TRUE under SQL three-valued logic. The +claim covers the grammar in the oracle's opening prose. It does not cover +comparison operators other than `eq`, joins, ordering, optimistic updates or +deletes, truncate, failed replay, or generated cleanup and restart histories. + +The joined result key oracle claims that an inner, left, or full two-source +join publishes one row per joined pair and that distinct pairs never share a +result key, for string keys, numeric keys, both infinities, and `NaN`. It does +not cover joins over subqueries, more than two sources, custom `getKey`, or +optimistic mutations. + +Both oracles use `mockSyncCollectionOptions` as a controlled provider. Their +claims are limited to how the Collection and compiler handle the sync +transactions that provider supplies; neither claims a real adapter's behavior. + +## RED and GREEN evidence + +Mutants ran on the reviewed tree through the oracle files alone. Each file was +restored after each run. + +| Mutant | Oracle | Outcome | +| --- | --- | --- | +| `eq` returns FALSE for a nullish operand | WHERE | Assertion failure in the pinned snapshot tests and both campaigns (5 tests). | +| Routing ignores a change's previous value | WHERE | Assertion failure in the change pinned tests and both campaigns (4 tests). | +| Routing continues while stale published rows await reconciliation | WHERE | Assertion failure in the restarted-source witness. | +| Routing treats `or` operands as conjuncts | WHERE | Assertion failure in the change pinned tests and a campaign (2 tests). | +| A dropped insert or update is recorded as sent | WHERE | Survived: equivalent in this oracle's domain. Routing withholds a dropped row from an `eq` subscription before any key is recorded, and unrouted predicates never skip the delete that would expose a stale record. | +| The same sent-key mutant | `collection-subscription.test.ts` page-offset witness | Survived the pre-review `eq` witness. After the review added an unrouted `or` predicate, the two cases beside a matching change fail by assertion (offset 3, expected 2). The unrouted `alone` case survives. | +| Result key joins the source keys with a comma | Join | Key-count invariant error or wrong-row assertion in both pinned delimiter and number histories and both campaigns (4 tests). | +| Result key uses plain `JSON.stringify` | Join | Wrong-row assertion in the pinned infinity history and both campaigns (3 tests). | +| Index compares source-key prefixes with `===` | Join | Key-count invariant error in the pinned `NaN` history and the fixed campaign (2 tests). | +| The same prefix comparison | Index refinement oracle | Assertion failure in three pinned `NaN` histories and both campaigns; the `-0` control passes. | + +The plain-JSON mutant is the pre-review implementation. The review found that +`JSON.stringify` prints `Infinity` and `-Infinity` as `null`, the marker for a +missing outer side. Adding both infinities to the key domain failed the fixed +and random campaigns with a full join over `-Infinity` on both sides. The repair +encodes a non-finite number as an object, which no source key can be. `NaN` +source keys failed before key encoding mattered, with the comma encoding as +well, because the join index compared source-key prefixes with `===`, so a +retracted `NaN`-keyed row never cancelled. The review fixed that comparison in +`@tanstack/db-ivm` and added an Index refinement oracle for it; the join oracle +now includes `NaN` keys and a pinned history in which a `NaN`-keyed pair leaves +and re-forms. + +After the repairs, `test:oracles` at `TANSTACK_DB_ORACLE_RUNS_MULTIPLIER=10` +passed 2,869 tests and the `@tanstack/db` suite passed 7,499 tests. Final +verification receipts belong in the pull request. + +## WHERE predicate publication oracle + +| Requirement | Outcome | +| --- | --- | +| ORC-001 Contract authority and limits | Pass. The evaluator contract in `src/query/compiler/evaluators.ts` states the operand rules. The opening prose lists the omissions and their owners. | +| ORC-002 Independent judgment | Pass. `expectedTruth` and `expectedVisible` evaluate plain objects without production comparison, normalization, prefilter, routing, or virtual-field helpers. | +| ORC-003 Distinguishable responsibilities | Pass. The contract, model, history grammar, production driver, and refinement check are separate marked sections. | +| ORC-004 Generated-history controls | Pass after repair. Reconstruction: the pinned snapshot with `eq($synced, false)` was outside the grammar because `false` was not a literal; the review added it, and every pinned history is now in the grammar. Ablation: removing `not` loses the nullish-comparison history, removing `or` loses the routing history, removing multi-operation transactions loses retraction inside one batch, removing insert reuse loses reinsertion, and removing the index axis loses the index path. Range: depth three, at most six rows, six transactions of three operations, and marginal values `NaN`, a valid Date, the normalization prefix, `null`, and a missing field. Exclusion: an update to an equivalent value and a second update to one key in one transaction are removed from the history. | +| ORC-005 Production path and observation | Pass. Public builder functions, `subscribeChanges`, and `currentStateAsChanges` run against a real Collection. Each consumer's exact key set is compared after the subscriptions attach and after each sync transaction commits. | +| ORC-006 Checker calibration | Pass. The checker test requires the two-valued answer to fail. The mutant table classifies each run. | +| ORC-007 Fixed/random replay | Pass. The fixed seed `44_500_301`, the unseeded campaign, and the replay entry share one property and budget. | +| ORC-008 Stateful-model minimality | Pass. The model keeps a row map across sync transactions. `$synced` and `$origin` separate a pending optimistic insert from a synced row, and the pinned `eq($synced, false)` snapshot distinguishes them. `$key` always equals the row id; it stays because the grammar reads it as an operand. The touched-key set distinguishes a subscription without initial state from one with it. | +| ORC-009 Vocabulary mapping | Pass after repair. The prose now uses live-query Collection and sync transaction. A subscriber is the callback of one subscription; a consumer is the oracle's label for one observed endpoint. | +| ORC-010 Failure fidelity and cleanup | Pass after repair. Cleanup ran in `finally`, so a cleanup failure replaced the check failure. `withOracleCleanup` now runs every cleanup step and keeps the check failure as the `cause` of an `AggregateError` when cleanup also fails. | +| ORC-011 Independent second formulation | Not triggered. No review has named a fault that the Kleene model and production could share. | +| ORC-012 Review evidence | This record. The coverage map links it. | +| ORC-013 Reusable boundary law | Pass. FALSE-for-UNKNOWN is rejected by the negated nullish snapshot. Ignoring the previous value is rejected by generated change histories. Treating `or` operands as conjuncts is rejected by the pinned `or` change history. Routing past stale rows is rejected by the restarted-source witness. | +| ORC-014 Controlled-premise handoff | Not triggered. The claim is limited to the controlled provider's sync transactions. | + +## Joined result key oracle + +| Requirement | Outcome | +| --- | --- | +| ORC-001 Contract authority and limits | Pass after repair. The prose now cites the live-query guide: joins behave like SQL joins, and a join result has a composite key of the parent keys. The guide does not fix the key format. | +| ORC-002 Independent judgment | Pass. `expectedPairs` is a nested loop over plain arrays and does not import the compiler or its key encoding. | +| ORC-003 Distinguishable responsibilities | Pass. The contract, model, history grammar, production driver, and refinement check are separate marked sections. | +| ORC-004 Generated-history controls | Pass after repair. Reconstruction: the pinned number history used right key `x`, which was outside the key domain; it now uses `c`, and all three pinned histories are in the grammar. Ablation: removing delimiter strings, number and string twins, or infinities loses a collision class; removing full joins loses unmatched rows on the right; removing synced changes loses unmatched rows created by a group move. Range: at most five rows a side, groups 0 through 2, and three changes. Exclusion: keys are unique within one side, and a change that keeps a row's group is skipped. | +| ORC-005 Production path and observation | Pass. `createLiveQueryCollection` compiles a real join. Published pairs and the key count are compared after preload and after each synced change. | +| ORC-006 Checker calibration | Pass. The comma and plain-JSON mutants fail, as classified above. | +| ORC-007 Fixed/random replay | Pass. The fixed seed `44_501_962`, the unseeded campaign, and the replay entry share one property and budget. | +| ORC-008 Stateful-model minimality | Not triggered. The model recomputes the pairs from the current rows. | +| ORC-009 Vocabulary mapping | Pass. A pair is one published row of the live-query Collection, described by its two source keys. | +| ORC-010 Failure fidelity and cleanup | Pass after repair, through `withOracleCleanup`. | +| ORC-011 Independent second formulation | Not triggered. No shared-fault hypothesis has been named. | +| ORC-012 Review evidence | This record. The coverage map links it. | +| ORC-013 Reusable boundary law | Pass. The law is that distinct pairs have distinct keys and a retracted pair cancels. The comma encoding is rejected by the delimiter and number pinned histories, plain JSON by the infinity pinned history, and `===` prefixes by the `NaN` pinned history. | +| ORC-014 Controlled-premise handoff | Not triggered. The claim is limited to how the compiler keys rows from the controlled provider. | + +## Open work + +- The sent-key mutant survives the unrouted `alone` page-offset case. +- Generated cleanup and restart histories for filtered subscriptions remain + with the lifecycle publication owner, as the coverage map records. diff --git a/packages/db-ivm/src/indexes.ts b/packages/db-ivm/src/indexes.ts index 07c0fa5d5..44cbd59db 100644 --- a/packages/db-ivm/src/indexes.ts +++ b/packages/db-ivm/src/indexes.ts @@ -69,7 +69,7 @@ class PrefixMap extends Map< const [currentValue, currentMultiplicity] = valueMapOrSingleValue const currentPrefix = getPrefix(currentValue) - if (currentPrefix !== prefix) { + if (!isSamePrefix(currentPrefix, prefix)) { throw new Error(`Mismatching prefixes, this should never happen`) } @@ -383,7 +383,7 @@ export class Index { // Check if they're the same value by prefix/suffix comparison if ( - currentPrefix === newPrefix && + isSamePrefix(currentPrefix, newPrefix) && (currentValue === newValue || hash(currentValue) === hash(newValue)) ) { const newMultiplicity = currentMultiplicity + multiplicity @@ -406,7 +406,7 @@ export class Index { // At least one has a prefix, use PrefixMap const prefixMap = new PrefixMap() - if (currentPrefix === newPrefix) { + if (isSamePrefix(currentPrefix, newPrefix)) { // Same prefix, different suffixes - need ValueMap within PrefixMap const valueMap = new ValueMap() valueMap.set(hash(currentValue), currentSingleValue) @@ -478,6 +478,11 @@ export class Index { * @param value - The value to extract the prefix from. * @returns The prefix and the suffix. */ +// Prefixes are Map keys, so they compare as a Map does: NaN equals NaN. +function isSamePrefix(a: unknown, b: unknown): boolean { + return a === b || (Number.isNaN(a) && Number.isNaN(b)) +} + function getPrefix(value: TValue): TPrefix | NO_PREFIX { // If the value is an array and the first element is a string or number, then the // first element is the prefix. This is used to distinguish between values without diff --git a/packages/db-ivm/tests/index-refinement-oracle.property.test.ts b/packages/db-ivm/tests/index-refinement-oracle.property.test.ts new file mode 100644 index 000000000..c5c300ee9 --- /dev/null +++ b/packages/db-ivm/tests/index-refinement-oracle.property.test.ts @@ -0,0 +1,275 @@ +/** + * # Does an Index hold exactly the summed multiplicity of each value? + * + * Law and source: an `Index` maps each key to a multiset of values. Adding + * `[value, m]` changes that value's multiplicity by `m`, and `get(key)` returns + * each value whose summed multiplicity is not zero. Two values are the same + * when they have the same structural identity, the identity that `hash` and + * `MultiSet.consolidate` use: `NaN` equals `NaN`, `-0` equals `0`, and values + * of different types differ. Join, reduce, orderBy, and top-K store their + * state in an `Index`, so this law carries every operator that keys rows by + * source key. + * + * Why an example can miss the failure: the Index stores an array value whose + * first element is a string, number, or bigint under that element, the + * prefix, and compares prefixes before hashing. Integer and string prefixes + * behave, so a prefix comparison that disagrees with the `Map` holding the + * prefixes shows only for `NaN`, which a `Map` treats as one key and `===` + * does not. + * + * Model: `expectedIndex` sums multiplicities in plain `Map`s under a string + * identity built from each value's type and printed value. It does not call + * the Index, `hash`, or the prefix helper. + * + * History grammar: up to twelve additions to two keys. Values are prefixed + * arrays `[prefix, payload]` with prefixes `NaN`, `0`, `-0`, `1`, `'1'`, + * `1n`, and `'a'`, or unprefixed numbers and strings. Multiplicities are + * -2, -1, 1, or 2, so values appear, grow, cancel, and reappear, and one key + * can hold several values under one prefix, several prefixes, or prefixed and + * unprefixed values together. + * + * Production driver: a fresh `Index` receives each addition through + * `addValue`. + * + * Refinement check: after each addition, `get` and `has` for both keys equal + * the model. + * + * Calibration: comparing prefixes with `===` kept a cancelled `NaN`-prefixed + * value, or threw `Mismatching prefixes`, and fails the pinned histories and + * both campaigns. + * + * Known omissions: `Index` joins, compaction, presence tracking, and structural + * payloads other than numbers and strings are outside this owner; the join + * operator tests and the incrementalization law own the operators. + */ +import { fc } from '@fast-check/vitest' +import { describe, expect, it } from 'vitest' +import { Index } from '../src/indexes.js' + +type Prefix = number | string | bigint +type Value = [Prefix, number | string] | number | string +type Addition = { key: `k1` | `k2`; value: Value; multiplicity: number } + +// --------------------------------------------------------------------------- +// Model +// --------------------------------------------------------------------------- + +// `NaN` equals `NaN` and `-0` equals `0`; the type keeps 1, '1', and 1n apart. +function scalarIdentity(value: Prefix): string { + if (typeof value === `number`) { + return `number:${Number.isNaN(value) ? `NaN` : String(value === 0 ? 0 : value)}` + } + return `${typeof value}:${String(value)}` +} + +function identity(value: Value): string { + return Array.isArray(value) + ? `[${scalarIdentity(value[0])},${scalarIdentity(value[1])}]` + : scalarIdentity(value) +} + +function expectedIndex( + additions: ReadonlyArray, +): Map> { + const keys = new Map>() + for (const { key, value, multiplicity } of additions) { + const values = keys.get(key) ?? new Map() + const id = identity(value) + const sum = (values.get(id) ?? 0) + multiplicity + if (sum === 0) values.delete(id) + else values.set(id, sum) + if (values.size === 0) keys.delete(key) + else keys.set(key, values) + } + return keys +} + +// --------------------------------------------------------------------------- +// History grammar +// --------------------------------------------------------------------------- + +const prefixes: ReadonlyArray = [ + Number.NaN, + 0, + -0, + 1, + `1`, + BigInt(1), + `a`, +] + +const valueArbitrary: fc.Arbitrary = fc.oneof( + { + weight: 3, + arbitrary: fc + .tuple(fc.constantFrom(...prefixes), fc.constantFrom(0, 1, `x`)) + .map(([prefix, payload]): Value => [prefix, payload]), + }, + fc.constantFrom(5, 6, `u`), +) + +// Reusing earlier values lets histories cancel what they added. +const historyArbitrary: fc.Arbitrary> = fc + .array( + fc.record({ + key: fc.constantFrom(`k1` as const, `k2` as const), + value: valueArbitrary, + multiplicity: fc.constantFrom(-2, -1, 1, 2), + repeat: fc.option(fc.nat({ max: 11 }), { nil: undefined }), + }), + { maxLength: 12 }, + ) + .map((steps) => { + const additions: Array = [] + for (const { key, value, multiplicity, repeat } of steps) { + const earlier = + repeat === undefined + ? undefined + : additions[repeat % (additions.length || 1)] + additions.push( + earlier === undefined + ? { key, value, multiplicity } + : { + key: earlier.key, + value: earlier.value, + multiplicity: -earlier.multiplicity, + }, + ) + } + return additions + }) + +const add = (value: Value, multiplicity: number): Addition => ({ + key: `k1`, + value, + multiplicity, +}) + +// Each history cancels or merges `NaN`-prefixed values in a different layout. +const pinnedHistories: ReadonlyArray<{ + name: string + additions: Array +}> = [ + { + name: `a single NaN-prefixed value cancels`, + additions: [add([Number.NaN, 0], 1), add([Number.NaN, 0], -1)], + }, + { + name: `two NaN-prefixed values share one prefix`, + additions: [ + add([Number.NaN, 0], 1), + add([Number.NaN, 1], 1), + add([Number.NaN, 0], -1), + ], + }, + { + name: `a NaN-prefixed value beside an unprefixed value`, + additions: [add(5, 1), add([Number.NaN, 0], 1), add([Number.NaN, 0], -1)], + }, + { + name: `negative zero and zero share a prefix`, + additions: [add([-0, 0], 1), add([0, 0], -1)], + }, +] + +// --------------------------------------------------------------------------- +// Production driver and refinement check +// --------------------------------------------------------------------------- + +function published( + index: Index, + key: string, +): Map { + const values = new Map() + for (const [value, multiplicity] of index.get(key)) { + const id = identity(value) + values.set(id, (values.get(id) ?? 0) + multiplicity) + } + return values +} + +function expectRefinement(additions: ReadonlyArray): void { + const index = new Index() + for (const [step, addition] of additions.entries()) { + index.addValue(addition.key, [addition.value, addition.multiplicity]) + const model = expectedIndex(additions.slice(0, step + 1)) + for (const key of [`k1`, `k2`]) { + const checkpoint = `after addition ${step} for ${key}` + expect(published(index, key), checkpoint).toEqual( + model.get(key) ?? new Map(), + ) + expect(index.has(key), `${checkpoint} has`).toBe(model.has(key)) + } + } +} + +// --------------------------------------------------------------------------- +// Campaigns. The fixed and random campaigns run the same property, grammar, +// check, and budget. `TANSTACK_DB_IVM_INDEX_SEED` and +// `TANSTACK_DB_IVM_INDEX_PATH` select a direct replay. + +const replaySeed = process.env.TANSTACK_DB_IVM_INDEX_SEED +const replayPath = process.env.TANSTACK_DB_IVM_INDEX_PATH +const campaigns = + replaySeed === undefined && replayPath === undefined + ? [ + { name: `20260930`, seed: 20260930 as number | undefined }, + { name: `random`, seed: undefined }, + ] + : [ + { + name: `replay`, + seed: replaySeed === undefined ? undefined : Number(replaySeed), + }, + ] + +describe(`Index refinement oracle`, () => { + if (replaySeed === undefined && replayPath === undefined) { + for (const { name, additions } of pinnedHistories) { + it(`matches the multiset model when ${name}`, () => + expectRefinement(additions)) + } + } + + for (const { name, seed } of campaigns) { + it(`matches the multiset model across generated additions (${name})`, () => { + if (replayPath !== undefined && replaySeed === undefined) + throw new Error(`TANSTACK_DB_IVM_INDEX_PATH requires a seed`) + if ( + replaySeed !== undefined && + (replaySeed.trim() === `` || + typeof seed !== `number` || + !Number.isSafeInteger(seed)) + ) + throw new Error(`TANSTACK_DB_IVM_INDEX_SEED must be an integer`) + fc.assert( + fc.property(historyArbitrary, (additions) => + expectRefinement(additions), + ), + { + numRuns: 300, + ...(seed === undefined ? {} : { seed }), + ...(replayPath === undefined ? {} : { path: replayPath }), + }, + ) + }) + } + + // Positive execution witness: generated histories reach a NaN prefix that + // cancels after another value joined its key. + it(`reaches cancelled NaN-prefixed values beside other values`, () => { + const sample = fc.sample(historyArbitrary, { seed: 20260930, numRuns: 300 }) + const reached = sample.filter((additions) => + additions.some( + (addition, step) => + Array.isArray(addition.value) && + Number.isNaN(addition.value[0]) && + addition.multiplicity < 0 && + additions + .slice(0, step) + .some((earlier) => earlier.value !== addition.value), + ), + ) + expect(reached.length).toBeGreaterThan(10) + }) +}) diff --git a/packages/db/package.json b/packages/db/package.json index 31dbf02f6..a9e6db7a1 100644 --- a/packages/db/package.json +++ b/packages/db/package.json @@ -22,7 +22,7 @@ "lint": "eslint . --fix", "test": "vitest --run", "test:facade-retention": "node --expose-gc --import tsx tests/facade-retention.probe.ts", - "test:oracles": "vitest --run --coverage.enabled=false tests/change-event-history-oracle.test.ts tests/index-suggestion-oracle.test.ts tests/paced-mutations-oracle.test.ts tests/db-client-hydration-authority-oracle.test.ts tests/collection-mutation-startup-oracle.test.ts tests/collection-cleanup-restart-oracle.test.ts tests/effect-disposal-oracle.test.ts tests/optimistic-transaction-oracle.property.test.ts tests/optimistic-settlement-boundaries.test.ts tests/optimistic-history-publication.test.ts tests/optimistic-history-outcomes.test.ts tests/collection-metadata-publication-oracle.property.test.ts tests/collection-state-retention-oracle.property.test.ts tests/collection-truncate-ownership-oracle.property.test.ts tests/collection-subscription-lifecycle-history.property.test.ts tests/collection-subscription-lifecycle-oracle.test.ts tests/collection-subscription-reentrancy-oracle.test.ts tests/collection-subscription-lifecycle-publication.property.test.ts tests/collection-subscription-replay-oracle.property.test.ts tests/d2-source-reconciliation-oracle.property.test.ts tests/live-query-observer-history.property.test.ts tests/query/cold-join-reconciliation-oracle.test.ts tests/query/identity-output-shape-oracle.test.ts tests/query/optimizer-semantics-oracle.test.ts tests/query/index-path-collision-oracle.test.ts tests/query/includes-collection-oracle.property.test.ts tests/query/includes-functional-projection-oracle.test.ts tests/query/includes-functional-input-boundary.test.ts tests/query/includes-context-transport-oracle.test.ts tests/query/includes-cross-formulation-oracle.property.test.ts tests/query/includes-optimistic-oracle.property.test.ts tests/query/includes-oracle.property.test.ts tests/query/includes-publication-oracle.test.ts tests/query/includes-query-shape-oracle.test.ts tests/query/subquery-user-value-oracle.test.ts tests/query/includes-temporal-oracle.test.ts tests/query/includes-work-counter-oracle.test.ts tests/query/load-subset-oracle.property.test.ts tests/query/load-subset-replay-refinement-oracle.test.ts tests/query/load-subset-source-readiness-refinement-oracle.test.ts tests/query/load-subset-transaction-refinement-oracle.test.ts tests/query/ordered-source-loader-state.test.ts tests/query/ordered-demand-retirement.test.ts tests/query/ordered-default-work.test.ts tests/query/ordered-lifecycle-oracle.property.test.ts tests/query/ordered-work-oracle.property.test.ts tests/query/pagination-oracle.property.test.ts tests/query/includes-space-oracle.test.ts tests/query/virtual-row-fields-oracle.test.ts", + "test:oracles": "vitest --run --coverage.enabled=false tests/change-event-history-oracle.test.ts tests/index-suggestion-oracle.test.ts tests/paced-mutations-oracle.test.ts tests/db-client-hydration-authority-oracle.test.ts tests/collection-mutation-startup-oracle.test.ts tests/collection-cleanup-restart-oracle.test.ts tests/effect-disposal-oracle.test.ts tests/optimistic-transaction-oracle.property.test.ts tests/optimistic-settlement-boundaries.test.ts tests/optimistic-history-publication.test.ts tests/optimistic-history-outcomes.test.ts tests/collection-metadata-publication-oracle.property.test.ts tests/collection-state-retention-oracle.property.test.ts tests/collection-truncate-ownership-oracle.property.test.ts tests/collection-subscription-lifecycle-history.property.test.ts tests/collection-subscription-lifecycle-oracle.test.ts tests/collection-subscription-reentrancy-oracle.test.ts tests/collection-subscription-lifecycle-publication.property.test.ts tests/collection-subscription-replay-oracle.property.test.ts tests/d2-source-reconciliation-oracle.property.test.ts tests/live-query-observer-history.property.test.ts tests/query/cold-join-reconciliation-oracle.test.ts tests/query/identity-output-shape-oracle.test.ts tests/query/optimizer-semantics-oracle.test.ts tests/query/index-path-collision-oracle.test.ts tests/query/includes-collection-oracle.property.test.ts tests/query/includes-functional-projection-oracle.test.ts tests/query/includes-functional-input-boundary.test.ts tests/query/includes-context-transport-oracle.test.ts tests/query/includes-cross-formulation-oracle.property.test.ts tests/query/includes-optimistic-oracle.property.test.ts tests/query/includes-oracle.property.test.ts tests/query/includes-publication-oracle.test.ts tests/query/includes-query-shape-oracle.test.ts tests/query/subquery-user-value-oracle.test.ts tests/query/includes-temporal-oracle.test.ts tests/query/includes-work-counter-oracle.test.ts tests/query/load-subset-oracle.property.test.ts tests/query/load-subset-replay-refinement-oracle.test.ts tests/query/load-subset-source-readiness-refinement-oracle.test.ts tests/query/load-subset-transaction-refinement-oracle.test.ts tests/query/ordered-source-loader-state.test.ts tests/query/ordered-demand-retirement.test.ts tests/query/ordered-default-work.test.ts tests/query/ordered-lifecycle-oracle.property.test.ts tests/query/ordered-work-oracle.property.test.ts tests/query/pagination-oracle.property.test.ts tests/query/includes-space-oracle.test.ts tests/query/virtual-row-fields-oracle.test.ts tests/query/where-predicate-publication-oracle.property.test.ts tests/query/join-result-key-oracle.property.test.ts", "bench:nested-includes": "vitest bench tests/query/includes-performance.bench.ts --run" }, "type": "module", diff --git a/packages/db/src/SortedMap.ts b/packages/db/src/SortedMap.ts index c87710834..c0b31b14c 100644 --- a/packages/db/src/SortedMap.ts +++ b/packages/db/src/SortedMap.ts @@ -233,13 +233,11 @@ export class SortedMap { * * @returns An iterator for the map's values */ - values(): IterableIterator { - return function* (this: SortedMap) { - this.restoreOrder() - for (const key of this.sortedKeys) { - yield this.map.get(key)! - } - }.call(this) + *values(): IterableIterator { + this.restoreOrder() + for (const key of this.sortedKeys) { + yield this.map.get(key)! + } } /** diff --git a/packages/db/src/collection/change-events.ts b/packages/db/src/collection/change-events.ts index f22bcda97..6409b883f 100644 --- a/packages/db/src/collection/change-events.ts +++ b/packages/db/src/collection/change-events.ts @@ -7,6 +7,8 @@ import { optimizeExpressionWithIndexes, } from '../utils/index-optimization.js' import { ensureIndexForField } from '../indexes/auto-index.js' +import { getPropRefPropertyPath } from '../query/ir.js' +import { isVirtualPropName } from '../virtual-props.js' import { makeComparator } from '../utils/comparison.js' import { buildCompareOptions } from '../query/compiler/order-by' import type { @@ -19,6 +21,14 @@ import type { CollectionImpl } from './index.js' import type { BasicExpression, OrderBy } from '../query/ir.js' import type { WithVirtualProps } from '../virtual-props.js' +/** + * Yields visible entries, enriched with virtual properties, whose stored row + * passes `prefilter`. + */ +export type StoredRowScan = ( + prefilter: (row: object) => boolean, +) => Iterable<[TKey, WithVirtualProps]> + /** * Returns the current state of the collection as an array of changes * @param collection - The collection to get changes from @@ -56,12 +66,23 @@ export function currentStateAsChanges< >( collection: CollectionLike, TKey>, options: CurrentStateAsChangesOptions = {}, + scanStoredRows?: StoredRowScan, ): Array, TKey>> | void { // Helper function to collect filtered results const collectFilteredResults = ( filterFn?: (value: WithVirtualProps) => boolean, ): Array, TKey>> => { const result: Array, TKey>> = [] + // Reject rows by one stored field before copying them to add virtual + // properties. Survivors still pass through the full predicate. + const prefilter = + scanStoredRows && options.where && compileEqualityPrefilter(options.where) + if (filterFn && scanStoredRows && prefilter) { + for (const [key, value] of scanStoredRows(prefilter)) { + if (filterFn(value)) result.push({ type: `insert`, key, value }) + } + return result + } for (const [key, value] of collection.entries()) { // If no filter function is provided, include all items if (filterFn?.(value) ?? true) { @@ -198,6 +219,121 @@ export function createFilterFunctionFromExpression( } } +/** A field and the string or boolean literal a top-level `eq` requires. */ +export type EqualityRoute = { + path: Array + /** Stable identity of `path`, for grouping routes by field. */ + pathKey: string + expected: string | boolean +} + +/** Read result for a route path whose property access threw. */ +export const UNREADABLE_ROUTE_VALUE: unique symbol = Symbol( + `unreadable route value`, +) + +/** + * Finds a cheap necessary condition for `expression` to be TRUE, or returns + * undefined when the expression has none. + * + * A top-level conjunct `eq(field, literal)` with a string or boolean literal is + * TRUE only when the field holds the identical string or boolean: equality + * normalization never maps another type onto a plain string or boolean. A + * row whose field holds anything else therefore fails the whole expression. + * + * With `storedRows`, the condition is read from a stored row instead of its + * enriched copy, so conjuncts on virtual fields are skipped: stored rows need + * not carry them. + */ +export function findEqualityRoute( + expression: BasicExpression, + { storedRows = false }: { storedRows?: boolean } = {}, +): EqualityRoute | undefined { + const conjuncts: Array = [] + const collect = (node: BasicExpression) => { + if (node.type === `func` && node.name === `and`) node.args.forEach(collect) + else conjuncts.push(node) + } + collect(expression) + + // A string literal usually rejects more rows than a boolean one. + let best: EqualityRoute | undefined + for (const conjunct of conjuncts) { + if (conjunct.type !== `func` || conjunct.name !== `eq`) continue + const [left, right] = conjunct.args + const ref = + left?.type === `ref` ? left : right?.type === `ref` ? right : undefined + const literal = + left?.type === `val` ? left : right?.type === `val` ? right : undefined + if (!ref || !literal) continue + const expected: unknown = literal.value + if (typeof expected !== `string` && typeof expected !== `boolean`) continue + + const path = getPropRefPropertyPath(ref) + if (storedRows && (path.length === 0 || isVirtualPropName(path[0]!))) { + continue + } + if (best === undefined || typeof best.expected === `boolean`) { + best = { path, pathKey: JSON.stringify(path), expected } + } + if (typeof expected === `string`) break + } + return best +} + +/** + * Reads a route field the way the single-row evaluator does. A throwing read + * returns UNREADABLE_ROUTE_VALUE so callers leave the decision to the full + * predicate. + */ +export function readRouteValue( + row: unknown, + path: ReadonlyArray, +): unknown { + try { + let value: unknown = row + for (const segment of path) { + if (value === null || value === undefined) return undefined + value = (value as Record)[segment] + } + return value + } catch { + return UNREADABLE_ROUTE_VALUE + } +} + +/** + * Compiles the route of `expression` as a row test that is false only when + * the full predicate must be false. + * + * The test reads a stored row instead of its enriched copy. The copy holds + * each enumerable own root property of the stored row and lacks the others, so its field is either the stored value or `undefined`, + * which never equals the literal. A read that throws passes the row to the + * full predicate. + */ +export function compileEqualityPrefilter( + expression: BasicExpression, +): ((row: object) => boolean) | undefined { + const route = findEqualityRoute(expression, { storedRows: true }) + if (route === undefined) return undefined + const { path, expected } = route + // Most routes name one top-level field; read it without walking a path. + if (path.length === 1) { + const field = path[0]! + return (row) => { + try { + return (row as Record)[field] === expected + } catch { + return true + } + } + } + return (row) => { + const value = readRouteValue(row, path) + return value === UNREADABLE_ROUTE_VALUE || value === expected + } +} + /** * Creates a filtered callback that only calls the original callback with changes that match the where clause * @param originalCallback - The original callback to filter diff --git a/packages/db/src/collection/changes.ts b/packages/db/src/collection/changes.ts index 5bf5be736..f6b1403f8 100644 --- a/packages/db/src/collection/changes.ts +++ b/packages/db/src/collection/changes.ts @@ -6,6 +6,7 @@ import { toExpression, } from '../query/builder/ref-proxy.js' import { CollectionSubscription } from './subscription.js' +import { readRouteValue } from './change-events.js' import type { StandardSchemaV1 } from '@standard-schema/spec' import type { ChangeMessage, SubscribeChangesOptions } from '../types' import type { CollectionLifecycleManager } from './lifecycle.js' @@ -259,9 +260,25 @@ export class CollectionChangesManager< const layoutListeners = [...this.layoutChangeListeners] const subscriptions = [...this.changeSubscriptions] withPublicationContext(() => { - const callbacks: Array<() => void> = subscriptions.map( - (subscription) => () => subscription.emitEvents(enrichedEvents), - ) + // An empty batch signals readiness to every subscriber. + const routed = + rawEvents.length > 0 + ? routeChanges(enrichedEvents, subscriptions) + : undefined + const callbacks: Array<() => void> = [] + for (const subscription of subscriptions) { + const own = routed?.get(subscription) + callbacks.push(() => { + // An earlier callback in this publication can end routing for this + // subscription, for example by leaving stale rows to reconcile. + if (own === undefined || !subscription.changeRoute) { + subscription.emitEvents(enrichedEvents) + } else if (own.length > 0) { + // A routed subscription with no candidate change cannot publish. + subscription.emitEvents(own) + } + }) + } if (rawEvents.length === 0) { callbacks.unshift(...layoutListeners) } @@ -425,3 +442,64 @@ export class CollectionChangesManager< this.deferral = undefined } } + +/** + * Gives each routed subscription only the changes whose value or previous + * value holds its route literal, in batch order. Subscriptions without a + * route are absent from the result and receive the whole batch; with no + * routed subscription the result is undefined. + */ +function routeChanges( + changes: Array>, + subscriptions: Array, +): Map>> | undefined { + const routes = subscriptions.map((subscription) => subscription.changeRoute) + if (routes.every((route) => route === undefined)) return undefined + const routed = new Map< + CollectionSubscription, + Array> + >() + const groups = new Map< + string, + { + path: Array + byLiteral: Map> + } + >() + for (const [index, subscription] of subscriptions.entries()) { + const route = routes[index] + if (!route) continue + routed.set(subscription, []) + let group = groups.get(route.pathKey) + if (!group) { + group = { path: route.path, byLiteral: new Map() } + groups.set(route.pathKey, group) + } + const peers = group.byLiteral.get(route.expected) + if (peers) peers.push(subscription) + else group.byLiteral.set(route.expected, [subscription]) + } + + const deliver = ( + targets: Array | undefined, + change: ChangeMessage, + ) => { + for (const subscription of targets ?? []) { + routed.get(subscription)!.push(change) + } + } + for (const change of changes) { + for (const group of groups.values()) { + const value = readRouteValue(change.value, group.path) + const previous = + change.previousValue === undefined + ? undefined + : readRouteValue(change.previousValue, group.path) + // The where filter reads these same values, and a read that throws makes + // its predicate false, so an unreadable value matches no literal here. + deliver(group.byLiteral.get(value), change) + if (previous !== value) deliver(group.byLiteral.get(previous), change) + } + } + return routed +} diff --git a/packages/db/src/collection/index.ts b/packages/db/src/collection/index.ts index 589b61980..239778486 100644 --- a/packages/db/src/collection/index.ts +++ b/packages/db/src/collection/index.ts @@ -1051,7 +1051,9 @@ export class CollectionImpl< public currentStateAsChanges( options: CurrentStateAsChangesOptions = {}, ): Array, TKey>> | void { - return currentStateAsChanges(this, options) + return currentStateAsChanges(this, options, (prefilter) => + this._state.entriesPassing(prefilter), + ) } /** diff --git a/packages/db/src/collection/state.ts b/packages/db/src/collection/state.ts index 2caa91d33..a7de166e1 100644 --- a/packages/db/src/collection/state.ts +++ b/packages/db/src/collection/state.ts @@ -500,6 +500,23 @@ export class CollectionStateManager< } } + /** + * Visible entries whose stored row passes `prefilter`, enriched with virtual + * properties. Rows that fail are never copied. + */ + public *entriesPassing( + prefilter: (row: object) => boolean, + ): IterableIterator<[TKey, WithVirtualProps]> { + // Without optimistic state, the visible rows are the synced rows in order. + const rows = + this.optimisticUpserts.size === 0 && this.optimisticDeletes.size === 0 + ? this.syncedData + : this.entries() + for (const [key, row] of rows) { + if (prefilter(row)) yield [key, this.enrichWithVirtualProps(row, key)] + } + } + /** * Get all entries (virtual derived state) */ diff --git a/packages/db/src/collection/subscription.ts b/packages/db/src/collection/subscription.ts index 9bf361f41..4f76ee879 100644 --- a/packages/db/src/collection/subscription.ts +++ b/packages/db/src/collection/subscription.ts @@ -12,7 +12,9 @@ import { LoadSubsetOperationAbortedError } from '../errors.js' import { createFilterFunctionFromExpression, createFilteredCallback, + findEqualityRoute, } from './change-events.js' +import type { EqualityRoute } from './change-events.js' import type { BasicExpression, OrderBy } from '../query/ir.js' import type { IndexReader } from '../indexes/base-index.js' import type { @@ -172,6 +174,9 @@ export class CollectionSubscription private filteredCallback: (changes: Array>) => boolean + /** Field and literal the where clause requires, if it has a cheap one. */ + private readonly equalityRoute: EqualityRoute | undefined + private orderByIndex: IndexReader | undefined // Status tracking @@ -234,6 +239,10 @@ export class CollectionSubscription this.callback = callbackWithSentKeysTracking + this.equalityRoute = options.whereExpression + ? findEqualityRoute(options.whereExpression) + : undefined + // Create a filtered callback if where clause is provided this.filteredCallback = options.whereExpression ? createFilteredCallback(this.callback, options) @@ -924,7 +933,15 @@ export class CollectionSubscription /** Create the record for a fresh, abortable acquisition attempt. */ private createSubsetAcquisitionRecord( demand: SubsetDemand, - ): SubsetAcquisitionRecord & { abortController: AbortController } { + ): SubsetAcquisitionRecord { + // Eager sync never passes subset options to an adapter, so an abortable + // acquisition would only allocate a controller and an AbortError. + if (this.collection.config.syncMode !== `on-demand`) { + return { + options: demand.requestOptions, + syncRunGeneration: this.collection._sync.getSyncRunGeneration(), + } + } const abortController = new AbortController() const requestSignal = demand.requestOptions.signal let removeRequestAbortListener: (() => void) | undefined @@ -1116,6 +1133,23 @@ export class CollectionSubscription return this.filteredCallback(newChanges) } + /** + * The route through which this subscription may receive only the changes + * whose value or previous value holds the route's literal. A change reaches + * the where filter only through those values, so the others cannot publish, + * and sent-key records cover published rows only. Stale published rows and + * truncate replay consume unfiltered changes, so no route applies then. + */ + get changeRoute(): EqualityRoute | undefined { + if ( + this.stalePublishedRows.size > 0 || + this.truncateReplayState !== undefined + ) { + return undefined + } + return this.equalityRoute + } + /** Keep direct snapshot reads private while an authoritative replay is open. */ private publishSnapshot(changes: Array>): void { if (!this.bufferPrivately(changes)) this.callback(changes) @@ -1591,15 +1625,21 @@ export class CollectionSubscription // 3. We're collecting all changes atomically, so filtering doesn't make sense const skipDeleteFilter = this.isBufferingForTruncate + // sentKeys records only published rows; trackSentKeys adds them after + // delivery. Keys inserted earlier in this batch are tracked locally, so a + // row the where clause drops cannot advance pagination or later look like + // a duplicate insert. + const insertedInBatch = new Set() const newChanges = [] for (const change of changes) { let newChange = change - const keyInSentKeys = this.sentKeys.has(change.key) + const keyInSentKeys = + this.sentKeys.has(change.key) || insertedInBatch.has(change.key) if (!keyInSentKeys) { if (change.type === `update`) { newChange = { ...change, type: `insert`, previousValue: undefined } - this.sentKeys.add(change.key) + insertedInBatch.add(change.key) } else if (change.type === `delete`) { // Filter out deletes for keys that have not been sent, // UNLESS we're buffering for truncate (where all deletes should pass through) @@ -1607,7 +1647,7 @@ export class CollectionSubscription continue } } else { - this.sentKeys.add(change.key) + insertedInBatch.add(change.key) } } else { // Key was already sent - handle based on change type @@ -1621,6 +1661,7 @@ export class CollectionSubscription // Remove from sentKeys so future inserts for this key are allowed // (e.g., after truncate + reinsert) this.sentKeys.delete(change.key) + insertedInBatch.delete(change.key) } } newChanges.push(newChange) diff --git a/packages/db/src/query/compiler/evaluators.ts b/packages/db/src/query/compiler/evaluators.ts index e901a2b56..6905d7b8f 100644 --- a/packages/db/src/query/compiler/evaluators.ts +++ b/packages/db/src/query/compiler/evaluators.ts @@ -260,8 +260,19 @@ function compileFunction(func: Func, isSingleRow: boolean): (data: any) => any { const argA = compiledArgs[0]! const argB = compiledArgs[1]! return (data) => { - const a = normalizeEqualityOperand(argA(data)) - const b = normalizeEqualityOperand(argB(data)) + const rawA = argA(data) + const rawB = argB(data) + // Same-type strings and booleans need no normalization; this is the + // hot path for predicate scans and change filtering. + const typeA = typeof rawA + if ( + typeA === typeof rawB && + (typeA === `string` || typeA === `boolean`) + ) { + return rawA === rawB + } + const a = normalizeEqualityOperand(rawA) + const b = normalizeEqualityOperand(rawB) // In 3-valued logic, any comparison with null/undefined returns UNKNOWN if (isUnknown(a) || isUnknown(b)) { return null diff --git a/packages/db/src/query/compiler/index.ts b/packages/db/src/query/compiler/index.ts index da5636515..3831788ed 100644 --- a/packages/db/src/query/compiler/index.ts +++ b/packages/db/src/query/compiler/index.ts @@ -48,6 +48,7 @@ import { isExpressionLike, } from '../ir.js' import { ensureIndexForField } from '../../indexes/auto-index.js' +import { createSourceRecord } from '../../utils/source-record.js' import { deepEquals } from '../../utils.js' import { normalizeValue } from '../../utils/comparison.js' import { @@ -401,7 +402,7 @@ export function compileQuery( mapNestedQueries(query, rawQuery, queryMapping) // Create a copy of the inputs map to avoid modifying the original - const allInputs = { ...inputs } + const allInputs = Object.assign(createSourceRecord(), inputs) const rawSources = collectCollectionSources(rawQuery) bindSourceInputs(rawSources, allInputs) @@ -1082,12 +1083,14 @@ export function compileQuery( } } - // Normalize every logical row before DISTINCT and ordering. Those operators - // track visibility by row key, so an insert-before-delete replacement with - // the same key would otherwise keep the old value and hide route or order - // changes. Joined contributors may differ in unselected namespaces; only - // the public value and its route/order inputs must be congruent. - if (!selectHasAggregates) { + // DISTINCT tracks visibility by selected value, so an insert-before-delete + // replacement with the same key must be normalized first; joined + // contributors may differ in unselected namespaces, and only the public + // value and its route/order inputs must be congruent. Ordering already + // batches each key's retractions before its insertions (topKBatch), and + // materialized relations reduce by public key, so other queries keep their + // compiled rows. + if (!selectHasAggregates && query.distinct) { pipeline = canonicalizeSelectedRows( pipeline, query, diff --git a/packages/db/src/query/compiler/joins.ts b/packages/db/src/query/compiler/joins.ts index 67816a3f2..9a0f0c23d 100644 --- a/packages/db/src/query/compiler/joins.ts +++ b/packages/db/src/query/compiler/joins.ts @@ -740,6 +740,14 @@ function processJoinSource( } } +// JSON prints non-finite numbers as null, the missing-side marker. An object +// never collides with a source key, which is a string or a number. +function encodeNonFiniteKey(_: string, value: unknown): unknown { + return typeof value === `number` && !Number.isFinite(value) + ? { number: String(value) } + : value +} + function getFirstFromAlias(query: QueryIR): string | undefined { return getFromSources(query.from)[0]?.alias } @@ -770,8 +778,13 @@ function processJoinResults( Object.assign(mergedNamespacedRow, joinedNamespacedRow) } - // We create a composite key that combines the main and joined keys - const resultKey = `[${mainKey},${joinedKey}]` + // Combine the main and joined keys without ambiguity: keys may contain + // delimiters, numbers and strings may print alike, and a missing outer + // side encodes as null, which no source key can be. + const resultKey = JSON.stringify( + [mainKey ?? null, joinedKey ?? null], + encodeNonFiniteKey, + ) return [resultKey, mergedNamespacedRow] as [string, NamespacedRow] }), diff --git a/packages/db/src/query/effect.ts b/packages/db/src/query/effect.ts index 308a535c5..d5d667a4c 100644 --- a/packages/db/src/query/effect.ts +++ b/packages/db/src/query/effect.ts @@ -1,4 +1,5 @@ import { D2, output } from '@tanstack/db-ivm' +import { createSourceRecord } from '../utils/source-record.js' import { createDeferred } from '../deferred.js' import { runAllCallbacks } from '../utils/callbacks.js' import { normalizeError } from '../utils/error.js' @@ -389,18 +390,14 @@ class EffectPipelineRunner { // Mutable objects passed to compileQuery by reference. // The join compiler captures these references and reads them later when // the graph runs, so they must be populated before the first graph run. - private readonly subscriptions: Record = {} - private readonly lazySourcesCallbacks: Record< - string, - LazyCollectionCallbacks - > = {} + private readonly subscriptions = createSourceRecord() + private readonly lazySourcesCallbacks = + createSourceRecord() private readonly lazySources = new Set() private readonly demand = new SubsetDemandController() // OrderBy optimization info populated by the compiler when limit is present - private readonly optimizableOrderByCollections: Record< - string, - OrderByOptimizationInfo - > = {} + private readonly optimizableOrderByCollections = + createSourceRecord() // Ordered subscription state for cursor-based loading private readonly orderedLoaders = new Map() @@ -464,12 +461,11 @@ class EffectPipelineRunner { /** Compile the D2 graph and query pipeline */ private compilePipeline(): void { this.graph = new D2() - this.inputs = Object.fromEntries( - this.collectionSources.map((source) => [ - source.sourceId, - this.graph!.newInput(), - ]), - ) + const inputs = createSourceRecord>() + for (const source of this.collectionSources) { + inputs[source.sourceId] = this.graph.newInput() + } + this.inputs = inputs const compilation = compileQuery( this.query, diff --git a/packages/db/src/query/live/collection-config-builder.ts b/packages/db/src/query/live/collection-config-builder.ts index 67e388089..f2b41d47b 100644 --- a/packages/db/src/query/live/collection-config-builder.ts +++ b/packages/db/src/query/live/collection-config-builder.ts @@ -11,6 +11,7 @@ import { createDeferred } from '../../deferred.js' import { deepEquals } from '../../utils.js' import { runAllCallbacks } from '../../utils/callbacks.js' import { normalizeError } from '../../utils/error.js' +import { createSourceRecord } from '../../utils/source-record.js' import { CollectionSubscriber } from './collection-subscriber.js' import { getCollectionBuilder } from './collection-registry.js' import { LIVE_QUERY_INTERNAL } from './internal.js' @@ -143,9 +144,9 @@ export class CollectionConfigBuilder< | undefined // Map of opaque source ID to subscription - readonly subscriptions: Record = {} + readonly subscriptions = createSourceRecord() // Map of opaque source ID to demand callbacks for that lazy source - lazySourcesCallbacks: Record = {} + lazySourcesCallbacks = createSourceRecord() // Set of opaque source IDs that are lazy (don't load initial state) readonly lazySources = new Set() private readonly activeDemands = new Map< @@ -164,7 +165,7 @@ export class CollectionConfigBuilder< private syncRunGeneration = 0 private windowOperationGeneration = 0 // Map of lexical source IDs to optimizable ORDER BY state - optimizableOrderByCollections: Record = {} + optimizableOrderByCollections = createSourceRecord() constructor( private readonly config: LiveQueryCollectionConfig, @@ -800,8 +801,8 @@ export class CollectionConfigBuilder< this.pendingOrderedLoads.clear() this.orderedLoadFailed = false this.windowFailed = false - this.optimizableOrderByCollections = {} - this.lazySourcesCallbacks = {} + this.optimizableOrderByCollections = createSourceRecord() + this.lazySourcesCallbacks = createSourceRecord() // Clear subscription references to prevent memory leaks // Note: Individual subscriptions are already unsubscribed via unsubscribeCallbacks @@ -815,12 +816,11 @@ export class CollectionConfigBuilder< */ private compileBasePipeline() { this.graphCache = new D2() - this.inputsCache = Object.fromEntries( - this.collectionSources.map((source) => [ - source.sourceId, - this.graphCache!.newInput(), - ]), - ) + const inputs = createSourceRecord>() + for (const source of this.collectionSources) { + inputs[source.sourceId] = this.graphCache.newInput() + } + this.inputsCache = inputs const compilation = compileQuery( this.query, @@ -893,20 +893,23 @@ export class CollectionConfigBuilder< }), ) - const bucketFacades = new BucketFacadeAdapter( - this.id, - this.bucketFacadesCache ?? [], - (count) => { - syncState.messagesCount += count - }, - ) - syncState.unsubscribeCallbacks.add(() => bucketFacades.cleanup()) + // A query whose pipeline was not materialized publishes its rows directly + // and pays for no facade state. + const facadeCompilations = this.bucketFacadesCache + const bucketFacades = facadeCompilations + ? new BucketFacadeAdapter(this.id, facadeCompilations, (count) => { + syncState.messagesCount += count + }) + : undefined + if (bucketFacades) { + syncState.unsubscribeCallbacks.add(() => bucketFacades.cleanup()) + } // Flush pending changes and reset the accumulator. // Called at the end of each graph run to commit all accumulated changes. syncState.flushPendingChanges = () => { const hasParentChanges = pendingChanges.size > 0 - const hasChildChanges = bucketFacades.hasPendingChanges() + const hasChildChanges = bucketFacades?.hasPendingChanges() ?? false if (!hasParentChanges && !hasChildChanges) { return @@ -926,6 +929,17 @@ export class CollectionConfigBuilder< return } + // A key has at most one result row, so one flush can add or remove at + // most one. Check before any state changes: anything else means an + // upstream operator broke multiplicity. + for (const [key, { inserts, deletes }] of pendingChanges) { + if (Math.abs(inserts - deletes) > 1) { + throw new Error( + `Live query result key ${String(key)} changed by ${inserts - deletes} rows in one flush; a key has at most one result row.`, + ) + } + } + let facadePublication: | ReturnType | undefined @@ -933,28 +947,30 @@ export class CollectionConfigBuilder< | ReturnType | undefined try { - facadePublication = bucketFacades.flush() + facadePublication = bucketFacades?.flush() rootPublication = hasParentChanges ? config.collection._deferPublication() : undefined - const changesToApply: Map> = new Map( - [...pendingChanges].map(([key, changes]) => { - const resolved: Changes = { - ...changes, - value: bucketFacades.resolve(changes.value), - } - if (changes.previousValue !== undefined) { - resolved.previousValue = bucketFacades.resolve( - changes.previousValue, - ) - } - return [key, resolved] - }), - ) + const changesToApply: Map> = bucketFacades + ? new Map( + [...pendingChanges].map(([key, changes]) => { + const resolved: Changes = { + ...changes, + value: bucketFacades.resolve(changes.value), + } + if (changes.previousValue !== undefined) { + resolved.previousValue = bucketFacades.resolve( + changes.previousValue, + ) + } + return [key, resolved] + }), + ) + : pendingChanges // New facades are not reachable until their root row is installed, so // make them ready first. A facade failure then leaves the root intact, // and the root commit is the final state change before publication. - facadePublication.prepare() + facadePublication?.prepare() if (hasParentChanges) { begin() let lookup: ((key: string | number) => boolean) | undefined @@ -980,7 +996,7 @@ export class CollectionConfigBuilder< let publicationError: unknown for (const publish of [ rootPublication?.publish, - facadePublication.publish, + facadePublication?.publish, ]) { if (!publish) continue try { @@ -1014,7 +1030,6 @@ export class CollectionConfigBuilder< ) { const { write, collection } = config const { deletes, inserts, value, orderByIndex } = changes - // Store the key of the result so that we can retrieve it in the // getKey function this.resultKeys.set(value, key) diff --git a/packages/db/src/query/live/materialized-pipeline.ts b/packages/db/src/query/live/materialized-pipeline.ts index d7f76e79c..69ea0f5b6 100644 --- a/packages/db/src/query/live/materialized-pipeline.ts +++ b/packages/db/src/query/live/materialized-pipeline.ts @@ -68,6 +68,15 @@ export type MaterializedCompilation = { facades: Array } +export type MaterializedRootCompilation = { + pipeline: ResultStream + /** + * Undefined when the compiled pipeline passes through unchanged. Its rows + * then carry no facade references or private route state to resolve. + */ + facades: Array | undefined +} + type RelationScope = `root` | `child` type BuiltRelations = WeakMap< @@ -86,12 +95,12 @@ export function materializeCompilation( compilation: CompilationResult, getRootKey?: (row: any) => unknown, reduceJoinedPublicKeys = false, -): MaterializedCompilation { +): MaterializedRootCompilation { if ( !compilation.includes?.length && !(getRootKey && reduceJoinedPublicKeys) ) { - return { pipeline: compilation.pipeline, facades: [] } + return { pipeline: compilation.pipeline, facades: undefined } } const built: BuiltRelations = new WeakMap() diff --git a/packages/db/src/utils/source-record.ts b/packages/db/src/utils/source-record.ts new file mode 100644 index 000000000..ee4bf5f2f --- /dev/null +++ b/packages/db/src/utils/source-record.ts @@ -0,0 +1,9 @@ +/** + * Creates a record keyed by source ids. Each source id is unique, so a plain + * object would receive a new hidden class for every compiled query. A record + * without a prototype avoids those shape transitions and cannot confuse an + * alias such as `constructor` with an inherited member. + */ +export function createSourceRecord(): Record { + return Object.create(null) +} diff --git a/packages/db/tests/collection-indexes.test.ts b/packages/db/tests/collection-indexes.test.ts index 8f874777d..43f0c3973 100644 --- a/packages/db/tests/collection-indexes.test.ts +++ b/packages/db/tests/collection-indexes.test.ts @@ -1012,7 +1012,7 @@ describe(`Collection Indexes`, () => { type: `index`, operation: `gte`, field: `age`, - value: { from: 25, fromInclusive: true }, + value: 25, }) }) }) @@ -2077,7 +2077,7 @@ describe(`Collection Indexes`, () => { type: `index`, operation: `gte`, field: `age`, - value: { from: 30, fromInclusive: true }, + value: 30, }, { type: `index`, operation: `eq`, field: `status`, value: `active` }, { type: `index`, operation: `eq`, field: `name`, value: `Alice` }, diff --git a/packages/db/tests/collection-subscription.test.ts b/packages/db/tests/collection-subscription.test.ts index c8ba4e652..17697f131 100644 --- a/packages/db/tests/collection-subscription.test.ts +++ b/packages/db/tests/collection-subscription.test.ts @@ -1700,6 +1700,119 @@ describe(`CollectionSubscription status tracking`, () => { } }) + // The next limited page starts after the rows this subscription published. + // A change its where clause drops is not published and cannot advance the + // offset, whether it arrives alone or beside a matching change. Change + // routing withholds dropped rows from an `eq` subscription before sent keys + // are recorded, so an `or` predicate, which is not routed, is also needed to + // reach that record. + const statusA = new Func(`eq`, [new PropRef([`status`]), new Value(`a`)]) + it.each( + [ + `alone`, + `beside a matching change`, + `as an update beside a matching change`, + ].flatMap((grouping) => [ + { grouping, predicate: `routed eq`, where: statusA }, + { + grouping, + predicate: `unrouted or`, + where: new Func(`or`, [ + statusA, + new Func(`eq`, [new PropRef([`status`]), new Value(`z`)]), + ]), + }, + ]), + )( + `does not advance the page offset for a dropped change that arrives $grouping ($predicate)`, + async ({ grouping, where }) => { + type Row = { id: string; rank: number; status: string } + const loads: Array = [] + let sync!: { + begin: () => void + write: (message: { type: `insert` | `update`; value: Row }) => void + commit: () => void + } + const collection = createCollection({ + id: `limited-offset-dropped-change-${grouping}-${String(where.name)}`, + getKey: ({ id }) => id, + syncMode: `on-demand`, + sync: { + sync: ({ begin, write, commit, markReady }) => { + sync = { begin, write, commit } + begin() + write({ type: `insert`, value: { id: `a1`, rank: 1, status: `a` } }) + write({ type: `insert`, value: { id: `b0`, rank: 5, status: `b` } }) + commit() + markReady() + return { + loadSubset: (options) => { + loads.push(options) + return true + }, + } + }, + }, + }) + const index = collection.createIndex((row) => row.rank, { + indexType: BTreeIndex, + }) + const published = new Set() + const subscription = collection.subscribeChanges( + (changes) => { + for (const change of changes) { + if (change.type === `delete`) published.delete(change.key) + else published.add(change.key) + } + }, + { whereExpression: where }, + ) + subscription.setOrderByIndex(index) + const requestPage = () => + subscription.requestLimitedSnapshot({ + orderBy: [ + { + expression: new PropRef([`rank`]), + compareOptions: { direction: `asc`, nulls: `first` }, + }, + ], + limit: 1, + }) + + try { + requestPage() + sync.begin() + if (grouping === `as an update beside a matching change`) { + sync.write({ + type: `update`, + value: { id: `b0`, rank: 5, status: `c` }, + }) + } else { + sync.write({ + type: `insert`, + value: { id: `b1`, rank: 2, status: `b` }, + }) + } + if (grouping !== `alone`) { + sync.write({ + type: `insert`, + value: { id: `a2`, rank: 3, status: `a` }, + }) + } + sync.commit() + requestPage() + + expect([...published].sort()).toEqual( + grouping === `alone` ? [`a1`] : [`a1`, `a2`], + ) + expect(loads.at(-1)?.offset).toBe(published.size) + } finally { + subscription.unsubscribe() + await collection.cleanup() + } + }, + ) + it(`does not observe limited adapter work after it unsubscribes`, async () => { type Row = { id: string; rank: number } const pending = createDeferred() diff --git a/packages/db/tests/oracle-config.ts b/packages/db/tests/oracle-config.ts index f7ba4d687..9b51b818b 100644 --- a/packages/db/tests/oracle-config.ts +++ b/packages/db/tests/oracle-config.ts @@ -46,6 +46,8 @@ const staticOracleProperties = [ `collection-state.same-key`, `query-identity.compiled-output`, `query-identity.equality-partition`, + `where-predicate.publication`, + `join-result-key.pairs`, `derived-publication.membership-work`, `collection-publication.metadata-cancellation`, `collection-publication.metadata-only`, diff --git a/packages/db/tests/query/compiler/basic.test.ts b/packages/db/tests/query/compiler/basic.test.ts index 0f2fa5ca8..75e26e4a7 100644 --- a/packages/db/tests/query/compiler/basic.test.ts +++ b/packages/db/tests/query/compiler/basic.test.ts @@ -55,7 +55,8 @@ describe(`Query2 Compiler`, () => { ) expect(materialized.pipeline).toBe(compilation.pipeline) - expect(materialized.facades).toEqual([]) + // Undefined facades mark a pass-through pipeline with nothing to resolve. + expect(materialized.facades).toBeUndefined() }) test(`compiles a simple FROM query`, () => { diff --git a/packages/db/tests/query/indexes.test.ts b/packages/db/tests/query/indexes.test.ts index 6abc065f6..9c47c0a89 100644 --- a/packages/db/tests/query/indexes.test.ts +++ b/packages/db/tests/query/indexes.test.ts @@ -13,7 +13,14 @@ import { length, or, } from '../../src/query/builder/functions' -import { mockSyncCollectionOptions, stripVirtualProps } from '../utils' +import { + createIndexUsageTracker, + expectIndexUsage, + mockSyncCollectionOptions, + stripVirtualProps, + withIndexTracking, +} from '../utils' +import type { IndexUsageStats } from '../utils' interface TestItem { id: string @@ -28,163 +35,6 @@ type TestItem2 = Omit & { id2: string } -// Index usage tracking utilities (copied from collection-indexes.test.ts) -interface IndexUsageStats { - rangeQueryCalls: number - fullScanCalls: number - indexesUsed: Array - queriesExecuted: Array<{ - type: `index` | `fullScan` - operation?: string - field?: string - value?: any - }> -} - -function createIndexUsageTracker(collection: any): { - stats: IndexUsageStats - restore: () => void -} { - const stats: IndexUsageStats = { - rangeQueryCalls: 0, - fullScanCalls: 0, - indexesUsed: [], - queriesExecuted: [], - } - - // Track rangeQuery calls on index objects (index usage) - const originalIndexes = new Map() - - // Mock the indexes getter to intercept index access - const originalIndexesGetter = Object.getOwnPropertyDescriptor( - Object.getPrototypeOf(collection), - `indexes`, - )?.get - Object.defineProperty(collection, `indexes`, { - get: function () { - const indexes = originalIndexesGetter?.call(collection) || new Map() - - // Mock each index's rangeQuery method - for (const [indexId, index] of indexes.entries()) { - if (!originalIndexes.has(indexId)) { - const originalLookup = index.lookup - originalIndexes.set(indexId, originalLookup) - - index.lookup = function (operation: string, value: any) { - stats.rangeQueryCalls++ - stats.indexesUsed.push(indexId) - stats.queriesExecuted.push({ - type: `index`, - operation, - field: index.expression?.path?.join(`.`), - value, - }) - return originalLookup.call(this, operation, value) - } - } - } - - return indexes - }, - configurable: true, - }) - - // Track full scan calls (entries() iteration) - const originalEntries = collection.entries - collection.entries = function* () { - // Only count as full scan if we're in a filtering context - // Check the call stack to see if we're inside createFilterFunction - const stack = new Error().stack || `` - if ( - stack.includes(`createFilterFunction`) || - stack.includes(`currentStateAsChanges`) - ) { - stats.fullScanCalls++ - stats.queriesExecuted.push({ - type: `fullScan`, - }) - } - yield* originalEntries.call(this) - } - - const restore = () => { - // Restore original indexes getter - if (originalIndexesGetter) { - Object.defineProperty(collection, `indexes`, { - get: originalIndexesGetter, - configurable: true, - }) - } - - // Restore original lookup methods on indexes - const indexes = originalIndexesGetter?.call(collection) || new Map() - for (const [indexId, originalLookup] of originalIndexes.entries()) { - const index = indexes.get(indexId) - if (index) { - index.lookup = originalLookup - } - } - - collection.entries = originalEntries - } - - return { stats, restore } -} - -// Helper to assert index usage -function expectIndexUsage( - stats: IndexUsageStats, - expectations: { - shouldUseIndex: boolean - shouldUseFullScan?: boolean - indexCallCount?: number - fullScanCallCount?: number - }, -) { - if (expectations.shouldUseIndex) { - expect(stats.rangeQueryCalls).toBeGreaterThan(0) - expect(stats.indexesUsed.length).toBeGreaterThan(0) - - if (expectations.indexCallCount !== undefined) { - expect(stats.rangeQueryCalls).toBe(expectations.indexCallCount) - } - } else { - expect(stats.rangeQueryCalls).toBe(0) - expect(stats.indexesUsed.length).toBe(0) - } - - if (expectations.shouldUseFullScan !== undefined) { - if (expectations.shouldUseFullScan) { - expect(stats.fullScanCalls).toBeGreaterThan(0) - - if (expectations.fullScanCallCount !== undefined) { - expect(stats.fullScanCalls).toBe(expectations.fullScanCallCount) - } - } else { - expect(stats.fullScanCalls).toBe(0) - } - } -} - -// Helper to run a test with index usage tracking (automatically handles setup/cleanup) -function withIndexTracking( - collection: any, - testFn: (tracker: { stats: IndexUsageStats }) => void | Promise, -): void | Promise { - const tracker = createIndexUsageTracker(collection) - - try { - const result = testFn(tracker) - if (result instanceof Promise) { - return result.finally(() => tracker.restore()) - } - tracker.restore() - } catch (error) { - tracker.restore() - throw error - } -} - const testData: Array = [ { id: `1`, diff --git a/packages/db/tests/query/join-result-key-oracle.property.test.ts b/packages/db/tests/query/join-result-key-oracle.property.test.ts new file mode 100644 index 000000000..7fcfafc4e --- /dev/null +++ b/packages/db/tests/query/join-result-key-oracle.property.test.ts @@ -0,0 +1,289 @@ +/** + * # Does every joined pair keep its own result row? + * + * Law and source: the live-query guide defines joins like SQL joins that + * combine matching rows into single result rows, and gives a join result a + * composite key of the parent keys (`docs/guides/live-queries.md`, "Joins" and + * the `getKey` option). A joined live-query Collection therefore publishes one + * row for each matching pair of source rows, plus one row for each unmatched + * row a left or full join keeps, and two distinct pairs never share a key. The + * guide does not fix the key's format. The compiler used to join the two source + * keys with a comma, so keys that contain delimiters, or a number and a string + * that print alike, could collide and drop or reject a valid row. + * + * Model: `expectedPairs` recomputes the pair set with nested loops over plain + * arrays. It does not import the compiler or its key encoding. + * + * History grammar: left and right rows draw keys from a domain of plain, + * comma-bearing, bracket-bearing, and quoted strings, numbers alongside the + * strings that print the same, both infinities, and `NaN`. Rows join on a small + * group value. Each history uses an inner, left, or full join, then applies up + * to three synced group changes to either side. + * + * Production driver: two `mockSyncCollectionOptions` Collections and a public + * `createLiveQueryCollection` join that selects both keys. + * + * Refinement check: after preload and after each synced change, the published + * rows equal the model's pair multiset, and the result key count equals the + * row count. + * + * Calibration: the fixed pair (`a,b`, `c`) versus (`a`, `b,c`) and the pair + * (1, `c`) versus (`1`, `c`) collided under the comma encoding; restoring it + * fails both pinned histories and both campaigns. Plain `JSON.stringify` + * printed `Infinity` and `-Infinity` as `null`; restoring it fails the pinned + * infinity history and both campaigns. Comparing the join index's source-key + * prefixes with `===` let a retracted `NaN`-keyed row survive; restoring it + * fails the pinned `NaN` history. + * + * Known omissions: joins over subqueries, more than two sources, custom + * `getKey`, and optimistic mutations are outside this owner. + * The mock sync source is a controlled provider; this oracle claims only the + * compiler's keying of the rows it supplies, not any real adapter's behavior. + */ +import { fc, test as fcTest } from '@fast-check/vitest' +import { describe, expect, it } from 'vitest' +import { createCollection } from '../../src/collection/index.js' +import { createLiveQueryCollection, eq } from '../../src/query/index.js' +import { + oraclePropertyOptions, + oracleRuns, + readOracleRunConfig, +} from '../oracle-config.js' +import { mockSyncCollectionOptions, withOracleCleanup } from '../utils.js' + +const property = `join-result-key.pairs` +const requestedReplayProperty = readOracleRunConfig().replayProperty + +type Key = string | number +type Row = { id: Key; g: number } +type JoinType = `inner` | `left` | `full` +type Change = { side: `left` | `right`; index: number; g: number } +type History = { + join: JoinType + left: Array + right: Array + changes: Array +} + +// Keys that a delimiter-joined encoding cannot tell apart, plus non-finite +// numbers, which JSON prints as `null`, the missing-side marker. A `NaN` key +// also checks that the join's source-key index treats `NaN` as one key. +const keyDomain: ReadonlyArray = [ + `a`, + `b`, + `c`, + `a,b`, + `b,c`, + `[a`, + `b]`, + `"a"`, + 1, + `1`, + 2, + `2`, + Number.POSITIVE_INFINITY, + Number.NEGATIVE_INFINITY, + Number.NaN, +] + +// --------------------------------------------------------------------------- +// Model +// --------------------------------------------------------------------------- + +type Pair = { l: Key | undefined; r: Key | undefined } + +// Nested-loop join over the current rows. +function expectedPairs( + join: JoinType, + left: ReadonlyArray, + right: ReadonlyArray, +): Array { + const pairs: Array = [] + const matchedRight = new Set() + for (const l of left) { + const matches = right.filter((r) => r.g === l.g) + for (const r of matches) { + pairs.push({ l: l.id, r: r.id }) + matchedRight.add(r) + } + if (matches.length === 0 && join !== `inner`) { + pairs.push({ l: l.id, r: undefined }) + } + } + if (join === `full`) { + for (const r of right) { + if (!matchedRight.has(r)) pairs.push({ l: undefined, r: r.id }) + } + } + return pairs.map(describePair).sort() +} + +function describePair({ l, r }: Pair): string { + const side = (key: Key | undefined) => + key === undefined ? `-` : `${typeof key}:${String(key)}` + return `${side(l)} | ${side(r)}` +} + +// --------------------------------------------------------------------------- +// History grammar +// --------------------------------------------------------------------------- + +const rowsArbitrary = fc.uniqueArray( + fc.record({ + id: fc.constantFrom(...keyDomain), + g: fc.integer({ min: 0, max: 2 }), + }), + { selector: (row) => `${typeof row.id}:${String(row.id)}`, maxLength: 5 }, +) + +const historyArbitrary: fc.Arbitrary = fc.record({ + join: fc.constantFrom(`inner`, `left`, `full`), + left: rowsArbitrary, + right: rowsArbitrary, + changes: fc.array( + fc.record({ + side: fc.constantFrom(`left`, `right`), + index: fc.nat({ max: 4 }), + g: fc.integer({ min: 0, max: 2 }), + }), + { maxLength: 3 }, + ), +}) + +// Under the comma encoding, (`a,b`, `c`) and (`a`, `b,c`) share `[a,b,c]`, +// and (1, `c`) and (`1`, `c`) share `[1,c]`. Under plain JSON, (`a`, Infinity) +// and (`a`, -Infinity) share `["a",null]`, and so do the unmatched right rows +// once `a` moves away. +const pinnedHistories: ReadonlyArray = [ + { + join: `inner`, + left: [ + { id: `a,b`, g: 1 }, + { id: `a`, g: 1 }, + ], + right: [ + { id: `c`, g: 1 }, + { id: `b,c`, g: 1 }, + ], + changes: [{ side: `right`, index: 0, g: 2 }], + }, + { + join: `left`, + left: [ + { id: 1, g: 1 }, + { id: `1`, g: 1 }, + ], + right: [{ id: `c`, g: 1 }], + changes: [{ side: `right`, index: 0, g: 0 }], + }, + { + join: `full`, + left: [{ id: `a`, g: 1 }], + right: [ + { id: Number.POSITIVE_INFINITY, g: 1 }, + { id: Number.NEGATIVE_INFINITY, g: 1 }, + ], + changes: [{ side: `left`, index: 0, g: 0 }], + }, + { + // The pair leaves and re-forms; the NaN-keyed left row must cancel. + join: `inner`, + left: [{ id: Number.NaN, g: 0 }], + right: [{ id: `a`, g: 0 }], + changes: [ + { side: `left`, index: 0, g: 1 }, + { side: `right`, index: 0, g: 1 }, + ], + }, +] + +// --------------------------------------------------------------------------- +// Production driver and refinement check +// --------------------------------------------------------------------------- + +let collectionSerial = 0 + +async function runHistory(history: History): Promise { + const serial = collectionSerial++ + const leftRows = history.left.map((row) => ({ ...row })) + const rightRows = history.right.map((row) => ({ ...row })) + const left = createCollection( + mockSyncCollectionOptions({ + id: `join-key-left-${serial}`, + getKey: (row) => row.id, + initialData: leftRows.map((row) => ({ ...row })), + }), + ) + const right = createCollection( + mockSyncCollectionOptions({ + id: `join-key-right-${serial}`, + getKey: (row) => row.id, + initialData: rightRows.map((row) => ({ ...row })), + }), + ) + const live = createLiveQueryCollection((q) => + q + .from({ l: left }) + .join({ r: right }, ({ l, r }) => eq(l.g, r.g), history.join) + .select(({ l, r }) => ({ l: l.id, r: r.id })), + ) + const check = (checkpoint: string) => { + const published = live.toArray + .map((row) => describePair({ l: row.l, r: row.r })) + .sort() + expect(published, checkpoint).toEqual( + expectedPairs(history.join, leftRows, rightRows), + ) + expect(live.size, `${checkpoint} key count`).toBe(published.length) + } + await withOracleCleanup(async () => { + await live.preload() + check(`after preload`) + for (const [step, change] of history.changes.entries()) { + const rows = change.side === `left` ? leftRows : rightRows + const collection = change.side === `left` ? left : right + const target = rows[change.index % Math.max(rows.length, 1)] + if (!target || target.g === change.g) continue + collection.utils.begin() + collection.utils.write({ + type: `update`, + value: { id: target.id, g: change.g }, + }) + collection.utils.commit() + target.g = change.g + check(`after change ${step}`) + } + }, [ + () => live.cleanup(), + () => Promise.all([left.cleanup(), right.cleanup()]), + ]) +} + +describe(`joined result key oracle`, () => { + if (requestedReplayProperty === undefined) { + for (const history of pinnedHistories) { + it(`keeps every ${history.join} join pair for ${history.left + .map((row) => + typeof row.id === `number` ? String(row.id) : JSON.stringify(row.id), + ) + .join(` `)}`, () => runHistory(history)) + } + + fcTest.prop([historyArbitrary], { + seed: 44_501_962, + numRuns: oracleRuns(80), + })(`publishes one row per joined pair (fixed)`, runHistory) + + fcTest.prop([historyArbitrary], oraclePropertyOptions(80, property))( + `publishes one row per joined pair (random)`, + runHistory, + ) + } else if (requestedReplayProperty === property) { + fcTest.prop([historyArbitrary], oraclePropertyOptions(80, property))( + `publishes one row per joined pair (replay)`, + runHistory, + ) + } else { + it.skip(`runs only when its replay property is selected`, () => {}) + } +}) diff --git a/packages/db/tests/query/join.test.ts b/packages/db/tests/query/join.test.ts index 9917dca4c..b6118c95c 100644 --- a/packages/db/tests/query/join.test.ts +++ b/packages/db/tests/query/join.test.ts @@ -394,7 +394,7 @@ function testJoinType(joinType: JoinType, autoIndex: `off` | `eager`) { }) // Initially Dave has null department - const daveBefore = joinQuery.get(`[4,undefined]`) + const daveBefore = joinQuery.get(`[4,null]`) expect(daveBefore).toMatchObject({ user_name: `Dave`, department_name: undefined, @@ -419,7 +419,7 @@ function testJoinType(joinType: JoinType, autoIndex: `off` | `eager`) { department_name: `Engineering`, }) - const daveAfter2 = joinQuery.get(`[4,undefined]`) + const daveAfter2 = joinQuery.get(`[4,null]`) expect(daveAfter2).toBeUndefined() }) } diff --git a/packages/db/tests/query/live-query-result-multiplicity.test.ts b/packages/db/tests/query/live-query-result-multiplicity.test.ts new file mode 100644 index 000000000..ee3aaca20 --- /dev/null +++ b/packages/db/tests/query/live-query-result-multiplicity.test.ts @@ -0,0 +1,51 @@ +/** + * A live-query result has at most one row per key. Queries that keep their + * compiled pipeline skip the keyed reduction that used to reject extra + * contributors, so the output boundary checks the invariant for every shape. + * No public query can produce the violation; the witness injects one. + */ +import { MultiSet } from '@tanstack/db-ivm' +import { expect, it } from 'vitest' +import { createCollection } from '../../src/collection/index.js' +import { createLiveQueryCollection } from '../../src/query/index.js' +import { getCollectionBuilder } from '../../src/query/live/collection-registry.js' +import { mockSyncCollectionOptions } from '../utils.js' + +type Row = { id: string } + +it(`rejects a flush that publishes two rows for one key before writing any row`, async () => { + const source = createCollection( + mockSyncCollectionOptions({ + id: `result-multiplicity-source`, + getKey: (row) => row.id, + initialData: [{ id: `a` }], + }), + ) + const live = createLiveQueryCollection((q) => q.from({ row: source })) + try { + await live.preload() + const syncState = getCollectionBuilder(live)?.currentSyncState + const input = syncState && Object.values(syncState.inputs)[0] + if (!syncState?.graph || !syncState.flushPendingChanges || !input) { + throw new Error(`Missing live query sync state`) + } + input.sendData( + new MultiSet([ + [[`c`, { id: `c` }], 1], + [[`b`, { id: `b` }], 1], + [[`b`, { id: `b` }], 1], + ]), + ) + syncState.graph.run() + expect(() => syncState.flushPendingChanges!()).toThrow( + `a key has at most one result row`, + ) + // The check runs before any write: the valid key in the same flush is not + // published, and no sync transaction is left open. + expect([...live.keys()]).toEqual([`a`]) + expect(live._state.pendingSyncedTransactions).toEqual([]) + } finally { + await live.cleanup() + await source.cleanup() + } +}) diff --git a/packages/db/tests/query/ordered-work-oracle.property.test.ts b/packages/db/tests/query/ordered-work-oracle.property.test.ts index 84fe6c0ac..f1cee83b2 100644 --- a/packages/db/tests/query/ordered-work-oracle.property.test.ts +++ b/packages/db/tests/query/ordered-work-oracle.property.test.ts @@ -1906,7 +1906,7 @@ describe(`ordered source work oracle`, () => { expect(changedDuringStart).toBe(true) await vi.runAllTimersAsync() expect(batches).toEqual([ - [{ type: `enter`, key: `[2,undefined]`, value: truth.get(2) }], + [{ type: `enter`, key: `[2,null]`, value: truth.get(2) }], ]) expect(requests.some((options) => options.refetch === true)).toBe( true, diff --git a/packages/db/tests/query/where-predicate-publication-oracle.property.test.ts b/packages/db/tests/query/where-predicate-publication-oracle.property.test.ts new file mode 100644 index 000000000..435808f43 --- /dev/null +++ b/packages/db/tests/query/where-predicate-publication-oracle.property.test.ts @@ -0,0 +1,1133 @@ +/** + * # Which rows does a WHERE clause publish, and to which subscribers? + * + * Law and source: a WHERE clause keeps a row only when its predicate is TRUE + * under SQL three-valued logic. The evaluator contract in + * `src/query/compiler/evaluators.ts` states the operand rules: `eq` with a + * null or undefined operand is UNKNOWN, `NaN` equals `NaN` (PostgreSQL float + * semantics), a valid Date compares as its millisecond timestamp, and other + * values of different types are unequal. `not`, `and`, and `or` follow Kleene + * logic, so `not(UNKNOWN)` is UNKNOWN. Virtual row fields (`$synced`, + * `$origin`, `$key`) are row fields for filtering; a pending optimistic insert + * is `$synced: false` and `$origin: 'local'`, and a synced row from this + * adapter is `$synced: true` and `$origin: 'remote'`. + * + * Filtered publication must agree with that predicate after every source + * transaction. `Collection.subscribeChanges` with a `whereExpression` turns an + * update into an insert when the row starts matching and into a delete when it + * stops. A subscriber that requested the initial state sees every TRUE row. + * A subscriber that did not sees a TRUE row once a later insert or update + * touches it. A live-query Collection publishes every TRUE row. + * + * Why an example can miss the failure: a direct `eq` filter treats FALSE and + * UNKNOWN alike, because both drop the row; the difference appears only under + * `not`. A subscription may skip a source batch that cannot satisfy its + * predicate. That shortcut is wrong for `or`, for literals that other types + * normalize onto (a Date equals its timestamp), and for a row that moves out of + * the predicate, yet non-negated single-row examples pass all of those. + * + * Model: `expectedTruth` is an independent Kleene evaluator over plain row + * objects, and `expectedVisible` adds the touched-row rule for a subscriber + * without initial state. Neither imports production comparison, + * normalization, prefilter, or virtual-field helpers. + * + * History grammar: rows carry field `v` from a value domain of equal and + * unequal strings, a string that begins with the internal normalization + * prefix, booleans, numbers, `NaN`, a valid Date, `null`, and a missing field. + * Predicates are `eq` leaves over `v`, a literal, or a virtual field, composed + * with `not`, `and`, and `or` up to depth three. A snapshot history may hold + * pending optimistic inserts and applies no sync transactions: the + * Collection holds synced commits while a user transaction persists, and the + * optimistic-history oracle owns that law. A change history starts from synced + * rows and applies up to six synced transactions of one to three operations. + * Operations insert a new or previously deleted key, or update or delete an + * existing key, so one transaction can move several rows across the predicate + * boundary. An update + * to an equivalent value, or a second update to one key in the same + * transaction, is dropped from the history: whether a net no-op publishes is + * change detection, which the change-event history oracle owns. Each history + * runs with and without a `BasicIndex` on `v`. + * + * Production driver: a `mockSyncCollectionOptions` Collection; a live-query + * Collection built with the public `eq`/`not`/`and`/`or` builder functions; + * direct `subscribeChanges` subscribers with and without `includeInitialState` + * using the equivalent IR (a subscriber is the callback of one subscription); peer subscribers whose `eq` predicates share one field + * with different literals, plus one on a virtual field, so published changes + * must be routed to exactly the subscribers they can satisfy; and the public + * `Collection.currentStateAsChanges` snapshot. + * + * Refinement check: each consumer's exact key set equals the model after the + * subscribers attach and after each sync transaction commits. Direct + * subscribers are reconstructed from their callback batches. Two fixed + * witnesses check boundaries outside the generated grammar: Collection + * readiness reaches a filtered subscriber as one empty batch, and an eager + * source restarted after cleanup retracts a vanished row even when its first + * batch cannot satisfy the predicate. A third witness checks that a same-key + * update which stays TRUE publishes its new row and exact subscriber payload. + * + * Calibration: pinned histories separate Kleene from two-valued logic, and a + * checker test requires the two-valued answer to fail. Hostile production + * mutants that returned FALSE for `eq(string, null)`, prefiltered on `or` + * operands as if they were conjuncts, prefiltered number literals past Date or + * NaN rows, skipped an empty + * Collection-readiness batch, or skipped a batch while stale published rows awaited + * reconciliation each passed the pre-existing `@tanstack/db` suite and fail + * here. + * + * Known omissions: comparison operators other than `eq`, `in`, `like`, + * Temporal and binary operands, joins, ordering, optimistic updates and + * deletes, truncate, and failed replay are outside this owner. The structural + * comparison, cold-join, and subscription replay oracles own those boundaries. + * Generated cleanup and restart histories for filtered subscribers remain + * open; the lifecycle publication owner needs a predicate dimension to reach + * them. The mock sync source is a controlled provider; this oracle claims only + * how the Collection filters and routes the sync transactions it supplies. + */ +import { fc, test as fcTest } from '@fast-check/vitest' +import { describe, expect, it } from 'vitest' +import { createCollection } from '../../src/collection/index.js' +import { BasicIndex } from '../../src/indexes/basic-index.js' +import { and, eq, not, or } from '../../src/query/builder/functions.js' +import { Func, PropRef, Value } from '../../src/query/ir.js' +import { createLiveQueryCollection } from '../../src/query/index.js' +import { + oraclePropertyOptions, + oracleRuns, + readOracleRunConfig, +} from '../oracle-config.js' +import { + mockSyncCollectionOptions, + mockSyncCollectionOptionsNoInitialState, + withOracleCleanup, +} from '../utils.js' +import type { BasicExpression } from '../../src/query/ir.js' +import type { ChangeMessage } from '../../src/types.js' + +const property = `where-predicate.publication` +const requestedReplayProperty = readOracleRunConfig().replayProperty + +// A string whose content equals an internal normalized-key prefix must still +// compare as an ordinary string. +const PREFIXED = `\u0000tanstack-db:string:a` +const MISSING = Symbol(`missing field`) +// A valid Date equals the number of its timestamp. +const DATE_ONE = new Date(1) + +type FieldValue = string | boolean | number | Date | null | typeof MISSING +type Row = { id: string; v?: unknown } +type Truth = boolean | null + +const fieldValues: ReadonlyArray = [ + `a`, + `b`, + PREFIXED, + true, + false, + 1, + 2, + Number.NaN, + DATE_ONE, + null, + MISSING, +] +const literalValues: ReadonlyArray = [ + `a`, + PREFIXED, + true, + false, + 1, + Number.NaN, + null, +] + +type Operand = + | { kind: `field` } + | { kind: `literal`; value: FieldValue } + | { kind: `virtual`; name: `$synced` | `$origin` | `$key` } + +type Predicate = + | { kind: `eq`; left: Operand; right: Operand } + | { kind: `not`; arg: Predicate } + | { kind: `and` | `or`; args: [Predicate, Predicate] } + +type SeedRow = { v: FieldValue; optimistic: boolean } +type Operation = + | { type: `insert`; v: FieldValue; reuse?: number } + | { type: `update`; target: number; v: FieldValue } + | { type: `delete`; target: number } +type History = + | { kind: `snapshot`; rows: Array; predicate: Predicate } + | { + kind: `changes` + rows: Array + predicate: Predicate + transactions: Array> + } + +// --------------------------------------------------------------------------- +// Model +// --------------------------------------------------------------------------- + +type ModelRow = { + id: string + v: FieldValue + $synced: boolean + $origin: `local` | `remote` + $key: string +} + +function operandValue(operand: Operand, row: ModelRow): unknown { + switch (operand.kind) { + case `field`: + return row.v === MISSING ? undefined : row.v + case `literal`: + return operand.value === MISSING ? undefined : operand.value + case `virtual`: + return row[operand.name] + } +} + +// SQL equality: nullish is UNKNOWN, NaN equals NaN, a Date compares as its +// timestamp, and other values compare by type and value. +function expectedEquality(a: unknown, b: unknown): Truth { + if (a === null || a === undefined || b === null || b === undefined) { + return null + } + const left = a instanceof Date ? a.getTime() : a + const right = b instanceof Date ? b.getTime() : b + if (Number.isNaN(left) && Number.isNaN(right)) return true + return typeof left === typeof right && left === right +} + +// Kleene logic: FALSE dominates AND, TRUE dominates OR, NOT keeps UNKNOWN. +function expectedTruth(predicate: Predicate, row: ModelRow): Truth { + switch (predicate.kind) { + case `eq`: + return expectedEquality( + operandValue(predicate.left, row), + operandValue(predicate.right, row), + ) + case `not`: { + const value = expectedTruth(predicate.arg, row) + return value === null ? null : !value + } + case `and`: { + const values = predicate.args.map((arg) => expectedTruth(arg, row)) + if (values.includes(false)) return false + return values.includes(null) ? null : true + } + case `or`: { + const values = predicate.args.map((arg) => expectedTruth(arg, row)) + if (values.includes(true)) return true + return values.includes(null) ? null : false + } + } +} + +// A subscriber without initial state sees a TRUE row only after a later +// insert or update touches it. +function expectedVisible( + predicate: Predicate, + rows: ReadonlyMap, + touched?: ReadonlySet, +): Array { + return [...rows.values()] + .filter( + (row) => + expectedTruth(predicate, row) === true && + (touched === undefined || touched.has(row.id)), + ) + .map((row) => row.id) + .sort() +} + +// Updates between these values are no-ops outside this grammar. +function equivalentFieldValues(a: FieldValue, b: FieldValue): boolean { + if (a instanceof Date && b instanceof Date) { + return a.getTime() === b.getTime() + } + return Object.is(a, b) +} + +function syncedModelRow(id: string, v: FieldValue): ModelRow { + return { id, v, $synced: true, $origin: `remote`, $key: id } +} + +// --------------------------------------------------------------------------- +// History grammar +// --------------------------------------------------------------------------- + +const fieldValueArbitrary = fc.constantFrom(...fieldValues) +const operandArbitrary: fc.Arbitrary = fc.oneof( + { weight: 4, arbitrary: fc.constant({ kind: `field` as const }) }, + { + weight: 4, + arbitrary: fc + .constantFrom(...literalValues) + .map((value): Operand => ({ kind: `literal`, value })), + }, + { + weight: 1, + arbitrary: fc + .constantFrom(`$synced` as const, `$origin` as const, `$key` as const) + .map((name): Operand => ({ kind: `virtual`, name })), + }, +) +const leafArbitrary: fc.Arbitrary = fc + .tuple(operandArbitrary, operandArbitrary) + .map(([left, right]) => ({ kind: `eq` as const, left, right })) + +const predicateArbitrary: fc.Arbitrary = fc.letrec<{ + predicate: Predicate +}>((tie) => ({ + predicate: fc.oneof( + { depthSize: `small`, maxDepth: 3 }, + leafArbitrary, + tie(`predicate`).map((arg) => ({ kind: `not` as const, arg })), + fc + .tuple( + fc.constantFrom(`and` as const, `or` as const), + tie(`predicate`), + tie(`predicate`), + ) + .map(([kind, left, right]) => ({ + kind, + args: [left, right] as [Predicate, Predicate], + })), + ), +})).predicate + +// A top-level `eq(v, string | boolean)` conjunct lets a subscription skip +// batches that cannot match and lets an unindexed snapshot reject stored rows +// before copying them. Histories favor that shape so generated histories +// exercise both shortcuts, not only the full filter. +const prefilterablePredicateArbitrary: fc.Arbitrary = fc + .tuple( + fc.constantFrom(`a`, PREFIXED, true), + fc.option(predicateArbitrary, { nil: undefined }), + ) + .map(([value, rest]): Predicate => { + const conjunct: Predicate = { + kind: `eq`, + left: { kind: `field` }, + right: { kind: `literal`, value }, + } + return rest === undefined + ? conjunct + : { kind: `and`, args: [conjunct, rest] } + }) + +// Histories favor one matching and one non-matching string so rows cross the +// predicate boundary often; the full domain stays reachable. +const changeValueArbitrary = fc.oneof( + { weight: 2, arbitrary: fc.constantFrom(`a`, `b`) }, + fieldValueArbitrary, +) + +const operationArbitrary: fc.Arbitrary = fc.oneof( + fc + .tuple( + changeValueArbitrary, + fc.option(fc.nat({ max: 7 }), { nil: undefined }), + ) + .map( + ([v, reuse]): Operation => + reuse === undefined + ? { type: `insert`, v } + : { type: `insert`, v, reuse }, + ), + { + weight: 2, + arbitrary: fc + .tuple(fc.nat({ max: 7 }), changeValueArbitrary) + .map(([target, v]): Operation => ({ type: `update`, target, v })), + }, + fc.nat({ max: 7 }).map((target): Operation => ({ type: `delete`, target })), +) + +const historyArbitrary: fc.Arbitrary = fc.oneof( + fc.record({ + kind: fc.constant(`snapshot` as const), + rows: fc.array( + fc.record({ v: changeValueArbitrary, optimistic: fc.boolean() }), + { minLength: 1, maxLength: 6 }, + ), + predicate: fc.oneof(predicateArbitrary, prefilterablePredicateArbitrary), + }), + { + weight: 2, + arbitrary: fc.record({ + kind: fc.constant(`changes` as const), + rows: fc.array(changeValueArbitrary, { maxLength: 5 }), + predicate: fc.oneof(predicateArbitrary, { + weight: 2, + arbitrary: prefilterablePredicateArbitrary, + }), + transactions: fc.array( + fc.array(operationArbitrary, { minLength: 1, maxLength: 3 }), + { minLength: 1, maxLength: 6 }, + ), + }), + }, +) + +const fieldEq = (value: FieldValue): Predicate => ({ + kind: `eq`, + left: { kind: `field` }, + right: { kind: `literal`, value }, +}) + +// Kleene and two-valued logic disagree on each snapshot history. The first is +// the negated nullish comparison that a FALSE-for-UNKNOWN fast path misses. +const pinnedHistories: ReadonlyArray = [ + { + kind: `snapshot`, + rows: [ + { v: `a`, optimistic: false }, + { v: `b`, optimistic: false }, + { v: null, optimistic: false }, + { v: MISSING, optimistic: false }, + ], + predicate: { kind: `not`, arg: fieldEq(`a`) }, + }, + { + kind: `snapshot`, + rows: [ + { v: `a`, optimistic: true }, + { v: null, optimistic: false }, + { v: Number.NaN, optimistic: false }, + { v: PREFIXED, optimistic: true }, + ], + predicate: { + kind: `or`, + args: [ + { + kind: `not`, + arg: { + kind: `eq`, + left: { kind: `literal`, value: true }, + right: { kind: `field` }, + }, + }, + { + kind: `eq`, + left: { kind: `virtual`, name: `$synced` }, + right: { kind: `literal`, value: false }, + }, + ], + }, + }, + { + kind: `snapshot`, + rows: [ + { v: PREFIXED, optimistic: false }, + { v: `a`, optimistic: false }, + { v: Number.NaN, optimistic: false }, + { v: 1, optimistic: false }, + { v: MISSING, optimistic: false }, + ], + predicate: { + kind: `and`, + args: [ + { kind: `not`, arg: fieldEq(PREFIXED) }, + { kind: `not`, arg: fieldEq(Number.NaN) }, + ], + }, + }, +] + +// An unindexed snapshot may reject stored rows by one field before copying +// them. A pending optimistic row is visible although it is not yet synced. +const pinnedPrefilteredSnapshot: History = { + kind: `snapshot`, + rows: [ + { v: `a`, optimistic: true }, + { v: `a`, optimistic: false }, + { v: `b`, optimistic: true }, + ], + predicate: fieldEq(`a`), +} + +// Each change history moves rows across a predicate whose equality operand a +// skip-if-no-match shortcut can misread. +const pinnedChangeHistories: ReadonlyArray = [ + { + // `or` operands are alternatives, not conjuncts. + kind: `changes`, + rows: [`b`, `a`], + predicate: { kind: `or`, args: [fieldEq(`a`), fieldEq(true)] }, + transactions: [ + [{ type: `update`, target: 0, v: true }], + [ + { type: `update`, target: 1, v: false }, + { type: `insert`, v: true }, + ], + ], + }, + { + // A valid Date equals its timestamp, so a number literal cannot prefilter + // by identity. + kind: `changes`, + rows: [2, 2], + predicate: fieldEq(1), + transactions: [ + [{ type: `update`, target: 0, v: DATE_ONE }], + [ + { type: `update`, target: 0, v: 2 }, + { type: `update`, target: 1, v: 1 }, + ], + ], + }, + { + // NaN equals NaN, although NaN is not identical to itself. + kind: `changes`, + rows: [2], + predicate: fieldEq(Number.NaN), + transactions: [[{ type: `update`, target: 0, v: Number.NaN }]], + }, + { + // A subscriber without initial state never saw the dropped `b` row. + // Deleting it and reinserting its key with a matching value must publish + // the new row, not treat it as a duplicate insert. + kind: `changes`, + rows: [], + predicate: fieldEq(`a`), + transactions: [ + [ + { type: `insert`, v: `a` }, + { type: `insert`, v: `b` }, + ], + [{ type: `delete`, target: 1 }], + [{ type: `insert`, v: `a`, reuse: 0 }], + ], + }, + { + // A row that moves out must be retracted, including inside one batch. + kind: `changes`, + rows: [`a`, `a`, `b`], + predicate: { + kind: `and`, + args: [fieldEq(`a`), { kind: `not`, arg: fieldEq(PREFIXED) }], + }, + transactions: [ + [ + { type: `update`, target: 0, v: `b` }, + { type: `update`, target: 2, v: `a` }, + ], + [ + { type: `delete`, target: 1 }, + { type: `insert`, v: `a` }, + { type: `update`, target: 0, v: `a` }, + ], + ], + }, +] + +// --------------------------------------------------------------------------- +// Production driver +// --------------------------------------------------------------------------- + +type RefLike = Record + +function builderOperand(operand: Operand, row: RefLike): unknown { + switch (operand.kind) { + case `field`: + return row.v + case `literal`: + return operand.value === MISSING ? undefined : operand.value + case `virtual`: + return row[operand.name] + } +} + +function builderPredicate( + predicate: Predicate, + row: RefLike, +): BasicExpression { + switch (predicate.kind) { + case `eq`: + return eq( + builderOperand(predicate.left, row) as never, + builderOperand(predicate.right, row) as never, + ) + case `not`: + return not(builderPredicate(predicate.arg, row)) + case `and`: + return and( + builderPredicate(predicate.args[0], row), + builderPredicate(predicate.args[1], row), + ) + case `or`: + return or( + builderPredicate(predicate.args[0], row), + builderPredicate(predicate.args[1], row), + ) + } +} + +function irOperand(operand: Operand): BasicExpression { + switch (operand.kind) { + case `field`: + return new PropRef([`v`]) + case `literal`: + return new Value(operand.value === MISSING ? undefined : operand.value) + case `virtual`: + return new PropRef([operand.name]) + } +} + +function irPredicate(predicate: Predicate): BasicExpression { + switch (predicate.kind) { + case `eq`: + return new Func(`eq`, [ + irOperand(predicate.left), + irOperand(predicate.right), + ]) + case `not`: + return new Func(`not`, [irPredicate(predicate.arg)]) + case `and`: + case `or`: + return new Func(predicate.kind, predicate.args.map(irPredicate)) + } +} + +function sourceRow(id: string, v: FieldValue): Row { + return v === MISSING ? { id } : { id, v } +} + +type Consumer = + | `live query` + | `subscriber with initial state` + | `subscriber without initial state` + | `direct snapshot` + | `peer ${string}` + +// Peers share one field with different literals, plus one on a second field, +// so each change must reach exactly the peers whose predicate it can satisfy. +const peerPredicates: ReadonlyArray = [ + ...[`a`, `b`, PREFIXED, true, false].map(fieldEq), + { + kind: `eq`, + left: { kind: `virtual`, name: `$synced` }, + right: { kind: `literal`, value: true }, + }, +] + +type Observation = { + checkpoint: string + consumer: Consumer + keys: Array +} + +// Rebuild a subscriber's visible key set from its callback batches. +function recordVisibleKeys(visible: Set) { + return (changes: Array>) => { + for (const change of changes) { + if (change.type === `delete`) visible.delete(change.key) + else visible.add(change.key) + } + } +} + +let collectionSerial = 0 + +async function observeHistory( + history: History, + indexed: boolean, +): Promise<{ observations: Array; model: Array }> { + const seeds: Array = + history.kind === `snapshot` + ? history.rows.map((seed, index) => ({ ...seed, id: `r${index}` })) + : history.rows.map((v, index) => ({ + v, + optimistic: false, + id: `r${index}`, + })) + const collection = createCollection( + mockSyncCollectionOptions({ + id: `where-publication-${collectionSerial++}`, + getKey: (row) => row.id, + initialData: seeds + .filter((seed) => !seed.optimistic) + .map((seed) => sourceRow(seed.id, seed.v)), + }), + ) + const live = createLiveQueryCollection((q) => + q + .from({ row: collection }) + .where(({ row }) => + builderPredicate(history.predicate, row as unknown as RefLike), + ), + ) + const modelRows = new Map( + seeds.map((seed) => [ + seed.id, + { + ...syncedModelRow(seed.id, seed.v), + $synced: !seed.optimistic, + $origin: seed.optimistic ? `local` : `remote`, + }, + ]), + ) + const touched = new Set() + const withInitial = new Set() + const withoutInitial = new Set() + const observations: Array = [] + const model: Array = [] + const record = ( + checkpoint: string, + consumer: Consumer, + keys: Iterable, + predicate: Predicate = history.predicate, + ) => { + observations.push({ + checkpoint, + consumer, + keys: [...keys].map(String).sort(), + }) + model.push({ + checkpoint, + consumer, + keys: expectedVisible( + predicate, + modelRows, + consumer === `subscriber without initial state` ? touched : undefined, + ), + }) + } + const peers = peerPredicates.map((predicate) => ({ + predicate, + visible: new Set(), + })) + const recordSubscribers = (checkpoint: string) => { + record(checkpoint, `live query`, live.keys()) + record(checkpoint, `subscriber with initial state`, withInitial) + record(checkpoint, `subscriber without initial state`, withoutInitial) + for (const peer of peers) { + record( + checkpoint, + `peer ${describePredicate(peer.predicate)}`, + peer.visible, + peer.predicate, + ) + } + } + + const subscriptions: Array<{ unsubscribe: () => void }> = [] + await withOracleCleanup(async () => { + await collection.stateWhenReady() + if (indexed) { + collection.createIndex((row) => row.v, { indexType: BasicIndex }) + } + // Pending optimistic inserts stay local: the adapter's onInsert awaits a + // sync acknowledgement that this history never sends. + for (const seed of seeds.filter((entry) => entry.optimistic)) { + collection.insert(sourceRow(seed.id, seed.v)) + } + await live.preload() + subscriptions.push( + collection.subscribeChanges(recordVisibleKeys(withInitial), { + includeInitialState: true, + whereExpression: irPredicate(history.predicate), + }), + collection.subscribeChanges(recordVisibleKeys(withoutInitial), { + whereExpression: irPredicate(history.predicate), + }), + ...peers.map((peer) => + collection.subscribeChanges(recordVisibleKeys(peer.visible), { + includeInitialState: true, + whereExpression: irPredicate(peer.predicate), + }), + ), + ) + recordSubscribers(`after subscribe`) + const snapshot = collection.currentStateAsChanges({ + where: irPredicate(history.predicate), + }) + if (!snapshot) throw new Error(`an unoptimized snapshot must return rows`) + record( + `after subscribe`, + `direct snapshot`, + snapshot.map((change) => change.key), + ) + + if (history.kind === `changes`) { + let nextKey = 0 + // Keys deleted by an earlier transaction may be reinserted. Reinsertion + // in the deleting transaction is a net change outside this grammar. + const deletedKeys: Array = [] + for (const [step, operations] of history.transactions.entries()) { + collection.utils.begin() + const updatedKeys = new Set() + const reusableKeys = [...deletedKeys] + for (const operation of operations) { + const keys = [...modelRows.keys()] + if (operation.type === `insert`) { + const reused = + operation.reuse !== undefined && reusableKeys.length > 0 + ? reusableKeys.splice( + operation.reuse % reusableKeys.length, + 1, + )[0]! + : undefined + if (reused !== undefined) { + deletedKeys.splice(deletedKeys.indexOf(reused), 1) + } + const id = reused ?? `n${nextKey++}` + collection.utils.write({ + type: `insert`, + value: sourceRow(id, operation.v), + }) + modelRows.set(id, syncedModelRow(id, operation.v)) + touched.add(id) + continue + } + // Targets name a current key; an empty Collection has none. + if (keys.length === 0) continue + const id = keys[operation.target % keys.length]! + if (operation.type === `delete`) { + collection.utils.write({ type: `delete`, value: { id } }) + modelRows.delete(id) + touched.delete(id) + deletedKeys.push(id) + continue + } + if ( + updatedKeys.has(id) || + equivalentFieldValues(modelRows.get(id)!.v, operation.v) + ) { + continue + } + updatedKeys.add(id) + // Synced updates merge fields, so a missing field is written as an + // explicit undefined; both are UNKNOWN operands. + collection.utils.write({ + type: `update`, + value: { id, v: operation.v === MISSING ? undefined : operation.v }, + }) + modelRows.set(id, syncedModelRow(id, operation.v)) + touched.add(id) + } + collection.utils.commit() + recordSubscribers(`after transaction ${step}`) + } + } + }, [ + () => { + for (const subscription of subscriptions) subscription.unsubscribe() + }, + () => live.cleanup(), + () => collection.cleanup(), + ]) + return { observations, model } +} + +async function runHistory(history: History): Promise { + for (const indexed of [false, true]) { + const { observations, model } = await observeHistory(history, indexed) + const path = `${indexed ? `index` : `scan`} path for ${describePredicate(history.predicate)}` + const listObservations = (list: Array) => + list.map((o) => `${o.checkpoint} ${o.consumer}: [${o.keys.join(`,`)}]`) + expect(listObservations(observations), path).toEqual( + listObservations(model), + ) + } +} + +function describeValue(value: FieldValue): string { + if (value === MISSING) return `` + if (typeof value === `string`) return JSON.stringify(value) + if (value instanceof Date) return `Date(${value.getTime()})` + return String(value) +} + +function describePredicate(predicate: Predicate): string { + const operand = (o: Operand) => + o.kind === `field` + ? `v` + : o.kind === `virtual` + ? o.name + : describeValue(o.value) + switch (predicate.kind) { + case `eq`: + return `eq(${operand(predicate.left)}, ${operand(predicate.right)})` + case `not`: + return `not(${describePredicate(predicate.arg)})` + case `and`: + case `or`: + return `${predicate.kind}(${predicate.args.map(describePredicate).join(`, `)})` + } +} + +// --------------------------------------------------------------------------- +// Refinement check and calibration +// --------------------------------------------------------------------------- + +describe(`WHERE predicate publication oracle`, () => { + if (requestedReplayProperty === undefined) { + it(`pinned snapshot histories distinguish Kleene logic from two-valued logic`, () => { + // A checker that folds UNKNOWN into FALSE before NOT must disagree with + // the model on each pinned snapshot, or the history does not calibrate. + const twoValued = (predicate: Predicate, row: ModelRow): boolean => { + switch (predicate.kind) { + case `eq`: + return ( + expectedEquality( + operandValue(predicate.left, row), + operandValue(predicate.right, row), + ) === true + ) + case `not`: + return !twoValued(predicate.arg, row) + case `and`: + return predicate.args.every((arg) => twoValued(arg, row)) + case `or`: + return predicate.args.some((arg) => twoValued(arg, row)) + } + } + for (const history of pinnedHistories) { + if (history.kind !== `snapshot`) throw new Error(`expected a snapshot`) + const rows = new Map( + history.rows.map((seed, index): [string, ModelRow] => [ + `r${index}`, + { + ...syncedModelRow(`r${index}`, seed.v), + $synced: !seed.optimistic, + $origin: seed.optimistic ? `local` : `remote`, + }, + ]), + ) + const naive = [...rows.values()] + .filter((row) => twoValued(history.predicate, row)) + .map((row) => row.id) + .sort() + expect(naive).not.toEqual(expectedVisible(history.predicate, rows)) + } + }) + + it(`Collection readiness reaches a filtered subscriber as one empty batch`, async () => { + // `markReady` notifies subscribers of Collection readiness with an + // empty batch. A subscriber + // whose predicate no row can satisfy must still receive it. + const collection = createCollection( + mockSyncCollectionOptionsNoInitialState({ + id: `where-publication-ready-${collectionSerial++}`, + getKey: (row) => row.id, + }), + ) + const batches: Array = [] + const subscription = collection.subscribeChanges( + (changes) => batches.push(changes.length), + { whereExpression: irPredicate(fieldEq(`a`)) }, + ) + await withOracleCleanup(() => { + collection.utils.begin() + collection.utils.commit() + collection.utils.markReady() + expect(batches).toEqual([0]) + }, [() => subscription.unsubscribe(), () => collection.cleanup()]) + }) + + it(`publishes the new payload for a same-key update that stays TRUE`, async () => { + type PayloadRow = { id: string; v: string; n: number } + const collection = createCollection( + mockSyncCollectionOptions({ + id: `where-publication-payload-${collectionSerial++}`, + getKey: (row) => row.id, + initialData: [{ id: `r0`, v: `a`, n: 1 }], + }), + ) + const live = createLiveQueryCollection((q) => + q.from({ row: collection }).where(({ row }) => eq(row.v, `a`)), + ) + const events: Array<{ + type: string + key: string | number + n: number + previousN: number | undefined + }> = [] + await collection.stateWhenReady() + await live.preload() + const subscription = collection.subscribeChanges( + (changes) => { + for (const change of changes) { + events.push({ + type: change.type, + key: change.key, + n: change.value.n, + previousN: change.previousValue?.n, + }) + } + }, + { + includeInitialState: true, + whereExpression: irPredicate(fieldEq(`a`)), + }, + ) + await withOracleCleanup(() => { + events.length = 0 + expect(live.get(`r0`)?.n).toBe(1) + collection.utils.begin() + collection.utils.write({ + type: `update`, + value: { id: `r0`, v: `a`, n: 2 }, + }) + collection.utils.commit() + expect(live.get(`r0`)?.n).toBe(2) + expect(events).toEqual([ + { type: `update`, key: `r0`, n: 2, previousN: 1 }, + ]) + }, [ + () => subscription.unsubscribe(), + () => live.cleanup(), + () => collection.cleanup(), + ]) + }) + + it(`a pending optimistic delete hides a row from a prefiltered unindexed scan`, async () => { + // The scan may read synced rows directly only when no optimistic state + // changes visibility. A pending delete removes a synced row. + const collection = createCollection( + mockSyncCollectionOptions({ + id: `where-publication-optimistic-delete-${collectionSerial++}`, + getKey: (row) => row.id, + initialData: [ + { id: `r0`, v: `a` }, + { id: `r1`, v: `a` }, + { id: `r2`, v: `b` }, + ], + }), + ) + const visible = new Set() + let subscription: { unsubscribe: () => void } | undefined + await withOracleCleanup(async () => { + await collection.stateWhenReady() + // The adapter's onDelete awaits a sync acknowledgement that this test + // never sends, so the delete stays pending. + collection.delete(`r1`) + subscription = collection.subscribeChanges(recordVisibleKeys(visible), { + includeInitialState: true, + whereExpression: irPredicate(fieldEq(`a`)), + }) + const snapshot = collection.currentStateAsChanges({ + where: irPredicate(fieldEq(`a`)), + }) + expect(snapshot?.map((change) => change.key)).toEqual([`r0`]) + expect([...visible]).toEqual([`r0`]) + }, [() => subscription?.unsubscribe(), () => collection.cleanup()]) + }) + + it(`a layout-only publication reaches a filtered subscriber as one empty batch`, async () => { + // An ordered live query whose rows move without changing their selected + // values publishes no rows, only a layout change. Every subscriber still + // learns that the layout changed, whatever its predicate. + type RankedRow = { id: string; rank: number } + const source = createCollection( + mockSyncCollectionOptions({ + id: `where-publication-layout-${collectionSerial++}`, + getKey: (row) => row.id, + initialData: [ + { id: `r0`, rank: 1 }, + { id: `r1`, rank: 2 }, + ], + }), + ) + const live = createLiveQueryCollection((q) => + q + .from({ row: source }) + .orderBy(({ row }) => row.rank) + .select(({ row }) => ({ id: row.id })), + ) + const batches: Array = [] + let subscription: { unsubscribe: () => void } | undefined + await withOracleCleanup(async () => { + await live.preload() + subscription = live.subscribeChanges( + (changes) => batches.push(changes.length), + { + whereExpression: new Func(`eq`, [ + new PropRef([`id`]), + new Value(`missing`), + ]), + }, + ) + source.utils.begin() + source.utils.write({ type: `update`, value: { id: `r0`, rank: 3 } }) + source.utils.commit() + expect([...live.keys()]).toEqual([`r1`, `r0`]) + expect(batches).toEqual([0]) + }, [ + () => subscription?.unsubscribe(), + () => live.cleanup(), + () => source.cleanup(), + ]) + }) + + it(`a restarted source retracts a vanished row even when its first batch cannot match`, async () => { + // Cleanup keeps the subscriber's published rows. The restarted eager + // source's first publication must retract rows it no longer holds, + // although no change in that batch satisfies the predicate. The check + // runs before ready, whose empty batch would also reconcile them. + let sourceRows: Array = [{ id: `r0`, v: `a` }] + let markSourceReady = () => {} + const collection = createCollection({ + id: `where-publication-restart-${collectionSerial++}`, + getKey: (row) => row.id, + startSync: true, + sync: { + sync: ({ begin, write, commit, markReady }) => { + begin() + for (const row of sourceRows) write({ type: `insert`, value: row }) + commit() + markSourceReady = markReady + }, + }, + }) + const visible = new Set() + markSourceReady() + await collection.stateWhenReady() + const subscription = collection.subscribeChanges( + recordVisibleKeys(visible), + { + includeInitialState: true, + whereExpression: irPredicate(fieldEq(`a`)), + }, + ) + await withOracleCleanup(async () => { + expect([...visible]).toEqual([`r0`]) + sourceRows = [{ id: `r1`, v: `b` }] + await collection.cleanup() + collection.startSyncImmediate() + expect([...visible]).toEqual([]) + markSourceReady() + expect([...visible]).toEqual([]) + }, [() => subscription.unsubscribe(), () => collection.cleanup()]) + }) + + for (const history of [ + ...pinnedHistories, + pinnedPrefilteredSnapshot, + ...pinnedChangeHistories, + ]) { + it(`publishes TRUE rows for ${history.kind} history ${describePredicate(history.predicate)}`, () => + runHistory(history)) + } + + fcTest.prop([historyArbitrary], { + seed: 44_500_301, + numRuns: oracleRuns(80), + })(`publishes exactly the TRUE rows (fixed)`, runHistory) + + fcTest.prop([historyArbitrary], oraclePropertyOptions(80, property))( + `publishes exactly the TRUE rows (random)`, + runHistory, + ) + } else if (requestedReplayProperty === property) { + fcTest.prop([historyArbitrary], oraclePropertyOptions(80, property))( + `publishes exactly the TRUE rows (replay)`, + runHistory, + ) + } else { + it.skip(`runs only when its replay property is selected`, () => {}) + } +}) diff --git a/packages/db/tests/query/where-prefilter-property-visibility.test.ts b/packages/db/tests/query/where-prefilter-property-visibility.test.ts new file mode 100644 index 000000000..7efb54a57 --- /dev/null +++ b/packages/db/tests/query/where-prefilter-property-visibility.test.ts @@ -0,0 +1,136 @@ +/** + * A stored-row prefilter may reject only rows that the enriched-row predicate + * cannot accept. Enrichment copies enumerable own root properties; the full + * predicate catches nested property reads that throw. These descriptor cases + * complement the plain-row grammar in the WHERE publication oracle. + */ +import { expect, it } from 'vitest' +import { createCollection } from '../../src/collection/index.js' +import { Func, PropRef, Value } from '../../src/query/ir.js' +import { mockSyncCollectionOptions } from '../utils.js' + +type Row = { id: string; v?: string; nested?: { v?: string } } +const where = new Func('eq', [new PropRef(['v']), new Value('a')]) +const nestedWhere = new Func('eq', [ + new PropRef(['nested', 'v']), + new Value('a'), +]) + +function rowWithGetter(kind: 'inherited' | 'non-enumerable') { + const getter = () => { + throw new Error('getter read') + } + if (kind === 'inherited') { + const prototype = Object.defineProperty({}, 'v', { get: getter }) + return Object.assign(Object.create(prototype) as Row, { id: kind }) + } + return Object.defineProperty({ id: kind } as Row, 'v', { + get: getter, + enumerable: false, + }) +} + +it.each(['inherited', 'non-enumerable'] as const)( + 'ignores a %s getter omitted from enriched rows', + async (kind) => { + const collection = createCollection( + mockSyncCollectionOptions({ + id: `prefilter-getter-${kind}`, + getKey: (item) => item.id, + initialData: [rowWithGetter(kind)], + }), + ) + try { + await collection.stateWhenReady() + expect(collection.currentStateAsChanges({ where })).toEqual([]) + } finally { + await collection.cleanup() + } + }, +) + +it('keeps an enumerable own getter as a visible field', async () => { + const row = Object.defineProperty({ id: 'enumerable' } as Row, 'v', { + get: () => 'a', + enumerable: true, + }) + const collection = createCollection( + mockSyncCollectionOptions({ + id: 'prefilter-getter-enumerable', + getKey: (item) => item.id, + initialData: [row], + }), + ) + try { + await collection.stateWhenReady() + expect( + collection.currentStateAsChanges({ where })?.map((x) => x.key), + ).toEqual(['enumerable']) + } finally { + await collection.cleanup() + } +}) + +it('lets the full predicate handle a throwing nested getter', async () => { + const nested = Object.create( + Object.defineProperty({}, 'v', { + get() { + throw new Error('nested getter read') + }, + }), + ) as { v?: string } + const collection = createCollection( + mockSyncCollectionOptions({ + id: 'prefilter-getter-nested', + getKey: (item) => item.id, + initialData: [{ id: 'nested', nested }], + }), + ) + try { + await collection.stateWhenReady() + expect(collection.currentStateAsChanges({ where: nestedWhere })).toEqual([]) + } finally { + await collection.cleanup() + } +}) + +it('lets the full filter handle a throwing nested getter in a change', async () => { + const nested = () => + Object.create( + Object.defineProperty({}, 'v', { + get() { + throw new Error('nested getter read') + }, + }), + ) as { v?: string } + const collection = createCollection( + mockSyncCollectionOptions({ + id: 'prefilter-getter-nested-change', + getKey: (item) => item.id, + initialData: [], + }), + ) + const batches: Array> = [] + let subscription: { unsubscribe: () => void } | undefined + try { + await collection.stateWhenReady() + subscription = collection.subscribeChanges( + (changes) => batches.push(changes.map((change) => change.key)), + { whereExpression: nestedWhere }, + ) + collection.utils.begin() + collection.utils.write({ + type: 'insert', + value: { id: 'thrower', nested: nested() }, + }) + collection.utils.write({ + type: 'insert', + value: { id: 'match', nested: { v: 'a' } }, + }) + collection.utils.commit() + expect(batches).toEqual([['match']]) + } finally { + subscription?.unsubscribe() + await collection.cleanup() + } +}) diff --git a/packages/db/tests/utils.ts b/packages/db/tests/utils.ts index 586000fc0..bf39b787b 100644 --- a/packages/db/tests/utils.ts +++ b/packages/db/tests/utils.ts @@ -17,6 +17,39 @@ export type OutputWithVirtual< T extends object, TKey extends string | number = string | number, > = WithVirtualProps +/** + * Runs an oracle check, then every cleanup step in order. A cleanup failure + * never replaces the check's own failure: the check failure is thrown alone or + * as the `cause` of an `AggregateError` that also holds each cleanup failure. + */ +export async function withOracleCleanup( + check: () => Promise | void, + cleanups: ReadonlyArray<() => unknown>, +): Promise { + const failures: Array = [] + let checkFailed = false + try { + await check() + } catch (error) { + checkFailed = true + failures.push(error) + } + for (const cleanup of cleanups) { + try { + await cleanup() + } catch (error) { + failures.push(error) + } + } + if (failures.length === 1) throw failures[0] + if (failures.length > 1) { + throw new AggregateError( + failures, + checkFailed ? `Oracle check and cleanup failed` : `Oracle cleanup failed`, + { cause: failures[0] }, + ) + } +} // Keep sync startup, writes, readiness, and load outcomes in the test itself. export function createOnDemandCollection( @@ -80,32 +113,38 @@ export function createIndexUsageTracker(collection: any): { queriesExecuted: [], } - // Track index method calls by patching all existing indexes + // Track index method calls. Indexes are patched when first read through the + // collection, so indexes created after tracking starts (such as auto-indexes + // added by a live query) are tracked too. const originalMethods = new Map() - - for (const [indexId, index] of collection.indexes) { + // A range lookup delegates to rangeQuery; record it once, as the lookup. + let lookupDepth = 0 + const patchIndex = (indexId: unknown, index: any) => { + if (originalMethods.has(indexId)) return // Track lookup calls (new unified method) const originalLookup = index.lookup.bind(index) index.lookup = function (operation: any, value: any) { - // Only track non-range operations to avoid double counting - // Range operations (gt, gte, lt, lte) are handled by rangeQuery tracking - if (![`gt`, `gte`, `lt`, `lte`].includes(operation)) { - stats.rangeQueryCalls++ - stats.indexesUsed.push(String(indexId)) - stats.queriesExecuted.push({ - type: `index`, - operation, - field: index.expression?.path?.join(`.`), - value, - }) + stats.rangeQueryCalls++ + stats.indexesUsed.push(String(indexId)) + stats.queriesExecuted.push({ + type: `index`, + operation, + field: index.expression?.path?.join(`.`), + value, + }) + lookupDepth++ + try { + return originalLookup(operation, value) + } finally { + lookupDepth-- } - return originalLookup(operation, value) } // Track rangeQuery calls (for compound range queries) - if (index.rangeQuery) { - const originalRangeQuery = index.rangeQuery.bind(index) + const originalRangeQuery = index.rangeQuery?.bind(index) + if (originalRangeQuery) { index.rangeQuery = function (options: any) { + if (lookupDepth > 0) return originalRangeQuery(options) stats.rangeQueryCalls++ stats.indexesUsed.push(String(indexId)) @@ -129,14 +168,30 @@ export function createIndexUsageTracker(collection: any): { } originalMethods.set(indexId, { + index, lookup: originalLookup, - rangeQuery: index.rangeQuery ? index.rangeQuery.bind(index) : undefined, + rangeQuery: originalRangeQuery, }) } + const originalIndexesGetter = Object.getOwnPropertyDescriptor( + Object.getPrototypeOf(collection), + `indexes`, + )?.get + const readIndexes = (): Map => + originalIndexesGetter?.call(collection) ?? new Map() + for (const [indexId, index] of readIndexes()) patchIndex(indexId, index) + Object.defineProperty(collection, `indexes`, { + get: () => { + const indexes = readIndexes() + for (const [indexId, index] of indexes) patchIndex(indexId, index) + return indexes + }, + configurable: true, + }) - // Track full scan calls (entries() iteration) - const originalEntries = collection.entries - collection.entries = function* () { + // Track full scan calls: filtered iteration through either the public + // entries() or the stored-row scan used by a prefiltered snapshot. + const recordFullScan = () => { // Only count as full scan if we're in a filtering context // Check the call stack to see if we're inside createFilterFunction const stack = new Error().stack || `` @@ -149,21 +204,28 @@ export function createIndexUsageTracker(collection: any): { type: `fullScan`, }) } + } + const originalEntries = collection.entries + collection.entries = function* () { + recordFullScan() yield* originalEntries.call(this) } + const state = collection._state + const originalEntriesPassing = state.entriesPassing + state.entriesPassing = function* (prefilter: (row: object) => boolean) { + recordFullScan() + yield* originalEntriesPassing.call(this, prefilter) + } const restore = () => { - // Restore original index methods - for (const [indexId, index] of collection.indexes) { - const original = originalMethods.get(indexId) - if (original) { - index.lookup = original.lookup - if (original.rangeQuery) { - index.rangeQuery = original.rangeQuery - } - } + // Remove the instance getter so the prototype getter applies again + delete collection.indexes + for (const { index, lookup, rangeQuery } of originalMethods.values()) { + index.lookup = lookup + if (rangeQuery) index.rangeQuery = rangeQuery } collection.entries = originalEntries + state.entriesPassing = originalEntriesPassing } return { stats, restore } diff --git a/packages/react-db/tests/useLiveQuery.test.tsx b/packages/react-db/tests/useLiveQuery.test.tsx index fdfbd3864..ad2574193 100644 --- a/packages/react-db/tests/useLiveQuery.test.tsx +++ b/packages/react-db/tests/useLiveQuery.test.tsx @@ -464,19 +464,19 @@ describe(`Query Collections`, () => { // Verify that we have the expected joined results - expect(result.current.state.get(`[1,1]`)).toMatchObject({ + expect(result.current.state.get(`["1","1"]`)).toMatchObject({ id: `1`, name: `John Doe`, title: `Issue 1`, }) - expect(result.current.state.get(`[2,2]`)).toMatchObject({ + expect(result.current.state.get(`["2","2"]`)).toMatchObject({ id: `2`, name: `Jane Doe`, title: `Issue 2`, }) - expect(result.current.state.get(`[3,1]`)).toMatchObject({ + expect(result.current.state.get(`["3","1"]`)).toMatchObject({ id: `3`, name: `John Doe`, title: `Issue 3`, @@ -500,7 +500,7 @@ describe(`Query Collections`, () => { await waitFor(() => { expect(result.current.state.size).toBe(4) }) - expect(result.current.state.get(`[4,2]`)).toMatchObject({ + expect(result.current.state.get(`["4","2"]`)).toMatchObject({ id: `4`, name: `Jane Doe`, title: `Issue 4`, @@ -523,7 +523,7 @@ describe(`Query Collections`, () => { await waitFor(() => { // The updated title should be reflected in the joined results - expect(result.current.state.get(`[2,2]`)).toMatchObject({ + expect(result.current.state.get(`["2","2"]`)).toMatchObject({ id: `2`, name: `Jane Doe`, title: `Updated Issue 2`, @@ -548,7 +548,7 @@ describe(`Query Collections`, () => { await new Promise((resolve) => setTimeout(resolve, 10)) // After deletion, issue 3 should no longer have a joined result - expect(result.current.state.get(`[3,1]`)).toBeUndefined() + expect(result.current.state.get(`["3","1"]`)).toBeUndefined() expect(result.current.state.size).toBe(3) }) @@ -835,8 +835,8 @@ describe(`Query Collections`, () => { useEffect(() => { renderStates.push({ stateSize: queryResult.state.size, - hasTempKey: queryResult.state.has(`[temp-key,1]`), - hasPermKey: queryResult.state.has(`[4,1]`), + hasTempKey: queryResult.state.has(`["temp-key","1"]`), + hasPermKey: queryResult.state.has(`["4","1"]`), timestamp: Date.now(), }) }, [queryResult.state]) @@ -915,12 +915,12 @@ describe(`Query Collections`, () => { await waitFor(() => { // Verify optimistic state is immediately reflected expect(result.current.state.size).toBe(4) - expect(result.current.state.get(`[temp-key,1]`)).toMatchObject({ + expect(result.current.state.get(`["temp-key","1"]`)).toMatchObject({ id: `temp-key`, name: `John Doe`, title: `New Issue`, }) - expect(result.current.state.get(`[4,1]`)).toBeUndefined() + expect(result.current.state.get(`["4","1"]`)).toBeUndefined() }) // Wait for the transaction to be committed @@ -928,7 +928,7 @@ describe(`Query Collections`, () => { await waitFor(() => { // Wait for the permanent key to appear - expect(result.current.state.get(`[4,1]`)).toBeDefined() + expect(result.current.state.get(`["4","1"]`)).toBeDefined() }) // Check if we had any render where the temp key was removed but the permanent key wasn't added yet @@ -941,8 +941,8 @@ describe(`Query Collections`, () => { // Verify the temporary key is replaced by the permanent one expect(result.current.state.size).toBe(4) - expect(result.current.state.get(`[temp-key,1]`)).toBeUndefined() - expect(result.current.state.get(`[4,1]`)).toMatchObject({ + expect(result.current.state.get(`["temp-key","1"]`)).toBeUndefined() + expect(result.current.state.get(`["4","1"]`)).toMatchObject({ id: `4`, name: `John Doe`, title: `New Issue`, diff --git a/packages/solid-db/tests/useLiveQuery.test.tsx b/packages/solid-db/tests/useLiveQuery.test.tsx index aba713cf1..43e1dff80 100644 --- a/packages/solid-db/tests/useLiveQuery.test.tsx +++ b/packages/solid-db/tests/useLiveQuery.test.tsx @@ -399,19 +399,19 @@ describe(`Query Collections`, () => { // Verify that we have the expected joined results - expect(result.state.get(`[1,1]`)).toMatchObject({ + expect(result.state.get(`["1","1"]`)).toMatchObject({ id: `1`, name: `John Doe`, title: `Issue 1`, }) - expect(result.state.get(`[2,2]`)).toMatchObject({ + expect(result.state.get(`["2","2"]`)).toMatchObject({ id: `2`, name: `Jane Doe`, title: `Issue 2`, }) - expect(result.state.get(`[3,1]`)).toMatchObject({ + expect(result.state.get(`["3","1"]`)).toMatchObject({ id: `3`, name: `John Doe`, title: `Issue 3`, @@ -433,7 +433,7 @@ describe(`Query Collections`, () => { await waitFor(() => { expect(result.state.size).toBe(4) }) - expect(result.state.get(`[4,2]`)).toMatchObject({ + expect(result.state.get(`["4","2"]`)).toMatchObject({ id: `4`, name: `Jane Doe`, title: `Issue 4`, @@ -454,7 +454,7 @@ describe(`Query Collections`, () => { await waitFor(() => { // The updated title should be reflected in the joined results - expect(result.state.get(`[2,2]`)).toMatchObject({ + expect(result.state.get(`["2","2"]`)).toMatchObject({ id: `2`, name: `Jane Doe`, title: `Updated Issue 2`, @@ -476,7 +476,7 @@ describe(`Query Collections`, () => { await waitFor(() => { // After deletion, issue 3 should no longer have a joined result - expect(result.state.get(`[3,1]`)).toBeUndefined() + expect(result.state.get(`["3","1"]`)).toBeUndefined() expect(result.state.size).toBe(3) }) }) @@ -879,8 +879,8 @@ describe(`Query Collections`, () => { createComputed(() => { renderStates.push({ stateSize: queryResult.state.size, - hasTempKey: queryResult.state.has(`[temp-key,1]`), - hasPermKey: queryResult.state.has(`[4,1]`), + hasTempKey: queryResult.state.has(`["temp-key","1"]`), + hasPermKey: queryResult.state.has(`["4","1"]`), timestamp: Date.now(), }) }) @@ -954,19 +954,19 @@ describe(`Query Collections`, () => { // Verify optimistic state is immediately reflected expect(result.state.size).toBe(4) }) - expect(result.state.get(`[temp-key,1]`)).toMatchObject({ + expect(result.state.get(`["temp-key","1"]`)).toMatchObject({ id: `temp-key`, name: `John Doe`, title: `New Issue`, }) - expect(result.state.get(`[4,1]`)).toBeUndefined() + expect(result.state.get(`["4","1"]`)).toBeUndefined() // Wait for the transaction to be committed await transaction.isPersisted.promise await waitFor(() => { // Wait for the permanent key to appear - expect(result.state.get(`[4,1]`)).toBeDefined() + expect(result.state.get(`["4","1"]`)).toBeDefined() }) // Check if we had any render where the temp key was removed but the permanent key wasn't added yet @@ -979,8 +979,8 @@ describe(`Query Collections`, () => { // Verify the temporary key is replaced by the permanent one expect(result.state.size).toBe(4) - expect(result.state.get(`[temp-key,1]`)).toBeUndefined() - expect(result.state.get(`[4,1]`)).toMatchObject({ + expect(result.state.get(`["temp-key","1"]`)).toBeUndefined() + expect(result.state.get(`["4","1"]`)).toMatchObject({ id: `4`, name: `John Doe`, title: `New Issue`, diff --git a/packages/svelte-db/tests/useLiveQuery.svelte.test.ts b/packages/svelte-db/tests/useLiveQuery.svelte.test.ts index a4f8d9e41..1dfd8694c 100644 --- a/packages/svelte-db/tests/useLiveQuery.svelte.test.ts +++ b/packages/svelte-db/tests/useLiveQuery.svelte.test.ts @@ -505,19 +505,19 @@ describe(`Query Collections`, () => { // Verify that we have the expected joined results expect(query.state.size).toBe(3) - expect(query.state.get(`[1,1]`)).toMatchObject({ + expect(query.state.get(`["1","1"]`)).toMatchObject({ id: `1`, name: `John Doe`, title: `Issue 1`, }) - expect(query.state.get(`[2,2]`)).toMatchObject({ + expect(query.state.get(`["2","2"]`)).toMatchObject({ id: `2`, name: `Jane Doe`, title: `Issue 2`, }) - expect(query.state.get(`[3,1]`)).toMatchObject({ + expect(query.state.get(`["3","1"]`)).toMatchObject({ id: `3`, name: `John Doe`, title: `Issue 3`, @@ -539,7 +539,7 @@ describe(`Query Collections`, () => { flushSync() expect(query.state.size).toBe(4) - expect(query.state.get(`[4,2]`)).toMatchObject({ + expect(query.state.get(`["4","2"]`)).toMatchObject({ id: `4`, name: `Jane Doe`, title: `Issue 4`, @@ -561,7 +561,7 @@ describe(`Query Collections`, () => { flushSync() // The updated title should be reflected in the joined results - expect(query.state.get(`[2,2]`)).toMatchObject({ + expect(query.state.get(`["2","2"]`)).toMatchObject({ id: `2`, name: `Jane Doe`, title: `Updated Issue 2`, @@ -583,7 +583,7 @@ describe(`Query Collections`, () => { flushSync() // After deletion, issue 3 should no longer have a joined result - expect(query.state.get(`[3,1]`)).toBeUndefined() + expect(query.state.get(`["3","1"]`)).toBeUndefined() expect(query.state.size).toBe(3) }) }) @@ -807,7 +807,7 @@ describe(`Query Collections`, () => { // Verify the new issue is reflected in the query expect(queryResult.state.size).toBe(4) - expect(queryResult.state.get(`[4,1]`)).toMatchObject({ + expect(queryResult.state.get(`["4","1"]`)).toMatchObject({ id: `4`, name: `John Doe`, title: `New Issue`, diff --git a/packages/vue-db/tests/useLiveQuery.test.ts b/packages/vue-db/tests/useLiveQuery.test.ts index 65ea1a274..58e43da6f 100644 --- a/packages/vue-db/tests/useLiveQuery.test.ts +++ b/packages/vue-db/tests/useLiveQuery.test.ts @@ -332,19 +332,19 @@ describe(`Query Collections`, () => { // Verify that we have the expected joined results expect(state.value.size).toBe(3) - expect(state.value.get(`[1,1]`)).toMatchObject({ + expect(state.value.get(`["1","1"]`)).toMatchObject({ id: `1`, name: `John Doe`, title: `Issue 1`, }) - expect(state.value.get(`[2,2]`)).toMatchObject({ + expect(state.value.get(`["2","2"]`)).toMatchObject({ id: `2`, name: `Jane Doe`, title: `Issue 2`, }) - expect(state.value.get(`[3,1]`)).toMatchObject({ + expect(state.value.get(`["3","1"]`)).toMatchObject({ id: `3`, name: `John Doe`, title: `Issue 3`, @@ -366,7 +366,7 @@ describe(`Query Collections`, () => { await waitForVueUpdate() expect(state.value.size).toBe(4) - expect(state.value.get(`[4,2]`)).toMatchObject({ + expect(state.value.get(`["4","2"]`)).toMatchObject({ id: `4`, name: `Jane Doe`, title: `Issue 4`, @@ -388,7 +388,7 @@ describe(`Query Collections`, () => { await waitForVueUpdate() // The updated title should be reflected in the joined results - expect(state.value.get(`[2,2]`)).toMatchObject({ + expect(state.value.get(`["2","2"]`)).toMatchObject({ id: `2`, name: `Jane Doe`, title: `Updated Issue 2`, @@ -410,7 +410,7 @@ describe(`Query Collections`, () => { await waitForVueUpdate() // After deletion, issue 3 should no longer have a joined result - expect(state.value.get(`[3,1]`)).toBeUndefined() + expect(state.value.get(`["3","1"]`)).toBeUndefined() expect(state.value.size).toBe(3) }) @@ -621,8 +621,8 @@ describe(`Query Collections`, () => { watchEffect(() => { renderStates.push({ stateSize: state.value.size, - hasTempKey: state.value.has(`[temp-key,1]`), - hasPermKey: state.value.has(`[4,1]`), + hasTempKey: state.value.has(`["temp-key","1"]`), + hasPermKey: state.value.has(`["4","1"]`), timestamp: Date.now(), }) }) @@ -694,12 +694,12 @@ describe(`Query Collections`, () => { // Verify optimistic state is immediately reflected (should be synchronous) expect(state.value.size).toBe(4) - expect(state.value.get(`[temp-key,1]`)).toMatchObject({ + expect(state.value.get(`["temp-key","1"]`)).toMatchObject({ id: `temp-key`, name: `John Doe`, title: `New Issue`, }) - expect(state.value.get(`[4,1]`)).toBeUndefined() + expect(state.value.get(`["4","1"]`)).toBeUndefined() // Wait for the transaction to be committed await transaction.isPersisted.promise @@ -708,8 +708,8 @@ describe(`Query Collections`, () => { // Verify the temporary key is replaced by the permanent one expect(state.value.size).toBe(4) - expect(state.value.get(`[temp-key,1]`)).toBeUndefined() - expect(state.value.get(`[4,1]`)).toMatchObject({ + expect(state.value.get(`["temp-key","1"]`)).toBeUndefined() + expect(state.value.get(`["4","1"]`)).toMatchObject({ id: `4`, name: `John Doe`, title: `New Issue`,