pub mod errors;
pub mod generated;
#[cfg(feature = "tonic")]
pub mod grpc;
use alloc::boxed::Box;
use alloc::collections::{BTreeMap, BTreeSet};
use alloc::string::String;
use alloc::sync::Arc;
use alloc::vec::Vec;
use miden_protocol::address::Address;
use miden_protocol::block::BlockNumber;
use miden_protocol::note::{
Note,
NoteDetails,
NoteDetailsCommitment,
NoteHeader,
NoteId,
NoteInclusionProof,
NoteTag,
};
use miden_protocol::utils::serde::Serializable;
use miden_tx::auth::TransactionAuthenticator;
use miden_tx::utils::serde::{
ByteReader,
ByteWriter,
Deserializable,
DeserializationError,
SliceReader,
};
pub use self::errors::NoteTransportError;
use crate::note::{NoteFile, NoteSyncHint};
use crate::store::{NoteFilter, SettingScope};
use crate::sync::NoteTagSource;
use crate::{Client, ClientError};
pub const NOTE_TRANSPORT_MAINNET_ENDPOINT: &str = "https://transport.mainnet.miden.io";
pub const NOTE_TRANSPORT_TESTNET_ENDPOINT: &str = "https://transport.miden.io";
pub const NOTE_TRANSPORT_DEVNET_ENDPOINT: &str = "https://transport.devnet.miden.io";
#[deprecated(since = "0.17.1", note = "note transport stores a cursor for each tag")]
pub const NOTE_TRANSPORT_CURSOR_STORE_SETTING: &str = "note_transport_cursor";
pub const NOTE_TRANSPORT_CURSORS_KEY: &str = "note_transport_cursors";
type NoteTransportCursors = BTreeMap<NoteTag, NoteTransportCursor>;
const MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST: usize = 128;
const MAX_NOTE_TRANSPORT_PAGES_PER_GROUP: usize = 32;
#[deprecated(since = "0.17.1", note = "note transport no longer keeps per-tag backfill state")]
pub const NOTE_TRANSPORT_COVERED_TAGS_KEY: &str = "note_transport_covered_tags";
impl<AUTH> Client<AUTH> {
pub const MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST: usize = MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST;
#[deprecated(
since = "0.17.1",
note = "use MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST; account tags are no longer limited"
)]
pub const MAX_ACCOUNT_TAGS: usize = Self::MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST;
pub fn is_note_transport_enabled(&self) -> bool {
self.note_transport_api.is_some()
}
pub(crate) fn get_note_transport_api(
&self,
) -> Result<Arc<dyn NoteTransportClient>, NoteTransportError> {
self.note_transport_api.clone().ok_or(NoteTransportError::Disabled)
}
pub async fn send_private_note_with_proof(
&mut self,
note: Note,
address: &Address,
inclusion_proof: NoteInclusionProof,
) -> Result<(), ClientError> {
let api = self.get_note_transport_api()?;
let _ = address;
api.send_note_with_proof(TransportNote::from(note), inclusion_proof).await?;
Ok(())
}
async fn load_note_transport_cursors(&self) -> Result<NoteTransportCursors, ClientError> {
let bytes = self
.store
.get_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_CURSORS_KEY))
.await
.map_err(ClientError::StoreError)?;
let Some(bytes) = bytes else {
return Ok(BTreeMap::new());
};
match NoteTransportCursors::read_from_bytes(&bytes) {
Ok(cursors) => Ok(cursors),
Err(err) => {
tracing::warn!(?err, "resetting unreadable note transport cursors");
Ok(BTreeMap::new())
},
}
}
async fn save_note_transport_cursors(
&self,
cursors: &NoteTransportCursors,
) -> Result<(), ClientError> {
let key = String::from(NOTE_TRANSPORT_CURSORS_KEY);
if cursors.is_empty() {
self.store.remove_setting(SettingScope::Client, key).await?;
} else {
self.store.set_setting(SettingScope::Client, key, cursors.to_bytes()).await?;
}
Ok(())
}
}
impl<AUTH> Client<AUTH>
where
AUTH: TransactionAuthenticator + Sync + 'static,
{
#[deprecated(since = "0.17.1", note = "note transport no longer performs per-tag backfills")]
pub const MAX_BACKFILL_TAGS_PER_SYNC: usize = 64;
pub async fn fetch_private_notes(&mut self) -> Result<(), ClientError> {
self.ensure_genesis_in_place().await?;
let mut update = self.fetch_transport_notes_in_chunks().await?;
let fetch_error = update.fetch_error.take();
self.apply_note_transport_update(update).await?;
if let Some(error) = fetch_error {
return Err(error);
}
Ok(())
}
async fn screen_transport_notes(
&self,
notes: &mut Vec<(NoteId, Note, Option<BlockNumber>)>,
) -> Result<(), ClientError> {
let account_tags = self.tracked_account_tags().await?;
let notes_to_screen: Vec<Note> = notes
.iter()
.filter(|(_, note, _)| account_tags.contains(¬e.metadata().tag()))
.map(|(_, note, _)| note.clone())
.collect();
let consumable = self.note_screener().get_batch_consumability(¬es_to_screen).await?;
notes.retain(|(_, note, _)| {
!account_tags.contains(¬e.metadata().tag()) || consumable.contains_key(¬e.id())
});
Ok(())
}
async fn tracked_account_tags(&self) -> Result<BTreeSet<NoteTag>, ClientError> {
let tags = self
.store
.get_note_tags()
.await?
.into_iter()
.filter(|record| matches!(record.source, NoteTagSource::Account(_)))
.map(|record| record.tag)
.collect();
Ok(tags)
}
async fn fetch_transport_notes_in_chunks(
&self,
) -> Result<NoteTransportLayerUpdate, ClientError> {
let api = self.get_note_transport_api()?;
let tags: Vec<_> = self.store.get_unique_note_tags().await?.into_iter().collect();
let stored = self.load_note_transport_cursors().await?;
let mut cursors: NoteTransportCursors = tags
.iter()
.filter_map(|tag| stored.get(tag).map(|cursor| (*tag, *cursor)))
.collect();
let mut update = NoteTransportLayerUpdate::default();
let mut notes = Vec::new();
for (start, group) in transport_request_groups(&tags, &cursors) {
let mut cursor = start;
for _ in 0..MAX_NOTE_TRANSPORT_PAGES_PER_GROUP {
let page = match api
.fetch_notes_page(&group, cursor)
.await
.and_then(|page| validate_transport_page(page, &group, cursor))
{
Ok(page) => page,
Err(error) => {
update.fetch_error.get_or_insert(error.into());
break;
},
};
for (id, note, block_hint) in page.notes {
update.id_by_commitment.insert(note.details_commitment(), id);
notes.push((id, note, block_hint));
}
cursor = page.cursor;
for tag in &group {
let position = cursors.entry(*tag).or_insert(cursor);
*position = advance_transport_cursor(*position, cursor);
}
if !page.has_more {
break;
}
}
}
update.note_files = self.prepare_transport_notes(notes).await?;
update.cursors = Some(cursors);
Ok(update)
}
async fn prepare_transport_notes(
&self,
mut notes: Vec<(NoteId, Note, Option<BlockNumber>)>,
) -> Result<Vec<NoteFile>, ClientError> {
const NOTE_LOOKBACK_BLOCKS: u32 = 20;
if notes.is_empty() {
return Ok(Vec::new());
}
self.drop_notes_resolved_locally(&mut notes).await?;
Box::pin(self.screen_transport_notes(&mut notes)).await?;
let sync_height = self.get_sync_height().await?;
let fallback_after_block_num =
BlockNumber::from(sync_height.as_u32().saturating_sub(NOTE_LOOKBACK_BLOCKS));
let mut note_files = Vec::with_capacity(notes.len());
for (_, note, block_hint) in notes {
let tag = note.metadata().tag();
let after_block_num = block_hint.unwrap_or(fallback_after_block_num);
note_files.push(NoteFile::ExpectedNote {
details: note.into(),
sync_hint: NoteSyncHint::new(after_block_num, tag),
});
}
Ok(note_files)
}
pub(crate) async fn fetch_note_transport_updates(
&self,
) -> Result<NoteTransportLayerUpdate, ClientError> {
if !self.is_note_transport_enabled() {
return Ok(NoteTransportLayerUpdate::default());
}
self.fetch_transport_notes_in_chunks().await
}
pub(crate) async fn apply_note_transport_update(
&mut self,
update: NoteTransportLayerUpdate,
) -> Result<(Vec<NoteId>, Vec<NoteDetailsCommitment>), ClientError> {
let NoteTransportLayerUpdate {
note_files,
id_by_commitment,
cursors,
fetch_error,
} = update;
let written = self.import_notes(¬e_files).await?;
let mut imported_ids: Vec<NoteId> = written
.iter()
.filter_map(|commitment| id_by_commitment.get(commitment).copied())
.collect();
if let Some(cursors) = cursors {
self.save_note_transport_cursors(&cursors).await?;
}
if let Some(error) = fetch_error {
tracing::warn!(?error, "note transport fetch failed; saved successful pages for retry");
}
imported_ids.sort_unstable();
imported_ids.dedup();
Ok((imported_ids, written))
}
async fn drop_notes_resolved_locally(
&self,
notes: &mut Vec<(NoteId, Note, Option<BlockNumber>)>,
) -> Result<(), ClientError> {
let commitments = notes.iter().map(|(_, note, _)| note.details_commitment()).collect();
let records: BTreeMap<_, _> = self
.get_input_notes(NoteFilter::DetailsCommitments(commitments))
.await?
.into_iter()
.map(|record| (record.details_commitment(), record))
.collect();
notes.retain(|(id, note, _)| {
let Some(record) = records.get(¬e.details_commitment()) else {
return true;
};
if record.is_processing() {
tracing::warn!(%id, "skipping delivery of a note being consumed locally");
return false;
}
!((record.is_committed() || record.is_consumed()) && record.id() == Some(*id))
});
Ok(())
}
}
#[derive(Default)]
pub(crate) struct NoteTransportLayerUpdate {
note_files: Vec<NoteFile>,
id_by_commitment: BTreeMap<NoteDetailsCommitment, NoteId>,
cursors: Option<NoteTransportCursors>,
pub(crate) fetch_error: Option<ClientError>,
}
fn transport_request_groups(
tags: &[NoteTag],
cursors: &NoteTransportCursors,
) -> Vec<(NoteTransportCursor, Vec<NoteTag>)> {
let mut by_nonce = BTreeMap::<Option<u64>, Vec<(NoteTransportCursor, NoteTag)>>::new();
for tag in tags {
let cursor = cursors.get(tag).copied().unwrap_or_else(NoteTransportCursor::init);
by_nonce
.entry(cursor.parts().map(|(nonce, _)| nonce))
.or_default()
.push((cursor, *tag));
}
let mut groups = Vec::new();
for mut entries in by_nonce.into_values() {
entries.sort_unstable();
for chunk in entries.chunks(MAX_NOTE_TAGS_PER_TRANSPORT_REQUEST) {
let start = chunk[0].0;
groups.push((start, chunk.iter().map(|(_, tag)| *tag).collect()));
}
}
groups
}
fn advance_transport_cursor(
current: NoteTransportCursor,
page_cursor: NoteTransportCursor,
) -> NoteTransportCursor {
match (current.parts(), page_cursor.parts()) {
(Some((nonce, sequence)), Some((page_nonce, page_sequence))) if nonce == page_nonce => {
NoteTransportCursor::from_parts(nonce, sequence.max(page_sequence))
},
_ => page_cursor,
}
}
struct ValidatedTransportPage {
notes: Vec<(NoteId, Note, Option<BlockNumber>)>,
cursor: NoteTransportCursor,
has_more: bool,
}
fn validate_transport_page(
page: NoteTransportPage,
tags: &[NoteTag],
request_cursor: NoteTransportCursor,
) -> Result<ValidatedTransportPage, NoteTransportError> {
let Some((nonce, sequence)) = page.cursor.parts() else {
return Err(NoteTransportError::Network(String::from("fetch response has no cursor")));
};
if request_cursor.parts().is_some_and(|(request_nonce, request_sequence)| {
nonce == request_nonce
&& (sequence < request_sequence
|| (!page.notes.is_empty() && sequence == request_sequence))
}) || (!page.notes.is_empty() && sequence == 0)
|| (page.has_more && page.notes.is_empty())
{
return Err(NoteTransportError::Network(String::from(
"fetch response has invalid pagination progress",
)));
}
let mut notes = Vec::with_capacity(page.notes.len());
for info in page.notes {
let note = rejoin_note(&info.header, &info.details_bytes)?;
if !tags.contains(¬e.metadata().tag()) {
return Err(NoteTransportError::UnrequestedTag(note.metadata().tag()));
}
notes.push((info.header.id(), note, info.block_hint));
}
Ok(ValidatedTransportPage {
notes,
cursor: page.cursor,
has_more: page.has_more,
})
}
#[derive(Clone, Copy, Debug, PartialEq, PartialOrd, Eq, Ord)]
pub struct NoteTransportCursor(Option<(u64, u64)>);
impl NoteTransportCursor {
pub fn init() -> Self {
Self(None)
}
pub fn from_parts(nonce: u64, sequence: u64) -> Self {
Self(Some((nonce, sequence)))
}
pub fn parts(&self) -> Option<(u64, u64)> {
self.0
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct TransportNote {
header: NoteHeader,
details: NoteDetails,
}
impl TransportNote {
pub fn new(header: NoteHeader, details: NoteDetails) -> Result<Self, NoteTransportError> {
validate_note_parts(&header, &details)?;
Ok(Self { header, details })
}
pub fn header(&self) -> &NoteHeader {
&self.header
}
pub fn details(&self) -> &NoteDetails {
&self.details
}
pub fn into_parts(self) -> (NoteHeader, NoteDetails) {
(self.header, self.details)
}
}
impl From<Note> for TransportNote {
fn from(note: Note) -> Self {
let header = *note.header();
let details = NoteDetails::from(note);
Self { header, details }
}
}
#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
pub trait NoteTransportClient: Send + Sync {
async fn send_note_with_proof(
&self,
note: TransportNote,
inclusion_proof: NoteInclusionProof,
) -> Result<(), NoteTransportError>;
async fn fetch_notes(
&self,
tag: &[NoteTag],
cursor: NoteTransportCursor,
) -> Result<(Vec<NoteInfo>, NoteTransportCursor), NoteTransportError>;
async fn fetch_notes_page(
&self,
tags: &[NoteTag],
cursor: NoteTransportCursor,
) -> Result<NoteTransportPage, NoteTransportError> {
let (notes, cursor) = self.fetch_notes(tags, cursor).await?;
let has_more = !notes.is_empty();
Ok(NoteTransportPage { notes, cursor, has_more })
}
}
pub struct NoteTransportPage {
pub notes: Vec<NoteInfo>,
pub cursor: NoteTransportCursor,
pub has_more: bool,
}
#[derive(Debug, Clone)]
pub struct NoteInfo {
pub header: NoteHeader,
pub details_bytes: Vec<u8>,
pub block_hint: Option<BlockNumber>,
}
impl NoteInfo {
pub fn new(header: NoteHeader, details_bytes: Vec<u8>) -> Self {
Self { header, details_bytes, block_hint: None }
}
}
impl Serializable for TransportNote {
fn write_into<W: ByteWriter>(&self, target: &mut W) {
self.header.write_into(target);
self.details.to_bytes().write_into(target);
}
}
impl Deserializable for TransportNote {
fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
let header = NoteHeader::read_from(source)?;
let details_bytes = Vec::<u8>::read_from(source)?;
let details = NoteDetails::read_from_bytes(&details_bytes)?;
Self::new(header, details)
.map_err(|error| DeserializationError::InvalidValue(format!("{error}")))
}
}
impl Serializable for NoteInfo {
fn write_into<W: ByteWriter>(&self, target: &mut W) {
self.header.write_into(target);
self.details_bytes.write_into(target);
self.block_hint.write_into(target);
}
}
impl Deserializable for NoteInfo {
fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
let header = NoteHeader::read_from(source)?;
let details_bytes = Vec::<u8>::read_from(source)?;
let block_hint = Option::<BlockNumber>::read_from(source)?;
Ok(NoteInfo { header, details_bytes, block_hint })
}
}
impl Serializable for NoteTransportCursor {
fn write_into<W: ByteWriter>(&self, target: &mut W) {
self.0.write_into(target);
}
}
impl Deserializable for NoteTransportCursor {
fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
Ok(Self(Option::<(u64, u64)>::read_from(source)?))
}
}
fn rejoin_note(header: &NoteHeader, details_bytes: &[u8]) -> Result<Note, NoteTransportError> {
let mut reader = SliceReader::new(details_bytes);
let details = NoteDetails::read_from(&mut reader)?;
validate_note_parts(header, &details)?;
let partial_metadata = *header.metadata().partial_metadata();
Ok(Note::new(
details.assets().clone(),
partial_metadata,
details.recipient().clone(),
))
}
pub(crate) fn validate_note_parts(
header: &NoteHeader,
details: &NoteDetails,
) -> Result<(), NoteTransportError> {
let header_commitment = header.details_commitment();
let details_commitment = details.commitment();
if header_commitment != details_commitment {
return Err(NoteTransportError::NoteDetailsMismatch {
header: header_commitment,
details: details_commitment,
});
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn transport_rejects_invalid_pagination() {
let cursor = NoteTransportCursor::from_parts(1, 10);
for (returned, has_more) in [
(NoteTransportCursor::init(), false),
(NoteTransportCursor::from_parts(1, 9), false),
(cursor, true),
] {
let page = NoteTransportPage {
notes: Vec::new(),
cursor: returned,
has_more,
};
assert!(validate_transport_page(page, &[], cursor).is_err());
}
}
}