miden-client 0.17.0

Client library that facilitates interaction with the Miden network
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
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::{InputNoteRecord, 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";
pub const NOTE_TRANSPORT_CURSOR_STORE_SETTING: &str = "note_transport_cursor";

/// Settings key for the note-transport backfill bookkeeping: a serialized `Vec<NoteTag>` of the
/// `User`- and `Account`-source tags whose full history has already been fetched up to the global
/// cursor. [`Client::sync_note_transport`] diffs the currently tracked tags against this set to
/// find tags added after the cursor advanced, and backfills only those. Reusing the settings k/v
/// avoids a Store-trait schema change while surviving process restarts.
pub const NOTE_TRANSPORT_COVERED_TAGS_KEY: &str = "note_transport_covered_tags";

/// Settings key for the durable relay outbox: a serialized `Vec<RelayOutboxEntry>` of private notes
/// whose transport delivery has not yet succeeded. [`Client::send_private_note_with_proof`] appends
/// (replacing any entry with the same note id) before relaying; [`Client::flush_relay_outbox`]
/// drains entries that re-send successfully. Reusing the settings k/v avoids a Store-trait schema
/// change while surviving process restarts.
pub const NOTE_TRANSPORT_OUTBOX_KEY: &str = "note_transport_outbox";

/// Client note transport methods.
impl<AUTH> Client<AUTH> {
    /// Check if note transport connection is configured
    pub fn is_note_transport_enabled(&self) -> bool {
        self.note_transport_api.is_some()
    }

    /// Returns the Note Transport client
    ///
    /// Errors if the note transport is not configured.
    pub(crate) fn get_note_transport_api(
        &self,
    ) -> Result<Arc<dyn NoteTransportClient>, NoteTransportError> {
        self.note_transport_api.clone().ok_or(NoteTransportError::Disabled)
    }

    /// Send a note through the note transport network together with its inclusion proof.
    ///
    /// The note will be end-to-end encrypted (unimplemented, currently plaintext) using the
    /// provided recipient's `address` details. The recipient will be able to retrieve this note
    /// through the note's [`NoteTag`].
    ///
    /// The transport carries the proof through [`NoteTransportClient::send_note_with_proof`]. The
    /// network verifies `inclusion_proof` against its node before it stores the note and relays the
    /// exact commitment block to the recipient. The proof exists once the transaction that created
    /// the note is committed and the sender has synced past it; see
    /// [`OutputNoteRecord::inclusion_proof`](crate::store::OutputNoteRecord::inclusion_proof).
    ///
    /// **Durability.** The note and its proof are persisted to the outbox before the transport
    /// call. If the call fails or is interrupted, the entry stays in the outbox and is retried on
    /// the next [`Client::flush_relay_outbox`] (which [`Client::sync_note_transport`] runs), so a
    /// transient transport failure does not drop the note. The receiver dedupes by note id, so a
    /// re-send after a partial success is harmless.
    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 note = TransportNote::from(note);
        let note_id = note.header().id();
        // The address is reserved for end-to-end encryption of the note details:
        // address.key().encrypt(note.details().to_bytes()).
        let _ = address;

        // Persist the payload before the network call so a failed or interrupted send leaves a
        // recoverable record rather than losing the only copy with the call frame. The proof
        // travels with the entry so a retried send relays the same value.
        let entry = RelayOutboxEntry { note, inclusion_proof };
        let mut outbox = self.load_relay_outbox().await?;
        // Replace any existing entry for this note id so the latest payload wins when a
        // still-pending note is re-sent.
        outbox.retain(|e| e.note.header().id() != note_id);
        outbox.push(entry.clone());
        self.save_relay_outbox(outbox).await?;

        entry.relay(api.as_ref()).await?;

        // Relay succeeded — drop the entry. A failed store write here is tolerable: the next flush
        // re-sends and the receiver dedupes by note id, so a stale entry never causes loss.
        let mut outbox = self.load_relay_outbox().await?;
        outbox.retain(|e| e.note.header().id() != note_id);
        self.save_relay_outbox(outbox).await?;

        Ok(())
    }

    /// Re-attempt every relay payload in the durable outbox. Each entry is a private note whose
    /// previous transport delivery failed. Successful re-sends are dropped; failures are kept for
    /// the next call. Every entry is attempted independently, so one persistently-failing note does
    /// not block the others.
    ///
    /// [`Client::sync_note_transport`] runs this automatically and ignores its error, so a relay
    /// failure can't block a sync. Callers driving retries themselves can invoke it directly and
    /// inspect the returned error.
    pub async fn flush_relay_outbox(&self) -> Result<(), ClientError> {
        let api = self.get_note_transport_api()?;

        let entries = self.load_relay_outbox().await?;
        if entries.is_empty() {
            return Ok(());
        }

        // Attempt every entry independently so a single persistently-failing note can't block the
        // rest. The outbox holds only the caller's own failed sends, so it stays small and this is
        // not a meaningful burst.
        let mut remaining = Vec::new();
        let mut last_err: Option<NoteTransportError> = None;

        for entry in entries {
            match entry.relay(api.as_ref()).await {
                Ok(()) => {},
                Err(err) => {
                    tracing::warn!(?err, "relay-outbox entry retry failed; will retry next sync");
                    remaining.push(entry);
                    last_err = Some(err);
                },
            }
        }

        self.save_relay_outbox(remaining).await?;

        if let Some(err) = last_err {
            return Err(err.into());
        }
        Ok(())
    }

    /// Load the durable relay outbox.
    ///
    /// Returns an empty `Vec` if the outbox key is absent. On deserialization failure (schema
    /// mismatch or storage corruption) the entry is dropped and an empty `Vec` is returned —
    /// leaving unreadable bytes in place would block every subsequent relay because each sync would
    /// re-read them.
    async fn load_relay_outbox(&self) -> Result<Vec<RelayOutboxEntry>, ClientError> {
        let bytes = self
            .store
            .get_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_OUTBOX_KEY))
            .await
            .map_err(ClientError::StoreError)?;
        let Some(bytes) = bytes else {
            return Ok(Vec::new());
        };
        match Vec::<RelayOutboxEntry>::read_from_bytes(&bytes) {
            Ok(entries) => Ok(entries),
            Err(err) => {
                tracing::warn!(?err, "dropping unreadable relay outbox; resetting to empty");
                self.store
                    .remove_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_OUTBOX_KEY))
                    .await
                    .map_err(ClientError::StoreError)?;
                Ok(Vec::new())
            },
        }
    }

    /// Persist the relay outbox, removing the key entirely when empty so the settings table doesn't
    /// accumulate empty-vec blobs.
    async fn save_relay_outbox(&self, entries: Vec<RelayOutboxEntry>) -> Result<(), ClientError> {
        let key = String::from(NOTE_TRANSPORT_OUTBOX_KEY);
        if entries.is_empty() {
            self.store
                .remove_setting(SettingScope::Client, key)
                .await
                .map_err(ClientError::StoreError)?;
            return Ok(());
        }
        let bytes = entries.to_bytes();
        self.store
            .set_setting(SettingScope::Client, key, bytes)
            .await
            .map_err(ClientError::StoreError)
    }

    /// The set of tracked tags eligible for history backfill.
    ///
    /// Only `User`- and `Account`-source tags qualify: those are the tags a consumer explicitly
    /// started tracking (via [`Client::add_note_tag`], account import, or address creation) and may
    /// therefore have historical private notes sitting below the global cursor. `Note`-source tags
    /// are created by transport delivery and note import, so backfilling them would re-fetch tags
    /// the fetch path itself just registered; `Subscription` tags are excluded for the same reason.
    async fn backfill_candidate_tags(&self) -> Result<BTreeSet<NoteTag>, ClientError> {
        let tags = self
            .store
            .get_note_tags()
            .await?
            .into_iter()
            .filter(|record| {
                matches!(record.source, NoteTagSource::User | NoteTagSource::Account(_))
            })
            .map(|record| record.tag)
            .collect();
        Ok(tags)
    }

    /// Load the set of tags whose history has already been fetched up to the global cursor.
    ///
    /// Returns an empty set when the key is absent (e.g. a store that predates the feature). On a
    /// deserialization failure the entry is dropped and an empty set is returned: re-treating every
    /// tracked tag as new only triggers a one-off backfill, which dedupes, whereas leaving
    /// unreadable bytes in place would fail every subsequent sync.
    async fn load_covered_tags(&self) -> Result<BTreeSet<NoteTag>, ClientError> {
        let bytes = self
            .store
            .get_setting(SettingScope::Client, String::from(NOTE_TRANSPORT_COVERED_TAGS_KEY))
            .await
            .map_err(ClientError::StoreError)?;
        let Some(bytes) = bytes else {
            return Ok(BTreeSet::new());
        };
        match BTreeSet::<NoteTag>::read_from_bytes(&bytes) {
            Ok(tags) => Ok(tags),
            Err(err) => {
                tracing::warn!(?err, "dropping unreadable covered-tags set; resetting to empty");
                self.store
                    .remove_setting(
                        SettingScope::Client,
                        String::from(NOTE_TRANSPORT_COVERED_TAGS_KEY),
                    )
                    .await
                    .map_err(ClientError::StoreError)?;
                Ok(BTreeSet::new())
            },
        }
    }

    /// Persist the covered-tags set, removing the key entirely when empty so the settings table
    /// doesn't accumulate empty-vec blobs.
    async fn save_covered_tags(&self, tags: &BTreeSet<NoteTag>) -> Result<(), ClientError> {
        let key = String::from(NOTE_TRANSPORT_COVERED_TAGS_KEY);
        if tags.is_empty() {
            self.store
                .remove_setting(SettingScope::Client, key)
                .await
                .map_err(ClientError::StoreError)?;
            return Ok(());
        }
        self.store
            .set_setting(SettingScope::Client, key, tags.to_bytes())
            .await
            .map_err(ClientError::StoreError)
    }
}

