Skip to content

feat(platform)!: subscribe to committed state transitions matching document, address, identity, token and contract filters - #5285

Draft
PastaPastaPasta wants to merge 25 commits into
v5.1-devfrom
feat/platform-subscriptions
Draft

PastaPastaPasta wants to merge 25 commits into
v5.1-devfrom
feat/platform-subscriptions

Conversation

@PastaPastaPasta

@PastaPastaPasta PastaPastaPasta commented Oct 5, 2026 •

Copy link
Copy Markdown
Member

Issue being fixed or feature implemented

Closes #5274.

Supersedes #5283, which was opened from a fork, where CI does not run.

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 DriveDocumentQueryFilter for 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.

  • documents: one data contract, optionally narrowed by:
    • a document type;
    • actions (create / replace / delete / transfer / update price / purchase), each with clauses on the new data, typed like getDocuments V1 where clauses;
    • $id of the document before the transition;
    • the transfer recipient or the buyer;
    • the new price;
    • the batch owner.
  • addresses, up to 256 platform addresses, and identities, up to 64, each on a side: SENDER, RECIPIENT or ANY. The roles are classified for every StateTransition variant in one exhaustive match, so a new variant won't compile until it is classified.
  • tokens: by token id and/or party.
  • data contracts: creation, updates, moderation and fee claims.

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:

    • an identifier sent as bytes vs base58 vs Identifier;
    • a u8 field sent as U64.

    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 $id passed validate() but can never match, because transition data has no $ fields. They are now rejected.

DAPI (rs-dapi/.../subscribe_to_state_transitions)

  • Served from the Tenderdash block store, not drive-abci. Each subscription walks committed heights in order with its own cursor. History and new blocks are read the same way. New-block websocket events only wake the tip tracker, so a dropped event delays a block but cannot lose it. No consensus code changes; Drive answers the RPC with unimplemented.
  • Never skips a height:
    • A block stored before its results are saved is retried.
    • blockchain returns only the 20 highest heights of a range, so pages are requested 20 at a time and checked complete.
    • A transaction this build cannot decode, or an unknown protocol version, ends the stream at that height instead of advancing past it.
  • Stream contract:
    • The first message is a checkpoint at start − 1.
    • Matches are sent in block order, and a checkpoint follows every block with a match.
    • A checkpoint repeats every 10 s otherwise, which keeps quiet streams inside proxy idle timeouts.
    • A checkpoint at h means everything from the first scanned block through h has been sent; resume from h + 1.
  • Cost:
    • Reads are single-flight: concurrent misses on one block or metas page share one Tenderdash call, and a block is decoded once for all subscribers.
    • The block cache is bounded by bytes.
  • Limits:
    • 1024 subscriptions per node.
    • 16 per client address (the last X-Forwarded-For entry; IPv6 per /64).
    • 8 subscriptions catching up on history at once. The permit is held only while reading, so a client that stops reading cannot keep it.
    • Start at most 50k blocks behind or 1k ahead of the tip.
    • A client that does not read for 60 s is dropped with RESOURCE_EXHAUSTED.
  • Contract updates: a data contract update in the stream rebinds the document filters on that contract.

SDK (rs-sdk/src/platform/subscriptions.rs)

Sdk::subscribe_to_state_transitions(filters, from_block_height) returns a subscription with next() / into_stream().

  • Checks what the node sends. It verifies each transition's hash and re-matches it against the filters, bound to data contracts fetched with proofs. Transitions that don't match are skipped and logged.
  • Resumes by itself from the last checkpoint, possibly on another node, with backoff. Delivery is at least once.
  • Follows updates of the contracts it filters on, so it matches against the same contract version as the node.

Wiring

  • Proto and regenerated clients.
  • build.rs and the metrics allowlist.
  • An Envoy route with the Core streams' timeouts. It replaces the route for subscribePlatformEvents, which has no RPC.
  • The endpoint doc 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

  • Four codebase explorations, then two rounds of design review by three independent reviewers. The reviewers were two Claude agents, one checking protocol correctness against the Tenderdash source and one checking codebase fit, plus a Codex GPT-6.1 reviewer for API and product.
  • An independent code review afterwards. Every finding above MINOR was fixed. Two findings in that category:
    • one client could hold the catch-up capacity;
    • every subscriber issued its own Tenderdash calls per block.

Not in this PR

  • wasm-sdk / js-evo-sdk bindings, an async iterator for browsers. The matcher and the Rust SDK already build for wasm32, and rs-dapi-client's grpc-web transport supports server streaming. The JS surface is a follow-up.
  • Firehose filters, all transitions of a kind. Explorers are better served by running indexers.

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:
    • roles for identities and addresses;
    • batch attribution across document and token transitions;
    • $id original-document matching;
    • rejected undecidable clauses;
    • limits;
    • the wire round-trip;
    • rebinding.
  • cargo test -p rs-dapi --lib: 344 passed. The 21 new tests run the full scan loop against an in-memory chain:
    • replay across pages, then live blocks, without gaps;
    • starting live;
    • checkpoints while idle;
    • start bounds;
    • per-address limits and IPv6 /64 bucketing;
    • a block whose results arrive late;
    • a client that stops reading must not starve another's catch-up;
    • an undecodable transaction ends the stream;
    • meta page validation, including heights below the store base and header-only metas.
  • cargo test -p dash-sdk --lib subscriptions. The new network-only test tests/fetch/state_transition_subscriptions.rs compiles under network-testing but has not been run against a live network.
  • cargo clippy is clean for every touched crate. dash-sdk checks for wasm32-unknown-unknown. The dashmate Envoy template test passes. check-grpc-coverage.py passes.
  • Not done: an end-to-end run on a local dashmate devnet.

Breaking Changes

  • Rust gRPC server API (dapi-grpc): platform_server::Platform gains a required associated type subscribeToStateTransitionsStream and method subscribe_to_state_transitions (tonic generates no default stubs). Downstream implementations of the trait no longer compile until they add both; rs-dapi serves it and drive-abci's QueryService returns unimplemented, which shows the minimal migration.
  • On the protobuf wire the RPC and messages are purely additive. The drive filter changes affect only DriveDocumentQueryFilter, which nothing else uses and which is not consensus code; no shipped protocol generation changes behaviour.

In-place changes to shipped generations

  • rs-platform-value Value's Display (display.rs): text longer than 20 bytes was truncated with split_at(20), which panics when byte 20 falls inside a multi-byte UTF-8 character. It now cuts at the last character boundary at or before byte 20. This cannot modify consensus anywhere it is used: output is byte-identical for every input that did not panic before (all text whose byte 20 is a character boundary, including all ASCII), and an input that used to panic would have aborted the process rather than produced a result to agree on.

Checklist:

  • I have performed a self-review of my own code
  • I have commented my code, particularly in hard-to-understand areas
  • I have added or updated relevant unit/integration/functional/e2e tests
  • I have added "!" to the title and described breaking changes in the corresponding section if my code contains any
  • I have made corresponding changes to the documentation if needed
  • If I added or changed GroveDB structure, I described it in the area's structure.rs, regenerated grovedb-structure.json, and checked the structure viewer link posted on this pull request

For repository code-owners and collaborators only

  • I have assigned this pull request to a milestone

🤖 Generated with Claude Code

@coderabbitai

coderabbitai Bot commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

Important

Draft PR not reviewed

Draft PRs are not automatically reviewed by default.

  • Trigger a manual review

To automatically review draft PRs, update your CodeRabbit configuration:

reviews:
  auto_review:
    drafts: true
  • Autopilot · Keep fixing CodeRabbit findings and required CI, and resolving merge conflicts

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@github-actions github-actions Bot added this to the v5.1.0 milestone Oct 5, 2026
@thepastaclaw

thepastaclaw commented Oct 5, 2026 •

Copy link
Copy Markdown
Collaborator

⛔ Final review complete — 1 blocking finding(s) (commit a1a1cf8) · triage: critical

PastaPastaPasta added a commit to thepastaclaw/review-system that referenced this pull request Oct 5, 2026
* fix(ingest): a review requested on a draft runs to the end

