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, ¬e_tags, &mut id_by_commitment).await?;
320
321 self.import_notes(¬e_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(¬e.metadata().tag()))
407 .map(|(note, _)| note.clone())
408 .collect();
409 let consumable = self.note_screener().get_batch_consumability(¬es_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(¬e.metadata().tag()) || consumable.contains_key(¬e.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 ¬e_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(¬e_info.header, ¬e_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, ¬e_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(¬e_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(¬e.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}