impl<AUTH> Client<AUTH>
where
    AUTH: TransactionAuthenticator + Sync + 'static,
{
    /// Per-sync cap on the number of newly tracked tags to backfill. Bounds the burst when many
    /// tags are registered at once (e.g. restoring many accounts or addresses). Deferred tags stay
    /// uncovered and are picked up on subsequent syncs.
    pub const MAX_BACKFILL_TAGS_PER_SYNC: usize = 64;

    /// Safety cap on the per-tag backfill drain. A well-behaved server eventually returns no
    /// forward cursor progress, ending the loop; this bound only guards against a server that
    /// advances the cursor indefinitely without ever returning an empty batch. It is far above any
    /// honest per-tag backlog, so reaching it signals a server bug rather than real history.
    const MAX_BACKFILL_ITERATIONS: usize = 1_000;

    /// Fetch notes for tracked note tags.
    ///
    /// The client will query the configured note transport node for all tracked note tags. To list
    /// tracked tags please use [`Client::get_note_tags`]. To add a new note tag please use
    /// [`Client::add_note_tag`]. Only notes directed at your addresses will be stored and readable
    /// given the use of end-to-end encryption (unimplemented). Fetched notes will be stored into
    /// the client's store.
    ///
    /// An internal pagination mechanism is employed to reduce the number of downloaded notes: this
    /// fetches only notes past the stored cursor. Historical notes for a newly tracked tag are
    /// recovered automatically by [`Client::sync_note_transport`], which backfills each new tag.
    pub async fn fetch_private_notes(&mut self) -> Result<(), ClientError> {
        self.ensure_genesis_in_place().await?;

        let note_tags: Vec<NoteTag> =
            self.store.get_unique_note_tags().await?.into_iter().collect();
        let cursor = self.store.get_note_transport_cursor().await?;

        let mut id_by_commitment = BTreeMap::new();
        let (note_files, new_cursor) =
            self.fetch_transport_notes(cursor, &note_tags, &mut id_by_commitment).await?;

        self.import_notes(&note_files).await?;
        self.store.update_note_transport_cursor(new_cursor).await?;

        Ok(())
    }

    /// Plans the backfill of historical private notes for tags added after the global cursor
    /// advanced.
    ///
    /// The global transport cursor is shared across all tracked tags and only moves forward, so a
    /// tag that starts being tracked late never sees its notes that already sit below the cursor.
    /// This diffs the tracked `User`/`Account` tags (see [`Self::backfill_candidate_tags`]) against
    /// the persisted covered set (see [`NOTE_TRANSPORT_COVERED_TAGS_KEY`]) and drains each newly
    /// tracked tag from the start, fetching only that tag's own history rather than re-scanning
    /// everything. Tags no longer tracked are dropped from the covered set so a later re-add
    /// backfills again instead of resuming from a stale mark. Imports dedupe, so the overlap with
    /// the steady-state stream is harmless.
    ///
    /// At most [`Self::MAX_BACKFILL_TAGS_PER_SYNC`] tags are backfilled per call; any remainder
    /// stays uncovered and is picked up on the next sync.
    ///
    /// Returns the pruned covered set, whether pruning changed it, and the tags to backfill. Reads
    /// only: persisting the covered set is left to the apply phase, which writes it after the
    /// imported notes so a crash re-backfills instead of skipping a tag whose notes were never
    /// written.
    async fn plan_backfill(&self) -> Result<(BTreeSet<NoteTag>, bool, Vec<NoteTag>), ClientError> {
        let candidates = self.backfill_candidate_tags().await?;
        let loaded = self.load_covered_tags().await?;

        // Drop tags no longer tracked. Keeping a removed tag marked covered would make a later
        // re-add skip its backlog, silently missing notes that arrived while it was untracked.
        let covered: BTreeSet<NoteTag> = loaded.intersection(&candidates).copied().collect();
        let pruned = covered.len() != loaded.len();

        let new_tags: Vec<NoteTag> = candidates
            .difference(&covered)
            .copied()
            .take(Self::MAX_BACKFILL_TAGS_PER_SYNC)
            .collect();

        Ok((covered, pruned, new_tags))
    }

    /// Drain a single tag's full history from the transport, paging until the cursor stops
    /// advancing. Uses a local cursor and never touches the global one, so it cannot regress
    /// steady-state progress. Returns the note files from every fetched page, in page order and
    /// none of them written.
    async fn backfill_tag(
        &self,
        tag: NoteTag,
        id_by_commitment: &mut BTreeMap<NoteDetailsCommitment, NoteId>,
    ) -> Result<Vec<NoteFile>, ClientError> {
        let mut note_files = Vec::new();
        let mut cursor = NoteTransportCursor::init();
        for _ in 0..Self::MAX_BACKFILL_ITERATIONS {
            let (page_files, new_cursor) =
                self.fetch_transport_notes(cursor, &[tag], id_by_commitment).await?;
            note_files.extend(page_files);
            // Terminate on any lack of forward progress. A well-behaved server returns `new_cursor
            // == cursor` when there are no new notes for this tag (since `rcursor = max(cursor,
            // max_seq_returned)`); using `<=` also handles implementations that return an `init()`
            // cursor on empty batches (see the in-tree mock transport).

            if new_cursor <= cursor {
                return Ok(note_files);
            }
            cursor = new_cursor;
        }

        Err(ClientError::NoteTransportError(NoteTransportError::PaginationDidNotTerminate(
            Self::MAX_BACKFILL_ITERATIONS,
        )))
    }

    /// Screens the transport-delivered notes carrying a tag derived from a tracked account,
    /// discarding those that no tracked account can consume. Notes carrying any other tag are kept
    /// as delivered.
    async fn screen_transport_notes(
        &self,
        notes: &mut Vec<(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(&note.metadata().tag()))
            .map(|(note, _)| note.clone())
            .collect();
        let consumable = self.note_screener().get_batch_consumability(&notes_to_screen).await?;

        // Discard the notes whose tag match the tracked accounts but are not consumable.
        notes.retain(|(note, _)| {
            !account_tags.contains(&note.metadata().tag()) || consumable.contains_key(&note.id())
        });

        Ok(())
    }

    /// Returns the tracked tags that were registered for an account, i.e. derived from its ID.
    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)
    }

    /// Fetches and returns one batch of notes from the note transport layer for the provided tags
    /// without applying any update to the store.
    ///
    /// The server paginates; this method issues one transport call and returns the note files
    /// together with the new cursor. The returned cursor equals the input cursor when the batch was
    /// empty (i.e. no new notes). Callers that want to drain a tag's full backlog should loop until
    /// `new_cursor == cursor` (see [`Client::backfill_tag`]). Callers that do steady-state polling
    /// (see [`Client::sync_state`] / [`Client::fetch_private_notes`]) should call this once per
    /// tick with the stored cursor.
    ///
    /// Each downloaded note's id is recorded in `id_by_commitment` so the caller can resolve the
    /// written records back to note ids once the final record set is known. Persistence of the
    /// returned cursor is left to the caller so that drain loops can guard against regression of an
    /// already-advanced stored cursor.
    async fn fetch_transport_notes(
        &self,
        cursor: NoteTransportCursor,
        tags: &[NoteTag],
        id_by_commitment: &mut BTreeMap<NoteDetailsCommitment, NoteId>,
    ) -> Result<(Vec<NoteFile>, NoteTransportCursor), ClientError> {
        // Fallback lookback window, in blocks, used only for notes the transport delivered without
        // block information. Scanning back from sync height handles the race where a note is
        // committed on-chain just before the NTL delivers its data. Without it,
        // check_expected_notes would scan from sync_height forward and miss the already-committed
        // note. A transport-provided block is deterministic and always preferred.
        const NOTE_LOOKBACK_BLOCKS: u32 = 20;

        let mut notes = Vec::new();
        // TODO: perhaps we should not need to map received IDs with details commitments, and
        // instead we may allow `InputNoteRecord` to optionally keep NoteIds. Then within
        // `import_note` we could match everything by ID and remove this map check
        let (note_infos, rcursor) =
            self.get_note_transport_api()?.fetch_notes(tags, cursor).await?;
        for note_info in &note_infos {
            // e2ee impl hint: for key in self.store.decryption_keys() try
            // key.decrypt(details_bytes_encrypted)
            //
            // An invalid delivery fails the fetch and the cursor stays on this page.
            let note = rejoin_note(&note_info.header, &note_info.details_bytes)?;
            let tag = note.metadata().tag();
            if !tags.contains(&tag) {
                return Err(NoteTransportError::UnrequestedTag(tag).into());
            }

            // The header carries the attachment-aware (on-chain) note id; the rejoined note has
            // empty attachments and would hash to a different id, so key off the header.
            id_by_commitment.insert(note.details_commitment(), note_info.header.id());

            notes.push((note, note_info.block_hint));
        }

        // Screen the transport-delivered notes to discard the ones that are not relevant to the
        // accounts tracked by the client. Boxed to avoid a `clippy::large_futures` warning, since
        // the sync future is already close to the size limit.
        Box::pin(self.screen_transport_notes(&mut notes)).await?;

        self.drop_notes_processed_locally(&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();
            // Prefer the transport-provided block, falling back to the lookback window when absent.
            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, rcursor))
    }

    /// Fetches the notes the Note Transport Layer holds for the tracked tags.
    ///
    /// Runs the per-tag backfill and fetches a page of notes. This performs no node call and writes
    /// nothing but the relay outbox, so it can run concurrently with the chain fetch. The caller
    /// imports the returned files and then persists the cursor and the covered-tag set.
    ///
    /// Returns empty data when note transport is not configured.
    pub(crate) async fn fetch_note_transport_updates(
        &self,
    ) -> Result<NoteTransportLayerUpdate, ClientError> {
        let mut note_transport_update = NoteTransportLayerUpdate::default();
        if !self.is_note_transport_enabled() {
            return Ok(note_transport_update);
        }

        // Drain any private notes whose previous relay attempt failed. A flush error is logged, not
        // propagated: a failing relay must not block the sync, and the entries stay durable for the
        // next attempt. This is the one write this phase performs; it touches only the outbox
        // setting, which is independent of everything the apply phase writes.
        if let Err(err) = self.flush_relay_outbox().await {
            tracing::warn!(?err, "relay outbox flush failed during sync; entries retained");
        }

        // Recover historical private notes for any tag added after the global cursor advanced. This
        // drains each newly tracked tag from the start, fetching only that tag's own history.
        let (mut covered, pruned, new_tags) = self.plan_backfill().await?;
        let backfilled = !new_tags.is_empty();
        for tag in new_tags {
            note_transport_update
                .note_files
                .extend(self.backfill_tag(tag, &mut note_transport_update.id_by_commitment).await?);
            covered.insert(tag);
        }
        if pruned || backfilled {
            note_transport_update.covered_tags = Some(covered);
        }

        let cursor = self.store.get_note_transport_cursor().await?;
        let note_tags: Vec<NoteTag> =
            self.store.get_unique_note_tags().await?.into_iter().collect();
        let (note_files, new_cursor) = self
            .fetch_transport_notes(cursor, &note_tags, &mut note_transport_update.id_by_commitment)
            .await?;
        note_transport_update.note_files.extend(note_files);
        note_transport_update.cursor = Some(new_cursor);

        Ok(note_transport_update)
    }

    /// Writes everything [`Client::fetch_note_transport_updates`] returned, in three steps:
    ///
    /// 1. Imports the fetched notes, which resolves their on-chain state and stores the records.
    /// 2. Saves the covered-tag set, when the backfill changed it.
    /// 3. Advances the stored note transport cursor, when a page was fetched.
    ///
    /// The notes are written before the covered-tag set and the cursor, so a crash between them
    /// re-fetches instead of skipping notes that were never written.
    ///
    /// Returns the ids of the imported notes and the details commitments of the records written.
    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,
            covered_tags,
            cursor,
        } = update;

        let written = self.import_notes(&note_files).await?;
        let mut imported_ids: Vec<NoteId> = written
            .iter()
            .filter_map(|commitment| id_by_commitment.get(commitment).copied())
            .collect();

        if let Some(covered_tags) = covered_tags {
            self.save_covered_tags(&covered_tags).await?;
        }

        if let Some(cursor) = cursor {
            self.store.update_note_transport_cursor(cursor).await?;
        }

        imported_ids.sort_unstable();
        imported_ids.dedup();

        Ok((imported_ids, written))
    }

    /// Drops deliveries of notes a local transaction is consuming; importing them would fail on the
    /// no-overwrite-while-processing guard.
    async fn drop_notes_processed_locally(
        &self,
        notes: &mut Vec<(Note, Option<BlockNumber>)>,
    ) -> Result<(), ClientError> {
        if notes.is_empty() {
            return Ok(());
        }

        let commitments = notes.iter().map(|(note, _)| note.details_commitment()).collect();
        let processing: BTreeSet<NoteDetailsCommitment> = self
            .get_input_notes(NoteFilter::DetailsCommitments(commitments))
            .await?
            .into_iter()
            .filter(InputNoteRecord::is_processing)
            .map(|record| record.details_commitment())
            .collect();

        if !processing.is_empty() {
            tracing::warn!(?processing, "skipping deliveries of notes being consumed locally");
            notes.retain(|(note, _)| !processing.contains(&note.details_commitment()));
        }
        Ok(())
    }
}