Ticking "Request normal/priority review" on a draft's gate comment
re-queued its head, and the next PR-ingest pass closed it as `pr_draft`
because the PR was still a draft, cancelling the run a few minutes after
it started and putting the draft body back with the box unticked
(dashpay/platform#5285, run 3512, 2026-10-05).

The draft branch of `ingest_repo` now keeps a head for the draft's
current commit when someone asked for it (any trigger but new_pr /
new_push); automatic heads and heads for older commits still close.
The deferred-comment pass no longer overwrites the worker's running,
done or failed body on such a draft, but still honours a box ticked on
an older draft body.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix(queue): a failed review on a draft gets the draft body's boxes back

Mentions and review requests are ignored on drafts, so the boxes are the
only way to retry there.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Preliminary review — Phase 1 blocker gate

Verified the supplied findings against head 5681b07. Seven in-scope suggestions remain, including incomplete SDK contract following, unknown-version handling, and historical schema binding; none violates the review policy's consensus-blocking invariants. This was a static review: the supplied CI snapshot shows Rust workspace tests skipped and PR Hygiene pending, so it does not establish that the subscription tests passed.

The Phase-1 findings confirmed by a fresh verifier are above the gate budget. Phase 2 is deferred until a fresh same-head revalidation brings them under it.

🟡 7 suggestion(s)

Review provenance

Source: reviewer 1: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 2: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); final verifier: gpt-6.1-sol (agent: sol-gate-verifier, role: verifier)

  • Triage: normal by gpt-6.1-sol (effort low) — This is a large, intricate streaming RPC and SDK change, but the diff changes client-facing query/filter handling rather than consensus rules, funds movement, cryptography, peer-to-peer deserialization, or storage migrations.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Fresh verifier: gpt-6.1-sol — verifier; agent sol-gate-verifier
  • Phase 2 reviewers: not run (deferred by the Phase-1 gate)
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-sdk/src/platform/subscriptions.rs`:
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:122-128: SDK silently drops its contract-follow filter when the request is full
  When a request contains 16 filters and includes document filters, this branch omits the internal contract-update filter but still accepts the subscription. DAPI calls follow() for every successfully executed transition it scans, including updates that match none of the caller's filters; the SDK can call follow() only for transitions it receives. Unless another caller filter happens to match those updates, the SDK retains its old contract binding while DAPI rebinds. This contradicts the method's guarantee that both sides follow contract updates and can cause local matching to diverge after an update. Reject requests that leave no room for the internal filter, or ensure the caller's existing filters cover all required contract updates. This is a client subscription correctness issue, not a consensus blocker.
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:279-280: SDK falls back to latest() for unknown protocol versions instead of skipping
  An unknown protocol version is replaced with PlatformVersion::latest(), after which the SDK evaluates filters and follows contract updates using tables that do not describe the advertised version. Hash verification and successful decoding do not establish that this substitution is valid. DAPI explicitly fails closed for unknown block versions, whereas the SDK can emit an event carrying an unsupported protocol version or mutate its filter bindings under guessed rules. Reject or skip the event with a diagnostic rather than selecting a different protocol version.

In `packages/dash-platform-queries/src/subscriptions/resolved.rs`:
- [SUGGESTION] packages/dash-platform-queries/src/subscriptions/resolved.rs:174-192: Historical replay starts with the current contract rather than the historical binding
  DAPI fetches the current contract before starting the scan, and follow() changes that binding only when a contract update is encountered. Consequently, the initial historical segment is evaluated using the current schema. Backward-compatible updates protect many existing-field filters, so ordinary field additions do not automatically imply incorrect matching. However, matching also calls data_as_stored(), which generates omitted generatedFrom properties using the bound schema. A generated property added after the replayed block can therefore be synthesized for an old transition and satisfy a clause even though that property was not generated when the transition executed. Bind replay to the contract version active at its starting height, or explicitly document that historical document clauses use the initial current-schema binding rather than execution-time schema semantics.
- [SUGGESTION] packages/dash-platform-queries/src/subscriptions/resolved.rs:297-312: Document clause matching deep-clones resolved clauses on the scan hot path
  For each relevant inner document transition, this branch clones the document type name and clones each attempted DocumentActionMatchClauses alternative into a temporary DriveDocumentQueryFilter. The clause structure owns collections and values, so these are deep clones rather than borrowed views. Block decoding is shared, but matching is performed separately for every subscription, multiplying these allocations across subscriptions, transitions, and alternatives. Prebuild reusable match state when resolving the filter or provide a borrowing matcher so evaluation does not clone the resolved clauses.

In `packages/dapi-grpc/protos/platform/v0/platform.proto`:
- [SUGGESTION] packages/dapi-grpc/protos/platform/v0/platform.proto:4247-4249: Document ActionMatch cannot distinguish a missing action from CREATE
  ActionMatch.action is a non-optional proto3 enum, and CREATE is zero. An omitted action therefore decodes as CREATE, and action_match_from_proto converts it to Some(DocumentAction::Create). The shared Rust validator explicitly rejects action alternatives without a named action, but wire requests cannot express that missing state. A client that sends clauses without naming an action silently receives a create-only subscription instead of an invalid-argument response. Make the field optional and validate its presence, or introduce a rejected unspecified enum value before publishing the new API.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/block_source.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/block_source.rs:202: block_results failures discard the cause before retrying
  Every block_results error becomes ReadMiss::NotYet, including persistent transport or configuration failures. Retrying is appropriate when results lag behind block storage, but after the unreadable deadline the client receives only a generic message and the underlying cause has not been recorded here. The block-fetch half preserves its error, making failures in the results half substantially harder to diagnose. Record the error before mapping it to the retryable state.

In `packages/rs-sdk/tests/fetch/state_transition_subscriptions.rs`:
- [SUGGESTION] packages/rs-sdk/tests/fetch/state_transition_subscriptions.rs:9-34: SDK event verification has no offline test coverage
  The new network test checks only the opening checkpoint. The unit test in subscriptions.rs checks retry-status classification, not accept(), so neither exercises hash rejection, decoding rejection, local filter matching, contract following, or removal of the internal filter's matches. These checks implement the advertised SDK verification behavior and can regress without either test failing. Add offline unit tests in the subscriptions module that construct response messages and verify rejected events, internal-update suppression, caller-visible match indexes, and unsupported protocol versions. The PR states that the network test was not run, and the supplied CI snapshot shows Rust workspace tests skipped.

Comment thread packages/rs-sdk/src/platform/subscriptions.rs Outdated
Comment thread packages/rs-sdk/src/platform/subscriptions.rs Outdated
Comment thread packages/dash-platform-queries/src/subscriptions/resolved.rs
Comment thread packages/dapi-grpc/protos/platform/v0/platform.proto
Comment thread packages/dash-platform-queries/src/subscriptions/resolved.rs
Comment thread packages/rs-sdk/tests/fetch/state_transition_subscriptions.rs
PastaPastaPasta added a commit that referenced this pull request Oct 5, 2026
… intact

Addresses review of #5285:
- The SDK refuses a request whose filters leave no room for the filter it
  needs to follow data contract updates, instead of silently not following
  them while the node does.
- A transition of a protocol version or kind the SDK does not know ends the
  subscription with an error rather than being judged under guessed rules;
  wrong-hash or non-matching transitions are still skipped.
- `ActionMatch.action` is `optional` on the wire, so a missing action is
  rejected instead of decoding as CREATE.
- DAPI logs why block results were unreadable before retrying.
- Documents that replayed history is matched against the current contract
  version followed through updates.
- Offline tests for the SDK's checks: delivery with caller filter indexes,
  wrong hash, non-matching and internal-filter-only transitions skipped,
  unknown version or transition rejected, resume height.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@PastaPastaPasta

Copy link
Copy Markdown
Member Author

Thanks. All seven were checked against the code; six are addressed in f2c4819:

  • Contract-follow filter dropped when the request is full (subscriptions.rs): the SDK now refuses a request with document filters that leaves no room for its contract-update filter, instead of going on without it.
  • Unknown protocol version falls back to latest(): the SDK now fails closed with an error, as DAPI does. A transition it cannot decode is also an error now rather than a skip; both mean the SDK needs upgrading. Wrong-hash and non-matching transitions are still skipped as node misbehaviour.
  • Historical replay uses the current contract: valid. Recovering the version in force at an arbitrary start height is not generally possible (contract history exists only for contracts that keep it), so the behaviour is now documented on ResolvedFilters::resolve and in the endpoint doc, including the generatedFrom case.
  • ActionMatch.action cannot be missing: the field is now optional Action action; an unset action is rejected with INVALID_ARGUMENT (new wire test). Clients regenerated.
  • block_results failures lose their cause: logged at debug before the height is retried.
  • No offline coverage of the SDK's verification: added unit tests that drive accept() with constructed responses: delivery with caller filter indexes, wrong hash, non-matching, internal-filter-only matches, unknown protocol version and undecodable transition, and the resume height.

Not changed: deep clones in document matching (resolved.rs). The clone happens only after a transition has passed the contract and document-type check, so only transitions on the subscribed type pay it, for at most 12 alternatives. Removing it needs DriveDocumentQueryFilter to borrow its clauses instead of owning them, a wider refactor of the drive type I'd rather leave to a follow-up.


🤖 Posted autonomously by Codex on behalf of pasta.

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

Static verification at f2c4819 confirms one blocking SDK trust-boundary defect and five non-blocking issues involving resource limits, failure handling, and rebinding coverage. Six prior findings are fixed; the bounded clause-cloning optimization can be deferred to a separate follow-up. No builds or tests were run; the supplied exact-head CI snapshot shows Rust workspace tests skipped and PR Hygiene pending.

🔴 1 blocking | 🟡 5 suggestion(s)

Review provenance

Source: reviewer 1: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 2: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 3: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 6: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 7: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: normal by gpt-6.1-sol (effort low) — This is a large, intricate streaming RPC and SDK change, but it does not clearly alter consensus, funds movement, cryptography, peer-facing network deserialization, or storage migrations.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort high); agent phase2-reviewer
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-sdk/src/platform/subscriptions.rs`:
- [BLOCKING] packages/rs-sdk/src/platform/subscriptions.rs:285-288: Authenticate contract updates before replacing proved filter bindings
  Initial document bindings come from DataContract::fetch(), whose proof-verification path authenticates the contract. This call to follow() instead installs the contract embedded in an unproved network event: the shared implementation converts the payload and rebinds filters by contract ID without authenticating its commitment. A malicious node can fabricate an update for a watched contract and compute the accompanying hash itself. Changing a generatedFrom rule can then make later document events match a schema that was never proved, while the internal update remains hidden from the caller. This violates the documented proved-schema matching guarantee, independently of the explicitly unproved completeness and block metadata. Treat updates as refresh notifications and obtain an authenticated binding through a proved query or verified contract history before rebinding. If the required historical version cannot be authenticated, fail closed or expose explicitly unverified semantics. Add a regression showing that a fabricated update cannot replace the proved binding.
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:260-263: Fail explicitly on an unsupported response envelope
  An absent or unrecognized response envelope produces Ok(None), even though unsupported protocol versions and undecodable transitions now return errors because the SDK cannot interpret them safely. Unknown protobuf oneof fields are discarded during decoding, so a newer unsupported envelope also reaches this branch as None. If a stream keeps delivering such responses, next() keeps waiting without yielding events or advancing its resume height; those messages also reset the consecutive-failure counter. Return an explicit unsupported-envelope error rather than treating the response as a locally rejected transition, and add an offline accept() assertion for this case.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs:322-326: Reacquire replay capacity after releasing it for client sends
  Releasing replay_permit before sending correctly prevents slow consumers from holding catch-up capacity. However, acquisition occurs only in the outer loop. After a matching block releases the permit, the inner page loop reads all remaining transaction-bearing heights without reacquiring it; the periodic-checkpoint branch has the same effect. Subscribers replaying distinct ranges can therefore exceed max_replaying concurrent historical reads, because single-flight sharing only combines reads of the same height. Reacquire capacity before subsequent historical reads while continuing to release it before client sends. Extend the stalled-client coverage with distinct replay ranges and a blocking/counting source that verifies the concurrency bound after a matching block.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/mod.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/mod.rs:71-82: Reject oversized filter lists before fetching their contracts
  The 1–16 filter bound is enforced by ResolvedFilters::resolve() only after this loop fetches every distinct document contract. from_proto() does not validate the request-level count, and nonexistent contracts return Ok(None), so they do not stop the loop. The 128 KiB transport limit still permits thousands of minimal document filters, turning a request that must be rejected into thousands of sequential Drive lookups. Admission bounds concurrent requests, not backend work within each request. Validate the count before conversion and contract fetching, and test that an oversized request performs no Drive calls. Apply early cardinality validation in the SDK as well, where contract fetching currently precedes rejection.

In `packages/dash-platform-queries/src/subscriptions/tests.rs`:
- [SUGGESTION] packages/dash-platform-queries/src/subscriptions/tests.rs:493-505: The rebinding test cannot detect a no-op implementation
  This test resolves its filter against contract and then rebinds to a clone of that same contract. The unconstrained document-type filter already matches the transition before rebinding, so removing rebind_data_contract()'s implementation would leave the assertion passing. The test also bypasses follow(), including its conversion of the update payload. Use distinct contract versions and assert a matching outcome that changes only after rebinding, then exercise an actual DataContractUpdate through follow(). The new SDK accept() tests cover useful verification branches, but their internal-filter fixture uses an identity filter and does not test actual contract following.

In `packages/dash-platform-queries/src/subscriptions/resolved.rs`:
- [SUGGESTION] packages/dash-platform-queries/src/subscriptions/resolved.rs:188-196: Report contract conversion failures instead of silently ignoring them
  Deserializing a StateTransition does not establish that its embedded contract can be constructed. DataContract::try_from_platform_versioned() still performs fallible document-type construction and required-since checks when full_validation is false. This if-let discards any error, leaving the previous binding without a diagnostic, so the SDK can continue after receiving an update it could not interpret. Handle the conversion failure explicitly: return an error where following must fail closed, or at minimum report the error and contract ID if retaining the previous binding is intentional. Keep this distinct from the documented historical-replay case where a successfully converted older contract cannot serve a particular filter.
Out-of-scope follow-up suggestions (1)

These are valid observations, but they are outside this PR's scope and should be handled in separate issues or author/maintainer-requested PRs rather than blocking this review.

  • Direct connections bypass the per-IP subscription limit — The bypass is conditional on exposing DAPI outside the repository's gateway boundary. packages/dashmate/docker-compose.yml and templates/dynamic-compose.yml.dot expose port 3010 only within the container network rather than publishing it; the Envoy route reaches that internal service, and use_remote_address is enabled. client_ip() explicitly documents that a missing forwarded header denotes a non-gateway request. Without evidence of an untrusted direct-access path in the supported deployment, requesting an additional direct-client policy would be speculative hardening rather than a demonstrated defect in this PR.
    • Follow-up: Consider creating a separate issue or author/maintainer-requested PR for this.

Comment thread packages/rs-sdk/src/platform/subscriptions.rs Outdated
Comment thread packages/dash-platform-queries/src/subscriptions/tests.rs
Comment thread packages/dash-platform-queries/src/subscriptions/resolved.rs Outdated
Comment thread packages/rs-sdk/src/platform/subscriptions.rs
PastaPastaPasta added a commit that referenced this pull request Oct 5, 2026
… sends

Addresses the second review of #5285:
- Blocking: the SDK rebound its filters to the contract embedded in a data
  contract update the node streamed, which a node could fabricate. An update
  is now only a signal: the contract is read again with a proof before the
  next message is checked. `ResolvedFilters::follow` is documented as for a
  party that trusts its block source (DAPI) and fails when the update's
  contract cannot be built; DAPI ends the stream at that height rather than
  silently keeping the old binding.
