Skip to main content

miden_client/note_transport/
mod.rs

1pub mod errors;
2pub mod generated;
3#[cfg(feature = "tonic")]
4pub mod grpc;
5
6use alloc::boxed::Box;
7use alloc::collections::{BTreeMap, BTreeSet};
8use alloc::string::String;
9use alloc::sync::Arc;
10use alloc::vec::Vec;
11
12use miden_protocol::address::Address;
13use miden_protocol::block::BlockNumber;
14use miden_protocol::note::{
15    Note,
16    NoteDetails,
17    NoteDetailsCommitment,
18    NoteHeader,
19    NoteId,
20    NoteInclusionProof,
21    NoteTag,
22};
23use miden_protocol::utils::serde::Serializable;
24use miden_tx::auth::TransactionAuthenticator;
25use miden_tx::utils::serde::{
26    ByteReader,
27    ByteWriter,
28    Deserializable,
29    DeserializationError,
30    SliceReader,
31};
32
33pub use self::errors::NoteTransportError;
34use crate::note::{NoteFile, NoteSyncHint};
35use crate::store::{InputNoteRecord, NoteFilter, SettingScope};
36use crate::sync::NoteTagSource;
37use crate::{Client, ClientError};
38
39pub const NOTE_TRANSPORT_MAINNET_ENDPOINT: &str = "https://transport.mainnet.miden.io";
40pub const NOTE_TRANSPORT_TESTNET_ENDPOINT: &str = "https://transport.miden.io";
41pub const NOTE_TRANSPORT_DEVNET_ENDPOINT: &str = "https://transport.devnet.miden.io";
42pub const NOTE_TRANSPORT_CURSOR_STORE_SETTING: &str = "note_transport_cursor";
43
44/// Settings key for the note-transport backfill bookkeeping: a serialized `Vec<NoteTag>` of the
45/// `User`- and `Account`-source tags whose full history has already been fetched up to the global
46/// cursor. [`Client::sync_note_transport`] diffs the currently tracked tags against this set to
47/// find tags added after the cursor advanced, and backfills only those. Reusing the settings k/v
48/// avoids a Store-trait schema change while surviving process restarts.
49pub const NOTE_TRANSPORT_COVERED_TAGS_KEY: &str = "note_transport_covered_tags";
50
51/// Settings key for the durable relay outbox: a serialized `Vec<RelayOutboxEntry>` of private notes
52/// whose transport delivery has not yet succeeded. [`Client::send_private_note_with_proof`] appends
53/// (replacing any entry with the same note id) before relaying; [`Client::flush_relay_outbox`]
54/// drains entries that re-send successfully. Reusing the settings k/v avoids a Store-trait schema
55/// change while surviving process restarts.
56pub const NOTE_TRANSPORT_OUTBOX_KEY: &str = "note_transport_outbox";
57
58/// Client note transport methods.
59impl<AUTH> Client<AUTH> {
60    /// Check if note transport connection is configured
61    pub fn is_note_transport_enabled(&self) -> bool {
62        self.note_transport_api.is_some()
63    }
64
65    /// Returns the Note Transport client
66    ///
67    /// Errors if the note transport is not configured.
68    pub(crate) fn get_note_transport_api(
69        &self,
70    ) -> Result<Arc<dyn NoteTransportClient>, NoteTransportError> {
71        self.note_transport_api.clone().ok_or(NoteTransportError::Disabled)
72    }
73
74    /// Send a note through the note transport network together with its inclusion proof.
75    ///
76    /// The note will be end-to-end encrypted (unimplemented, currently plaintext) using the
77    /// provided recipient's `address` details. The recipient will be able to retrieve this note
78    /// through the note's [`NoteTag`].
79    ///
80    /// The transport carries the proof through [`NoteTransportClient::send_note_with_proof`]. The
81    /// network verifies `inclusion_proof` against its node before it stores the note and relays the
82    /// exact commitment block to the recipient. The proof exists once the transaction that created
83    /// the note is committed and the sender has synced past it; see
84    /// [`OutputNoteRecord::inclusion_proof`](crate::store::OutputNoteRecord::inclusion_proof).
85    ///
86    /// **Durability.** The note and its proof are persisted to the outbox before the transport
87    /// call. If the call fails or is interrupted, the entry stays in the outbox and is retried on
88    /// the next [`Client::flush_relay_outbox`] (which [`Client::sync_note_transport`] runs), so a
89    /// transient transport failure does not drop the note. The receiver dedupes by note id, so a
90    /// re-send after a partial success is harmless.
91    pub async fn send_private_note_with_proof(
92        &mut self,
93        note: Note,
94        address: &Address,
95        inclusion_proof: NoteInclusionProof,
96    ) -> Result<(), ClientError> {
97        let api = self.get_note_transport_api()?;
98
99        let note = TransportNote::from(note);
100        let note_id = note.header().id();
101        // The address is reserved for end-to-end encryption of the note details:
102        // address.key().encrypt(note.details().to_bytes()).
103        let _ = address;
104
105        // Persist the payload before the network call so a failed or interrupted send leaves a
106        // recoverable record rather than losing the only copy with the call frame. The proof
107        // travels with the entry so a retried send relays the same value.
108        let entry = RelayOutboxEntry { note, inclusion_proof };
109        let mut outbox = self.load_relay_outbox().await?;
110        // Replace any existing entry for this note id so the latest payload wins when a
111        // still-pending note is re-sent.
112        outbox.retain(|e| e.note.header().id() != note_id);
113        outbox.push(entry.clone());
114        self.save_relay_outbox(outbox).await?;
115
116        entry.relay(api.as_ref()).await?;
117
118        // Relay succeeded — drop the entry. A failed store write here is tolerable: the next flush
119        // re-sends and the receiver dedupes by note id, so a stale entry never causes loss.
120        let mut outbox = self.load_relay_outbox().await?;
121        outbox.retain(|e| e.note.header().id() != note_id);
122        self.save_relay_outbox(outbox).await?;
123
124        Ok(())
125    }
126
127    /// Re-attempt every relay payload in the durable outbox. Each entry is a private note whose
128    /// previous transport delivery failed. Successful re-sends are dropped; failures are kept for
129    /// the next call. Every entry is attempted independently, so one persistently-failing note does
130    /// not block the others.
131    ///
132    /// [`Client::sync_note_transport`] runs this automatically and ignores its error, so a relay
133    /// failure can't block a sync. Callers driving retries themselves can invoke it directly and
134    /// inspect the returned error.
135    pub async fn flush_relay_outbox(&self) -> Result<(), ClientError> {
136        let api = self.get_note_transport_api()?;
137
138        let entries = self.load_relay_outbox().await?;
139        if entries.is_empty() {
140            return Ok(());
141        }
142
143        // Attempt every entry independently so a single persistently-failing note can't block the
144        // rest. The outbox holds only the caller's own failed sends, so it stays small and this is
145        // not a meaningful burst.
146        let mut remaining = Vec::new();
147        let mut last_err: Option<NoteTransportError> = None;
148
149        for entry in entries {
150            match entry.relay(api.as_ref()).await {
151                Ok(()) => {},
152                Err(err) => {
153                    tracing::warn!(?err, "relay-outbox entry retry failed; will retry next sync");
154                    remaining.push(entry);
155                    last_err = Some(err);
156                },
157            }
158        }
159
160        self.save_relay_outbox(remaining).await?;
161
162        if let Some(err) = last_err {
163            return Err(err.into());
164        }
165        Ok(())
166    }
167
168    /// Load the durable relay outbox.
169    ///
170    /// Returns an empty `Vec` if the outbox key is absent. On deserialization failure (schema
171    /// mismatch or storage corruption) the entry is dropped and an empty `Vec` is returned —
172    /// leaving unreadable bytes in place would block every subsequent relay because each sync would
173    /// re-read them.
174    async fn load_relay_outbox(&self) -> Result<Vec<RelayOutboxEntry>, ClientError> {
175        let bytes = self
176            .store
177            .get_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_OUTBOX_KEY))
178            .await
179            .map_err(ClientError::StoreError)?;
180        let Some(bytes) = bytes else {
181            return Ok(Vec::new());
182        };
183        match Vec::<RelayOutboxEntry>::read_from_bytes(&bytes) {
184            Ok(entries) => Ok(entries),
185            Err(err) => {
186                tracing::warn!(?err, "dropping unreadable relay outbox; resetting to empty");
187                self.store
188                    .remove_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_OUTBOX_KEY))
189                    .await
190                    .map_err(ClientError::StoreError)?;
191                Ok(Vec::new())
192            },
193        }
194    }
195
196    /// Persist the relay outbox, removing the key entirely when empty so the settings table doesn't
197    /// accumulate empty-vec blobs.
198    async fn save_relay_outbox(&self, entries: Vec<RelayOutboxEntry>) -> Result<(), ClientError> {
199        let key = String::from(NOTE_TRANSPORT_OUTBOX_KEY);
200        if entries.is_empty() {
201            self.store
202                .remove_setting(SettingScope::Client, key)
203                .await
204                .map_err(ClientError::StoreError)?;
205            return Ok(());
206        }
207        let bytes = entries.to_bytes();
208        self.store
209            .set_setting(SettingScope::Client, key, bytes)
210            .await
211            .map_err(ClientError::StoreError)
212    }
213
214    /// The set of tracked tags eligible for history backfill.
215    ///
216    /// Only `User`- and `Account`-source tags qualify: those are the tags a consumer explicitly
217    /// started tracking (via [`Client::add_note_tag`], account import, or address creation) and may
218    /// therefore have historical private notes sitting below the global cursor. `Note`-source tags
219    /// are created by transport delivery and note import, so backfilling them would re-fetch tags
220    /// the fetch path itself just registered; `Subscription` tags are excluded for the same reason.
221    async fn backfill_candidate_tags(&self) -> Result<BTreeSet<NoteTag>, ClientError> {
222        let tags = self
223            .store
224            .get_note_tags()
225            .await?
226            .into_iter()
227            .filter(|record| {
228                matches!(record.source, NoteTagSource::User | NoteTagSource::Account(_))
229            })
230            .map(|record| record.tag)
231            .collect();
232        Ok(tags)
233    }
234
235    /// Load the set of tags whose history has already been fetched up to the global cursor.
236    ///
237    /// Returns an empty set when the key is absent (e.g. a store that predates the feature). On a
238    /// deserialization failure the entry is dropped and an empty set is returned: re-treating every
239    /// tracked tag as new only triggers a one-off backfill, which dedupes, whereas leaving
240    /// unreadable bytes in place would fail every subsequent sync.
241    async fn load_covered_tags(&self) -> Result<BTreeSet<NoteTag>, ClientError> {
242        let bytes = self
243            .store
244            .get_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_COVERED_TAGS_KEY))
245            .await
246            .map_err(ClientError::StoreError)?;
247        let Some(bytes) = bytes else {
248            return Ok(BTreeSet::new());
249        };
250        match BTreeSet::<NoteTag>::read_from_bytes(&bytes) {
251            Ok(tags) => Ok(tags),
252            Err(err) => {
253                tracing::warn!(?err, "dropping unreadable covered-tags set; resetting to empty");
254                self.store
255                    .remove_setting(
256                        SettingScope::Client,
257                        String::from(NOTE_TRANSPORT_COVERED_TAGS_KEY),
258                    )
259                    .await
260                    .map_err(ClientError::StoreError)?;
261                Ok(BTreeSet::new())
262            },
263        }
264    }
265
266    /// Persist the covered-tags set, removing the key entirely when empty so the settings table
267    /// doesn't accumulate empty-vec blobs.
268    async fn save_covered_tags(&self, tags: &BTreeSet<NoteTag>) -> Result<(), ClientError> {
269        let key = String::from(NOTE_TRANSPORT_COVERED_TAGS_KEY);
270        if tags.is_empty() {
271            self.store
272                .remove_setting(SettingScope::Client, key)
273                .await
274                .map_err(ClientError::StoreError)?;
275            return Ok(());
276        }
277        self.store
278            .set_setting(SettingScope::Client, key, tags.to_bytes())
279            .await
280            .map_err(ClientError::StoreError)
281    }
282}
283
284impl<AUTH> Client<AUTH>
285where
286    AUTH: TransactionAuthenticator + Sync + 'static,
287{
288    /// Per-sync cap on the number of newly tracked tags to backfill. Bounds the burst when many
289    /// tags are registered at once (e.g. restoring many accounts or addresses). Deferred tags stay
290    /// uncovered and are picked up on subsequent syncs.
291    pub const MAX_BACKFILL_TAGS_PER_SYNC: usize = 64;
292
293    /// Safety cap on the per-tag backfill drain. A well-behaved server eventually returns no
294    /// forward cursor progress, ending the loop; this bound only guards against a server that
295    /// advances the cursor indefinitely without ever returning an empty batch. It is far above any
296    /// honest per-tag backlog, so reaching it signals a server bug rather than real history.
297    const MAX_BACKFILL_ITERATIONS: usize = 1_000;
298
299    /// Fetch notes for tracked note tags.
300    ///
301    /// The client will query the configured note transport node for all tracked note tags. To list
302    /// tracked tags please use [`Client::get_note_tags`]. To add a new note tag please use
303    /// [`Client::add_note_tag`]. Only notes directed at your addresses will be stored and readable
304    /// given the use of end-to-end encryption (unimplemented). Fetched notes will be stored into
305    /// the client's store.
306    ///
307    /// An internal pagination mechanism is employed to reduce the number of downloaded notes: this
308    /// fetches only notes past the stored cursor. Historical notes for a newly tracked tag are
309    /// recovered automatically by [`Client::sync_note_transport`], which backfills each new tag.
310    pub async fn fetch_private_notes(&mut self) -> Result<(), ClientError> {
311        self.ensure_genesis_in_place().await?;
312
313        let note_tags: Vec<NoteTag> =
314            self.store.get_unique_note_tags().await?.into_iter().collect();
315        let cursor = self.store.get_note_transport_cursor().await?;
316
317        let mut id_by_commitment = BTreeMap::new();
318        let (note_files, new_cursor) =
319            self.fetch_transport_notes(cursor, &note_tags, &mut id_by_commitment).await?;
320
321        self.import_notes(&note_files).await?;
322        self.store.update_note_transport_cursor(new_cursor).await?;
323
324        Ok(())
325    }
326
327    /// Plans the backfill of historical private notes for tags added after the global cursor
328    /// advanced.
329    ///
330    /// The global transport cursor is shared across all tracked tags and only moves forward, so a
331    /// tag that starts being tracked late never sees its notes that already sit below the cursor.
332    /// This diffs the tracked `User`/`Account` tags (see [`Self::backfill_candidate_tags`]) against
333    /// the persisted covered set (see [`NOTE_TRANSPORT_COVERED_TAGS_KEY`]) and drains each newly
334    /// tracked tag from the start, fetching only that tag's own history rather than re-scanning
335    /// everything. Tags no longer tracked are dropped from the covered set so a later re-add
336    /// backfills again instead of resuming from a stale mark. Imports dedupe, so the overlap with
337    /// the steady-state stream is harmless.
338    ///
339    /// At most [`Self::MAX_BACKFILL_TAGS_PER_SYNC`] tags are backfilled per call; any remainder
340    /// stays uncovered and is picked up on the next sync.
341    ///
342    /// Returns the pruned covered set, whether pruning changed it, and the tags to backfill. Reads
343    /// only: persisting the covered set is left to the apply phase, which writes it after the
344    /// imported notes so a crash re-backfills instead of skipping a tag whose notes were never
345    /// written.
346    async fn plan_backfill(&self) -> Result<(BTreeSet<NoteTag>, bool, Vec<NoteTag>), ClientError> {
347        let candidates = self.backfill_candidate_tags().await?;
348        let loaded = self.load_covered_tags().await?;
349
350        // Drop tags no longer tracked. Keeping a removed tag marked covered would make a later
351        // re-add skip its backlog, silently missing notes that arrived while it was untracked.
352        let covered: BTreeSet<NoteTag> = loaded.intersection(&candidates).copied().collect();
353        let pruned = covered.len() != loaded.len();
354
355        let new_tags: Vec<NoteTag> = candidates
356            .difference(&covered)
357            .copied()
358            .take(Self::MAX_BACKFILL_TAGS_PER_SYNC)
359            .collect();
360
361        Ok((covered, pruned, new_tags))
362    }
363
364    /// Drain a single tag's full history from the transport, paging until the cursor stops
365    /// advancing. Uses a local cursor and never touches the global one, so it cannot regress
366    /// steady-state progress. Returns the note files from every fetched page, in page order and
367    /// none of them written.
368    async fn backfill_tag(
369        &self,
370        tag: NoteTag,
371        id_by_commitment: &mut BTreeMap<NoteDetailsCommitment, NoteId>,
372    ) -> Result<Vec<NoteFile>, ClientError> {
373        let mut note_files = Vec::new();
374        let mut cursor = NoteTransportCursor::init();
375        for _ in 0..Self::MAX_BACKFILL_ITERATIONS {
376            let (page_files, new_cursor) =
377                self.fetch_transport_notes(cursor, &[tag], id_by_commitment).await?;
378            note_files.extend(page_files);
379            // Terminate on any lack of forward progress. A well-behaved server returns `new_cursor
380            // == cursor` when there are no new notes for this tag (since `rcursor = max(cursor,
381            // max_seq_returned)`); using `<=` also handles implementations that return an `init()`
382            // cursor on empty batches (see the in-tree mock transport).
383
384            if new_cursor <= cursor {
385                return Ok(note_files);
386            }
387            cursor = new_cursor;
388        }
389
390        Err(ClientError::NoteTransportError(NoteTransportError::PaginationDidNotTerminate(
391            Self::MAX_BACKFILL_ITERATIONS,
392        )))
393    }
394
395    /// Screens the transport-delivered notes carrying a tag derived from a tracked account,
396    /// discarding those that no tracked account can consume. Notes carrying any other tag are kept
397    /// as delivered.
398    async fn screen_transport_notes(
399        &self,
400        notes: &mut Vec<(Note, Option<BlockNumber>)>,
401    ) -> Result<(), ClientError> {
402        let account_tags = self.tracked_account_tags().await?;
403
404        let notes_to_screen: Vec<Note> = notes
405            .iter()
406            .filter(|(note, _)| account_tags.contains(&note.metadata().tag()))
407            .map(|(note, _)| note.clone())
408            .collect();
409        let consumable = self.note_screener().get_batch_consumability(&notes_to_screen).await?;
410
411        // Discard the notes whose tag match the tracked accounts but are not consumable.
412        notes.retain(|(note, _)| {
413            !account_tags.contains(&note.metadata().tag()) || consumable.contains_key(&note.id())
414        });
415
416        Ok(())
417    }
418
419    /// Returns the tracked tags that were registered for an account, i.e. derived from its ID.
420    async fn tracked_account_tags(&self) -> Result<BTreeSet<NoteTag>, ClientError> {
421        let tags = self
422            .store
423            .get_note_tags()
424            .await?
425            .into_iter()
426            .filter(|record| matches!(record.source, NoteTagSource::Account(_)))
427            .map(|record| record.tag)
428            .collect();
429        Ok(tags)
430    }
431
432    /// Fetches and returns one batch of notes from the note transport layer for the provided tags
433    /// without applying any update to the store.
434    ///
435    /// The server paginates; this method issues one transport call and returns the note files
436    /// together with the new cursor. The returned cursor equals the input cursor when the batch was
437    /// empty (i.e. no new notes). Callers that want to drain a tag's full backlog should loop until
438    /// `new_cursor == cursor` (see [`Client::backfill_tag`]). Callers that do steady-state polling
439    /// (see [`Client::sync_state`] / [`Client::fetch_private_notes`]) should call this once per
440    /// tick with the stored cursor.
441    ///
442    /// Each downloaded note's id is recorded in `id_by_commitment` so the caller can resolve the
443    /// written records back to note ids once the final record set is known. Persistence of the
444    /// returned cursor is left to the caller so that drain loops can guard against regression of an
445    /// already-advanced stored cursor.
446    async fn fetch_transport_notes(
447        &self,
448        cursor: NoteTransportCursor,
449        tags: &[NoteTag],
450        id_by_commitment: &mut BTreeMap<NoteDetailsCommitment, NoteId>,
451    ) -> Result<(Vec<NoteFile>, NoteTransportCursor), ClientError> {
452        // Fallback lookback window, in blocks, used only for notes the transport delivered without
453        // block information. Scanning back from sync height handles the race where a note is
454        // committed on-chain just before the NTL delivers its data. Without it,
455        // check_expected_notes would scan from sync_height forward and miss the already-committed
456        // note. A transport-provided block is deterministic and always preferred.
457        const NOTE_LOOKBACK_BLOCKS: u32 = 20;
458
459        let mut notes = Vec::new();
460        // TODO: perhaps we should not need to map received IDs with details commitments, and
461        // instead we may allow `InputNoteRecord` to optionally keep NoteIds. Then within
462        // `import_note` we could match everything by ID and remove this map check
463        let (note_infos, rcursor) =
464            self.get_note_transport_api()?.fetch_notes(tags, cursor).await?;
465        for note_info in &note_infos {
466            // e2ee impl hint: for key in self.store.decryption_keys() try
467            // key.decrypt(details_bytes_encrypted)
468            //
469            // An invalid delivery fails the fetch and the cursor stays on this page.
470            let note = rejoin_note(&note_info.header, &note_info.details_bytes)?;
471            let tag = note.metadata().tag();
472            if !tags.contains(&tag) {
473                return Err(NoteTransportError::UnrequestedTag(tag).into());
474            }
475
476            // The header carries the attachment-aware (on-chain) note id; the rejoined note has
477            // empty attachments and would hash to a different id, so key off the header.
478            id_by_commitment.insert(note.details_commitment(), note_info.header.id());
479
480            notes.push((note, note_info.block_hint));
481        }
482
483        // Screen the transport-delivered notes to discard the ones that are not relevant to the
484        // accounts tracked by the client. Boxed to avoid a `clippy::large_futures` warning, since
485        // the sync future is already close to the size limit.
486        Box::pin(self.screen_transport_notes(&mut notes)).await?;
487
488        self.drop_notes_processed_locally(&mut notes).await?;
489
490        let sync_height = self.get_sync_height().await?;
491        let fallback_after_block_num =
492            BlockNumber::from(sync_height.as_u32().saturating_sub(NOTE_LOOKBACK_BLOCKS));
493
494        let mut note_files = Vec::with_capacity(notes.len());
495        for (note, block_hint) in notes {
496            let tag = note.metadata().tag();
497            // Prefer the transport-provided block, falling back to the lookback window when absent.
498            let after_block_num = block_hint.unwrap_or(fallback_after_block_num);
499            note_files.push(NoteFile::ExpectedNote {
500                details: note.into(),
501                sync_hint: NoteSyncHint::new(after_block_num, tag),
502            });
503        }
504
505        Ok((note_files, rcursor))
506    }
507
508    /// Fetches the notes the Note Transport Layer holds for the tracked tags.
509    ///
510    /// Runs the per-tag backfill and fetches a page of notes. This performs no node call and writes
511    /// nothing but the relay outbox, so it can run concurrently with the chain fetch. The caller
512    /// imports the returned files and then persists the cursor and the covered-tag set.
513    ///
514    /// Returns empty data when note transport is not configured.
515    pub(crate) async fn fetch_note_transport_updates(
516        &self,
517    ) -> Result<NoteTransportLayerUpdate, ClientError> {
518        let mut note_transport_update = NoteTransportLayerUpdate::default();
519        if !self.is_note_transport_enabled() {
520            return Ok(note_transport_update);
521        }
522
523        // Drain any private notes whose previous relay attempt failed. A flush error is logged, not
524        // propagated: a failing relay must not block the sync, and the entries stay durable for the
525        // next attempt. This is the one write this phase performs; it touches only the outbox
526        // setting, which is independent of everything the apply phase writes.
527        if let Err(err) = self.flush_relay_outbox().await {
528            tracing::warn!(?err, "relay outbox flush failed during sync; entries retained");
529        }
530
531        // Recover historical private notes for any tag added after the global cursor advanced. This
532        // drains each newly tracked tag from the start, fetching only that tag's own history.
533        let (mut covered, pruned, new_tags) = self.plan_backfill().await?;
534        let backfilled = !new_tags.is_empty();
535        for tag in new_tags {
536            note_transport_update
537                .note_files
538                .extend(self.backfill_tag(tag, &mut note_transport_update.id_by_commitment).await?);
539            covered.insert(tag);
540        }
541        if pruned || backfilled {
542            note_transport_update.covered_tags = Some(covered);
543        }
544
545        let cursor = self.store.get_note_transport_cursor().await?;
546        let note_tags: Vec<NoteTag> =
547            self.store.get_unique_note_tags().await?.into_iter().collect();
548        let (note_files, new_cursor) = self
549            .fetch_transport_notes(cursor, &note_tags, &mut note_transport_update.id_by_commitment)
550            .await?;
551        note_transport_update.note_files.extend(note_files);
552        note_transport_update.cursor = Some(new_cursor);
553
554        Ok(note_transport_update)
555    }
556
557    /// Writes everything [`Client::fetch_note_transport_updates`] returned, in three steps:
558    ///
559    /// 1. Imports the fetched notes, which resolves their on-chain state and stores the records.
560    /// 2. Saves the covered-tag set, when the backfill changed it.
561    /// 3. Advances the stored note transport cursor, when a page was fetched.
562    ///
563    /// The notes are written before the covered-tag set and the cursor, so a crash between them
564    /// re-fetches instead of skipping notes that were never written.
565    ///
566    /// Returns the ids of the imported notes and the details commitments of the records written.
567    pub(crate) async fn apply_note_transport_update(
568        &mut self,
569        update: NoteTransportLayerUpdate,
570    ) -> Result<(Vec<NoteId>, Vec<NoteDetailsCommitment>), ClientError> {
571        let NoteTransportLayerUpdate {
572            note_files,
573            id_by_commitment,
574            covered_tags,
575            cursor,
576        } = update;
577
578        let written = self.import_notes(&note_files).await?;
579        let mut imported_ids: Vec<NoteId> = written
580            .iter()
581            .filter_map(|commitment| id_by_commitment.get(commitment).copied())
582            .collect();
583
584        if let Some(covered_tags) = covered_tags {
585            self.save_covered_tags(&covered_tags).await?;
586        }
587
588        if let Some(cursor) = cursor {
589            self.store.update_note_transport_cursor(cursor).await?;
590        }
591
592        imported_ids.sort_unstable();
593        imported_ids.dedup();
594
595        Ok((imported_ids, written))
596    }
597
598    /// Drops deliveries of notes a local transaction is consuming; importing them would fail on the
599    /// no-overwrite-while-processing guard.
600    async fn drop_notes_processed_locally(
601        &self,
602        notes: &mut Vec<(Note, Option<BlockNumber>)>,
603    ) -> Result<(), ClientError> {
604        if notes.is_empty() {
605            return Ok(());
606        }
607
608        let commitments = notes.iter().map(|(note, _)| note.details_commitment()).collect();
609        let processing: BTreeSet<NoteDetailsCommitment> = self
610            .get_input_notes(NoteFilter::DetailsCommitments(commitments))
611            .await?
612            .into_iter()
613            .filter(InputNoteRecord::is_processing)
614            .map(|record| record.details_commitment())
615            .collect();
616
617        if !processing.is_empty() {
618            tracing::warn!(?processing, "skipping deliveries of notes being consumed locally");
619            notes.retain(|(note, _)| !processing.contains(&note.details_commitment()));
620        }
621        Ok(())
622    }
623}
624
625// NOTE TRANSPORT FETCH
626// ================================================================================================
627
628/// What the note transport fetch returned, before anything is written.
629///
630/// Built by [`Client::fetch_note_transport_updates`] and consumed by
631/// [`Client::apply_note_transport_update`].
632#[derive(Default)]
633pub(crate) struct NoteTransportLayerUpdate {
634    /// Notes to import, backfill pages first and then the steady-state page.
635    note_files: Vec<NoteFile>,
636    /// Note ids by details commitment, taken from the note headers the transport returned. Used to
637    /// resolve the written records back to ids.
638    id_by_commitment: BTreeMap<NoteDetailsCommitment, NoteId>,
639    /// Covered-tag set to persist, `None` when it did not change.
640    covered_tags: Option<BTreeSet<NoteTag>>,
641    /// New global cursor, from the steady-state page. `None` when no page was fetched.
642    cursor: Option<NoteTransportCursor>,
643}
644
645/// Note transport cursor
646///
647/// Identifies a position in the note transport service's stored-note sequence.
648#[derive(Clone, Copy, Debug, PartialEq, PartialOrd, Eq, Ord)]
649pub struct NoteTransportCursor(Option<(u64, u64)>);
650
651impl NoteTransportCursor {
652    /// Returns the cursor that starts from the first retained note.
653    pub fn init() -> Self {
654        Self(None)
655    }
656
657    /// Builds a cursor from the nonce and sequence returned by the transport service.
658    pub fn from_parts(nonce: u64, sequence: u64) -> Self {
659        Self(Some((nonce, sequence)))
660    }
661
662    /// Returns the nonce and sequence, or `None` for the initial cursor.
663    pub fn parts(&self) -> Option<(u64, u64)> {
664        self.0
665    }
666}
667
668/// The part of a note that the note transport network sends to a recipient.
669///
670/// The transport sends the original header and details. It does not send note attachments. The
671/// constructor verifies that the header commits to the details.
672#[derive(Clone, Debug, PartialEq, Eq)]
673pub struct TransportNote {
674    header: NoteHeader,
675    details: NoteDetails,
676}
677
678impl TransportNote {
679    /// Creates a transport note from matching note parts.
680    pub fn new(header: NoteHeader, details: NoteDetails) -> Result<Self, NoteTransportError> {
681        validate_note_parts(&header, &details)?;
682        Ok(Self { header, details })
683    }
684
685    /// Returns the note header.
686    pub fn header(&self) -> &NoteHeader {
687        &self.header
688    }
689
690    /// Returns the note details.
691    pub fn details(&self) -> &NoteDetails {
692        &self.details
693    }
694
695    /// Returns the note header and details.
696    pub fn into_parts(self) -> (NoteHeader, NoteDetails) {
697        (self.header, self.details)
698    }
699}
700
701impl From<Note> for TransportNote {
702    fn from(note: Note) -> Self {
703        let header = *note.header();
704        let details = NoteDetails::from(note);
705        Self { header, details }
706    }
707}
708
709/// The main transport client trait for sending and receiving private notes.
710#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
711#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
712pub trait NoteTransportClient: Send + Sync {
713    /// Sends a note together with its inclusion proof.
714    ///
715    /// The transport carries the proof to the network. The network verifies `inclusion_proof`
716    /// before it stores the note and relays the exact commitment block to the recipient.
717    async fn send_note_with_proof(
718        &self,
719        note: TransportNote,
720        inclusion_proof: NoteInclusionProof,
721    ) -> Result<(), NoteTransportError>;
722
723    /// Fetches notes for the given tags.
724    ///
725    /// Downloads notes for the given tags. Returns notes after the provided cursor (pagination),
726    /// and an updated cursor.
727    async fn fetch_notes(
728        &self,
729        tag: &[NoteTag],
730        cursor: NoteTransportCursor,
731    ) -> Result<(Vec<NoteInfo>, NoteTransportCursor), NoteTransportError>;
732}
733
734/// Information about a note fetched from the note transport network
735#[derive(Debug, Clone)]
736pub struct NoteInfo {
737    /// Note header.
738    pub header: NoteHeader,
739    /// Serialized note details.
740    pub details_bytes: Vec<u8>,
741    /// Block from which the recipient starts scanning for the note's on-chain commitment. This is
742    /// either an unverified sender hint or the exact block verified by a proof-aware transport.
743    /// `None` applies the recipient's default lookback window.
744    pub block_hint: Option<BlockNumber>,
745}
746
747impl NoteInfo {
748    /// Builds a [`NoteInfo`] without a block hint (`block_hint` is `None`).
749    ///
750    /// Use the [`NoteInfo::block_hint`] field directly to attach a hint.
751    pub fn new(header: NoteHeader, details_bytes: Vec<u8>) -> Self {
752        Self { header, details_bytes, block_hint: None }
753    }
754}
755
756// RELAY OUTBOX
757// ================================================================================================
758
759/// A private note whose transport delivery has not yet succeeded.
760#[derive(Debug, Clone, PartialEq, Eq)]
761struct RelayOutboxEntry {
762    note: TransportNote,
763    inclusion_proof: NoteInclusionProof,
764}
765
766impl RelayOutboxEntry {
767    /// Sends the note and its inclusion proof through the transport.
768    async fn relay(&self, api: &dyn NoteTransportClient) -> Result<(), NoteTransportError> {
769        api.send_note_with_proof(self.note.clone(), self.inclusion_proof.clone()).await
770    }
771}
772
773// SERIALIZATION
774// ================================================================================================
775
776impl Serializable for TransportNote {
777    fn write_into<W: ByteWriter>(&self, target: &mut W) {
778        self.header.write_into(target);
779        self.details.to_bytes().write_into(target);
780    }
781}
782
783impl Deserializable for TransportNote {
784    fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
785        let header = NoteHeader::read_from(source)?;
786        let details_bytes = Vec::<u8>::read_from(source)?;
787        let details = NoteDetails::read_from_bytes(&details_bytes)?;
788        Self::new(header, details)
789            .map_err(|error| DeserializationError::InvalidValue(format!("{error}")))
790    }
791}
792
793impl Serializable for RelayOutboxEntry {
794    fn write_into<W: ByteWriter>(&self, target: &mut W) {
795        self.note.write_into(target);
796        self.inclusion_proof.write_into(target);
797    }
798}
799
800impl Deserializable for RelayOutboxEntry {
801    fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
802        let note = TransportNote::read_from(source)?;
803        let inclusion_proof = NoteInclusionProof::read_from(source)?;
804        Ok(Self { note, inclusion_proof })
805    }
806}
807
808impl Serializable for NoteInfo {
809    fn write_into<W: ByteWriter>(&self, target: &mut W) {
810        self.header.write_into(target);
811        self.details_bytes.write_into(target);
812        self.block_hint.write_into(target);
813    }
814}
815
816impl Deserializable for NoteInfo {
817    fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
818        let header = NoteHeader::read_from(source)?;
819        let details_bytes = Vec::<u8>::read_from(source)?;
820        let block_hint = Option::<BlockNumber>::read_from(source)?;
821        Ok(NoteInfo { header, details_bytes, block_hint })
822    }
823}
824
825impl Serializable for NoteTransportCursor {
826    fn write_into<W: ByteWriter>(&self, target: &mut W) {
827        self.0.write_into(target);
828    }
829}
830
831impl Deserializable for NoteTransportCursor {
832    fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
833        Ok(Self(Option::<(u64, u64)>::read_from(source)?))
834    }
835}
836
837fn rejoin_note(header: &NoteHeader, details_bytes: &[u8]) -> Result<Note, NoteTransportError> {
838    let mut reader = SliceReader::new(details_bytes);
839    let details = NoteDetails::read_from(&mut reader)?;
840    validate_note_parts(header, &details)?;
841    // The transport wire format only carries `NoteHeader` + serialized `NoteDetails`, not the
842    // attachments collection. We rejoin with empty attachments; this matches the original note only
843    // when it had no attachments in the first place.
844    let partial_metadata = *header.metadata().partial_metadata();
845    Ok(Note::new(
846        details.assets().clone(),
847        partial_metadata,
848        details.recipient().clone(),
849    ))
850}
851
852/// Checks that the note header commits to the supplied details.
853pub(crate) fn validate_note_parts(
854    header: &NoteHeader,
855    details: &NoteDetails,
856) -> Result<(), NoteTransportError> {
857    let header_commitment = header.details_commitment();
858    let details_commitment = details.commitment();
859    if header_commitment != details_commitment {
860        return Err(NoteTransportError::NoteDetailsMismatch {
861            header: header_commitment,
862            details: details_commitment,
863        });
864    }
865
866    Ok(())
867}
868
869#[cfg(test)]
870mod tests {
871    use miden_protocol::account::AccountId;
872    use miden_protocol::asset::FungibleAsset;
873    use miden_protocol::crypto::merkle::SparseMerklePath;
874    use miden_protocol::note::NoteType;
875    use miden_protocol::testing::account_id::{
876        ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET,
877        ACCOUNT_ID_REGULAR_PUBLIC_ACCOUNT_IMMUTABLE_CODE,
878        ACCOUNT_ID_SENDER,
879    };
880    use miden_standards::note::P2idNote;
881    use rand::SeedableRng;
882    use rand_chacha::ChaCha20Rng;
883
884    use super::*;
885    use crate::rng::draw_word;
886
887    #[test]
888    fn relay_outbox_entry_round_trips() {
889        let sender = AccountId::try_from(ACCOUNT_ID_SENDER).unwrap();
890        let target = AccountId::try_from(ACCOUNT_ID_REGULAR_PUBLIC_ACCOUNT_IMMUTABLE_CODE).unwrap();
891        let faucet = AccountId::try_from(ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET).unwrap();
892        let mut rng = ChaCha20Rng::seed_from_u64(0);
893        let note: Note = P2idNote::builder()
894            .sender(sender)
895            .target(target)
896            .asset(FungibleAsset::new(faucet, 100).unwrap())
897            .note_type(NoteType::Private)
898            .serial_number(draw_word(&mut rng))
899            .build()
900            .unwrap()
901            .into();
902
903        let inclusion_proof =
904            NoteInclusionProof::new(BlockNumber::from(7), 3, SparseMerklePath::default()).unwrap();
905        let entry = RelayOutboxEntry {
906            note: TransportNote::from(note),
907            inclusion_proof,
908        };
909
910        assert_eq!(RelayOutboxEntry::read_from_bytes(&entry.to_bytes()).unwrap(), entry);
911    }
912}