// NOTE TRANSPORT FETCH
// ================================================================================================

/// What the note transport fetch returned, before anything is written.
///
/// Built by [`Client::fetch_note_transport_updates`] and consumed by
/// [`Client::apply_note_transport_update`].
#[derive(Default)]
pub(crate) struct NoteTransportLayerUpdate {
    /// Notes to import, backfill pages first and then the steady-state page.
    note_files: Vec<NoteFile>,
    /// Note ids by details commitment, taken from the note headers the transport returned. Used to
    /// resolve the written records back to ids.
    id_by_commitment: BTreeMap<NoteDetailsCommitment, NoteId>,
    /// Covered-tag set to persist, `None` when it did not change.
    covered_tags: Option<BTreeSet<NoteTag>>,
    /// New global cursor, from the steady-state page. `None` when no page was fetched.
    cursor: Option<NoteTransportCursor>,
}

/// Note transport cursor
///
/// Identifies a position in the note transport service's stored-note sequence.
#[derive(Clone, Copy, Debug, PartialEq, PartialOrd, Eq, Ord)]
pub struct NoteTransportCursor(Option<(u64, u64)>);

impl NoteTransportCursor {
    /// Returns the cursor that starts from the first retained note.
    pub fn init() -> Self {
        Self(None)
    }

