Repository navigation
feat(truapi): observe core store queries as streams - #1125
Merged
Merged
Conversation
`Db::observe(sql, |q| …)` streams a query's result: once when first polled, then again after every commit that changes a table the query reads. This is the store's counterpart of Room's Flow queries, which the durable engine and its domains use to follow their own tables. The writer's `update_hook` collects touched tables and `Db::write` publishes them after `commit()` returns, so a woken observer always reads the committed rows and a rolled-back write wakes nobody. Wakes go through a `channel(0)` per observer, so a burst of commits costs one re-query. The tables a query reads are found with SQLite's authorizer, once per SQL constant, and matched against the schema's tables, which covers joins, subqueries, views and `count(*)`. The query closure gets an `ObservedStatement` rather than the connection, so it cannot run SQL the detection did not see. The statement records the raw values of every row it reads, and a re-query whose rows equal the last emission's emits nothing, without requiring `PartialEq` on the result. Every connection carries an authorizer that answers `Ignore` for deletes, which makes SQLite delete row by row instead of truncating, so an unconditional `DELETE FROM t` still reaches `update_hook`. A schema test rejects `WITHOUT ROWID` core tables, which the hook never reports. A failed re-query yields `Err` and keeps observing; the stream ends when the database closes or the query's tables cannot be resolved. Refs #964 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Review feedback: a function should not embed several self-contained pieces of logic. Table resolution becomes `names_read_by` (what the authorizer reports) and `schema_tables` (the real tables), intersected by `resolve_tables`. The same split applies to waking one observer (`Observer::notify`), registering one (`Invalidation::register`), the per-SQL table cache (`Db::tables`), the dedupe re-query loop (`Db::next_emission`), installing the writer's hooks (`track_changes`) and running a write transaction (`in_transaction`). Behaviour is unchanged. Refs #964 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
An unchanged re-query no longer loops inside the stream step waiting for the next wake. `Db::requery` runs the query once and returns `None` when the rows equal the last emission's, and the stream drops those with `filter_map`. Each wake now runs exactly one query, and waiting for a wake happens in one place. Refs #964 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Contributor
Bundle size reportCompared with
WebAssembly modules
No file changed size. Commit: 789d871 |
Contributor
|
CI Status: 24 required jobs green, 21 passed and 3 skipped by path filter. All job results
Signing credentials: failure as of 2026-10-02, a release may fail Commit |
Contributor
iOS simulator previewBuilt from gh run download 36868114828 --name simulator-preview-ef091824a
unzip polkadot-app-*.app.zip
xcrun simctl install booted polkadot-app.app
xcrun simctl launch booted io.parity.polkadotapp.developOr download it in a browser, which arrives as a zip wrapping An arm64 simulator slice, so it needs an Apple Silicon Mac and does not |
pgherveou
approved these changes
Oct 1, 2026
A dropped stream's observer stayed registered until the next write to one of its tables or the database closing, keeping its channel and the waiting task's waker alive. The stream now owns a `Registration` whose `Drop` unregisters the observer, and that is the only place observers are removed: `publish` and `close` only wake them. Review feedback also renames the shared authorizer to `authorize_for_change_tracking` and inlines the schema-tables SQL into `schema_tables`. Refs #964 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
An observed query could run its statement several times, appending every run's rows to one flat snapshot. Two different results could then record the same rows: buckets `(["a"], ["b"])` and `([], ["a", "b"])` both snapshot `["a", "b"]`, so dedupe dropped a real change. The query closure now takes `ObservedStatement` by value, and `query_map`, `query_optional` and `query_row` consume it, so a second run does not compile. Reading several parameter sets takes one observer each, each with its own snapshot. The snapshot is owned by the refresh and lent to the statement. Refs #964 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
valentunn
enabled auto-merge
October 2, 2026 09:58
github-merge-queue
Bot
removed this pull request from the merge queue due to no response for status checks
Oct 2, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Part 2 of #964
Summary
Adds ability to observe database. Consumers can specify a sql query which results they want to observe. They get a stream where each emission = real change of the results of that query.
Implementation idea:
authorizerthat records table read authorizations. We filter out sqlite-internal tablesNotable design decision
DELETE * FROMqueries because sqlite might utilize trimming optimization. We fix that by disallowing such, with a side effect that each row will now be deleted on its own. An alternative solution would be to use per-table invalidation triggers, but this is more complex and schema-invasive. We can always switch to it later if we find a need to. For now I dont anticipate any that large delete queries where the difference would matter