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#[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>;
48const MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST: usize = 128;
50const MAX_NOTE_TRANSPORT_PAGES_PER_GROUP: usize = 32;
52
53#[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
57impl<AUTH> Client<AUTH> {
59 pub const MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST: usize = MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST;
61
62 #[deprecated(
66 since = "0.17.1",
67 note = "use MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST; account tags are no longer limited"
68 )]
69 pub const MAX_ACCOUNT_TAGS: usize = Self::MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST;
70
71 pub fn is_note_transport_enabled(&self) -> bool {
73 self.note_transport_api.is_some()
74 }
75
76 pub(crate) fn get_note_transport_api(
80 &self,
81 ) -> Result<Arc<dyn NoteTransportClient>, NoteTransportError> {
82 self.note_transport_api.clone().ok_or(NoteTransportError::Disabled)
83 }
84
85 pub async fn send_private_note_with_proof(
104 &mut self,
105 note: Note,
106 address: &Address,
107 inclusion_proof: NoteInclusionProof,
108 ) -> Result<(), ClientError> {
109 let api = self.get_note_transport_api()?;
110
111 let _ = address;
114
115 api.send_note_with_proof(TransportNote::from(note), inclusion_proof).await?;
116
117 Ok(())
118 }
119
120 async fn load_note_transport_cursors(&self) -> Result<NoteTransportCursors, ClientError> {
125 let bytes = self
126 .store
127 .get_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_CURSORS_KEY))
128 .await
129 .map_err(ClientError::StoreError)?;
130 let Some(bytes) = bytes else {
131 return Ok(BTreeMap::new());
132 };
133
134 match NoteTransportCursors::read_from_bytes(&bytes) {
135 Ok(cursors) => Ok(cursors),
136 Err(err) => {
137 tracing::warn!(?err, "resetting unreadable note transport cursors");
138 Ok(BTreeMap::new())
139 },
140 }
141 }
142
143 async fn save_note_transport_cursors(
145 &self,
146 cursors: &NoteTransportCursors,
147 ) -> Result<(), ClientError> {
148 let key = String::from(NOTE_TRANSPORT_CURSORS_KEY);
149 if cursors.is_empty() {
150 self.store.remove_setting(SettingScope::Client, key).await?;
151 } else {
152 self.store.set_setting(SettingScope::Client, key, cursors.to_bytes()).await?;
153 }
154 Ok(())
155 }
156}
157
158impl<AUTH> Client<AUTH>
159where
160 AUTH: TransactionAuthenticator + Sync + 'static,
161{
162 #[deprecated(since = "0.17.1", note = "note transport no longer performs per-tag backfills")]
164 pub const MAX_BACKFILL_TAGS_PER_SYNC: usize = 64;
165
166 pub async fn fetch_private_notes(&mut self) -> Result<(), ClientError> {
179 self.ensure_genesis_in_place().await?;
180
181 let mut update = self.fetch_transport_notes_in_chunks().await?;
182 let fetch_error = update.fetch_error.take();
183 self.apply_note_transport_update(update).await?;
184 if let Some(error) = fetch_error {
185 return Err(error);
186 }
187
188 Ok(())
189 }
190
191 async fn screen_transport_notes(
195 &self,
196 notes: &mut Vec<(NoteId, Note, Option<BlockNumber>)>,
197 ) -> Result<(), ClientError> {
198 let account_tags = self.tracked_account_tags().await?;
199
200 let notes_to_screen: Vec<Note> = notes
201 .iter()
202 .filter(|(_, note, _)| account_tags.contains(¬e.metadata().tag()))
203 .map(|(_, note, _)| note.clone())
204 .collect();
205 let consumable = self.note_screener().get_batch_consumability(¬es_to_screen).await?;
206
207 notes.retain(|(_, note, _)| {
209 !account_tags.contains(¬e.metadata().tag()) || consumable.contains_key(¬e.id())
210 });
211
212 Ok(())
213 }
214
215 async fn tracked_account_tags(&self) -> Result<BTreeSet<NoteTag>, ClientError> {
217 let tags = self
218 .store
219 .get_note_tags()
220 .await?
221 .into_iter()
222 .filter(|record| matches!(record.source, NoteTagSource::Account(_)))
223 .map(|record| record.tag)
224 .collect();
225 Ok(tags)
226 }
227
228 async fn fetch_transport_notes_in_chunks(
236 &self,
237 ) -> Result<NoteTransportLayerUpdate, ClientError> {
238 let api = self.get_note_transport_api()?;
239 let tags: Vec<_> = self.store.get_unique_note_tags().await?.into_iter().collect();
240 let stored = self.load_note_transport_cursors().await?;
241 let mut cursors: NoteTransportCursors = tags
242 .iter()
243 .filter_map(|tag| stored.get(tag).map(|cursor| (*tag, *cursor)))
244 .collect();
245
246 let mut update = NoteTransportLayerUpdate::default();
247 let mut notes = Vec::new();
248 for (start, group) in transport_request_groups(&tags, &cursors) {
249 let mut cursor = start;
250 for _ in 0..MAX_NOTE_TRANSPORT_PAGES_PER_GROUP {
251 let page = match api
252 .fetch_notes_page(&group, cursor)
253 .await
254 .and_then(|page| validate_transport_page(page, &group, cursor))
255 {
256 Ok(page) => page,
257 Err(error) => {
258 update.fetch_error.get_or_insert(error.into());
259 break;
260 },
261 };
262 for (id, note, block_hint) in page.notes {
263 update.id_by_commitment.insert(note.details_commitment(), id);
264 notes.push((id, note, block_hint));
265 }
266 cursor = page.cursor;
267 for tag in &group {
268 let position = cursors.entry(*tag).or_insert(cursor);
269 *position = advance_transport_cursor(*position, cursor);
270 }
271 if !page.has_more {
272 break;
273 }
274 }
275 }
276
277 update.note_files = self.prepare_transport_notes(notes).await?;
278 update.cursors = Some(cursors);
279 Ok(update)
280 }
281
282 async fn prepare_transport_notes(
284 &self,
285 mut notes: Vec<(NoteId, Note, Option<BlockNumber>)>,
286 ) -> Result<Vec<NoteFile>, ClientError> {
287 const NOTE_LOOKBACK_BLOCKS: u32 = 20;
293
294 if notes.is_empty() {
295 return Ok(Vec::new());
296 }
297
298 self.drop_notes_resolved_locally(&mut notes).await?;
299
300 Box::pin(self.screen_transport_notes(&mut notes)).await?;
304
305 let sync_height = self.get_sync_height().await?;
306 let fallback_after_block_num =
307 BlockNumber::from(sync_height.as_u32().saturating_sub(NOTE_LOOKBACK_BLOCKS));
308
309 let mut note_files = Vec::with_capacity(notes.len());
310 for (_, note, block_hint) in notes {
311 let tag = note.metadata().tag();
312 let after_block_num = block_hint.unwrap_or(fallback_after_block_num);
314 note_files.push(NoteFile::ExpectedNote {
315 details: note.into(),
316 sync_hint: NoteSyncHint::new(after_block_num, tag),
317 });
318 }
319
320 Ok(note_files)
321 }
322
323 pub(crate) async fn fetch_note_transport_updates(
332 &self,
333 ) -> Result<NoteTransportLayerUpdate, ClientError> {
334 if !self.is_note_transport_enabled() {
335 return Ok(NoteTransportLayerUpdate::default());
336 }
337
338 self.fetch_transport_notes_in_chunks().await
339 }
340
341 pub(crate) async fn apply_note_transport_update(
351 &mut self,
352 update: NoteTransportLayerUpdate,
353 ) -> Result<(Vec<NoteId>, Vec<NoteDetailsCommitment>), ClientError> {
354 let NoteTransportLayerUpdate {
355 note_files,
356 id_by_commitment,
357 cursors,
358 fetch_error,
359 } = update;
360
361 let written = self.import_notes(¬e_files).await?;
362 let mut imported_ids: Vec<NoteId> = written
363 .iter()
364 .filter_map(|commitment| id_by_commitment.get(commitment).copied())
365 .collect();
366
367 if let Some(cursors) = cursors {
368 self.save_note_transport_cursors(&cursors).await?;
369 }
370 if let Some(error) = fetch_error {
371 tracing::warn!(?error, "note transport fetch failed; saved successful pages for retry");
372 }
373
374 imported_ids.sort_unstable();
375 imported_ids.dedup();
376
377 Ok((imported_ids, written))
378 }
379
380 async fn drop_notes_resolved_locally(
388 &self,
389 notes: &mut Vec<(NoteId, Note, Option<BlockNumber>)>,
390 ) -> Result<(), ClientError> {
391 let commitments = notes.iter().map(|(_, note, _)| note.details_commitment()).collect();
392 let records: BTreeMap<_, _> = self
393 .get_input_notes(NoteFilter::DetailsCommitments(commitments))
394 .await?
395 .into_iter()
396 .map(|record| (record.details_commitment(), record))
397 .collect();
398 notes.retain(|(id, note, _)| {
399 let Some(record) = records.get(¬e.details_commitment()) else {
400 return true;
401 };
402 if record.is_processing() {
403 tracing::warn!(%id, "skipping delivery of a note being consumed locally");
404 return false;
405 }
406 !((record.is_committed() || record.is_consumed()) && record.id() == Some(*id))
407 });
408 Ok(())
409 }
410}
411
412#[derive(Default)]
420pub(crate) struct NoteTransportLayerUpdate {
421 note_files: Vec<NoteFile>,
423 id_by_commitment: BTreeMap<NoteDetailsCommitment, NoteId>,
426 cursors: Option<NoteTransportCursors>,
428 pub(crate) fetch_error: Option<ClientError>,
430}
431
432fn transport_request_groups(
439 tags: &[NoteTag],
440 cursors: &NoteTransportCursors,
441) -> Vec<(NoteTransportCursor, Vec<NoteTag>)> {
442 let mut by_nonce = BTreeMap::<Option<u64>, Vec<(NoteTransportCursor, NoteTag)>>::new();
443 for tag in tags {
444 let cursor = cursors.get(tag).copied().unwrap_or_else(NoteTransportCursor::init);
445 by_nonce
446 .entry(cursor.parts().map(|(nonce, _)| nonce))
447 .or_default()
448 .push((cursor, *tag));
449 }
450
451 let mut groups = Vec::new();
452 for mut entries in by_nonce.into_values() {
453 entries.sort_unstable();
454 for chunk in entries.chunks(MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST) {
455 let start = chunk[0].0;
456 groups.push((start, chunk.iter().map(|(_, tag)| *tag).collect()));
457 }
458 }
459 groups
460}
461
462fn advance_transport_cursor(
467 current: NoteTransportCursor,
468 page_cursor: NoteTransportCursor,
469) -> NoteTransportCursor {
470 match (current.parts(), page_cursor.parts()) {
471 (Some((nonce, sequence)), Some((page_nonce, page_sequence))) if nonce == page_nonce => {
472 NoteTransportCursor::from_parts(nonce, sequence.max(page_sequence))
473 },
474 _ => page_cursor,
475 }
476}
477
478struct ValidatedTransportPage {
479 notes: Vec<(NoteId, Note, Option<BlockNumber>)>,
480 cursor: NoteTransportCursor,
481 has_more: bool,
482}
483
484fn validate_transport_page(
486 page: NoteTransportPage,
487 tags: &[NoteTag],
488 request_cursor: NoteTransportCursor,
489) -> Result<ValidatedTransportPage, NoteTransportError> {
490 let Some((nonce, sequence)) = page.cursor.parts() else {
491 return Err(NoteTransportError::Network(String::from("fetch response has no cursor")));
492 };
493 if request_cursor.parts().is_some_and(|(request_nonce, request_sequence)| {
494 nonce == request_nonce
495 && (sequence < request_sequence
496 || (!page.notes.is_empty() && sequence == request_sequence))
497 }) || (!page.notes.is_empty() && sequence == 0)
498 || (page.has_more && page.notes.is_empty())
499 {
500 return Err(NoteTransportError::Network(String::from(
501 "fetch response has invalid pagination progress",
502 )));
503 }
504 let mut notes = Vec::with_capacity(page.notes.len());
505 for info in page.notes {
506 let note = rejoin_note(&info.header, &info.details_bytes)?;
507 if !tags.contains(¬e.metadata().tag()) {
508 return Err(NoteTransportError::UnrequestedTag(note.metadata().tag()));
509 }
510 notes.push((info.header.id(), note, info.block_hint));
512 }
513 Ok(ValidatedTransportPage {
514 notes,
515 cursor: page.cursor,
516 has_more: page.has_more,
517 })
518}
519
520#[derive(Clone, Copy, Debug, PartialEq, PartialOrd, Eq, Ord)]
527pub struct NoteTransportCursor(Option<(u64, u64)>);
528
529impl NoteTransportCursor {
530 pub fn init() -> Self {
532 Self(None)
533 }
534
535 pub fn from_parts(nonce: u64, sequence: u64) -> Self {
537 Self(Some((nonce, sequence)))
538 }
539
540 pub fn parts(&self) -> Option<(u64, u64)> {
542 self.0
543 }
544}
545
546#[derive(Clone, Debug, PartialEq, Eq)]
551pub struct TransportNote {
552 header: NoteHeader,
553 details: NoteDetails,
554}
555
556impl TransportNote {
557 pub fn new(header: NoteHeader, details: NoteDetails) -> Result<Self, NoteTransportError> {
559 validate_note_parts(&header, &details)?;
560 Ok(Self { header, details })
561 }
562
563 pub fn header(&self) -> &NoteHeader {
565 &self.header
566 }
567
568 pub fn details(&self) -> &NoteDetails {
570 &self.details
571 }
572
573 pub fn into_parts(self) -> (NoteHeader, NoteDetails) {
575 (self.header, self.details)
576 }
577}
578
579impl From<Note> for TransportNote {
580 fn from(note: Note) -> Self {
581 let header = *note.header();
582 let details = NoteDetails::from(note);
583 Self { header, details }
584 }
585}
586
587#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
589#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
590pub trait NoteTransportClient: Send + Sync {
591 async fn send_note_with_proof(
596 &self,
597 note: TransportNote,
598 inclusion_proof: NoteInclusionProof,
599 ) -> Result<(), NoteTransportError>;
600
601 async fn fetch_notes(
606 &self,
607 tag: &[NoteTag],
608 cursor: NoteTransportCursor,
609 ) -> Result<(Vec<NoteInfo>, NoteTransportCursor), NoteTransportError>;
610
611 async fn fetch_notes_page(
615 &self,
616 tags: &[NoteTag],
617 cursor: NoteTransportCursor,
618 ) -> Result<NoteTransportPage, NoteTransportError> {
619 let (notes, cursor) = self.fetch_notes(tags, cursor).await?;
620 let has_more = !notes.is_empty();
621 Ok(NoteTransportPage { notes, cursor, has_more })
622 }
623}
624
625pub struct NoteTransportPage {
627 pub notes: Vec<NoteInfo>,
629 pub cursor: NoteTransportCursor,
631 pub has_more: bool,
633}
634
635#[derive(Debug, Clone)]
637pub struct NoteInfo {
638 pub header: NoteHeader,
640 pub details_bytes: Vec<u8>,
642 pub block_hint: Option<BlockNumber>,
646}
647
648impl NoteInfo {
649 pub fn new(header: NoteHeader, details_bytes: Vec<u8>) -> Self {
653 Self { header, details_bytes, block_hint: None }
654 }
655}
656
657impl Serializable for TransportNote {
661 fn write_into<W: ByteWriter>(&self, target: &mut W) {
662 self.header.write_into(target);
663 self.details.to_bytes().write_into(target);
664 }
665}
666
667impl Deserializable for TransportNote {
668 fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
669 let header = NoteHeader::read_from(source)?;
670 let details_bytes = Vec::<u8>::read_from(source)?;
671 let details = NoteDetails::read_from_bytes(&details_bytes)?;
672 Self::new(header, details)
673 .map_err(|error| DeserializationError::InvalidValue(format!("{error}")))
674 }
675}
676
677impl Serializable for NoteInfo {
678 fn write_into<W: ByteWriter>(&self, target: &mut W) {
679 self.header.write_into(target);
680 self.details_bytes.write_into(target);
681 self.block_hint.write_into(target);
682 }
683}
684
685impl Deserializable for NoteInfo {
686 fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
687 let header = NoteHeader::read_from(source)?;
688 let details_bytes = Vec::<u8>::read_from(source)?;
689 let block_hint = Option::<BlockNumber>::read_from(source)?;
690 Ok(NoteInfo { header, details_bytes, block_hint })
691 }
692}
693
694impl Serializable for NoteTransportCursor {
695 fn write_into<W: ByteWriter>(&self, target: &mut W) {
696 self.0.write_into(target);
697 }
698}
699
700impl Deserializable for NoteTransportCursor {
701 fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
702 Ok(Self(Option::<(u64, u64)>::read_from(source)?))
703 }
704}
705
706fn rejoin_note(header: &NoteHeader, details_bytes: &[u8]) -> Result<Note, NoteTransportError> {
707 let mut reader = SliceReader::new(details_bytes);
708 let details = NoteDetails::read_from(&mut reader)?;
709 validate_note_parts(header, &details)?;
710 let partial_metadata = *header.metadata().partial_metadata();
714 Ok(Note::new(
715 details.assets().clone(),
716 partial_metadata,
717 details.recipient().clone(),
718 ))
719}
720
721pub(crate) fn validate_note_parts(
723 header: &NoteHeader,
724 details: &NoteDetails,
725) -> Result<(), NoteTransportError> {
726 let header_commitment = header.details_commitment();
727 let details_commitment = details.commitment();
728 if header_commitment != details_commitment {
729 return Err(NoteTransportError::NoteDetailsMismatch {
730 header: header_commitment,
731 details: details_commitment,
732 });
733 }
734
735 Ok(())
736}
737
738#[cfg(test)]
739mod tests {
740 use super::*;
741
742 #[test]
743 fn transport_rejects_invalid_pagination() {
744 let cursor = NoteTransportCursor::from_parts(1, 10);
745 for (returned, has_more) in [
746 (NoteTransportCursor::init(), false),
747 (NoteTransportCursor::from_parts(1, 9), false),
748 (cursor, true),
749 ] {
750 let page = NoteTransportPage {
751 notes: Vec::new(),
752 cursor: returned,
753 has_more,
754 };
755 assert!(validate_transport_page(page, &[], cursor).is_err());
756 }
757 }
758}