    /// Builds a cursor from the nonce and sequence returned by the transport service.
    pub fn from_parts(nonce: u64, sequence: u64) -> Self {
        Self(Some((nonce, sequence)))
    }

    /// Returns the nonce and sequence, or `None` for the initial cursor.
    pub fn parts(&self) -> Option<(u64, u64)> {
        self.0
    }
}

/// The part of a note that the note transport network sends to a recipient.
///
/// The transport sends the original header and details. It does not send note attachments. The
/// constructor verifies that the header commits to the details.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct TransportNote {
    header: NoteHeader,
    details: NoteDetails,
}

impl TransportNote {
    /// Creates a transport note from matching note parts.
    pub fn new(header: NoteHeader, details: NoteDetails) -> Result<Self, NoteTransportError> {
        validate_note_parts(&header, &details)?;
        Ok(Self { header, details })
    }

    /// Returns the note header.
    pub fn header(&self) -> &NoteHeader {
        &self.header
    }

    /// Returns the note details.
    pub fn details(&self) -> &NoteDetails {
        &self.details
    }

    /// Returns the note header and 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 }
    }
}

/// The main transport client trait for sending and receiving private notes.
#[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 {
    /// Sends a note together with its inclusion proof.
    ///
    /// The transport carries the proof to the network. The network verifies `inclusion_proof`
    /// before it stores the note and relays the exact commitment block to the recipient.
    async fn send_note_with_proof(
        &self,
        note: TransportNote,
        inclusion_proof: NoteInclusionProof,
    ) -> Result<(), NoteTransportError>;

    /// Fetches notes for the given tags.
    ///
    /// Downloads notes for the given tags. Returns notes after the provided cursor (pagination),
    /// and an updated cursor.
    async fn fetch_notes(
        &self,
        tag: &[NoteTag],
        cursor: NoteTransportCursor,
    ) -> Result<(Vec<NoteInfo>, NoteTransportCursor), NoteTransportError>;
}

