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::{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";
42/// Settings key for the unused aggregate transport cursor.
43#[deprecated(since = "0.17.1", note = "note transport stores a cursor for each tag")]
44pub const NOTE_TRANSPORT_CURSOR_STORE_SETTING: &str = "note_transport_cursor";
45pub const NOTE_TRANSPORT_CURSORS_KEY: &str = "note_transport_cursors";
46
47type NoteTransportCursors = BTreeMap<NoteTag, NoteTransportCursor>;
48/// Maximum number of note tags in one transport fetch request. The service rejects larger requests.
49const MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST: usize = 128;
50/// Limits the pages for each request group so other groups can progress in the same sync.
51const MAX_NOTE_TRANSPORT_PAGES_PER_GROUP: usize = 32;
52
53/// Legacy settings key for note transport backfill state.
54#[deprecated(since = "0.17.1", note = "note transport no longer keeps per-tag backfill state")]
55pub const NOTE_TRANSPORT_COVERED_TAGS_KEY: &str = "note_transport_covered_tags";
56
57/// Settings key for the durable relay outbox: a serialized `Vec<RelayOutboxEntry>` of private notes
58/// whose transport delivery has not yet succeeded. [`Client::send_private_note_with_proof`] appends
59/// (replacing any entry with the same note id) before relaying; [`Client::flush_relay_outbox`]
60/// drains entries that re-send successfully. Reusing the settings k/v avoids a Store-trait schema
61/// change while surviving process restarts.
62pub const NOTE_TRANSPORT_OUTBOX_KEY: &str = "note_transport_outbox";
63
64/// Client note transport methods.
65impl<AUTH> Client<AUTH> {
66    /// Maximum number of note tags in one transport fetch request.
67    pub const MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST: usize = MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST;
68
69    /// Legacy name for [`Self::MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST`].
70    ///
71    /// This value is not a limit on the number of account tags that the client can track.
72    #[deprecated(
73        since = "0.17.1",
74        note = "use MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST; account tags are no longer limited"
75    )]
76    pub const MAX_ACCOUNT_TAGS: usize = Self::MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST;
77
78    /// Check if note transport connection is configured
79    pub fn is_note_transport_enabled(&self) -> bool {
80        self.note_transport_api.is_some()
81    }
82
83    /// Returns the Note Transport client
84    ///
85    /// Errors if the note transport is not configured.
86    pub(crate) fn get_note_transport_api(
87        &self,
88    ) -> Result<Arc<dyn NoteTransportClient>, NoteTransportError> {
89        self.note_transport_api.clone().ok_or(NoteTransportError::Disabled)
90    }
91
92    /// Send a note through the note transport network together with its inclusion proof.
93    ///
94    /// The note will be end-to-end encrypted (unimplemented, currently plaintext) using the
95    /// provided recipient's `address` details. The recipient will be able to retrieve this note
96    /// through the note's [`NoteTag`].
97    ///
98    /// The transport carries the proof through [`NoteTransportClient::send_note_with_proof`]. The
99    /// network verifies `inclusion_proof` against its node before it stores the note and relays the
100    /// exact commitment block to the recipient. The proof exists once the transaction that created
101    /// the note is committed and the sender has synced past it; see
102    /// [`OutputNoteRecord::inclusion_proof`](crate::store::OutputNoteRecord::inclusion_proof).
103    ///
104    /// **Durability.** The note and its proof are persisted to the outbox before the transport
105    /// call. If the call fails or is interrupted, the entry stays in the outbox and is retried on
106    /// the next [`Client::flush_relay_outbox`] (which [`Client::sync_note_transport`] runs), so a
107    /// transient transport failure does not drop the note. The receiver dedupes by note id, so a
108    /// re-send after a partial success is harmless.
109    pub async fn send_private_note_with_proof(
110        &mut self,
111        note: Note,
112        address: &Address,
113        inclusion_proof: NoteInclusionProof,
114    ) -> Result<(), ClientError> {
115        let api = self.get_note_transport_api()?;
116
117        let note = TransportNote::from(note);
118        let note_id = note.header().id();
119        // The address is reserved for end-to-end encryption of the note details:
120        // address.key().encrypt(note.details().to_bytes()).
121        let _ = address;
122
123        // Persist the payload before the network call so a failed or interrupted send leaves a
124        // recoverable record rather than losing the only copy with the call frame. The proof
125        // travels with the entry so a retried send relays the same value.
126        let entry = RelayOutboxEntry { note, inclusion_proof };
127        let mut outbox = self.load_relay_outbox().await?;
128        // Replace any existing entry for this note id so the latest payload wins when a
129        // still-pending note is re-sent.
130        outbox.retain(|e| e.note.header().id() != note_id);
131        outbox.push(entry.clone());
132        self.save_relay_outbox(outbox).await?;
133
134        entry.relay(api.as_ref()).await?;
135
136        // Relay succeeded — drop the entry. A failed store write here is tolerable: the next flush
137        // re-sends and the receiver dedupes by note id, so a stale entry never causes loss.
138        let mut outbox = self.load_relay_outbox().await?;
139        outbox.retain(|e| e.note.header().id() != note_id);
140        self.save_relay_outbox(outbox).await?;
141
142        Ok(())
143    }
144
145    /// Re-attempt every relay payload in the durable outbox. Each entry is a private note whose
146    /// previous transport delivery failed. Successful re-sends are dropped; failures are kept for
147    /// the next call. Every entry is attempted independently, so one persistently-failing note does
148    /// not block the others.
149    ///
150    /// [`Client::sync_note_transport`] runs this automatically and ignores its error, so a relay
151    /// failure can't block a sync. Callers driving retries themselves can invoke it directly and
152    /// inspect the returned error.
153    pub async fn flush_relay_outbox(&self) -> Result<(), ClientError> {
154        let api = self.get_note_transport_api()?;
155
156        let entries = self.load_relay_outbox().await?;
157        if entries.is_empty() {
158            return Ok(());
159        }
160
161        // Attempt every entry independently so a single persistently-failing note can't block the
162        // rest. The outbox holds only the caller's own failed sends, so it stays small and this is
163        // not a meaningful burst.
164        let mut remaining = Vec::new();
165        let mut last_err: Option<NoteTransportError> = None;
166
167        for entry in entries {
168            match entry.relay(api.as_ref()).await {
169                Ok(()) => {},
170                Err(err) => {
171                    tracing::warn!(?err, "relay-outbox entry retry failed; will retry next sync");
172                    remaining.push(entry);
173                    last_err = Some(err);
174                },
175            }
176        }
177
178        self.save_relay_outbox(remaining).await?;
179
180        if let Some(err) = last_err {
181            return Err(err.into());
182        }
183        Ok(())
184    }
185
186    /// Load the durable relay outbox.
187    ///
188    /// Returns an empty `Vec` if the outbox key is absent. On deserialization failure (schema
189    /// mismatch or storage corruption) the entry is dropped and an empty `Vec` is returned —
190    /// leaving unreadable bytes in place would block every subsequent relay because each sync would
191    /// re-read them.
192    async fn load_relay_outbox(&self) -> Result<Vec<RelayOutboxEntry>, ClientError> {
193        let bytes = self
194            .store
195            .get_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_OUTBOX_KEY))
196            .await
197            .map_err(ClientError::StoreError)?;
198        let Some(bytes) = bytes else {
199            return Ok(Vec::new());
200        };
201        match Vec::<RelayOutboxEntry>::read_from_bytes(&bytes) {
202            Ok(entries) => Ok(entries),
203            Err(err) => {
204                tracing::warn!(?err, "dropping unreadable relay outbox; resetting to empty");
205                self.store
206                    .remove_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_OUTBOX_KEY))
207                    .await
208                    .map_err(ClientError::StoreError)?;
209                Ok(Vec::new())
210            },
211        }
212    }
213
214    /// Persist the relay outbox, removing the key entirely when empty so the settings table doesn't
215    /// accumulate empty-vec blobs.
216    async fn save_relay_outbox(&self, entries: Vec<RelayOutboxEntry>) -> Result<(), ClientError> {
217        let key = String::from(NOTE_TRANSPORT_OUTBOX_KEY);
218        if entries.is_empty() {
219            self.store
220                .remove_setting(SettingScope::Client, key)
221                .await
222                .map_err(ClientError::StoreError)?;
223            return Ok(());
224        }
225        let bytes = entries.to_bytes();
226        self.store
227            .set_setting(SettingScope::Client, key, bytes)
228            .await
229            .map_err(ClientError::StoreError)
230    }
231
232    /// Loads the cursor for each tag used by the transport fetch.
233    ///
234    /// A missing or unreadable value resets all tags so the next fetch safely reads their retained
235    /// history.
236    async fn load_note_transport_cursors(&self) -> Result<NoteTransportCursors, ClientError> {
237        let bytes = self
238            .store
239            .get_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_CURSORS_KEY))
240            .await
241            .map_err(ClientError::StoreError)?;
242        let Some(bytes) = bytes else {
243            return Ok(BTreeMap::new());
244        };
245
246        match NoteTransportCursors::read_from_bytes(&bytes) {
247            Ok(cursors) => Ok(cursors),
248            Err(err) => {
249                tracing::warn!(?err, "resetting unreadable note transport cursors");
250                Ok(BTreeMap::new())
251            },
252        }
253    }
254
255    /// Saves the cursor for each tag used by the transport fetch.
256    async fn save_note_transport_cursors(
257        &self,
258        cursors: &NoteTransportCursors,
259    ) -> Result<(), ClientError> {
260        let key = String::from(NOTE_TRANSPORT_CURSORS_KEY);
261        if cursors.is_empty() {
262            self.store.remove_setting(SettingScope::Client, key).await?;
263        } else {
264            self.store.set_setting(SettingScope::Client, key, cursors.to_bytes()).await?;
265        }
266        Ok(())
267    }
268}
269
270impl<AUTH> Client<AUTH>
271where
272    AUTH: TransactionAuthenticator + Sync + 'static,
273{
274    /// Legacy per-sync cap for tag backfills.
275    #[deprecated(since = "0.17.1", note = "note transport no longer performs per-tag backfills")]
276    pub const MAX_BACKFILL_TAGS_PER_SYNC: usize = 64;
277
278    /// Fetch notes for tracked note tags.
279    ///
280    /// The client queries the configured note transport node for all tracked tags. To list tracked
281    /// tags, use [`Client::get_note_tags`]. To add a user-source tag, use [`Client::add_note_tag`].
282    /// Fetched notes are stored in the client store.
283    ///
284    /// An internal pagination mechanism starts each request from the lowest cursor of its tags.
285    /// Tags without a stored cursor start from the first retained note. The service can deliver a
286    /// note again after a tag is added or removed; the import drops these duplicates.
287    ///
288    /// A failed request returns an error after successful pages are imported and their cursors are
289    /// saved. Histories that exceed the page budget continue on the next call.
290    pub async fn fetch_private_notes(&mut self) -> Result<(), ClientError> {
291        self.ensure_genesis_in_place().await?;
292
293        let mut update = self.fetch_transport_notes_in_chunks().await?;
294        let fetch_error = update.fetch_error.take();
295        self.apply_note_transport_update(update).await?;
296        if let Some(error) = fetch_error {
297            return Err(error);
298        }
299
300        Ok(())
301    }
302
303    /// Screens the transport-delivered notes carrying a tag derived from a tracked account,
304    /// discarding those that no tracked account can consume. Notes carrying any other tag are kept
305    /// as delivered.
306    async fn screen_transport_notes(
307        &self,
308        notes: &mut Vec<(NoteId, Note, Option<BlockNumber>)>,
309    ) -> Result<(), ClientError> {
310        let account_tags = self.tracked_account_tags().await?;
311
312        let notes_to_screen: Vec<Note> = notes
313            .iter()
314            .filter(|(_, note, _)| account_tags.contains(&note.metadata().tag()))
315            .map(|(_, note, _)| note.clone())
316            .collect();
317        let consumable = self.note_screener().get_batch_consumability(&notes_to_screen).await?;
318
319        // Discard the notes whose tag match the tracked accounts but are not consumable.
320        notes.retain(|(_, note, _)| {
321            !account_tags.contains(&note.metadata().tag()) || consumable.contains_key(&note.id())
322        });
323
324        Ok(())
325    }
326
327    /// Returns the tracked tags that were registered for an account, i.e. derived from its ID.
328    async fn tracked_account_tags(&self) -> Result<BTreeSet<NoteTag>, ClientError> {
329        let tags = self
330            .store
331            .get_note_tags()
332            .await?
333            .into_iter()
334            .filter(|record| matches!(record.source, NoteTagSource::Account(_)))
335            .map(|record| record.tag)
336            .collect();
337        Ok(tags)
338    }
339
340    /// Fetches bounded pages for request groups of at most
341    /// [`Self::MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST`] tags.
342    ///
343    /// Each group starts from its lowest cursor. The import drops notes delivered again. Each tag
344    /// keeps the higher of its stored cursor and the page cursor when their database nonces match.
345    ///
346    /// A failed page keeps the pages fetched before it and leaves the other groups to proceed.
347    async fn fetch_transport_notes_in_chunks(
348        &self,
349    ) -> Result<NoteTransportLayerUpdate, ClientError> {
350        let api = self.get_note_transport_api()?;
351        let tags: Vec<_> = self.store.get_unique_note_tags().await?.into_iter().collect();
352        let stored = self.load_note_transport_cursors().await?;
353        let mut cursors: NoteTransportCursors = tags
354            .iter()
355            .filter_map(|tag| stored.get(tag).map(|cursor| (*tag, *cursor)))
356            .collect();
357
358        let mut update = NoteTransportLayerUpdate::default();
359        let mut notes = Vec::new();
360        for (start, group) in transport_request_groups(&tags, &cursors) {
361            let mut cursor = start;
362            for _ in 0..MAX_NOTE_TRANSPORT_PAGES_PER_GROUP {
363                let page = match api
364                    .fetch_notes_page(&group, cursor)
365                    .await
366                    .and_then(|page| validate_transport_page(page, &group, cursor))
367                {
368                    Ok(page) => page,
369                    Err(error) => {
370                        update.fetch_error.get_or_insert(error.into());
371                        break;
372                    },
373                };
374                for (id, note, block_hint) in page.notes {
375                    update.id_by_commitment.insert(note.details_commitment(), id);
376                    notes.push((id, note, block_hint));
377                }
378                cursor = page.cursor;
379                for tag in &group {
380                    let position = cursors.entry(*tag).or_insert(cursor);
381                    *position = advance_transport_cursor(*position, cursor);
382                }
383                if !page.has_more {
384                    break;
385                }
386            }
387        }
388
389        update.note_files = self.prepare_transport_notes(notes).await?;
390        update.cursors = Some(cursors);
391        Ok(update)
392    }
393
394    /// Screens fetched notes and prepares them for import with one set of store reads.
395    async fn prepare_transport_notes(
396        &self,
397        mut notes: Vec<(NoteId, Note, Option<BlockNumber>)>,
398    ) -> Result<Vec<NoteFile>, ClientError> {
399        // Fallback lookback window, in blocks, used only for notes the transport delivered without
400        // block information. Scanning back from sync height handles the race where a note is
401        // committed on-chain just before the NTL delivers its data. Without it,
402        // check_expected_notes would scan from sync_height forward and miss the already-committed
403        // note. A transport-provided block is deterministic and always preferred.
404        const NOTE_LOOKBACK_BLOCKS: u32 = 20;
405
406        if notes.is_empty() {
407            return Ok(Vec::new());
408        }
409
410        self.drop_notes_resolved_locally(&mut notes).await?;
411
412        // Screen the transport-delivered notes to discard the ones that are not relevant to the
413        // accounts tracked by the client. Boxed to avoid a `clippy::large_futures` warning, since
414        // the sync future is already close to the size limit.
415        Box::pin(self.screen_transport_notes(&mut notes)).await?;
416
417        let sync_height = self.get_sync_height().await?;
418        let fallback_after_block_num =
419            BlockNumber::from(sync_height.as_u32().saturating_sub(NOTE_LOOKBACK_BLOCKS));
420
421        let mut note_files = Vec::with_capacity(notes.len());
422        for (_, note, block_hint) in notes {
423            let tag = note.metadata().tag();
424            // Prefer the transport-provided block, falling back to the lookback window when absent.
425            let after_block_num = block_hint.unwrap_or(fallback_after_block_num);
426            note_files.push(NoteFile::ExpectedNote {
427                details: note.into(),
428                sync_hint: NoteSyncHint::new(after_block_num, tag),
429            });
430        }
431
432        Ok(note_files)
433    }
434
435    /// Fetches the notes the Note Transport Layer holds for the tracked tags.
436    ///
437    /// Fetches bounded pages for each group of tracked tags. This performs no node call and writes
438    /// nothing but the relay outbox, so it can run concurrently with the chain fetch. The caller
439    /// imports the returned files and then persists all tag cursors. A failed request preserves
440    /// successful pages in the returned update.
441    ///
442    /// Returns empty data when note transport is not configured.
443    pub(crate) async fn fetch_note_transport_updates(
444        &self,
445    ) -> Result<NoteTransportLayerUpdate, ClientError> {
446        if !self.is_note_transport_enabled() {
447            return Ok(NoteTransportLayerUpdate::default());
448        }
449
450        // Drain any private notes whose previous relay attempt failed. A flush error is logged, not
451        // propagated: a failing relay must not block the sync, and the entries stay durable for the
452        // next attempt. This is the one write this phase performs; it touches only the outbox
453        // setting, which is independent of everything the apply phase writes.
454        if let Err(err) = self.flush_relay_outbox().await {
455            tracing::warn!(?err, "relay outbox flush failed during sync; entries retained");
456        }
457
458        self.fetch_transport_notes_in_chunks().await
459    }
460
461    /// Writes everything [`Client::fetch_note_transport_updates`] returned, in two steps:
462    ///
463    /// 1. Imports the fetched notes, which resolves their on-chain state and stores the records.
464    /// 2. Saves the cursor for each tag.
465    ///
466    /// The notes are written before the cursors, so a crash between them re-fetches instead of
467    /// skipping notes that were never written.
468    ///
469    /// Returns the ids of the imported notes and the details commitments of the records written.
470    pub(crate) async fn apply_note_transport_update(
471        &mut self,
472        update: NoteTransportLayerUpdate,
473    ) -> Result<(Vec<NoteId>, Vec<NoteDetailsCommitment>), ClientError> {
474        let NoteTransportLayerUpdate {
475            note_files,
476            id_by_commitment,
477            cursors,
478            fetch_error,
479        } = update;
480
481        let written = self.import_notes(&note_files).await?;
482        let mut imported_ids: Vec<NoteId> = written
483            .iter()
484            .filter_map(|commitment| id_by_commitment.get(commitment).copied())
485            .collect();
486
487        if let Some(cursors) = cursors {
488            self.save_note_transport_cursors(&cursors).await?;
489        }
490        if let Some(error) = fetch_error {
491            tracing::warn!(?error, "note transport fetch failed; saved successful pages for retry");
492        }
493
494        imported_ids.sort_unstable();
495        imported_ids.dedup();
496
497        Ok((imported_ids, written))
498    }
499
500    /// Drops deliveries of notes whose local record does not need them.
501    ///
502    /// A note that a local transaction is consuming cannot be overwritten, so its import would
503    /// fail. A note that is already committed or consumed gains nothing from a second import and
504    /// would cost a node request. A request from a group's lowest cursor can deliver such notes
505    /// again. Expected, unverified, and invalid notes pass through because their records can need a
506    /// new inclusion proof. A resolved record must match the ID in the transport header.
507    async fn drop_notes_resolved_locally(
508        &self,
509        notes: &mut Vec<(NoteId, Note, Option<BlockNumber>)>,
510    ) -> Result<(), ClientError> {
511        let commitments = notes.iter().map(|(_, note, _)| note.details_commitment()).collect();
512        let records: BTreeMap<_, _> = self
513            .get_input_notes(NoteFilter::DetailsCommitments(commitments))
514            .await?
515            .into_iter()
516            .map(|record| (record.details_commitment(), record))
517            .collect();
518        notes.retain(|(id, note, _)| {
519            let Some(record) = records.get(&note.details_commitment()) else {
520                return true;
521            };
522            if record.is_processing() {
523                tracing::warn!(%id, "skipping delivery of a note being consumed locally");
524                return false;
525            }
526            !((record.is_committed() || record.is_consumed()) && record.id() == Some(*id))
527        });
528        Ok(())
529    }
530}
531
532// NOTE TRANSPORT FETCH
533// ================================================================================================
534
535/// What the note transport fetch returned, before anything is written.
536///
537/// Built by [`Client::fetch_note_transport_updates`] and consumed by
538/// [`Client::apply_note_transport_update`].
539#[derive(Default)]
540pub(crate) struct NoteTransportLayerUpdate {
541    /// Notes to import.
542    note_files: Vec<NoteFile>,
543    /// Note ids by details commitment, taken from the note headers the transport returned. Used to
544    /// resolve the written records back to ids.
545    id_by_commitment: BTreeMap<NoteDetailsCommitment, NoteId>,
546    /// Cursor for each tag used by the steady-state fetch.
547    cursors: Option<NoteTransportCursors>,
548    /// First request error. Successful pages remain available for import.
549    pub(crate) fetch_error: Option<ClientError>,
550}
551
552/// Splits the tracked tags into request groups and returns the cursor each group starts from.
553///
554/// Tags without a cursor start from the first retained note. Tags with a cursor are grouped by
555/// database nonce, because sequences from different databases do not compare. Each nonce group is
556/// sorted by sequence before it is split, so the tags in one request sit close together and the
557/// lowest cursor limits repeated deliveries.
558fn transport_request_groups(
559    tags: &[NoteTag],
560    cursors: &NoteTransportCursors,
561) -> Vec<(NoteTransportCursor, Vec<NoteTag>)> {
562    let mut by_nonce = BTreeMap::<Option<u64>, Vec<(NoteTransportCursor, NoteTag)>>::new();
563    for tag in tags {
564        let cursor = cursors.get(tag).copied().unwrap_or_else(NoteTransportCursor::init);
565        by_nonce
566            .entry(cursor.parts().map(|(nonce, _)| nonce))
567            .or_default()
568            .push((cursor, *tag));
569    }
570
571    let mut groups = Vec::new();
572    for mut entries in by_nonce.into_values() {
573        entries.sort_unstable();
574        for chunk in entries.chunks(MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST) {
575            let start = chunk[0].0;
576            groups.push((start, chunk.iter().map(|(_, tag)| *tag).collect()));
577        }
578    }
579    groups
580}
581
582/// Returns the position of a tag after a page that ends at `page_cursor`.
583///
584/// A page from the same database never moves a tag backwards. A page from a different database
585/// replaces a position that database does not recognize.
586fn advance_transport_cursor(
587    current: NoteTransportCursor,
588    page_cursor: NoteTransportCursor,
589) -> NoteTransportCursor {
590    match (current.parts(), page_cursor.parts()) {
591        (Some((nonce, sequence)), Some((page_nonce, page_sequence))) if nonce == page_nonce => {
592            NoteTransportCursor::from_parts(nonce, sequence.max(page_sequence))
593        },
594        _ => page_cursor,
595    }
596}
597
598struct ValidatedTransportPage {
599    notes: Vec<(NoteId, Note, Option<BlockNumber>)>,
600    cursor: NoteTransportCursor,
601    has_more: bool,
602}
603
604/// Validates a complete page before its notes or cursor can advance local progress.
605fn validate_transport_page(
606    page: NoteTransportPage,
607    tags: &[NoteTag],
608    request_cursor: NoteTransportCursor,
609) -> Result<ValidatedTransportPage, NoteTransportError> {
610    let Some((nonce, sequence)) = page.cursor.parts() else {
611        return Err(NoteTransportError::Network(String::from("fetch response has no cursor")));
612    };
613    if request_cursor.parts().is_some_and(|(request_nonce, request_sequence)| {
614        nonce == request_nonce
615            && (sequence < request_sequence
616                || (!page.notes.is_empty() && sequence == request_sequence))
617    }) || (!page.notes.is_empty() && sequence == 0)
618        || (page.has_more && page.notes.is_empty())
619    {
620        return Err(NoteTransportError::Network(String::from(
621            "fetch response has invalid pagination progress",
622        )));
623    }
624    let mut notes = Vec::with_capacity(page.notes.len());
625    for info in page.notes {
626        let note = rejoin_note(&info.header, &info.details_bytes)?;
627        if !tags.contains(&note.metadata().tag()) {
628            return Err(NoteTransportError::UnrequestedTag(note.metadata().tag()));
629        }
630        // The header ID includes attachments that the transport does not send.
631        notes.push((info.header.id(), note, info.block_hint));
632    }
633    Ok(ValidatedTransportPage {
634        notes,
635        cursor: page.cursor,
636        has_more: page.has_more,
637    })
638}
639
640/// Note transport cursor
641///
642/// Identifies a position in the note transport service's stored-note sequence.
643///
644/// The sequence is global across tags in one service database. Compare sequences only when their
645/// nonces match.
646#[derive(Clone, Copy, Debug, PartialEq, PartialOrd, Eq, Ord)]
647pub struct NoteTransportCursor(Option<(u64, u64)>);
648
649impl NoteTransportCursor {
650    /// Returns the cursor that starts from the first retained note.
651    pub fn init() -> Self {
652        Self(None)
653    }
654
655    /// Builds a cursor from the nonce and sequence returned by the transport service.
656    pub fn from_parts(nonce: u64, sequence: u64) -> Self {
657        Self(Some((nonce, sequence)))
658    }
659
660    /// Returns the nonce and sequence, or `None` for the initial cursor.
661    pub fn parts(&self) -> Option<(u64, u64)> {
662        self.0
663    }
664}
665
666/// The part of a note that the note transport network sends to a recipient.
667///
668/// The transport sends the original header and details. It does not send note attachments. The
669/// constructor verifies that the header commits to the details.
670#[derive(Clone, Debug, PartialEq, Eq)]
671pub struct TransportNote {
672    header: NoteHeader,
673    details: NoteDetails,
674}
675
676impl TransportNote {
677    /// Creates a transport note from matching note parts.
678    pub fn new(header: NoteHeader, details: NoteDetails) -> Result<Self, NoteTransportError> {
679        validate_note_parts(&header, &details)?;
680        Ok(Self { header, details })
681    }
682
683    /// Returns the note header.
684    pub fn header(&self) -> &NoteHeader {
685        &self.header
686    }
687
688    /// Returns the note details.
689    pub fn details(&self) -> &NoteDetails {
690        &self.details
691    }
692
693    /// Returns the note header and details.
694    pub fn into_parts(self) -> (NoteHeader, NoteDetails) {
695        (self.header, self.details)
696    }
697}
698
699impl From<Note> for TransportNote {
700    fn from(note: Note) -> Self {
701        let header = *note.header();
702        let details = NoteDetails::from(note);
703        Self { header, details }
704    }
705}
706
707/// The main transport client trait for sending and receiving private notes.
708#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
709#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
710pub trait NoteTransportClient: Send + Sync {
711    /// Sends a note together with its inclusion proof.
712    ///
713    /// The transport carries the proof to the network. The network verifies `inclusion_proof`
714    /// before it stores the note and relays the exact commitment block to the recipient.
715    async fn send_note_with_proof(
716        &self,
717        note: TransportNote,
718        inclusion_proof: NoteInclusionProof,
719    ) -> Result<(), NoteTransportError>;
720
721    /// Fetches notes for the given tags.
722    ///
723    /// Downloads notes for the given tags. Returns notes after the provided cursor (pagination),
724    /// and an updated cursor.
725    async fn fetch_notes(
726        &self,
727        tag: &[NoteTag],
728        cursor: NoteTransportCursor,
729    ) -> Result<(Vec<NoteInfo>, NoteTransportCursor), NoteTransportError>;
730
731    /// Fetches a page and reports whether another page is available.
732    ///
733    /// Transports without a continuation flag use an empty page to confirm the end of the history.
734    async fn fetch_notes_page(
735        &self,
736        tags: &[NoteTag],
737        cursor: NoteTransportCursor,
738    ) -> Result<NoteTransportPage, NoteTransportError> {
739        let (notes, cursor) = self.fetch_notes(tags, cursor).await?;
740        let has_more = !notes.is_empty();
741        Ok(NoteTransportPage { notes, cursor, has_more })
742    }
743}
744
745/// A page of notes from the transport service.
746pub struct NoteTransportPage {
747    /// Notes in global sequence order.
748    pub notes: Vec<NoteInfo>,
749    /// Position of the last returned note. An empty page keeps the request position.
750    pub cursor: NoteTransportCursor,
751    /// Indicates that the requested tags have another page.
752    pub has_more: bool,
753}
754
755/// Information about a note fetched from the note transport network
756#[derive(Debug, Clone)]
757pub struct NoteInfo {
758    /// Note header.
759    pub header: NoteHeader,
760    /// Serialized note details.
761    pub details_bytes: Vec<u8>,
762    /// Block from which the recipient starts scanning for the note's on-chain commitment. This is
763    /// either an unverified sender hint or the exact block verified by a proof-aware transport.
764    /// `None` applies the recipient's default lookback window.
765    pub block_hint: Option<BlockNumber>,
766}
767
768impl NoteInfo {
769    /// Builds a [`NoteInfo`] without a block hint (`block_hint` is `None`).
770    ///
771    /// Use the [`NoteInfo::block_hint`] field directly to attach a hint.
772    pub fn new(header: NoteHeader, details_bytes: Vec<u8>) -> Self {
773        Self { header, details_bytes, block_hint: None }
774    }
775}
776
777// RELAY OUTBOX
778// ================================================================================================
779
780/// A private note whose transport delivery has not yet succeeded.
781#[derive(Debug, Clone, PartialEq, Eq)]
782struct RelayOutboxEntry {
783    note: TransportNote,
784    inclusion_proof: NoteInclusionProof,
785}
786
787impl RelayOutboxEntry {
788    /// Sends the note and its inclusion proof through the transport.
789    async fn relay(&self, api: &dyn NoteTransportClient) -> Result<(), NoteTransportError> {
790        api.send_note_with_proof(self.note.clone(), self.inclusion_proof.clone()).await
791    }
792}
793
794// SERIALIZATION
795// ================================================================================================
796
797impl Serializable for TransportNote {
798    fn write_into<W: ByteWriter>(&self, target: &mut W) {
799        self.header.write_into(target);
800        self.details.to_bytes().write_into(target);
801    }
802}
803
804impl Deserializable for TransportNote {
805    fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
806        let header = NoteHeader::read_from(source)?;
807        let details_bytes = Vec::<u8>::read_from(source)?;
808        let details = NoteDetails::read_from_bytes(&details_bytes)?;
809        Self::new(header, details)
810            .map_err(|error| DeserializationError::InvalidValue(format!("{error}")))
811    }
812}
813
814impl Serializable for RelayOutboxEntry {
815    fn write_into<W: ByteWriter>(&self, target: &mut W) {
816        self.note.write_into(target);
817        self.inclusion_proof.write_into(target);
818    }
819}
820
821impl Deserializable for RelayOutboxEntry {
822    fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
823        let note = TransportNote::read_from(source)?;
824        let inclusion_proof = NoteInclusionProof::read_from(source)?;
825        Ok(Self { note, inclusion_proof })
826    }
827}
828
829impl Serializable for NoteInfo {
830    fn write_into<W: ByteWriter>(&self, target: &mut W) {
831        self.header.write_into(target);
832        self.details_bytes.write_into(target);
833        self.block_hint.write_into(target);
834    }
835}
836
837impl Deserializable for NoteInfo {
838    fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
839        let header = NoteHeader::read_from(source)?;
840        let details_bytes = Vec::<u8>::read_from(source)?;
841        let block_hint = Option::<BlockNumber>::read_from(source)?;
842        Ok(NoteInfo { header, details_bytes, block_hint })
843    }
844}
845
846impl Serializable for NoteTransportCursor {
847    fn write_into<W: ByteWriter>(&self, target: &mut W) {
848        self.0.write_into(target);
849    }
850}
851
852impl Deserializable for NoteTransportCursor {
853    fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
854        Ok(Self(Option::<(u64, u64)>::read_from(source)?))
855    }
856}
857
858fn rejoin_note(header: &NoteHeader, details_bytes: &[u8]) -> Result<Note, NoteTransportError> {
859    let mut reader = SliceReader::new(details_bytes);
860    let details = NoteDetails::read_from(&mut reader)?;
861    validate_note_parts(header, &details)?;
862    // The transport wire format only carries `NoteHeader` + serialized `NoteDetails`, not the
863    // attachments collection. We rejoin with empty attachments; this matches the original note only
864    // when it had no attachments in the first place.
865    let partial_metadata = *header.metadata().partial_metadata();
866    Ok(Note::new(
867        details.assets().clone(),
868        partial_metadata,
869        details.recipient().clone(),
870    ))
871}
872
873/// Checks that the note header commits to the supplied details.
874pub(crate) fn validate_note_parts(
875    header: &NoteHeader,
876    details: &NoteDetails,
877) -> Result<(), NoteTransportError> {
878    let header_commitment = header.details_commitment();
879    let details_commitment = details.commitment();
880    if header_commitment != details_commitment {
881        return Err(NoteTransportError::NoteDetailsMismatch {
882            header: header_commitment,
883            details: details_commitment,
884        });
885    }
886
887    Ok(())
888}
889
890#[cfg(test)]
891mod tests {
892    use miden_protocol::account::AccountId;
893    use miden_protocol::asset::FungibleAsset;
894    use miden_protocol::crypto::merkle::SparseMerklePath;
895    use miden_protocol::note::NoteType;
896    use miden_protocol::testing::account_id::{
897        ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET,
898        ACCOUNT_ID_REGULAR_PUBLIC_ACCOUNT_IMMUTABLE_CODE,
899        ACCOUNT_ID_SENDER,
900    };
901    use miden_standards::note::P2idNote;
902    use rand::SeedableRng;
903    use rand_chacha::ChaCha20Rng;
904
905    use super::*;
906    use crate::rng::draw_word;
907
908    #[test]
909    fn transport_rejects_invalid_pagination() {
910        let cursor = NoteTransportCursor::from_parts(1, 10);
911        for (returned, has_more) in [
912            (NoteTransportCursor::init(), false),
913            (NoteTransportCursor::from_parts(1, 9), false),
914            (cursor, true),
915        ] {
916            let page = NoteTransportPage {
917                notes: Vec::new(),
918                cursor: returned,
919                has_more,
920            };
921            assert!(validate_transport_page(page, &[], cursor).is_err());
922        }
923    }
924
925    #[test]
926    fn relay_outbox_entry_round_trips() {
927        let sender = AccountId::try_from(ACCOUNT_ID_SENDER).unwrap();
928        let target = AccountId::try_from(ACCOUNT_ID_REGULAR_PUBLIC_ACCOUNT_IMMUTABLE_CODE).unwrap();
929        let faucet = AccountId::try_from(ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET).unwrap();
930        let mut rng = ChaCha20Rng::seed_from_u64(0);
931        let note: Note = P2idNote::builder()
932            .sender(sender)
933            .target(target)
934            .asset(FungibleAsset::new(faucet, 100).unwrap())
935            .note_type(NoteType::Private)
936            .serial_number(draw_word(&mut rng))
937            .build()
938            .unwrap()
939            .into();
940
941        let inclusion_proof =
942            NoteInclusionProof::new(BlockNumber::from(7), 3, SparseMerklePath::default()).unwrap();
943        let entry = RelayOutboxEntry {
944            note: TransportNote::from(note),
945            inclusion_proof,
946        };
947
948        assert_eq!(RelayOutboxEntry::read_from_bytes(&entry.to_bytes()).unwrap(), entry);
949    }
950}