Repository navigation
feat(platform): subscribe to committed state transitions matching document, address, identity, token and contract filters - #5283
Closed
PastaPastaPasta wants to merge 11 commits into
Conversation
…system-field clauses `DriveDocumentQueryFilter` compared clause values and transition data with `Value`'s derived equality, so the same value in two accepted encodings never matched: an identifier given as bytes or base58 text against `Value::Identifier`, a `u8` property sent as `U64`. Values a clause reads are now brought to the form the document type stores before comparing, and `canonicalize_clause_values` does the same once for the clause operands (element-wise for IN and BETWEEN*, identifiers for the recipient/buyer clause, `U64` for the price clause). Schema properties go through their property type's codec directly, so values longer than an index key still compare. `validate()` now rejects clauses on `$`-prefixed system fields other than the primary-key `$id` clauses: document data carries no such fields, so those clauses could never match. The filter has no consumer yet, so neither change affects consensus. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
A server-streaming Platform RPC that delivers committed, successfully executed state transitions matching any of a request's filters: document transitions on a contract (narrowed by type, action, where clauses on the new data, `$id` of the original, recipient/buyer, price, batch owner), platform addresses, identities and tokens on a chosen side (sender/recipient/any), and data contracts. Clauses reuse the getDocuments V1 typed WhereClause. The stream sends each matching transition with its block height, time, protocol version, position, hash and bytes, and checkpoints that tell the client where to resume (`from_block_height`). Drive answers it with `unimplemented`: DAPI serves it from the Tenderdash block store. The Platform module now allows the camel-case stream type tonic generates, as the Core module already does. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
`dash_platform_queries::subscriptions` holds the transport-free filter model, its wire conversion and the matcher, so DAPI (which evaluates the filters) and clients (which re-check what DAPI sends) run the same code: - `StateTransitionFilter` and builders, `subscribe_request`; - `ResolvedFilters::resolve` validates limits and binds document filters to their contracts through `DriveDocumentQueryFilter`, rejecting clauses that cannot be decided from a transition (original-document clauses other than `$id`, except an indexOnly delete); - `ResolvedFilters::matches` reports the matched filters and, for batches, the matched inner transition positions; - every `StateTransition` variant is classified by the parties it names, so a new variant does not compile until it is. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ck store Each subscription walks committed blocks in height order with its own cursor, reading history and new blocks the same way through a shared, byte-bounded block cache; new-block websocket events only wake the tip tracker, so a dropped event delays a block but cannot lose it. - Heights are never skipped: a block whose results are not saved yet is retried, `blockchain` pages are requested 20 heights at a time and checked complete, and an undecodable transaction or unknown protocol version ends the stream at that height instead of advancing past it. - A checkpoint follows every block with a match and repeats every 10s, so clients can resume and quiet streams outlive proxy idle timeouts. - Limits: 1024 subscriptions per node, 16 per client address (last X-Forwarded-For entry; IPv6 per /64), 8 concurrent catch-ups, starts at most 50k blocks behind or 1k ahead of the tip, and a client that does not read for 60s is dropped with RESOURCE_EXHAUSTED. - A data contract update in the stream rebinds the document filters on that contract. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
`Sdk::subscribe_to_state_transitions(filters, from_block_height)` returns a `StateTransitionSubscription` (`next()`, or `into_stream()`) yielding matching transitions and checkpoints. - Every transition the node sends is checked: its hash, and a re-match against the filters bound to data contracts fetched with proofs; a transition that does not match is skipped and logged. - Broken streams resume from the last checkpoint, possibly on another node, with backoff; only progress clears the failure count, so a node failing right after every start is given up on after five attempts. Delivery is at least once. - A data contract update in the stream rebinds the local filters, as on the node. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Gives the stream the Core streams' timeouts (300s idle, 600s lifetime) and documents the endpoint. Replaces the route of `subscribePlatformEvents`, an RPC that does not exist. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…by slow clients - Reads are single-flight: concurrent misses on one block or one metas page share a Tenderdash call, short pages at the tip are cached by their exact range, and a block's transactions are decoded once for every subscriber. - The replay permit covers only reads from Tenderdash and is released before anything is sent, so a client that stops reading cannot hold the node's catch-up capacity. - A scan stops as soon as its client goes away, including while waiting for a block's results; start-height bounds use saturating arithmetic. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…suming - The subscription also asks for updates of the contracts its document filters name and rebinds its local filters on them, as the node does, so the two never check against different contract versions; those updates are not reported unless the caller asked for them. - A stream that ends cleanly counts as a failure with backoff, and only a stream that got past its opening checkpoint clears the failure count, so a node that keeps ending streams cannot cause a tight reconnect loop. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
ResolvedFilters::follow replaces the identical rebinding the node and the client each carried; drops unused len/is_empty. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Contributor
|
Important Draft PR not reviewedDraft PRs are not automatically reviewed by default.
To automatically review draft PRs, update your CodeRabbit configuration: reviews:
auto_review:
drafts: true
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
Collaborator
|
🕓 Review not started yet because this PR is a draft.
Commit 5681b07. Normal review starts when eligible; priority review starts as soon as a slot is available. |
5 of 7 tasks
Member
Author
|
Superseded by #5285: same branch and commits, opened from the dashpay/platform repository so CI runs. Closing this fork-based PR. 🤖 Posted autonomously by Codex on behalf of pasta. |
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.
Issue being fixed or feature implemented
Closes #5274.
Applications that react to Platform activity poll for it today, for example Dash Forge's inbox, CI runners and webhook relay. A wallet watching its addresses and a token payment gateway waiting for a transfer do the same. Polling cost grows with users times feeds, runs into the gateway's per-IP limit, and notices changes late.
This adds
subscribeToStateTransitions, a server-streaming Platform RPC. A client describes what it wants to hear about, and DAPI streams each committed, successfully executed state transition that matches. It covers the requests "tell me when this address gets funds", "when a document appears on this contract/type", "when a document matching this query is created", "when this document changes", "when this identity receives credits or tokens" and "when this contract is updated". A client resumes from a checkpoint without gaps.The closed #2795 was an earlier attempt: drive-abci events for committed blocks and transaction hashes, with no entity filters and no resume. Drive already had
DriveDocumentQueryFilterfor this purpose (#2761, #2781, #5114), but nothing used it.What was done?
Filters (
dash_platform_queries::subscriptions, shared by DAPI and the SDK)A request carries 1–16 filters. A transition is delivered when it matches any filter, and every constraint within one filter must hold.
$idof the document before the transition;SENDER,RECIPIENTorANY. The roles are classified for everyStateTransitionvariant in one exhaustive match, so a new variant won't compile until it is classified.Matching reads only what a transition says, so some filters are refused up front rather than accepted and then silently never matching. A clause on the document as it was before the transition is undecidable unless it is on
$id. The exception is the delete of an indexOnly document, which carries the document's values. The endpoint doc lists the remaining limits: requested vs credited amounts, mints to the configured destination, shielded parties, and group actions that are proposed but not yet executed.Drive filter fixes (
rs-drive/src/query/filter.rs)The filter had no consumer and two bugs that made valid filters silently never match:
Different encodings of the same value never compared equal. Two examples:
Identifier;u8field sent asU64.Clause values and the transition values a clause reads are now both brought to the form the schema stores them in. Long values still compare, because the per-type codec is used directly instead of the index-key path, which caps values at 255 bytes.
Clauses on
$system fields other than$idpassedvalidate()but can never match, because transition data has no$fields. They are now rejected.DAPI (
rs-dapi/.../subscribe_to_state_transitions)unimplemented.blockchainreturns only the 20 highest heights of a range, so pages are requested 20 at a time and checked complete.hmeans everything from the first scanned block throughhhas been sent; resume fromh + 1.X-Forwarded-Forentry; IPv6 per /64).RESOURCE_EXHAUSTED.SDK (
rs-sdk/src/platform/subscriptions.rs)Sdk::subscribe_to_state_transitions(filters, from_block_height)returns a subscription withnext()/into_stream().Wiring
build.rsand the metrics allowlist.subscribePlatformEvents, which has no RPC.packages/dapi/doc/endpoints/streams/subscribeToStateTransitions.md.Trust model (stated in the docs)
A client can check every delivered transition: its hash, and that it matches its filters against proved contracts. A client cannot check completeness: a node may leave a match out, and block heights and times are the node's word. Act on an event by reading the changed state with a proved query, and catch up from a proved read's
metadata.height + 1.Design process
Not in this PR
wasm32, andrs-dapi-client's grpc-web transport supports server streaming. The JS surface is a follow-up.How Has This Been Tested?
cargo test -p drive --lib query::filter: 67 passed. New tests cover the identifier encodings, long values, rejected system fields and recipients given as bytes.cargo test -p dash-platform-queries --lib subscriptions: 12 passed. They cover:$idoriginal-document matching;cargo test -p rs-dapi --lib: 344 passed. The 21 new tests run the full scan loop against an in-memory chain:cargo test -p dash-sdk --lib subscriptions. The new network-only testtests/fetch/state_transition_subscriptions.rscompiles undernetwork-testingbut has not been run against a live network.cargo clippyis clean for every touched crate.dash-sdkchecks forwasm32-unknown-unknown. The dashmate Envoy template test passes.check-grpc-coverage.pypasses.Breaking Changes
None. The RPC and messages are new. The drive filter changes affect only
DriveDocumentQueryFilter, which nothing else uses and which is not consensus code; no shipped generation changes behaviour.Checklist:
structure.rs, regeneratedgrovedb-structure.json, and checked the structure viewer link posted on this pull requestFor repository code-owners and collaborators only
🤖 Generated with Claude Code