/// Information about a note fetched from the note transport network
#[derive(Debug, Clone)]
pub struct NoteInfo {
    /// Note header.
    pub header: NoteHeader,
    /// Serialized note details.
    pub details_bytes: Vec<u8>,
    /// Block from which the recipient starts scanning for the note's on-chain commitment. This is
    /// either an unverified sender hint or the exact block verified by a proof-aware transport.
    /// `None` applies the recipient's default lookback window.
    pub block_hint: Option<BlockNumber>,
}

impl NoteInfo {
    /// Builds a [`NoteInfo`] without a block hint (`block_hint` is `None`).
    ///
    /// Use the [`NoteInfo::block_hint`] field directly to attach a hint.
    pub fn new(header: NoteHeader, details_bytes: Vec<u8>) -> Self {
        Self { header, details_bytes, block_hint: None }
    }
}

// RELAY OUTBOX
// ================================================================================================

/// A private note whose transport delivery has not yet succeeded.
#[derive(Debug, Clone, PartialEq, Eq)]
struct RelayOutboxEntry {
    note: TransportNote,
    inclusion_proof: NoteInclusionProof,
}

impl RelayOutboxEntry {
    /// Sends the note and its inclusion proof through the transport.
    async fn relay(&self, api: &dyn NoteTransportClient) -> Result<(), NoteTransportError> {
        api.send_note_with_proof(self.note.clone(), self.inclusion_proof.clone()).await
    }
}