- The SDK fails on a response envelope it does not understand instead of
  waiting on it.
- DAPI takes the replay permit again before each historical read after
  releasing it to send, so catch-up reads stay within `max_replaying`.
- DAPI and the SDK check the filter count before fetching any contract.
- Tests: rebinding through a real DataContractUpdate with a distinct
  version, updates of unbound contracts ignored, a fabricated update leaving
  the proved binding in place and queueing a proved refresh, unknown
  envelopes rejected, and catch-up reads bounded after sending matches.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

Verified the supplied findings against head 9085160: one blocking CPU-exhaustion issue and three correctness suggestions remain. Twelve prior findings are fixed; the bounded document-matching allocation optimization remains intentionally deferred. This was static verification only: the supplied CI snapshot shows Rust workspace tests skipped, with NPM release validation and PR Hygiene still pending.

🔴 1 blocking | 🟡 3 suggestion(s)

Review provenance

Source: reviewer 1: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 2: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 3: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 6: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); reviewer 7: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 8: gpt-6-astra (agent: phase2-reviewer, role: general); reviewer 9: gpt-6-astra (agent: phase2-reviewer, role: architecture-layering); reviewer 10: gpt-6-astra (agent: phase2-reviewer, role: platform-versioning); reviewer 11: gpt-6-astra (agent: phase2-reviewer, role: rust-quality); reviewer 12: gpt-6-astra (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: critical by gpt-6.1-sol (effort low) — This intricate streaming feature adds peer-facing network deserialization in packages/dash-platform-queries/src/subscriptions/proto.rs, including StateTransitionFilter::from_proto and document clause decoding, meeting the critical-surface criterion.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Single stage: Phase 1 and Phase 2 reviewed this head side by side, with no blocker gate between them (triage tier)
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — general (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — security-auditor (completed, effort xhigh); agent phase2-reviewer
  • Model comparison: every Phase-2 reviewer also ran on gpt-6-astra; the verifier weighed both sets without knowing which model wrote which
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-drive/src/query/filter.rs`:
- [BLOCKING] packages/rs-drive/src/query/filter.rs:827-829: Bound Base58 identifier operands before canonicalization
  An unauthenticated subscription request can contain a document DELETE filter for an existing contract/type with an original-document `$id EQUAL` operand consisting of 120,000 `z` characters. It fits the 128 KiB transport limit and passes the structural checks. This call then reaches serialize_value_for_key_v0(), which invokes Value::to_identifier_bytes() and bs58::decode(text).into_vec() before checking that the decoded identifier is 32 bytes. The pinned bs58 0.5.1 decoder processes each character against the growing decoded buffer, making this input quadratic CPU work. DAPI resolves filters synchronously on a Tokio worker before replay throttling applies; concurrent admitted requests can monopolize workers, and asynchronous timeouts cannot interrupt the decode. Bound identifier text before invoking the codec, or decode into a fixed 32-byte destination. Cover system identifiers, schema identifier properties, and list operands, while preserving genuinely long string-property operands.
- [SUGGESTION] packages/rs-drive/src/query/filter.rs:818-825: Preserve empty byte-array operands during canonicalization
  The index-key codec is not lossless for empty byte arrays: encode_value_for_tree_keys(Value::Bytes(vec![])) returns an empty buffer, and decode_value_for_tree_keys() interprets that buffer as Value::Null. A valid equality operand therefore becomes null and is rejected by validate_against_schema(), which requires Value::Bytes for byte-array equality. IN operands containing an empty byte array fail for the same reason. The existing withByteArrays fixture permits this value because byteArrayField has maxItems but no positive minItems. Preserve the byte-array type during normalization instead of decoding its empty index-key sentinel, and add EQUAL and IN regressions containing an empty byte array.

In `packages/rs-sdk/src/platform/subscriptions.rs`:
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:268-270: Retain pending contract refreshes across cancellation
  pop_first() removes the refresh obligation before awaiting DataContract::fetch(). If a caller drops a pending next() future through tokio::select! or a timeout, the subscription survives but the popped identifier is lost. The update notification has already been consumed, so a subsequent next() can read later events and checkpoints using the previous binding without retrying the required proved refresh. Keep the identifier queued until fetching and rebinding complete. Add an offline cancellation regression that suspends the fetch, drops the next() future, and verifies that the following call still refreshes the contract before reading another message.
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:306: Reject checkpoints with no representable resume height
  Checkpoint heights are unproved node input, and protobuf permits block_height = u64::MAX. Adding one panics with overflow checks enabled and wraps to zero without them, producing a resume cursor that DAPI explicitly rejects on reconnect. Unauthenticated checkpoint semantics do not require accepting metadata that panics the SDK or corrupts its cursor. Use checked_add() before updating resume_from, and add an offline accept() regression asserting that the maximum-height checkpoint returns an error without changing the previous cursor.

Comment thread packages/rs-drive/src/query/filter.rs Outdated
Comment thread packages/rs-drive/src/query/filter.rs Outdated
Comment thread packages/rs-sdk/src/platform/subscriptions.rs Outdated
Comment thread packages/rs-sdk/src/platform/subscriptions.rs Outdated
PastaPastaPasta added a commit that referenced this pull request Oct 5, 2026
…ding it

Addresses the third review of #5285:
- Blocking: base58 decoding is quadratic in its input, and filter operands
  come from unauthenticated subscription requests. Identifier operands given
  as text longer than any 32-byte identifier (44 characters) are refused
  before decoding, for system identifier fields and identifier properties;
  long string-property operands are unaffected.
- An empty byte-array operand no longer becomes null through the index-key
  round trip (its index form is empty), so `== []` and `IN [..., []]` work.
- SDK: a contract stays queued for its proved refresh until it is rebound,
  so a cancelled or failed `next()` retries it; a checkpoint at the maximum
  height is rejected instead of overflowing the resume height.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

Verified the combined reviewer findings against head be0b69d. Of the 17 prior findings, 15 are fixed, one remains blocking because derived identifier fields bypass the Base58 input bound, and one performance optimization remains intentionally deferred; three additional suggestions are confirmed. This was static verification only: the supplied CI snapshot shows Rust workspace tests skipped, with NPM release build validation and PR Hygiene pending.

🔴 1 blocking | 🟡 3 suggestion(s)

1 carried-forward finding(s) already raised on this PR; not re-posting as new inline comments.

Review provenance

Source: reviewer 1: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 2: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 3: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); reviewer 6: gpt-6-astra (agent: phase2-reviewer, role: general); reviewer 7: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 8: gpt-6-astra (agent: phase2-reviewer, role: architecture-layering); reviewer 9: gpt-6-astra (agent: phase2-reviewer, role: platform-versioning); reviewer 10: gpt-6-astra (agent: phase2-reviewer, role: rust-quality); reviewer 11: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 12: gpt-6-astra (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: critical by gpt-6.1-sol (effort low) — The change is intricate and adds peer-facing network deserialization in packages/dash-platform-queries/src/subscriptions/proto.rs::StateTransitionFilter::from_proto alongside a new public streaming RPC and decoding of remotely supplied block and transition data.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Single stage: Phase 1 and Phase 2 reviewed this head side by side, with no blocker gate between them (triage tier)
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — general (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — security-auditor (completed, effort xhigh); agent phase2-reviewer
  • Model comparison: every Phase-2 reviewer also ran on gpt-6-astra; the verifier weighed both sets without knowing which model wrote which
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-drive/src/query/filter.rs`:
- [SUGGESTION] packages/rs-drive/src/query/filter.rs:833-842: Preserve string values without round-tripping index sentinels
  The string index-key codec is not lossless: both empty text and a one-character U+0000 string encode as `[0]`, which `decode_value_for_tree_keys()` converts to empty text. Document serialization preserves their distinction with a length prefix, and the fixture's unrestricted `niceDocument.name` permits both. Because this normalizer runs on operands and transition values, an EQUAL filter for empty text also matches a transition containing only U+0000, and SDK re-matching repeats the same false positive. An IN list containing both strings becomes a duplicate list and is rejected; a STARTS_WITH operand containing only U+0000 becomes an empty prefix and is rejected. Preserve typed string values directly rather than using the index-key representation as a lossless canonicalizer. Add EQUAL, IN, and STARTS_WITH regressions distinguishing empty text from U+0000 without changing the shipped index encoding.
- [SUGGESTION] packages/rs-drive/src/query/filter.rs:871-882: Canonicalize packed U8 IN operands as candidate lists
  `WhereClause::in_values()`, schema validation, and `WhereOperator::eval()` explicitly support `Value::Bytes` as a packed IN list for U8 properties. The subscription wire decoder also preserves `BytesValue` as `Value::Bytes`. This match handles only `Value::Array` element-wise, so a valid U8 operand such as `IN Value::Bytes(vec![1, 2])` falls through to scalar conversion. The U8 property codec then calls `to_integer()` on the entire byte vector and rejects it, causing both SDK resolution and DAPI to refuse a clause that the existing validator and evaluator support. Handle packed byte IN lists before scalar conversion while retaining their U8-only field restriction and keeping byte-array scalar operands distinct. Cover native resolution, wire round-trip, and matching against U8 values.
- [BLOCKING] packages/rs-drive/src/query/filter.rs:817-832: Bound Base58 identifier operands before canonicalization
  (existing thread: https://github.com/dashpay/platform/pull/5285#discussion_r4186329460)
  The new bound protects the listed system identifiers and flattened identifier properties, but derived identifier index fields still bypass it. DPP supports a field such as `postId.$ownerId` with effective type `Identifier`; that field is absent from `flattened_properties()`, so `identifier` is false here. The fallback at lines 844–846 calls `serialize_value_for_key_v0()`, which resolves `derived_index_property_type(key)` and reaches `to_identifier_bytes()` → `bs58::decode(text).into_vec()`. A client can submit a CREATE filter with a new-document equality clause on that field and 120,000 `z` characters against a contract declaring the derived index. The request fits the 128 KiB transport cap, and `action_clauses()` canonicalizes it before schema validation rejects the unsupported field. The pinned decoder iterates over its growing output for each input character, so this still performs quadratic synchronous work on a DAPI worker before replay throttling. Reject derived-index-only fields before invoking their codec, since transition matching cannot read their values, or include their effective types in the pre-decode identifier bound. Add EQUAL and IN regressions using a derived identifier index.

In `packages/dapi-grpc/protos/platform/v0/platform.proto`:
- [SUGGESTION] packages/dapi-grpc/protos/platform/v0/platform.proto:166-167: Record the Rust server-trait break in release metadata
  The RPC is additive on the protobuf wire, but it changes the exported Rust server API. `dapi-grpc` generates servers without default stubs, so tonic adds a required `subscribeToStateTransitionsStream` associated type and `subscribe_to_state_transitions` method to `platform_server::Platform`. Existing downstream service or mock implementations no longer compile until they add these items; the new implementations in `QueryService` and `PlatformServiceImpl` show the migration. The PR's Breaking Changes section currently says there are none. Document this server-side Rust source break and apply the repository's breaking-change title/release convention so downstream implementers receive notice. This is package API release treatment, not a PlatformVersion activation change.

Comment thread packages/rs-drive/src/query/filter.rs Outdated
Comment thread packages/rs-drive/src/query/filter.rs
Comment thread packages/dapi-grpc/protos/platform/v0/platform.proto
PastaPastaPasta added a commit that referenced this pull request Oct 5, 2026
Addresses the fourth review of #5285:
- Blocking: a derived index property such as `postId.$ownerId` resolves to
  an identifier without being a schema property, so its operand skipped the
  base58 length bound and reached the quadratic decoder. A filter only reads
  schema properties plus `$id` and `$ownerId`; any other field is now refused
  before a codec runs, and the system identifier fields stay length-bounded.
- String operands and values are compared as text instead of through the
  index-key form, which maps `""` and `"\0"` alike; EQUAL, IN and STARTS_WITH
  now tell them apart.
- A `u8` field's IN candidates packed as bytes, which validation and
  evaluation support, are kept as they are instead of being rejected.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@PastaPastaPasta PastaPastaPasta changed the title feat(platform): subscribe to committed state transitions matching document, address, identity, token and contract filters feat(platform)!: subscribe to committed state transitions matching document, address, identity, token and contract filters Oct 5, 2026

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

At fbcd07a, verification confirms one blocking untrusted-input panic and one numeric range-normalization defect. Nineteen prior findings are fixed; the bounded document-clause cloning optimization remains intentionally deferred. Validation was static only: the supplied CI snapshot shows Rust workspace tests skipped, with NPM release validation and PR Hygiene pending.

🔴 1 blocking | 🟡 1 suggestion(s)

Review provenance

Source: reviewer 1: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 2: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 3: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 4: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 6: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 7: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); reviewer 8: gpt-6-astra (agent: phase2-reviewer, role: general); reviewer 9: gpt-6-astra (agent: phase2-reviewer, role: architecture-layering); reviewer 10: gpt-6-astra (agent: phase2-reviewer, role: platform-versioning); reviewer 11: gpt-6-astra (agent: phase2-reviewer, role: rust-quality); reviewer 12: gpt-6-astra (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: critical by gpt-6.1-sol (effort low) — The intricate streaming feature adds peer-facing network deserialization and validation in packages/dash-platform-queries/src/subscriptions/proto.rs, notably StateTransitionFilter::from_proto, alongside new decoding of remotely supplied blocks and transitions.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Single stage: Phase 1 and Phase 2 reviewed this head side by side, with no blocker gate between them (triage tier)
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — general (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6-astra — security-auditor (completed, effort xhigh); agent phase2-reviewer
  • Model comparison: every Phase-2 reviewer also ran on gpt-6-astra; the verifier weighed both sets without knowing which model wrote which
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-drive/src/query/filter.rs`:
- [BLOCKING] packages/rs-drive/src/query/filter.rs:865-866: Prevent malformed event values from panicking the SDK matcher
  A subscription node can fabricate a binary-decodable document CREATE for the watched contract/type, supply its correct SHA-256 hash, and encode an integer property as Value::Text("a".repeat(19) + "é"). SDK accept() deserializes the transition without schema or signature validation and passes its properties into this new canonicalization path. The integer codec calls Value::to_integer(), whose wrong-type error formats the value through Display. In packages/rs-platform-value/src/display.rs, that formatter calls text.split_at(20): byte 20 is inside the final UTF-8 character, so it panics before the codec returns an error. The .ok()? does not catch a panic, which escapes subscription.next(), terminating the consuming task or aborting a panic=abort client. Although the formatter predates this PR, the new matcher exposes it to fabricated network events; malformed subscription operands also reach it through canonical_operand(). Make value truncation UTF-8-safe and add regressions for both a rejected Unicode operand and a fabricated SDK event with a wrong-typed Unicode property.

In `packages/dash-platform-queries/src/subscriptions/resolved.rs`:
- [SUGGESTION] packages/dash-platform-queries/src/subscriptions/resolved.rs:556-561: Canonicalize numeric bounds before grouping range clauses
  These clauses are extracted and grouped before schema canonicalization at lines 612–614. Grouping paired bounds calls WhereClause::less_than(), which accepts only matching Value variants. On a U8 property, a CREATE filter containing stars >= Value::I64(1) and stars <= Value::U64(5) therefore fails with RangeClausesNotGroupable, although both values individually canonicalize to U8 and the equivalent BETWEEN operand succeeds. The wire decoder preserves int64_value and uint64_value as I64 and U64, so this affects wire requests as well as native SDK requests using different integer widths. Normalize raw operands against the document schema before extraction/grouping, retaining the identifier-length and schema-field guards rather than changing shipped grouping implementations. Add native and wire-resolution regressions for mixed-encoding bounds.

Comment thread packages/rs-drive/src/query/filter.rs
Comment thread packages/dash-platform-queries/src/subscriptions/resolved.rs Outdated

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

At head 7cfbb22, the 22 prior findings reconcile to 21 fixed and one intentionally deferred performance refactor. Four non-blocking issues remain: null values can satisfy numeric ranges, accepted native numeric operands can fail after wire encoding, and cache accounting plus eager response copies understate subscription memory use. Verification was static only; the supplied CI snapshot shows Rust workspace and Dashmate E2E tests skipped, with NPM release validation and PR Hygiene pending.

🟡 4 suggestion(s)

Review provenance

Source: reviewer 1: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 2: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 3: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 6: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 7: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); reviewer 8: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 9: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 10: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 11: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 12: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: critical by gpt-6.1-sol (effort low) — This large, intricate change introduces peer-facing network deserialization in StateTransitionFilter::from_proto in packages/dash-platform-queries/src/subscriptions/proto.rs and a new public streaming RPC with request validation and resource admission.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Single stage: Phase 1 and Phase 2 reviewed this head side by side, with no blocker gate between them (triage tier)
  • Fresh final gate: an independent Phase-2 review ran after iterative findings were reconciled
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-drive/src/query/filter.rs`:
- [SUGGESTION] packages/rs-drive/src/query/filter.rs:880-888: Reject null event values before evaluating numeric ranges
  The property codec returns an empty encoding for Value::Null and decodes that encoding back to Null. Because the empty-encoding guard excludes null inputs, this path returns Some(Null), which evaluate_clauses passes to WhereOperator::eval(). Its ordered comparisons use Value's derived PartialOrd, where Null sorts after every numeric variant. A CREATE filter such as stars > U64(3) on the U8 rating field therefore matches a transition carrying stars: Null. data_as_stored() does not validate this property, and SDK accept() can deliver a fabricated, binary-decodable event with a correct hash through this path. Prevent null from participating in numeric ordered comparisons, while preserving any intentionally supported null semantics elsewhere, and add matcher and offline SDK accept() regressions. This is newly exposed subscription-matching behavior, not consensus execution.

In `packages/dash-platform-queries/src/subscriptions/proto.rs`:
- [SUGGESTION] packages/dash-platform-queries/src/subscriptions/proto.rs:164-168: Preserve normalized operand semantics when encoding subscriptions
  The SDK resolves a clone of the filters but retains and transmits their original operands. Native resolution accepts an UPDATE_PRICE equality operand Value::U128(50), normalizing it to U64(50). This encoder instead reuses value_to_proto(), which emits Text("50") for U128; the subscription decoder preserves Text, and DAPI's strict to_integer::<u64>() rejects it. The document-clause encoder has the same discrepancy for a U8 field queried with U128(3), so this is not limited to unsupported 128-bit schema fields. Encode accepted operands from their normalized representation, or add subscription-specific integer encoding that preserves representable numeric values without converting them to text. Add native/wire round-trip resolution tests for both price and document clauses.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs:415-425: Defer transaction payload copies until responses are sent
  match_block() constructs all matching responses before the first send, and tx.bytes.as_ref().clone() deep-copies each serialized transaction from its shared Arc<Vec<u8>>. The 32-message channel bound does not bound these unsent copies: the send loop's iterator retains every remaining response while awaiting a slow consumer. Subscribers matching the same block therefore retain approximately subscriber count times matching payload bytes in private buffers, outside the shared block-cache budget, until delivery or the send deadline. Precompute lightweight transaction indexes and match descriptors while preserving contract-follow order, then materialize each payload when sending it. Add a stalled-reader regression with more matches than STREAM_BUFFER to check that unsent responses do not eagerly duplicate the entire block.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/block_source.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/block_source.rs:128-131: Include decoded allocations in the block-cache memory budget
  Each cached transaction retains both its serialized bytes and its decoded StateTransition, but this weight charges only one additional serialized length for the decoded representation. Compact Value arrays can substantially exceed that estimate: protocol 14 permits 1024-element boolean typed-array properties, whose transition encoding uses two bytes per boolean while the decoded Vec allocates a full Value slot per element. Several such properties fit within the 20 KiB transition limit, yielding decoded allocations many times larger than their charged weight. BlockSource caches successful transactions before applying subscriber filters, so a historical subscription with a non-matching identity filter can populate this amplified cache; replay concurrency and per-IP admission limits do not correct the retained-memory accounting. Include decoded collection capacities and nested allocations in the weight, or use a demonstrably conservative bound, and add a compact-array regression. The current nominal 64 MiB budget is not a conservative bound on retained block memory.

Comment thread packages/rs-drive/src/query/filter.rs Outdated
Comment thread packages/dash-platform-queries/src/subscriptions/proto.rs
PastaPastaPasta and others added 12 commits October 5, 2026 14:20
…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>
… intact

Addresses review of #5285:
- The SDK refuses a request whose filters leave no room for the filter it
  needs to follow data contract updates, instead of silently not following
  them while the node does.
- A transition of a protocol version or kind the SDK does not know ends the
  subscription with an error rather than being judged under guessed rules;
  wrong-hash or non-matching transitions are still skipped.
- `ActionMatch.action` is `optional` on the wire, so a missing action is
  rejected instead of decoding as CREATE.
- DAPI logs why block results were unreadable before retrying.
- Documents that replayed history is matched against the current contract
  version followed through updates.
- Offline tests for the SDK's checks: delivery with caller filter indexes,
  wrong hash, non-matching and internal-filter-only transitions skipped,
  unknown version or transition rejected, resume height.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
PastaPastaPasta and others added 7 commits October 5, 2026 14:21
… sends

Addresses the second review of #5285:
- Blocking: the SDK rebound its filters to the contract embedded in a data
  contract update the node streamed, which a node could fabricate. An update
  is now only a signal: the contract is read again with a proof before the
  next message is checked. `ResolvedFilters::follow` is documented as for a
  party that trusts its block source (DAPI) and fails when the update's
  contract cannot be built; DAPI ends the stream at that height rather than
  silently keeping the old binding.
- The SDK fails on a response envelope it does not understand instead of
  waiting on it.
- DAPI takes the replay permit again before each historical read after
  releasing it to send, so catch-up reads stay within `max_replaying`.
- DAPI and the SDK check the filter count before fetching any contract.
- Tests: rebinding through a real DataContractUpdate with a distinct
  version, updates of unbound contracts ignored, a fabricated update leaving
  the proved binding in place and queueing a proved refresh, unknown
  envelopes rejected, and catch-up reads bounded after sending matches.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ding it

Addresses the third review of #5285:
- Blocking: base58 decoding is quadratic in its input, and filter operands
  come from unauthenticated subscription requests. Identifier operands given
  as text longer than any 32-byte identifier (44 characters) are refused
  before decoding, for system identifier fields and identifier properties;
  long string-property operands are unaffected.
- An empty byte-array operand no longer becomes null through the index-key
  round trip (its index form is empty), so `== []` and `IN [..., []]` work.
- SDK: a contract stays queued for its proved refresh until it is rebound,
  so a cancelled or failed `next()` retries it; a checkpoint at the maximum
  height is rejected instead of overflowing the resume height.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Addresses the fourth review of #5285:
- Blocking: a derived index property such as `postId.$ownerId` resolves to
  an identifier without being a schema property, so its operand skipped the
  base58 length bound and reached the quadratic decoder. A filter only reads
  schema properties plus `$id` and `$ownerId`; any other field is now refused
  before a codec runs, and the system identifier fields stay length-bounded.
- String operands and values are compared as text instead of through the
  index-key form, which maps `""` and `"\0"` alike; EQUAL, IN and STARTS_WITH
  now tell them apart.
- A `u8` field's IN candidates packed as bytes, which validation and
  evaluation support, are kept as they are instead of being rejected.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…aying it

`Value`'s `Display` truncated text longer than 20 bytes with `split_at(20)`,
which panics when byte 20 falls inside a multi-byte character. Any error
that formats such a value panicked instead of returning, and the new
subscription matcher formats values from unauthenticated requests and
streamed transitions, so a crafted one could abort an SDK subscription. The
text is now cut at the last character boundary at or before byte 20; ASCII
output is unchanged, and only inputs that used to panic display differently.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Range bounds were grouped into one range before being canonicalized, and
grouping only compares like `Value` variants, so `stars >= I64(1)` with
`stars <= U64(5)` on a u8 field failed as not groupable although the
equivalent BETWEEN worked. `drive::query::filter::canonicalize_where_clause`
canonicalizes one clause (keeping the schema-field and identifier-length
guards), and the subscription matcher applies it before extraction.

Also adds regressions that a wrong-typed Unicode value, in a filter operand
or in a transition a node fabricates, is refused or skipped rather than
panicking the matcher.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…matching right

Addresses the non-blocking findings of the review of #5285 at 7cfbb22:
- An ordered comparison never matches a null value (`Value` orders null
  after every number, so `stars > 3` matched `stars: null`); equality and
  IN are unchanged.
- The SDK sends 128-bit integer operands that fit 64 bits as 64-bit
  integers, as the node resolves them, instead of the shared encoder's text.
- DAPI copies a matching transition's bytes only as it sends it, instead of
  materializing every response of a block before the first send.
- The block cache weighs each transaction by a conservative bound on its
  decoded form (two `Value` slots per encoded byte), with a 128 MiB budget.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The rebase onto v5.1-dev brought token shielded pools (#4760), and the
exhaustive classification of transitions stopped compiling until the new
kinds were placed:
- Batched token unshields name their recipient; shields, shielded transfers
  and the pool variants of mint, burn, claim and direct purchase send tokens
  to holders that are not public.
- The shielded token transitions that pay their fee from the pool
  (transfer, unshield, purchase) name their token, and the unshield its
  recipient; token filters match them by token and party.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@PastaPastaPasta
PastaPastaPasta force-pushed the feat/platform-subscriptions branch from 7cfbb22 to 0c67483 Compare October 5, 2026 19:27

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

At head 0c67483, 24 of the 26 prior findings are fixed, one remains partially unfixed, and the clause-cloning optimization remains intentionally deferred. Verification also confirms a blocking memory-retention issue in stalled subscriptions and a numeric-null regression coverage gap. This was a static review only; the supplied head-specific CI snapshot shows Rust workspace tests skipped and other checks pending, queued, or running.

🔴 1 blocking | 🟡 2 suggestion(s)

1 carried-forward finding(s) already raised on this PR; not re-posting as new inline comments.

Review provenance

Source: reviewer 1: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 2: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 3: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 6: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); reviewer 7: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: critical by gpt-6.1-sol (effort low) — The intricate new subscription endpoint changes peer-facing network deserialization in packages/dash-platform-queries/src/subscriptions/proto.rs through StateTransitionFilter::from_proto and its document-clause decoders, meeting the critical-surface bar.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Single stage: Phase 1 and Phase 2 reviewed this head side by side, with no blocker gate between them (triage tier)
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs`:
- [BLOCKING] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs:327-333: Account for decoded blocks pinned by stalled subscriptions
  The scan retains its Arc<CommittedBlock>, including every decoded transaction, throughout the send loop. Cache eviction cannot free a block still owned by a scan, and the replay permit is released before a backpressured send suspends it. Consequently, unauthenticated subscriptions replaying distinct transaction-heavy heights can each pin a different decoded block while additional subscriptions continue reading; the 128 MiB cache budget and eight-reader limit do not bound these allocations. Compact boolean arrays expand into Vec<Value> allocations many times their encoded size, so even the gateway's default 100 concurrent requests can retain several GiB on a sufficiently transaction-heavy history before the send deadlines expire. Charge in-flight block ownership against a shared memory budget, including blocks not retained by the cache, or retain only bounded serialized payload references after matching and release the decoded block before awaiting sends. Add a stalled-reader regression across distinct heavy heights and cache eviction.

In `packages/rs-drive/src/query/filter.rs`:
- [SUGGESTION] packages/rs-drive/src/query/filter.rs:3747-3753: Make the null regression exercise a numeric field
  This regression uses the string property niceDocument.name. canonical_property_value() already rejects Null for string properties because its String branch accepts only Text, so removing the new null guard from evaluate_clauses() would leave the test passing. The numeric path that required the guard instead encodes Null as an empty buffer and decodes it back to Null, allowing Value's derived ordering to satisfy a greater-than bound. Use the U8 rating.stars fixture with stars > U64(3) and a transition carrying stars: Null, so removing the guard makes the regression fail.

In `packages/dash-platform-queries/src/subscriptions/proto.rs`:
- [SUGGESTION] packages/dash-platform-queries/src/subscriptions/proto.rs:145: Preserve normalized operand semantics when encoding subscriptions
  (existing thread: https://github.com/dashpay/platform/pull/5285#discussion_r4187638043)
  The new narrowing fixes the reported U128 examples, but nonnegative I128 operands above i64::MAX still lose numeric semantics even when they fit u64. For example, an UPDATE_PRICE equality operand Value::I128(1_i128 << 63) resolves natively because price canonicalization calls to_integer::<u64>() and produces U64. Here, i64::try_from() fails and preserves I128, so the shared encoder emits decimal Text; the subscription decoder preserves that Text and DAPI rejects it as a non-integer. The SDK transmits the original filters, not its resolved copy, so a locally accepted subscription fails when opened. Fall back to unsigned-64 narrowing when signed narrowing fails, and extend the native/wire regression to i64::MAX + 1, including list operands.

Comment thread packages/rs-drive/src/query/filter.rs Outdated
Addresses the review of #5285 at 0c67483:
- Blocking: a scan kept its decoded block alive while waiting for a stalled
  client, outside the block cache's budget, so subscriptions replaying
  different heavy heights could each pin a decoded block. After matching, a
  scan now keeps only the serialized payloads (shared with the block) and
  drops the decoded block before any send waits. Regression: a stalled
  reader over many matching blocks leaves no block held beyond the cache.
- A signed 128-bit operand that fits only an unsigned 64-bit integer is
  sent as `u64` too.
- The null-ordering regression now uses a numeric field, whose codec
  round-trips null, so it fails without the guard.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

Static verification of the complete PR range at a181760 found four non-consensus suggestions and no remaining blockers. Of the 28 prior findings, 26 are fixed, operand serialization remains partly unresolved, and the bounded document-clause cloning optimization remains intentionally deferred. No builds or tests were run; the supplied head-specific CI snapshot reports Rust workspace tests skipped and several checks pending or queued.

🟡 4 suggestion(s)

1 carried-forward finding(s) already raised on this PR; not re-posting as new inline comments.

Review provenance

Source: reviewer 1: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 2: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 3: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 6: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 7: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); reviewer 8: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 9: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 10: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 11: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 12: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: critical by gpt-6.1-sol (effort low) — The change is intricate and adds peer-facing network deserialization in packages/dash-platform-queries/src/subscriptions/proto.rs (StateTransitionFilter::from_proto), alongside a new remotely accessible streaming endpoint with substantial input validation and resource-management logic.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Single stage: Phase 1 and Phase 2 reviewed this head side by side, with no blocker gate between them (triage tier)
  • Fresh final gate: an independent Phase-2 review ran after iterative findings were reconciled
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/dash-platform-queries/src/subscriptions/resolved.rs`:
- [SUGGESTION] packages/dash-platform-queries/src/subscriptions/resolved.rs:619-621: Bound price IN candidate lists during subscription resolution
  Price IN operands bypass the candidate bounds used by ordinary document clauses. This branch forwards the price clause to DriveDocumentQueryFilter, whose price validation checks only that array elements are integers; it does not call WhereClause::in_values(), which rejects empty lists, more than 100 candidates, and duplicates. A wire request with 20,000 repeated small uint64 candidates fits the 128 KiB request limit and resolves successfully. Every relevant document transition then deep-clones that list, and an UPDATE_PRICE with an absent price scans it through array.contains(). Replay admission limits concurrent readers, not this per-transition matching work. Validate price IN cardinality and duplicates at the subscription boundary, consistently with ordinary IN clauses, and add oversized native and wire request regressions. This does not require the deferred borrowing-matcher refactor.

In `packages/rs-sdk/src/platform/subscriptions.rs`:
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:202-206: Evict a failing stream's node before reconnecting
  Failures after a stream opens bypass the executor's address-failure handling. open() discards ExecutionResponse.address through into_inner(), and next() clears only the stream before retrying. AddressList deliberately uses a sticky active set, so with with_active_set_size(1), a node that sends the opening checkpoint and then OUT_OF_RANGE for pruned history is selected again on every reopen. The five-failure budget expires before its 5–7.5-minute rotation slot, even when a healthy standby retains the requested block. Preserve the selected address and evict it from rotation for node-specific resumable failures, without treating ordinary stream-lifetime expiry as node misbehavior. Add deterministic offline reconnect coverage with a failing active node and a healthy standby; verify the checkpoint-derived replacement height, replay of an interrupted block, backoff, and the failure limit for repeated opening-checkpoint-only failures.

In `packages/dash-platform-queries/src/subscriptions/participants.rs`:
- [SUGGESTION] packages/dash-platform-queries/src/subscriptions/participants.rs:204-209: Classify a document purchaser as a recipient
  document_recipient() handles only Transfer, so an identity filter with Role::Recipient never matches a document PURCHASE by that identity. The buyer is available without a state lookup: the batch transformer passes its owner_id as purchaser_id, and the purchase action assigns that identity as the document's new owner. The document purchase owner clause already uses the batch owner, and token_recipient() likewise classifies a DirectPurchase buyer as Recipient. Currently participants() records the document buyer only as Sender, while matches_batched() relies on this helper for recipient matching. Pass the batch owner into document_recipient(), return it for Purchase, and add a regression asserting that the buyer's recipient filter matches and reports the purchase's batch position.

In `packages/dash-platform-queries/src/subscriptions/proto.rs`:
- [SUGGESTION] packages/dash-platform-queries/src/subscriptions/proto.rs:170-175: Preserve normalized operand semantics when encoding subscriptions
  (existing thread: https://github.com/dashpay/platform/pull/5285#discussion_r4187638043)
  The I128 unsigned fallback fixes the reported price boundary, but the SDK still resolves a clone and serializes the original operands. For an F64 document property, EQUAL U128(1_u128 << 64) resolves natively because the property's codec converts it to Float. narrowed() cannot represent it as u64, so the shared encoder emits decimal Text; the subscription decoder preserves Text, which the F64 codec rejects. There is another instance of the same mismatch: a byte-array equality operand Array([U8(1), U8(2)]) resolves to Bytes natively, but serialization and decoding produce Array([U64(1), U64(2)]), which to_binary_bytes() rejects. Encode subscription operands from their schema-normalized representation using the existing Drive canonicalization, while leaving the separately documented getDocuments encoder behavior unchanged. Add native/wire resolution regressions for wide F64 operands, including IN, and byte-array operands supplied as U8 arrays.

Comment thread packages/dash-platform-queries/src/subscriptions/resolved.rs
Comment thread packages/rs-sdk/src/platform/subscriptions.rs Outdated
Comment thread packages/dash-platform-queries/src/subscriptions/participants.rs Outdated
…ion node

Addresses the non-blocking findings of the review of #5285 at a181760:
- The SDK sends clause operands in the form the schema stores them
  (`canonical_operands`), since the wire carries a few primitive types: a
  `u128` for a float field or `u8`s for a byte array no longer arrive as
  text and `u64`s the node refuses.
- When a stream fails because of its node (missing history, an unreadable
  or undecodable block, Tenderdash unreachable), the SDK bans that node
  before reconnecting, so a sticky active set does not reconnect to it until
  the failure budget runs out; ended lifetimes and slow-client errors do not
  ban.
- A document purchase makes the batch owner a recipient, as a token direct
  purchase already did.
- Price IN candidates are bounded and distinct, like any IN clause's.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

At head 8d4bba6, verification of the complete diff confirms two non-blocking subscription issues: resumed stream-opening refusals can repeatedly select the same unsuitable node, and heartbeat ticks discard replay semaphore queue positions. All 31 prior findings were revalidated: 30 are addressed, and the bounded clause-cloning optimization remains intentionally deferred; no blocking finding remains. This was static verification only: the supplied exact-head CI snapshot shows Rust workspace tests skipped and several other checks pending, queued, or running, so passing validation at this head is not established.

🟡 2 suggestion(s)

Review provenance

Source: reviewer 1: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 2: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 3: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 6: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 7: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); reviewer 8: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 9: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 10: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 11: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 12: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: normal by gpt-6.1-sol (effort low) — This is a large, intricate streaming API and SDK change, but its filtering and committed-block replay do not change consensus, funds movement, cryptography, peer-to-peer deserialization, or storage migrations.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Fresh final gate: an independent Phase-2 review ran after iterative findings were reconciled
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — general (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort high); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort high); agent phase2-reviewer
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-sdk/src/platform/subscriptions.rs`:
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:258-267: Fail over when a node refuses a resumed subscription before opening
  open() discards ExecutionError.address on failure, so subscription-specific failover applies only after a stream opens. A replacement node more than 1,000 blocks behind a valid checkpoint-derived resume cursor refuses the request with OUT_OF_RANGE in SubscriptionService::start(). The generic gRPC CanRetry implementation excludes OUT_OF_RANGE and FAILED_PRECONDITION, so the executor neither retries elsewhere nor bans that address. reopen() nevertheless treats these statuses as resumable and retries them; with an active-set size of one, it repeatedly selects the same unsuitable node until the five-failure budget expires, even when a healthy standby can serve the cursor. Retain the failed address and evict or ban it for node-specific opening failures during reconnection, while keeping genuinely invalid initial requests terminal. Add deterministic coverage in which a lagging replacement refuses the open and a healthy standby receives the unchanged resume cursor.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs:377-385: Preserve replay semaphore queue position across heartbeat ticks
  Each loop iteration constructs a new acquire_owned() future inside tokio::select!. When the checkpoint timer wins, that acquisition is cancelled; Tokio explicitly documents that cancellation loses its FIFO queue position. Historical reads can remain in flight beyond the 10-second heartbeat interval—the Tenderdash request timeout is 30 seconds—so an older waiter can repeatedly move behind newer subscriptions while emitting unchanged checkpoints. Under sustained arrivals, this can starve its replay despite permits periodically becoming available. Preserve the acquisition's queue position across heartbeat ticks, and ensure a backpressured checkpoint send cannot retain an acquired or reserved permit. Add a paused-time regression with staggered waiters that crosses a heartbeat interval before releasing capacity and asserts that the oldest healthy waiter still acquires first.

Comment thread packages/rs-sdk/src/platform/subscriptions.rs Outdated
Addresses the non-blocking findings of the review of #5285 at 8d4bba6:
- A subscription waiting for replay capacity created a new semaphore
  acquisition on every loop, so each heartbeat put it at the back of the
  queue and older waiters could be overtaken indefinitely. One acquisition is
  now kept for the whole wait, and its heartbeats never wait on the client
  (skipped when the buffer is full), so a permit granted meanwhile is not held
  behind a blocked send. Regression with staggered waiters crossing several
  heartbeats; it fails with the old per-loop acquisition.
- SDK: a node that refuses to reopen a subscription for its own reasons
  (lagging past the resume height, missing history) is banned before the
  next attempt, as one failing an open stream already was; a refusal of the
  first open stays the caller's to handle.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

Verified the complete PR range at the specified head: 32 prior findings are fixed, and the bounded clause-cloning optimization remains intentionally deferred. Two non-consensus suggestions remain: floating-point canonicalization corrupts negative zero, and replay-capacity waits bypass slow-reader cleanup. Validation was static only; the supplied head-specific CI snapshot shows Rust workspace tests skipped, with NPM release validation and PR Hygiene pending.

🟡 2 suggestion(s)

Review provenance

Source: reviewer 1: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 2: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 3: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 6: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); reviewer 7: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 8: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 9: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 10: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 11: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 12: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: critical by gpt-6.1-sol (effort low) — This intricate cross-package feature adds peer-facing network deserialization and validation in packages/dash-platform-queries/src/subscriptions/proto.rs::StateTransitionFilter::from_proto and a new public streaming RPC, meeting the critical-surface threshold.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Single stage: Phase 1 and Phase 2 reviewed this head side by side, with no blocker gate between them (triage tier)
  • Fresh final gate: an independent Phase-2 review ran after iterative findings were reconciled
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-drive/src/query/filter.rs`:
- [SUGGESTION] packages/rs-drive/src/query/filter.rs:886-895: Preserve negative zero when canonicalizing floating-point values
  The index-key round trip is not lossless for F64. DocumentPropertyType::encode_float() uses `val < 0.0`, so -0.0 takes the nonnegative branch and encodes as eight zero bytes; decode_float() complements those bytes and returns NaN. Consequently, a valid CREATE filter with EQUAL Float(0.0) on an unindexed number property misses a transition containing Float(-0.0), although the two finite values compare equal. Ordered comparisons involving that transition value also fail. Document serialization preserves the float's original bits, and data_as_stored() does not normalize it, so this defect is exposed by the new subscription canonicalizer rather than requiring a change to consensus storage. Handle F64 directly with Value::to_float(), preserving the existing null behavior, and add equality, range, and native/wire regressions with negative zero as both an operand and an event value. Leave the shipped index codec unchanged.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs:407-409: Preserve slow-reader cleanup while waiting for replay capacity
  A full response buffer is treated as success here, and replay_permit() resets last_checkpoint after each attempt without applying SEND_DEADLINE. This is reachable when a block's final checkpoint fills the buffer and the scan queues for its next historical read. If the client then stops draining but remains connected, the replay wait can continue beyond 60 seconds while retaining the subscription's Admission guard, per-IP slot, and buffered responses. FIFO acquisition preserves fairness but does not bound waiting time when many subscriptions are queued behind slow reads. Keep the pinned acquisition and nonblocking heartbeat, but track sustained full-buffer backpressure and enter the RESOURCE_EXHAUSTED cleanup path once SEND_DEADLINE expires. Add a paused-time regression with a full channel and withheld replay capacity that verifies termination without a permit grant.

Comment thread packages/rs-drive/src/query/filter.rs Outdated
…ts waiting to replay

Addresses the non-blocking findings of the review of #5285 at 5adc408:
- Float operands and values are compared as floats instead of through the
  index-key form, which decodes -0.0 as NaN; `== 0.0` now matches `-0.0` and
  ordered comparisons treat it as zero. The index codec is unchanged.
- A subscription waiting for replay capacity whose client keeps its buffer
  full is dropped with RESOURCE_EXHAUSTED after the same deadline as a
  blocked send, instead of holding its admission and per-IP slots until
  capacity frees.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

At head 575006e, the subscription changes preserve the consensus boundary and address 34 prior findings; the bounded document-clause cloning optimization remains intentionally deferred. Three non-blocking suggestions remain: stalled-stream recovery, terminal SDK error handling, and avoidable text allocations during matching. Validation was static only: the supplied CI snapshot skips Rust workspace and functional tests, with NPM release validation and PR Hygiene still pending.

🟡 3 suggestion(s)

Review provenance

Source: reviewer 1: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 2: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 3: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 6: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); reviewer 7: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 8: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 9: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 10: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 11: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 12: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: critical by gpt-6.1-sol (effort low) — The large, intricate change introduces peer-facing network deserialization in packages/dash-platform-queries/src/subscriptions/proto.rs, notably StateTransitionFilter::from_proto and document_filter_from_proto, alongside a new public streaming RPC.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Single stage: Phase 1 and Phase 2 reviewed this head side by side, with no blocker gate between them (triage tier)
  • Fresh final gate: an independent Phase-2 review ran after iterative findings were reconciled
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-sdk/src/platform/subscriptions.rs`:
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:195-196: Bound message waits so stalled streams can reconnect
  Recovery starts only after stream.message().await returns an error or EOF, but this await has no message-idle deadline. STREAMING_TIMEOUT bounds opening the RPC, not consuming its Streaming response; the executor and RequestSettings documentation explicitly confirm that distinction. HTTP/2 keepalive detects an unresponsive connection, not a stalled subscription handler whose connection still answers PINGs. Although the configured gateway bounds stream lifetime, a direct connection or another proxy need not enforce that bound, leaving next() pending indefinitely rather than resuming from its checkpoint. Add a portable message-idle watchdog with enough allowance for legitimate block reads, and route expiration through node eviction, backoff, and checkpoint-based reopening. Cover a stream that sends its opening checkpoint and then remains pending.
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:203-205: Latch terminal interpretation errors before returning from next()
  The ? returns an interpretation error after consuming the response but leaves self.stream open and records no terminal state. Calling next() again therefore reads the following response: an unsupported transition at height h can produce an error, followed by a checkpoint that advances resume_from to h + 1. This contradicts next()'s documented terminal-error behavior and differs from into_stream(), which stops after an error. Record a terminal state for interpretation failures and explicitly terminal transport or retry-budget outcomes, so later calls cannot consume messages or advance the cursor. Keep transient contract-refresh failures retryable. Add a regression for an interpretation error followed by a checkpoint and another next() call.

In `packages/rs-drive/src/query/filter.rs`:
- [SUGGESTION] packages/rs-drive/src/query/filter.rs:879-880: Borrow already-canonical text during event matching
  This branch clones the entire Value::Text, allocating another String. Owned output is useful when resolving operands, but evaluate_clauses() also reaches this helper for each string clause evaluated against transition data, compares the temporary value, and immediately drops it. Every relevant subscriber therefore copies the text independently, including the long unindexed values supported by this PR; sharing decoded blocks does not remove these allocations. Add a borrowed fast path for already-canonical text during evaluation, or return a borrowed/owned representation while retaining owned normalization where needed. This change is local to the canonicalization/evaluation helpers and does not require the intentionally deferred refactor of DriveDocumentQueryFilter's owned clause sets.

Comment thread packages/rs-sdk/src/platform/subscriptions.rs Outdated
Comment thread packages/rs-sdk/src/platform/subscriptions.rs Outdated
Comment thread packages/rs-drive/src/query/filter.rs Outdated
…ription ended

Addresses the non-blocking findings of the review of #5285 at 575006e:
- A stream that delivers nothing for 90 seconds (the node checkpoints every
  10) counts as stalled, even if its connection still answers: the SDK bans
  the node and resumes from the last checkpoint elsewhere. The deadline uses
  the transport's portable sleep, so it also works in wasm.
- Once `next()` returns a terminal error, every later call fails without
  reading the stream or moving the resume height; a failed proved re-read of
  an updated contract stays retryable.
- Drive: matching borrows text and float values that are already in their
  stored form instead of copying them for every comparison.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

Verified the complete change range at head 181b582: subscription processing remains outside consensus execution, and no blocking issue was confirmed. All 38 prior findings were revalidated—37 are fixed and one remains intentionally deferred—but four in-scope suggestions remain concerning admission, cleanup, cancellation, and historical decoding. Validation was static only; the supplied head-specific CI snapshot shows Rust workspace tests skipped and PR Hygiene pending.

🟡 4 suggestion(s)

Review provenance

Source: reviewer 1: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 2: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 3: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 6: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 7: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); reviewer 8: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 9: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 10: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 11: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 12: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: critical by gpt-6.1-sol (effort low) — This intricate cross-package streaming feature introduces peer-facing network deserialization in packages/dash-platform-queries/src/subscriptions/proto.rs (StateTransitionFilter::from_proto), alongside substantial request validation and replay logic.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Single stage: Phase 1 and Phase 2 reviewed this head side by side, with no blocker gate between them (triage tier)
  • Fresh final gate: an independent Phase-2 review ran after iterative findings were reconciled
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/mod.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/mod.rs:153-163: Apply per-address admission to direct gRPC connections
  This helper trusts X-Forwarded-For without checking the connection peer and returns None when the header is absent. SubscriptionService::admit(None) then omits per-address accounting. When rs-dapi's configurable bind address exposes it directly, one client can occupy all 1,024 node slots instead of its allotted 16, either by omitting the header or by rotating forged header values. Extract the connection peer with request.remote_addr(), trust forwarded addresses only from the gateway's trusted peer boundary, and otherwise use the peer itself. SourceKey::of_request() in shielded_proof_failure_budget.rs already implements this distinction; preserve the subscription-specific IPv6 /64 bucketing. Add extraction tests for direct connections with both absent and forged forwarding headers.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs:273-277: Release slow-client admission before waiting to send the terminal status
  A blocked send first waits SEND_DEADLINE, then returns Stop::Fail(slow_client()). run() immediately attempts another send into the same full channel with a fresh 60-second deadline. The spawning task retains its Admission guard until run() returns, so an undrained client holds the node-wide and per-address slots for roughly 120 seconds, rather than the documented 60-second slow-reader interval. The slow-client failure from replay_permit() incurs the same additional cleanup wait. Separate scan termination and admission release from best-effort terminal-status delivery, and do not restart the slow-client cleanup deadline while reporting that failure. Add paused-time coverage through SubscriptionService::start() that checks capacity reclamation; the existing replay_permit() test cannot detect the guard's extended lifetime.

In `packages/rs-sdk/src/platform/subscriptions.rs`:
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:227-229: Preserve the message-idle deadline across cancellation of next()
  The 90-second timer is local to this step() invocation, while the subscription and its stream survive cancellation. A caller multiplexing subscription.next() with tokio::select! drops both pending futures whenever another branch wins; the following call starts a fresh timer without receiving any message. Consequently, a stream that sends its opening checkpoint and then stalls can remain stuck indefinitely if another branch wins periodically at intervals below 90 seconds. Store an absolute idle deadline or equivalent watchdog state in StateTransitionSubscription, resetting it only when a message arrives or a replacement stream opens. Add a regression that repeatedly cancels pending reads and verifies that the original silence deadline still initiates recovery; an injected pending read can exercise the watchdog without constructing a mocked tonic Streaming response.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/block_source.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/block_source.rs:241-242: Use DPP's protocol-aware decoder for subscription transitions
  The unscoped decoder inherits platform-value's default 256-container nesting limit, but protocol versions through 12 set max_document_value_depth to None. Consensus decoding uses StateTransition::deserialize_from_bytes_untrusted_in_version(), which installs the historical policy. Legacy contract schemas can carry deeply nested array metadata: the open v0 document meta-schema admits unused keys, and validate_max_depth_v0() enforces its ceiling for maps but not array-only nesting. An otherwise successful historical transition containing that metadata is therefore recorded as undecodable here, and match_block() ends the subscription before checking whether the transaction matches any filter. SDK accept() repeats the unscoped decode despite already resolving PlatformVersion. Resolve the block's version before decoding and use DPP's protocol-aware entry point in both adapters. Add coverage showing that pre-limit historical decoding succeeds while protocol-13-and-later decoding enforces its limit; no shipped codec change is needed.
Out-of-scope follow-up suggestions (1)

These are valid observations, but they are outside this PR's scope and should be handled in separate issues or author/maintainer-requested PRs rather than blocking this review.

  • Isolate unary request capacity from long-lived streams — The existing Envoy rs_dapi cluster shares a default max_requests budget of 100 between long-lived Core streams and unary Platform requests. Saturating that streaming budget can therefore deny ordinary queries and broadcasts. This is a concrete gateway availability concern, but the shared budget and existing streaming exposure predate this PR.
    • Follow-up: Track a separate gateway-capacity issue to reserve unary capacity or give streaming routes an independent upstream request budget.

Comment on lines +153 to +163
fn client_ip<T>(request: &Request<T>) -> Option<IpAddr> {
request
.metadata()
.get("x-forwarded-for")?
.to_str()
.ok()?
.rsplit(',')
.next()?
.trim()
.parse()
.ok()

@thepastaclaw thepastaclaw Oct 6, 2026 •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

✅ Resolved at a1a1cf8c; see the replies below.

🟡 Suggestion: Apply per-address admission to direct gRPC connections

This helper trusts X-Forwarded-For without checking the connection peer and returns None when the header is absent. SubscriptionService::admit(None) then omits per-address accounting. When rs-dapi's configurable bind address exposes it directly, one client can occupy all 1,024 node slots instead of its allotted 16, either by omitting the header or by rotating forged header values. Extract the connection peer with request.remote_addr(), trust forwarded addresses only from the gateway's trusted peer boundary, and otherwise use the peer itself. SourceKey::of_request() in shielded_proof_failure_budget.rs already implements this distinction; preserve the subscription-specific IPv6 /64 bucketing. Add extraction tests for direct connections with both absent and forged forwarding headers.

source: gpt-6.1-sol (phase2-reviewer: general)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in a1a1cf8. client_ip now starts from the connection peer. It trusts the last X-Forwarded-For entry only when the peer is loopback or private (a gateway), using the same is_internal / last_forwarded_address helpers as the shielded proof budget. A client that connects directly is counted by its own address whatever header it sends, so admit always gets an address when a peer is known. Tests: should_read_the_client_address_the_gateway_forwards and should_use_the_peer_of_a_direct_connection_with_or_without_a_header.


🤖 Posted autonomously by Codex on behalf of pasta.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Resolved (re-reviewed at a1a1cf8c): Your client_ip() now starts with remote_addr() and ignores forwarding metadata from public direct peers, using the existing gateway-boundary helpers. The absent-header and forged-header cases are covered; canonicalizing mapped addresses in quota bucketing is a separate remaining issue.

Comment on lines +273 to +277
async fn run(mut self) {
let outcome = self.scan().await;
if let Err(Stop::Fail(status)) = outcome {
debug!(code = ?status.code(), message = status.message(), "ending state transition subscription");
let _ = timeout(SEND_DEADLINE, self.sender.send(Err(status))).await;

@thepastaclaw thepastaclaw Oct 6, 2026 •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

✅ Resolved at a1a1cf8c; see the replies below.

🟡 Suggestion: Release slow-client admission before waiting to send the terminal status

A blocked send first waits SEND_DEADLINE, then returns Stop::Fail(slow_client()). run() immediately attempts another send into the same full channel with a fresh 60-second deadline. The spawning task retains its Admission guard until run() returns, so an undrained client holds the node-wide and per-address slots for roughly 120 seconds, rather than the documented 60-second slow-reader interval. The slow-client failure from replay_permit() incurs the same additional cleanup wait. Separate scan termination and admission release from best-effort terminal-status delivery, and do not restart the slow-client cleanup deadline while reporting that failure. Add paused-time coverage through SubscriptionService::start() that checks capacity reclamation; the existing replay_permit() test cannot detect the guard's extended lifetime.

source: gpt-6.1-sol (phase2-reviewer: general)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in a1a1cf8. A blocked send now ends the scan with Stop::SlowClient. run() gives its admission back before it reports anything. For a slow client it uses try_send for the status (the buffer is already full, so it is only delivered if there happens to be room). Other failures still get a bounded SEND_DEADLINE send, but the slot is already free. Regression: should_free_a_slow_clients_slot_once_its_deadline_passes (paused time, one slot: a second subscriber is admitted right after the first deadline). I checked that it fails against the previous ordering.


🤖 Posted autonomously by Codex on behalf of pasta.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Resolved (re-reviewed at a1a1cf8c): Your run() drops Admission before terminal delivery, and Stop::SlowClient uses try_send rather than another send deadline. The start()-level paused-time regression checks that the sole subscription slot is available just after the first slow-reader deadline.

Comment on lines +227 to +229
let idle = sleep(MESSAGE_IDLE_DEADLINE);
futures::pin_mut!(idle);
match futures::future::select(message, idle).await {

@thepastaclaw thepastaclaw Oct 6, 2026 •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

✅ Resolved at a1a1cf8c; see the replies below.

🟡 Suggestion: Preserve the message-idle deadline across cancellation of next()

The 90-second timer is local to this step() invocation, while the subscription and its stream survive cancellation. A caller multiplexing subscription.next() with tokio::select! drops both pending futures whenever another branch wins; the following call starts a fresh timer without receiving any message. Consequently, a stream that sends its opening checkpoint and then stalls can remain stuck indefinitely if another branch wins periodically at intervals below 90 seconds. Store an absolute idle deadline or equivalent watchdog state in StateTransitionSubscription, resetting it only when a message arrives or a replacement stream opens. Add a regression that repeatedly cancels pending reads and verifies that the original silence deadline still initiates recovery; an injected pending read can exercise the watchdog without constructing a mocked tonic Streaming response.

source: gpt-6.1-sol (phase2-reviewer: rust-quality)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in a1a1cf8. The silence timer is now a field of StateTransitionSubscription (idle). The first read after a message or after a new stream opens arms it. It is cleared when a message arrives and when a replacement stream opens, so cancelling next() no longer restarts it. The select is factored into read_or_idle. Regression: should_keep_the_idle_deadline_across_cancelled_reads (paused time) cancels a pending read twice at 40s each, and the third read sees the original 90s deadline fire.


🤖 Posted autonomously by Codex on behalf of pasta.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Resolved (re-reviewed at a1a1cf8c): Your idle field retains the same timer across cancelled read_or_idle calls and is reset after a message or replacement stream. The paused-time regression explicitly cancels two 40-second reads before observing the original 90-second deadline.

Comment on lines +241 to +242
state_transition: StateTransition::deserialize_from_bytes_untrusted(&bytes)
.map_err(|e| e.to_string()),

@thepastaclaw thepastaclaw Oct 6, 2026 •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

✅ Resolved at a1a1cf8c; see the replies below.

🟡 Suggestion: Use DPP's protocol-aware decoder for subscription transitions

The unscoped decoder inherits platform-value's default 256-container nesting limit, but protocol versions through 12 set max_document_value_depth to None. Consensus decoding uses StateTransition::deserialize_from_bytes_untrusted_in_version(), which installs the historical policy. Legacy contract schemas can carry deeply nested array metadata: the open v0 document meta-schema admits unused keys, and validate_max_depth_v0() enforces its ceiling for maps but not array-only nesting. An otherwise successful historical transition containing that metadata is therefore recorded as undecodable here, and match_block() ends the subscription before checking whether the transaction matches any filter. SDK accept() repeats the unscoped decode despite already resolving PlatformVersion. Resolve the block's version before decoding and use DPP's protocol-aware entry point in both adapters. Add coverage showing that pre-limit historical decoding succeeds while protocol-13-and-later decoding enforces its limit; no shipped codec change is needed.

source: gpt-6.1-sol (phase2-reviewer: architecture-layering)

@PastaPastaPasta PastaPastaPasta Oct 6, 2026 •

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in a1a1cf8. DAPI (BlockSource) and the SDK (accept) both decode with StateTransition::deserialize_from_bytes_untrusted_in_version under the version from the block header, as consensus does. On DAPI, a block whose protocol version this node does not know is not decoded under some other version. Its transactions count as undecodable: the scan ends with failed_precondition and resume from <height> on an upgraded node, and the SDK moves to another node. I did not add a separate subscription test for the pre-v13 and v13+ depth behaviour. The rule lives in DPP and is covered there (document_batch_rejects_excessive_value_depth_during_decode, document_batch_value_depth_limits_align_between_decode_and_validation). The subscription code only chooses which version to pass, and it is now the same choice consensus makes.


🤖 Posted autonomously by Codex on behalf of pasta.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Resolved (re-reviewed at a1a1cf8c): Your DAPI and SDK adapters now pass the resolved block protocol version to deserialize_from_bytes_untrusted_in_version, and unknown versions are not decoded under substitute tables. The cited DPP tests cover the decoder's depth policy; the separately reported SDK issue concerns unsafe lifetimes of fabricated legacy values, not this corrected version selection.

…slots on time

Addresses the non-blocking findings of the review of #5285 at 181b582:
- The per-address subscription limit reads the client from the connection
  peer, trusting `X-Forwarded-For` only from a gateway (loopback or private)
  peer, as the shielded proof budget does; a client connecting directly is
  counted by its own address whatever header it sends.
- A scan gives its admission back before reporting why it ended, and a slow
  client's status is sent only if its buffer has room, so the slot is freed
  at the 60-second deadline instead of after a second one.
- DAPI and the SDK decode a transition under its block's protocol version
  with DPP's protocol-aware decoder, as consensus did (older versions put no
  depth limit on document values).
- SDK: the stream-silence deadline lives in the subscription, so cancelling
  `next()` (e.g. in a `select!`) no longer restarts it.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

@thepastaclaw thepastaclaw left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review — Final validation — Phase 1 + Phase 2

Revalidated the exact head and all 42 prior findings: 41 are fixed, and the bounded document-clause borrowing refactor remains intentionally deferred. Four new in-scope issues remain, including a blocking SDK denial-of-service path through node-selected historical decoding. This was static verification only; the supplied CI snapshot skips Rust workspace tests and leaves NPM release validation and PR Hygiene pending.

🔴 1 blocking | 🟡 3 suggestion(s)

Review provenance

Source: reviewer 1: muse-spark-1.3-contributor (agent: phase1-reviewer, role: general); reviewer 2: gpt-6.1-sol (agent: phase2-reviewer, role: general); reviewer 3: gpt-6.1-sol (agent: phase2-reviewer, role: architecture-layering); reviewer 4: gpt-6.1-sol (agent: phase2-reviewer, role: platform-versioning); reviewer 5: gpt-6.1-sol (agent: phase2-reviewer, role: rust-quality); reviewer 6: muse-spark-1.3-contributor (agent: phase1-reviewer, role: rust-quality); reviewer 7: gpt-6.1-sol (agent: phase2-reviewer, role: security-auditor); final verifier: gpt-6.1-sol (agent: sol-verifier, role: final-verifier)

  • Triage: critical by gpt-6.1-sol (effort low) — The change is large and intricate and adds peer-facing network deserialization and validation in packages/dash-platform-queries/src/subscriptions/proto.rs::StateTransitionFilter::from_proto and the new DAPI subscription request handler.
  • Phase 1 reviewers: muse-spark-1.3-contributor — general (completed, effort xhigh); agent phase1-reviewer, muse-spark-1.3-contributor — rust-quality (completed, effort xhigh); agent phase1-reviewer
  • Phase 1 model: muse-spark-1.3-contributor — not quota-gated; passed over gemini-3.8-flash-high (antigravity below 15% reserve: weekly 15% left, 5h 100% left), glm-5.3-flash (not used above high effort; tier asks max)
  • Single stage: Phase 1 and Phase 2 reviewed this head side by side, with no blocker gate between them (triage tier)
  • Fresh verifier: gpt-6.1-sol — final-verifier; agent sol-verifier
  • Phase 2 reviewers: gpt-6.1-sol — general (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — architecture-layering (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — platform-versioning (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — rust-quality (completed, effort xhigh); agent phase2-reviewer, gpt-6.1-sol — security-auditor (completed, effort xhigh); agent phase2-reviewer
🤖 Prompt for all review comments with AI agents
These findings are from an automated code review. Verify each finding against the current code and only fix it if needed.

In `packages/rs-sdk/src/platform/subscriptions.rs`:
- [BLOCKING] packages/rs-sdk/src/platform/subscriptions.rs:439-446: Make node-selected historical decoding safe for recursive value lifetimes
  The node controls this unproved protocol_version. Protocol 12 uses SYSTEM_LIMITS_V2, whose max_document_value_depth is None, so this call disables the defensive Value nesting ceiling. A malicious node can supply a correctly hashed, fabricated document CREATE containing many thousands of singleton arrays within StateTransition's 100,000-byte decoding budget. Value decoding is iterative and refunds each container's provisional allocation charge before descending, so the byte budget does not prevent this nesting. Subsequent handling is not stack-safe: Value has recursive clone/formatting paths and recursive drop glue. In particular, an unrelated identity filter can reject the batch without inspecting its document data, after which accept() recursively drops the entire decoded tree and can exhaust the client's stack. DPP's serialization tests explicitly document recursive consumers and destruction, use enlarged stacks, and forget an intentionally deep fixture. The hash and proved contract bindings do not authenticate this fabricated transition or its advertised protocol version. Preserve historical interpretation, but make decoding cleanup and subsequent value handling stack-bounded, or enforce a separate documented client resource ceiling during decoding before constructing an unsafe tree. Do not change shipped consensus decoding policy. Add crafted historical-response coverage for both matching and non-matching disposal paths; a post-decode depth check alone cannot safely dispose of the rejected tree.
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:357-359: Retain reconnect backoff across cancelled next() calls
  Unlike the repaired idle watchdog, this retry sleep belongs to the temporary reopen() future. After step() consumes a resumable failure, it clears self.stream and increments consecutive_failures. If the caller multiplexes next() with another branch that repeatedly wins before the one-second first retry delay, cancelling next() drops the sleep. Each subsequent call starts the full delay again, so no replacement RPC is attempted even after arbitrarily many calls and ample elapsed time. Keep the pending retry timer or absolute deadline in StateTransitionSubscription, consuming it when the delay completes and resetting it only for a new retry attempt. Add paused-time coverage that repeatedly cancels reads during backoff and verifies that a replacement open eventually starts; the timer/open boundary can be injected without constructing tonic Streaming.
- [SUGGESTION] packages/rs-sdk/src/platform/subscriptions.rs:244-247: Cancel the detached WASM watchdog when its read completes
  On wasm32, transport::sleep() calls into_send_sleep(), whose into_send() spawns a local task owning the gloo timer and returns only a oneshot receiver. Once a pending message read has polled that timer, receiving the message clears self.idle and drops the receiver, but does not cancel the spawned task or its 90-second JavaScript timer. Repeated pending-then-completed reads therefore accumulate detached watchdog tasks for the following 90 seconds, making retained work scale with message rate rather than remaining one watchdog per subscription. This is newly exposed by using the existing sleep adapter on the subscription hot path; the Rust SDK's WASM path is supported even though JavaScript bindings are deferred. Use a cancellation-aware portable timer or a reusable watchdog, and verify that replacing or dropping it reclaims the WASM timer and task.

In `packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs`:
- [SUGGESTION] packages/rs-dapi/src/services/platform_service/subscribe_to_state_transitions/stream.rs:241: Canonicalize IPv4-mapped peers before applying IPv6 quotas
  The new direct-peer extraction can pass IPv4-mapped IPv6 addresses from remote_addr() on a dual-stack IPv6 listener. Both ::ffff:198.51.100.1 and ::ffff:203.0.113.7 enter the IPv6 branch and become :: when their lower four segments are cleared. All such direct IPv4 clients consequently share one 16-subscription bucket, allowing one client to exhaust admission for the others despite available node-wide capacity. Canonicalize the address before both family selection and returning the IPv4 key, as SourceKey::from already does. Add regressions showing that distinct mapped IPv4 peers remain distinct and that a mapped address shares its quota with its native IPv4 representation.

Comment on lines +439 to +446
let platform_version = PlatformVersion::get(matched.protocol_version)
.map_err(|e| Error::Protocol(e.into()))?;
// Decoded under the rules of the block's version, as consensus decoded it.
let state_transition =
StateTransition::deserialize_from_bytes_untrusted_in_version(
&matched.state_transition,
platform_version,
)?;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔴 Blocking: Make node-selected historical decoding safe for recursive value lifetimes

The node controls this unproved protocol_version. Protocol 12 uses SYSTEM_LIMITS_V2, whose max_document_value_depth is None, so this call disables the defensive Value nesting ceiling. A malicious node can supply a correctly hashed, fabricated document CREATE containing many thousands of singleton arrays within StateTransition's 100,000-byte decoding budget. Value decoding is iterative and refunds each container's provisional allocation charge before descending, so the byte budget does not prevent this nesting. Subsequent handling is not stack-safe: Value has recursive clone/formatting paths and recursive drop glue. In particular, an unrelated identity filter can reject the batch without inspecting its document data, after which accept() recursively drops the entire decoded tree and can exhaust the client's stack. DPP's serialization tests explicitly document recursive consumers and destruction, use enlarged stacks, and forget an intentionally deep fixture. The hash and proved contract bindings do not authenticate this fabricated transition or its advertised protocol version. Preserve historical interpretation, but make decoding cleanup and subsequent value handling stack-bounded, or enforce a separate documented client resource ceiling during decoding before constructing an unsafe tree. Do not change shipped consensus decoding policy. Add crafted historical-response coverage for both matching and non-matching disposal paths; a post-decode depth check alone cannot safely dispose of the rejected tree.

source: gpt-6.1-sol (phase2-reviewer: general, architecture-layering, platform-versioning, rust-quality, security-auditor)

/// IPv4 addresses count individually; IPv6 addresses by /64, the smallest block a client
/// usually controls.
fn per_ip_key(ip: IpAddr) -> IpAddr {
match ip {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Suggestion: Canonicalize IPv4-mapped peers before applying IPv6 quotas

The new direct-peer extraction can pass IPv4-mapped IPv6 addresses from remote_addr() on a dual-stack IPv6 listener. Both ::ffff:198.51.100.1 and ::ffff:203.0.113.7 enter the IPv6 branch and become :: when their lower four segments are cleared. All such direct IPv4 clients consequently share one 16-subscription bucket, allowing one client to exhaust admission for the others despite available node-wide capacity. Canonicalize the address before both family selection and returning the IPv4 key, as SourceKey::from already does. Add regressions showing that distinct mapped IPv4 peers remain distinct and that a mapped address shares its quota with its native IPv4 representation.

Suggested change
match ip {
let ip = ip.to_canonical();
match ip {

source: gpt-6.1-sol (phase2-reviewer: security-auditor)

Comment on lines +357 to +359
loop {
if self.consecutive_failures > 0 {
sleep(FIRST_RETRY_DELAY * 2u32.pow(self.consecutive_failures - 1)).await;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Suggestion: Retain reconnect backoff across cancelled next() calls

Unlike the repaired idle watchdog, this retry sleep belongs to the temporary reopen() future. After step() consumes a resumable failure, it clears self.stream and increments consecutive_failures. If the caller multiplexes next() with another branch that repeatedly wins before the one-second first retry delay, cancelling next() drops the sleep. Each subsequent call starts the full delay again, so no replacement RPC is attempted even after arbitrarily many calls and ample elapsed time. Keep the pending retry timer or absolute deadline in StateTransitionSubscription, consuming it when the delay completes and resetting it only for a new retry attempt. Add paused-time coverage that repeatedly cancels reads during backoff and verifies that a replacement open eventually starts; the timer/open boundary can be injected without constructing tonic Streaming.

source: gpt-6.1-sol (phase2-reviewer: general, architecture-layering, platform-versioning, rust-quality, security-auditor)

Comment on lines +244 to +247
let idle = self
.idle
.get_or_insert_with(|| Box::pin(sleep(MESSAGE_IDLE_DEADLINE)));
let message = read_or_idle(idle, stream.message())

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Suggestion: Cancel the detached WASM watchdog when its read completes

On wasm32, transport::sleep() calls into_send_sleep(), whose into_send() spawns a local task owning the gloo timer and returns only a oneshot receiver. Once a pending message read has polled that timer, receiving the message clears self.idle and drops the receiver, but does not cancel the spawned task or its 90-second JavaScript timer. Repeated pending-then-completed reads therefore accumulate detached watchdog tasks for the following 90 seconds, making retained work scale with message rate rather than remaining one watchdog per subscription. This is newly exposed by using the existing sleep adapter on the subscription hot path; the Rust SDK's WASM path is supported even though JavaScript bindings are deferred. Use a cancellation-aware portable timer or a reusable watchdog, and verify that replacing or dropping it reclaims the WASM timer and task.

source: gpt-6.1-sol (phase2-reviewer: general, architecture-layering, platform-versioning, rust-quality, security-auditor)

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants