Skip to content

feat: Add subscribe_resolved (latest entry per key + blob-complete gating) - #87

Open
cbenhagen wants to merge 5 commits into
n0-computer:mainfrom
cbenhagen:feat/subscribe-resolved
Open

cbenhagen wants to merge 5 commits into
n0-computer:mainfrom
cbenhagen:feat/subscribe-resolved

Conversation

@cbenhagen

Copy link
Copy Markdown
Contributor

Description

Adds a subscription that correlates live sync events with Query::single_latest_per_key() (latest entry per key, timestamp/hash tie-break), waits until that entry’s blob is locally complete (or the entry is empty), then emits the resolved key + entry (optional bytes). This closes the gap between per-author InsertLocal / InsertRemote and ContentReady { hash }, which carries no key—the store has no hash→key index, so apps otherwise re-query or track pending state themselves.

Problem / baseline

  • LiveEvent inserts are per-author; they are not necessarily the latest entry for that key.
  • Latest-for-key is defined by Query::single_latest_per_key() and existing store selection rules.
  • ContentReady { hash } does not name a key; mapping completion to keys is non-trivial without re-querying.

Behavior

When the latest entry for a key in scope may have changed, re-resolve via the store; emit only when the blob is Complete (or the entry is empty). Resolution always uses the query, not the raw insert event.

Implementation

  • Dirty keys on scoped inserts; flush with get_one / get_many under single_latest_per_key.
  • Blob status gating; pending key→hash; reconcile on ContentReady; drop stale pending if the winner changes.
  • Blob reads use a small internal ResolveBlobs trait (production delegates to iroh-blobs); if get_bytes fails after status is Complete, the subscription stream yields Err and the resolver keeps running.
  • Forward LiveEvent on a dedicated task into a bounded mpsc so slow resolution does not fill Engine::subscribe’s bounded channel.

API

  • Module subscribe_resolved: subscribe_resolved_with, ResolvedFetcher, ResolvedKeyValue, ResolvedSubscribeOpts.
  • Engine::subscribe_resolved, Doc::subscribe_resolved.
  • Store helpers / docs for single_latest_per_key with key prefix vs author filter.

Tests

  • Integration (tests/subscribe_resolved.rs, one case in tests/sync.rs): two authors / latest-per-key; sync across nodes + blob; burst inserts with slowed resolution; include_empty + delete; initial_snapshot; Query::author() after grouping; whole-replica scope (KeyFilter::Any / dirty_all flush).
  • Unit (subscribe_resolved module): mock ResolveBlobs for blob status error then recovery, get_bytes error on the stream, and “no Ok while status keeps failing.”

Caveats

  • Empty / tombstone-style winners: use .include_empty() on the scope Query, not ResolvedSubscribeOpts.
  • Blob not Complete: no item for that key until complete or winner changes; get_bytes failures after Complete surface as stream Err.
  • Inserts are not assumed to be the winner; always resolve through the store.

Out of scope

  • Engine-native “resolved key” events in LiveActor (possible follow-up).
  • New irpc for remote blob bytes (v1 assumes client-side blob API where needed).

Breaking Changes

None; additive API.

Notes & open questions

None.

Change checklist

  • Self-review.
  • Documentation updates following the style guide, if relevant.
  • Tests if relevant.
  • All breaking changes documented.

@n0bot n0bot Bot added this to iroh Mar 29, 2026
@github-project-automation github-project-automation Bot moved this to 🚑 Needs Triage in iroh Mar 29, 2026
@dignifiedquire dignifiedquire moved this from 🚑 Needs Triage to 👀 In review in iroh Apr 7, 2026
@cbenhagen
cbenhagen force-pushed the feat/subscribe-resolved branch from 0cdc11f to aa504c0 Compare May 25, 2026 07:01
cbenhagen added 5 commits May 25, 2026 09:03
…tion

Subscribe streams that follow Query::single_latest_per_key and emit when each
key's latest entry has a complete local blob (or is empty). Live events are
forwarded on a dedicated task so resolution work does not stall the bounded
Engine::subscribe channel.

- subscribe_resolved_with, Engine::subscribe_resolved, Doc::subscribe_resolved
- Store helpers for single-latest queries; tests: local two authors, sync + blob, burst inserts
On wasm32-unknown-unknown, n0_future::boxed::BoxFuture is local (!Send), so
Engine::subscribe streams are not Send. Use a target-gated bound so wasm builds
match n0_future::task::spawn (spawn_local) while native still requires Send.
- Runnable example: router wiring, prefix query, include_content,
  initial_snapshot, while-let stream loop, snapshot + new key + same-key update
- README: point to cargo run --example subscribe_resolved as main Doc API sample
Register cfg(wasm_browser) in unexpected_cfgs so custom cfgs from build.rs
match the linter. Use wasm_browser for SendUnlessWasmBrowser instead of
duplicating target triples, matching actor.rs. Replace tokio::time::sleep
with n0_future::time::sleep for resolution_delay and extend module docs for
Wasm/Tokio expectations.
Use explicit Engine::subscribe path and plain code for the private SendUnlessWasmBrowser bound (deny(broken_intra_doc_links)).
@cbenhagen
cbenhagen force-pushed the feat/subscribe-resolved branch from aa504c0 to 1229382 Compare May 25, 2026 07:03

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

Status: 👀 In review

Development

Successfully merging this pull request may close these issues.

2 participants