// SERIALIZATION
// ================================================================================================

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 RelayOutboxEntry {
    fn write_into<W: ByteWriter>(&self, target: &mut W) {
        self.note.write_into(target);
        self.inclusion_proof.write_into(target);
    }
}

impl Deserializable for RelayOutboxEntry {
    fn read_from<R: ByteReader>(source: &mut R) -> Result<Self, DeserializationError> {
        let note = TransportNote::read_from(source)?;
        let inclusion_proof = NoteInclusionProof::read_from(source)?;
        Ok(Self { note, inclusion_proof })
    }
}

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)?;
    // The transport wire format only carries `NoteHeader` + serialized `NoteDetails`, not the
    // attachments collection. We rejoin with empty attachments; this matches the original note only
    // when it had no attachments in the first place.
    let partial_metadata = *header.metadata().partial_metadata();
    Ok(Note::new(
        details.assets().clone(),
        partial_metadata,
        details.recipient().clone(),
    ))
}

/// Checks that the note header commits to the supplied details.
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 miden_protocol::account::AccountId;
    use miden_protocol::asset::FungibleAsset;
    use miden_protocol::crypto::merkle::SparseMerklePath;
    use miden_protocol::note::NoteType;
    use miden_protocol::testing::account_id::{
        ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET,
        ACCOUNT_ID_REGULAR_PUBLIC_ACCOUNT_IMMUTABLE_CODE,
        ACCOUNT_ID_SENDER,
    };
    use miden_standards::note::P2idNote;
    use rand::SeedableRng;
    use rand_chacha::ChaCha20Rng;

    use super::*;
    use crate::rng::draw_word;

    #[test]
    fn relay_outbox_entry_round_trips() {
        let sender = AccountId::try_from(ACCOUNT_ID_SENDER).unwrap();
        let target = AccountId::try_from(ACCOUNT_ID_REGULAR_PUBLIC_ACCOUNT_IMMUTABLE_CODE).unwrap();
        let faucet = AccountId::try_from(ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET).unwrap();
        let mut rng = ChaCha20Rng::seed_from_u64(0);
        let note: Note = P2idNote::builder()
            .sender(sender)
            .target(target)
            .asset(FungibleAsset::new(faucet, 100).unwrap())
            .note_type(NoteType::Private)
            .serial_number(draw_word(&mut rng))
            .build()
            .unwrap()
            .into();

        let inclusion_proof =
            NoteInclusionProof::new(BlockNumber::from(7), 3, SparseMerklePath::default()).unwrap();
        let entry = RelayOutboxEntry {
            note: TransportNote::from(note),
            inclusion_proof,
        };

        assert_eq!(RelayOutboxEntry::read_from_bytes(&entry.to_bytes()).unwrap(), entry);
    }
}