1use alloc::boxed::Box;
2use alloc::collections::{BTreeMap, BTreeSet};
3use alloc::sync::Arc;
4use alloc::vec::Vec;
5use core::cmp::Ordering;
6
7use async_trait::async_trait;
8use futures::{StreamExt, TryStreamExt};
9use miden_protocol::Word;
10use miden_protocol::account::{Account, AccountHeader, AccountId, StorageSlotType};
11use miden_protocol::block::account_tree::AccountIdKey;
12use miden_protocol::block::{BlockHeader, BlockNumber, ValidatorConfig};
13use miden_protocol::crypto::merkle::MerklePath;
14use miden_protocol::crypto::merkle::mmr::{InOrderIndex, MmrDelta, PartialMmr};
15use miden_protocol::note::{NoteId, NoteTag, Nullifier};
16use miden_protocol::protocol_config::ProtocolConfig;
17use tracing::info;
18
19use super::state_sync_update::{TransactionUpdateTracker, build_account_patch};
20use super::{
21 AccountUpdates,
22 NoteObserver,
23 PartialBlockchainUpdates,
24 PublicAccountUpdate,
25 StateSyncUpdate,
26};
27use crate::ClientError;
28use crate::note::{NoteConsumption, NoteUpdateTracker};
29use crate::rpc::domain::account::{
30 AccountDetails,
31 AccountProof,
32 AccountStorageMapDetails,
33 GetAccountRequest,
34 StorageMapFetch,
35 VaultFetch,
36};
37use crate::rpc::domain::note::{CommittedNote, FetchedNote, ResolvedSyncNotesBlock, SyncedNote};
38use crate::rpc::domain::sync::{ChainMmrInfo, SyncTarget};
39use crate::rpc::domain::transaction::TransactionRecord as RpcTransactionRecord;
40use crate::rpc::{AccountStateAt, NodeRpcClient, NoteContentFetch, RpcError};
41use crate::store::input_note_states::UnverifiedNoteState;
42use crate::store::{InputNoteRecord, OutputNoteRecord, StoreError};
43use crate::transaction::TransactionRecord;
44
45const MAX_CONCURRENT_ACCOUNT_FETCHES: usize = 4;
47
48enum PublicAccountSync {
53 Apply(Box<PublicAccountUpdate>),
55 Superseded,
57 Ignore,
59}
60
61struct FetchedSyncData {
67 mmr_delta: MmrDelta,
69 chain_tip_header: BlockHeader,
71 note_blocks: Vec<ResolvedSyncNotesBlock>,
74 transactions: Vec<RpcTransactionRecord>,
76 protocol_config: Option<ProtocolConfig>,
79}
80
81struct RecoverableConsumedNote {
86 nullifier: Nullifier,
87 consumer: AccountId,
88 block_num: BlockNumber,
89}
90
91struct RelevantNoteBlock {
96 block_header: BlockHeader,
97 mmr_path: MerklePath,
98 observer_requires_block: bool,
99}
100
101#[derive(Default)]
103struct NoteBlockRelevance {
104 has_client_note: bool,
105 observer_requires_block: bool,
106}
107
108impl NoteBlockRelevance {
109 fn is_relevant(&self) -> bool {
110 self.has_client_note || self.observer_requires_block
111 }
112}
113
114pub struct StateSyncInput {
130 pub accounts: Vec<AccountHeader>,
132 pub note_tags: BTreeSet<NoteTag>,
134 pub input_notes: Vec<InputNoteRecord>,
136 pub output_notes: Vec<OutputNoteRecord>,
142 pub uncommitted_transactions: Vec<TransactionRecord>,
144}
145
146#[allow(clippy::large_enum_variant)]
151pub enum NoteUpdateAction {
152 Commit(CommittedNote),
155 Insert(InputNoteRecord),
157 Discard,
159}
160
161#[async_trait(?Send)]
162pub trait OnNoteReceived {
163 async fn on_note_received(
174 &self,
175 committed_note: CommittedNote,
176 public_note: Option<InputNoteRecord>,
177 ) -> Result<NoteUpdateAction, ClientError>;
178}
179#[derive(Clone)]
186pub struct StateSync {
187 rpc_api: Arc<dyn NodeRpcClient>,
189 note_screener: Arc<dyn OnNoteReceived>,
192 note_observers: Vec<Arc<dyn NoteObserver>>,
195 tx_discard_delta: Option<u32>,
198 sync_nullifiers: bool,
201 validator_config: ValidatorConfig,
204}
205
206impl StateSync {
207 pub fn new(
219 rpc_api: Arc<dyn NodeRpcClient>,
220 note_screener: Arc<dyn OnNoteReceived>,
221 tx_discard_delta: Option<u32>,
222 validator_config: ValidatorConfig,
223 ) -> Self {
224 Self {
225 rpc_api,
226 note_screener,
227 note_observers: Vec::new(),
228 tx_discard_delta,
229 sync_nullifiers: true,
230 validator_config,
231 }
232 }
233
234 #[must_use]
238 pub fn with_note_observer(mut self, observer: Arc<dyn NoteObserver>) -> Self {
239 self.note_observers.push(observer);
240 self
241 }
242
243 pub fn disable_nullifier_sync(&mut self) {
248 self.sync_nullifiers = false;
249 }
250
251 pub fn enable_nullifier_sync(&mut self) {
253 self.sync_nullifiers = true;
254 }
255
256 pub(crate) async fn run_apply_hooks(
261 &self,
262 state_sync_update: &StateSyncUpdate,
263 ) -> Result<(), ClientError> {
264 for observer in &self.note_observers {
265 crate::errors::log_observer_failure(
266 observer.name(),
267 "NoteObserver::apply",
268 observer.apply(state_sync_update).await,
269 );
270 }
271 Ok(())
272 }
273
274 pub async fn sync_state(
297 &self,
298 current_partial_mmr: &mut PartialMmr,
299 input: StateSyncInput,
300 ) -> Result<StateSyncUpdate, ClientError> {
301 let block_num = block_num_from_forest(current_partial_mmr)?;
302
303 let mut chain_sync_data = self.fetch_state(block_num, input).await?;
304 self.derive_state_updates(&mut chain_sync_data).await?;
305 self.fetch_nullifiers(&mut chain_sync_data).await?;
306
307 let mut working_mmr = current_partial_mmr.clone();
309 let update = Self::build_update(chain_sync_data, &mut working_mmr)?;
310 *current_partial_mmr = working_mmr;
311
312 Ok(update)
313 }
314
315 pub async fn fetch_state(
324 &self,
325 block_from: BlockNumber,
326 input: StateSyncInput,
327 ) -> Result<ChainSyncData, ClientError> {
328 let StateSyncInput {
329 accounts,
330 note_tags,
331 input_notes,
332 output_notes,
333 uncommitted_transactions,
334 } = input;
335
336 let note_tags = Arc::new(note_tags);
337 let account_ids: Vec<AccountId> = accounts.iter().map(AccountHeader::id).collect();
338
339 let note_updates = NoteUpdateTracker::new(input_notes, output_notes);
340 let transaction_updates = TransactionUpdateTracker::new(uncommitted_transactions);
341 let mut account_updates = AccountUpdates::default();
342
343 let Some(sync_data) = self.fetch_sync_data(block_from, &account_ids, ¬e_tags).await?
344 else {
345 return Ok(ChainSyncData {
347 block_from,
348 advance: None,
349 superseded_states: Vec::new(),
350 note_updates,
351 transaction_updates,
352 account_updates,
353 });
354 };
355
356 let FetchedSyncData {
357 mmr_delta,
358 chain_tip_header,
359 note_blocks,
360 transactions,
361 protocol_config,
362 } = sync_data;
363
364 let new_commitments = derive_account_commitments(&transactions);
365 let superseded_states = self
366 .account_state_sync(
367 &mut account_updates,
368 &accounts,
369 &new_commitments,
370 block_from,
371 &chain_tip_header,
372 )
373 .await?;
374
375 Ok(ChainSyncData {
376 block_from,
377 advance: Some(ChainAdvance {
378 chain_tip_header,
379 mmr_delta,
380 note_blocks_awaiting_screening: note_blocks,
381 transactions,
382 relevant_note_blocks: Vec::new(),
383 protocol_config,
384 }),
385 superseded_states,
386 note_updates,
387 transaction_updates,
388 account_updates,
389 })
390 }
391
392 pub async fn derive_state_updates(
398 &self,
399 chain_sync_data: &mut ChainSyncData,
400 ) -> Result<(), ClientError> {
401 let ChainSyncData {
402 advance,
403 superseded_states,
404 note_updates,
405 transaction_updates,
406 ..
407 } = chain_sync_data;
408
409 let Some(advance) = advance.as_mut() else {
410 return Ok(());
411 };
412
413 for superseded_state in core::mem::take(superseded_states) {
415 transaction_updates.apply_superseded_account_state(superseded_state);
416 }
417
418 advance.relevant_note_blocks = self
419 .screen_note_blocks(
420 core::mem::take(&mut advance.note_blocks_awaiting_screening),
421 note_updates,
422 )
423 .await?;
424
425 self.apply_transactions_and_nullifiers(
426 &advance.chain_tip_header,
427 &advance.transactions,
428 note_updates,
429 transaction_updates,
430 )?;
431
432 self.recover_consumed_public_notes(note_updates, &advance.transactions).await?;
433
434 Ok(())
435 }
436
437 pub fn build_update(
443 chain_sync_data: ChainSyncData,
444 partial_mmr: &mut PartialMmr,
445 ) -> Result<StateSyncUpdate, ClientError> {
446 let ChainSyncData {
447 block_from,
448 advance,
449 note_updates,
450 transaction_updates,
451 account_updates,
452 ..
453 } = chain_sync_data;
454
455 let mut partial_blockchain_updates = PartialBlockchainUpdates::default();
456
457 let Some(ChainAdvance {
458 chain_tip_header,
459 mmr_delta,
460 note_blocks_awaiting_screening,
461 relevant_note_blocks,
462 protocol_config,
463 ..
464 }) = advance
465 else {
466 return Ok(StateSyncUpdate::from_parts(
468 block_from,
469 partial_blockchain_updates,
470 note_updates,
471 transaction_updates,
472 account_updates,
473 None,
474 ));
475 };
476 if !note_blocks_awaiting_screening.is_empty() {
478 return Err(ClientError::UnscreenedNoteBlocks);
479 }
480
481 let chain_tip = chain_tip_header.block_num();
482
483 Self::advance_mmr(
484 mmr_delta,
485 &chain_tip_header,
486 partial_mmr,
487 &mut partial_blockchain_updates,
488 )?;
489
490 let blocks_with_unspent_notes: BTreeSet<BlockNumber> =
491 note_updates.unspent_input_note_block_numbers().collect();
492
493 Self::validate_and_track_note_blocks(
494 relevant_note_blocks,
495 &blocks_with_unspent_notes,
496 partial_mmr,
497 &mut partial_blockchain_updates,
498 )?;
499
500 Ok(StateSyncUpdate::from_parts(
501 chain_tip,
502 partial_blockchain_updates,
503 note_updates,
504 transaction_updates,
505 account_updates,
506 protocol_config,
507 ))
508 }
509
510 pub async fn fetch_nullifiers(
519 &self,
520 chain_sync_data: &mut ChainSyncData,
521 ) -> Result<(), ClientError> {
522 if !self.sync_nullifiers {
523 return Ok(());
524 }
525
526 let Some(chain_tip) = chain_sync_data
527 .advance
528 .as_ref()
529 .map(|advance| advance.chain_tip_header.block_num())
530 else {
531 return Ok(());
532 };
533
534 self.nullifiers_state_sync(
535 &mut chain_sync_data.note_updates,
536 &mut chain_sync_data.transaction_updates,
537 chain_tip,
538 chain_sync_data.block_from,
539 )
540 .await
541 }
542
543 async fn recover_consumed_public_notes(
548 &self,
549 note_updates: &mut NoteUpdateTracker,
550 transactions: &[RpcTransactionRecord],
551 ) -> Result<(), ClientError> {
552 let mut recoverable_consumed_notes: BTreeMap<NoteId, RecoverableConsumedNote> =
553 BTreeMap::new();
554 for tx in transactions {
555 for (nullifier, note_id) in tx.trusted_consumed_note_refs() {
556 recoverable_consumed_notes.insert(
557 note_id,
558 RecoverableConsumedNote {
559 nullifier,
560 consumer: tx.transaction_header.account_id(),
561 block_num: tx.block_num,
562 },
563 );
564 }
565 }
566 recoverable_consumed_notes.retain(|note_id, _| !note_updates.tracks_note(*note_id));
569
570 let note_ids: Vec<NoteId> = recoverable_consumed_notes.keys().copied().collect();
571 if note_ids.is_empty() {
572 return Ok(());
573 }
574
575 for fetched in self.rpc_api.get_notes_by_id(¬e_ids).await? {
576 match fetched {
577 FetchedNote::Public(note, _) => {
578 let Some(reference) = recoverable_consumed_notes.get(¬e.id()) else {
579 continue;
580 };
581 if note.nullifier() != reference.nullifier {
584 return Err(RpcError::InvalidResponse(format!(
585 "node returned note {} whose nullifier doesn't match the consumed reference",
586 note.id()
587 ))
588 .into());
589 }
590 note_updates.insert_consumed_public_note(
591 note,
592 reference.consumer,
593 reference.block_num,
594 )?;
595 },
596 FetchedNote::Private(note_id, ..) => {
597 return Err(RpcError::InvalidResponse(format!(
598 "node returned private note {note_id} for a public consumed-note reference"
599 ))
600 .into());
601 },
602 }
603 }
604
605 Ok(())
606 }
607
608 async fn fetch_sync_data(
618 &self,
619 current_block_num: BlockNumber,
620 account_ids: &[AccountId],
621 note_tags: &Arc<BTreeSet<NoteTag>>,
622 ) -> Result<Option<FetchedSyncData>, ClientError> {
623 let chain_mmr_info = self
625 .rpc_api
626 .sync_chain_mmr(current_block_num, SyncTarget::CommittedChainTip)
627 .await?;
628 let chain_tip = chain_mmr_info.block_to;
629
630 Self::validate_chain_mmr_response(
632 &chain_mmr_info,
633 current_block_num,
634 &self.validator_config,
635 )?;
636
637 if chain_tip == current_block_num {
639 info!(block_num = %current_block_num, "Already at chain tip, nothing to sync.");
640 return Ok(None);
641 }
642
643 info!(
644 block_from = %current_block_num,
645 block_to = %chain_tip,
646 "Syncing state.",
647 );
648
649 let (note_blocks, transaction_records) = futures::try_join!(
652 self.fetch_note_blocks(current_block_num + 1, chain_tip, note_tags),
653 self.fetch_transactions(current_block_num + 1, chain_tip, account_ids),
654 )?;
655
656 Self::validate_note_blocks_range(¬e_blocks, current_block_num, chain_tip)?;
658
659 let note_count: usize = note_blocks.iter().map(|b| b.notes.len()).sum();
660 info!(
661 blocks_with_notes = note_blocks.len(),
662 notes = note_count,
663 "Fetched note sync data.",
664 );
665
666 Self::validate_transaction_records_range(
667 &transaction_records,
668 current_block_num,
669 chain_tip,
670 )?;
671
672 Ok(Some(FetchedSyncData {
673 mmr_delta: chain_mmr_info.mmr_delta,
674 chain_tip_header: chain_mmr_info.block_header,
675 note_blocks,
676 transactions: transaction_records,
677 protocol_config: chain_mmr_info.protocol_config,
678 }))
679 }
680
681 async fn fetch_note_blocks(
687 &self,
688 block_from: BlockNumber,
689 block_to: BlockNumber,
690 note_tags: &BTreeSet<NoteTag>,
691 ) -> Result<Vec<ResolvedSyncNotesBlock>, ClientError> {
692 if note_tags.is_empty() {
693 return Ok(Vec::new());
694 }
695
696 self.rpc_api
697 .sync_notes_with_content(
698 block_from,
699 block_to,
700 note_tags,
701 NoteContentFetch::PublicDetailsAndAttachments,
702 )
703 .await
704 .map_err(ClientError::RpcError)
705 }
706
707 async fn fetch_transactions(
709 &self,
710 block_from: BlockNumber,
711 block_to: BlockNumber,
712 account_ids: &[AccountId],
713 ) -> Result<Vec<RpcTransactionRecord>, ClientError> {
714 if account_ids.is_empty() {
715 return Ok(Vec::new());
716 }
717
718 self.rpc_api
719 .sync_transactions(block_from, block_to, account_ids.to_vec())
720 .await
721 .map_err(ClientError::RpcError)
722 }
723
724 fn validate_chain_mmr_response(
727 chain_mmr_info: &ChainMmrInfo,
728 current_block_num: BlockNumber,
729 validator_config: &ValidatorConfig,
730 ) -> Result<(), ClientError> {
731 if chain_mmr_info.block_header.block_num() != chain_mmr_info.block_to {
732 return Err(ClientError::ChainValidationError(format!(
733 "sync_chain_mmr block_header.block_num ({}) does not match block_to ({})",
734 chain_mmr_info.block_header.block_num(),
735 chain_mmr_info.block_to
736 )));
737 }
738 if chain_mmr_info.block_from != current_block_num {
739 return Err(ClientError::ChainValidationError(format!(
740 "sync_chain_mmr block_from mismatch: expected {current_block_num}, got {}",
741 chain_mmr_info.block_from
742 )));
743 }
744 if chain_mmr_info.block_to < current_block_num {
745 return Err(ClientError::ChainValidationError(format!(
746 "sync_chain_mmr block_to ({}) is behind current block {current_block_num}",
747 chain_mmr_info.block_to
748 )));
749 }
750
751 chain_mmr_info
753 .block_signatures
754 .verify_against(chain_mmr_info.block_header.commitment(), validator_config)
755 .map_err(|err| {
756 ClientError::ChainValidationError(format!(
757 "chain tip block header {} does not carry valid validator signatures: {err}",
758 chain_mmr_info.block_header.block_num()
759 ))
760 })?;
761
762 Ok(())
763 }
764
765 fn validate_note_blocks_range(
768 note_blocks: &[ResolvedSyncNotesBlock],
769 current_block_num: BlockNumber,
770 chain_tip: BlockNumber,
771 ) -> Result<(), ClientError> {
772 for block in note_blocks {
773 let block_num = block.block_header.block_num();
774 if block_num <= current_block_num || block_num > chain_tip {
775 return Err(ClientError::ChainValidationError(format!(
776 "sync_notes returned block {block_num} outside requested range ({current_block_num}, {chain_tip}]"
777 )));
778 }
779 }
780 Ok(())
781 }
782
783 fn validate_transaction_records_range(
786 records: &[RpcTransactionRecord],
787 current_block_num: BlockNumber,
788 chain_tip: BlockNumber,
789 ) -> Result<(), ClientError> {
790 for record in records {
791 let block_num = record.block_num;
792 if block_num <= current_block_num || block_num > chain_tip {
793 return Err(ClientError::ChainValidationError(format!(
794 "sync_transactions returned block {block_num} outside requested range ({current_block_num}, {chain_tip}]"
795 )));
796 }
797 }
798 Ok(())
799 }
800
801 fn advance_mmr(
808 mmr_delta: MmrDelta,
809 chain_tip_header: &BlockHeader,
810 current_partial_mmr: &mut PartialMmr,
811 partial_blockchain_updates: &mut PartialBlockchainUpdates,
812 ) -> Result<(), ClientError> {
813 let mut new_authentication_nodes =
814 current_partial_mmr.apply(mmr_delta).map_err(StoreError::MmrError)?;
815 let new_peaks = current_partial_mmr.peaks();
816
817 let peaks_commitment = new_peaks.hash_peaks();
821 if peaks_commitment != chain_tip_header.chain_commitment() {
822 return Err(ClientError::ChainValidationError(format!(
823 "MMR peaks commitment is {} and does not match block header chain commitment {}",
824 peaks_commitment.to_hex(),
825 chain_tip_header.chain_commitment().to_hex()
826 )));
827 }
828
829 partial_blockchain_updates.new_peaks = new_peaks;
830
831 new_authentication_nodes.append(
835 &mut current_partial_mmr
836 .add(chain_tip_header.commitment(), false)
837 .map_err(StoreError::MmrError)?,
838 );
839
840 partial_blockchain_updates.insert(chain_tip_header.clone(), false);
841 partial_blockchain_updates.extend_authentication_nodes(new_authentication_nodes);
842
843 Ok(())
844 }
845
846 async fn screen_note_blocks(
854 &self,
855 note_blocks: Vec<ResolvedSyncNotesBlock>,
856 note_updates: &mut NoteUpdateTracker,
857 ) -> Result<Vec<RelevantNoteBlock>, ClientError> {
858 let mut relevant_blocks = Vec::new();
859
860 for block in note_blocks {
861 let relevance =
862 self.note_state_sync(note_updates, block.notes, &block.block_header).await?;
863
864 if relevance.is_relevant() {
865 relevant_blocks.push(RelevantNoteBlock {
866 block_header: block.block_header,
867 mmr_path: block.mmr_path,
868 observer_requires_block: relevance.observer_requires_block,
869 });
870 }
871 }
872
873 Ok(relevant_blocks)
874 }
875
876 fn validate_and_track_note_blocks(
886 relevant_blocks: Vec<RelevantNoteBlock>,
887 blocks_with_unspent_notes: &BTreeSet<BlockNumber>,
888 partial_mmr: &mut PartialMmr,
889 partial_blockchain_updates: &mut PartialBlockchainUpdates,
890 ) -> Result<(), ClientError> {
891 let nodes_before: BTreeSet<InOrderIndex> = partial_mmr.nodes().map(|(k, _)| *k).collect();
892
893 for RelevantNoteBlock {
894 block_header,
895 mmr_path,
896 observer_requires_block,
897 } in relevant_blocks
898 {
899 let block_pos = block_header.block_num().as_usize();
900 let was_tracked = partial_mmr.is_tracked(block_pos);
901
902 partial_mmr
905 .track(block_pos, block_header.commitment(), &mmr_path)
906 .map_err(StoreError::MmrError)?;
907
908 if observer_requires_block
909 || blocks_with_unspent_notes.contains(&block_header.block_num())
910 {
911 partial_blockchain_updates.insert(block_header, true);
912 } else if !was_tracked {
913 partial_mmr.untrack(block_pos);
914 }
915 }
916
917 partial_blockchain_updates.extend_authentication_nodes(
919 partial_mmr
920 .nodes()
921 .filter(|(index, _)| !nodes_before.contains(index))
922 .map(|(index, value)| (*index, *value)),
923 );
924
925 Ok(())
926 }
927
928 fn apply_transactions_and_nullifiers(
932 &self,
933 chain_tip_header: &BlockHeader,
934 transactions: &[RpcTransactionRecord],
935 note_updates: &mut NoteUpdateTracker,
936 transaction_updates: &mut TransactionUpdateTracker,
937 ) -> Result<(), ClientError> {
938 note_updates.extend_nullifiers(compute_ordered_nullifiers(transactions));
939
940 for record in transactions {
941 transaction_updates
942 .apply_transaction_inclusion(record, u64::from(chain_tip_header.timestamp())); }
944 transaction_updates
945 .apply_sync_height_update(chain_tip_header.block_num(), self.tx_discard_delta);
946
947 for transaction in transactions {
948 note_updates.apply_output_note_inclusion_proofs(&transaction.output_notes)?;
952
953 Self::mark_erased_notes_as_consumed(note_updates, transaction);
955 }
956
957 Ok(())
958 }
959
960 fn mark_erased_notes_as_consumed(
966 note_updates: &mut NoteUpdateTracker,
967 transaction: &RpcTransactionRecord,
968 ) {
969 for note_header in &transaction.erased_output_notes {
970 let _ = note_updates.mark_erased_note_as_consumed(note_header, transaction.block_num);
972 }
973 }
974
975 async fn account_state_sync(
988 &self,
989 account_updates: &mut AccountUpdates,
990 accounts: &[AccountHeader],
991 account_commitment_updates: &[(AccountId, Word)],
992 block_from: BlockNumber,
993 chain_tip_header: &BlockHeader,
994 ) -> Result<Vec<Word>, ClientError> {
995 let (public_accounts, private_accounts): (Vec<_>, Vec<_>) =
998 accounts.iter().partition(|header| !header.id().is_private());
999
1000 let superseded_states = self
1001 .sync_public_accounts(
1002 account_updates,
1003 account_commitment_updates,
1004 &public_accounts,
1005 block_from,
1006 chain_tip_header,
1007 )
1008 .await?;
1009
1010 let diverging_private_accounts: Vec<&AccountHeader> = private_accounts
1013 .into_iter()
1014 .filter(|header| {
1015 let local_commitment = header.to_commitment();
1016 account_commitment_updates
1017 .iter()
1018 .any(|(id, digest)| *id == header.id() && *digest != local_commitment)
1019 })
1020 .collect();
1021
1022 let proven_commitments: Vec<Option<Word>> =
1023 futures::stream::iter(diverging_private_accounts.iter().map(|header| {
1024 self.verify_private_account_mismatch(
1025 header.id(),
1026 header.to_commitment(),
1027 chain_tip_header,
1028 )
1029 }))
1030 .buffered(MAX_CONCURRENT_ACCOUNT_FETCHES)
1031 .try_collect()
1032 .await?;
1033
1034 let mismatched_private_accounts: Vec<(AccountId, Word)> = diverging_private_accounts
1035 .iter()
1036 .zip(proven_commitments)
1037 .filter_map(|(header, proven_commitment)| {
1038 proven_commitment.map(|commitment| (header.id(), commitment))
1039 })
1040 .collect();
1041
1042 account_updates.extend(AccountUpdates::new(Vec::new(), mismatched_private_accounts));
1043
1044 Ok(superseded_states)
1045 }
1046
1047 async fn verify_private_account_mismatch(
1058 &self,
1059 account_id: AccountId,
1060 local_commitment: Word,
1061 chain_tip_header: &BlockHeader,
1062 ) -> Result<Option<Word>, ClientError> {
1063 let chain_tip = chain_tip_header.block_num();
1064 let (proof_block_num, proof) = self
1065 .rpc_api
1066 .get_account(account_id, GetAccountRequest::new().at(AccountStateAt::Block(chain_tip)))
1067 .await?;
1068
1069 if proof_block_num != chain_tip {
1070 return Err(ClientError::ChainValidationError(format!(
1071 "get_account returned a proof at block {proof_block_num}, expected chain tip {chain_tip}"
1072 )));
1073 }
1074
1075 let (witness, _) = proof.into_parts();
1076 let witness_id = witness.id();
1077 let proven_commitment = witness.state_commitment();
1078 if witness.into_proof().compute_root() != chain_tip_header.account_root() {
1081 return Err(ClientError::ChainValidationError(format!(
1082 "account witness for {account_id} does not verify against the chain tip account root"
1083 )));
1084 }
1085
1086 if witness_id != account_id
1089 || proven_commitment == Word::empty()
1090 || proven_commitment == local_commitment
1091 {
1092 return Ok(None);
1093 }
1094
1095 Ok(Some(proven_commitment))
1096 }
1097
1098 async fn sync_public_accounts(
1107 &self,
1108 account_updates: &mut AccountUpdates,
1109 commitment_updates: &[(AccountId, Word)],
1110 current_public_accounts: &[&AccountHeader],
1111 block_from: BlockNumber,
1112 chain_tip_header: &BlockHeader,
1113 ) -> Result<Vec<Word>, ClientError> {
1114 let local_headers: BTreeMap<AccountId, &AccountHeader> =
1115 current_public_accounts.iter().map(|header| (header.id(), *header)).collect();
1116
1117 let diverging_accounts: Vec<(AccountId, &AccountHeader)> = commitment_updates
1120 .iter()
1121 .filter_map(|(id, commitment)| {
1122 let local_header = local_headers.get(id).copied()?;
1123 (local_header.to_commitment() != *commitment).then_some((*id, local_header))
1124 })
1125 .collect();
1126
1127 let synced_accounts: Vec<PublicAccountSync> =
1130 futures::stream::iter(diverging_accounts.iter().map(|(id, local_header)| {
1131 self.sync_public_account(*id, local_header, block_from, chain_tip_header)
1132 }))
1133 .buffered(MAX_CONCURRENT_ACCOUNT_FETCHES)
1134 .try_collect()
1135 .await?;
1136
1137 let mut superseded_states = Vec::new();
1139 for ((_, local_header), synced_account) in diverging_accounts.iter().zip(synced_accounts) {
1140 match synced_account {
1141 PublicAccountSync::Apply(public_update) => {
1142 account_updates.extend(AccountUpdates::new(vec![*public_update], Vec::new()));
1143 },
1144 PublicAccountSync::Superseded => {
1145 superseded_states.push(local_header.to_commitment());
1146 },
1147 PublicAccountSync::Ignore => {},
1148 }
1149 }
1150
1151 Ok(superseded_states)
1152 }
1153
1154 async fn sync_public_account(
1163 &self,
1164 account_id: AccountId,
1165 local_header: &AccountHeader,
1166 block_from: BlockNumber,
1167 chain_tip_header: &BlockHeader,
1168 ) -> Result<PublicAccountSync, ClientError> {
1169 let target_block_num = chain_tip_header.block_num();
1170
1171 let (proof_block_num, proof) = self
1174 .rpc_api
1175 .get_account(
1176 account_id,
1177 GetAccountRequest::new()
1178 .at(AccountStateAt::Block(target_block_num))
1179 .with_storage(StorageMapFetch::All)
1180 .with_vault(VaultFetch::Always),
1181 )
1182 .await
1183 .map_err(ClientError::RpcError)?;
1184
1185 let details =
1186 Self::validate_account_proof(proof, proof_block_num, account_id, chain_tip_header)?;
1187
1188 match details
1189 .header
1190 .nonce()
1191 .as_canonical_u64()
1192 .cmp(&local_header.nonce().as_canonical_u64())
1193 {
1194 Ordering::Less => return Ok(PublicAccountSync::Ignore),
1197 Ordering::Equal => return Ok(PublicAccountSync::Superseded),
1199 Ordering::Greater => {},
1201 }
1202
1203 let vault_oversized = details.vault_details.too_many_assets;
1204 let any_map_oversized = details
1205 .storage_details
1206 .map_details
1207 .iter()
1208 .any(AccountStorageMapDetails::is_limit_exceeded);
1209
1210 let public_update = if vault_oversized || any_map_oversized {
1214 self.build_patch_update(account_id, &details, block_from, proof_block_num)
1216 .await?
1217 } else {
1218 let account = Account::try_from(&details).map_err(ClientError::RpcError)?;
1220 PublicAccountUpdate::Full(account)
1221 };
1222
1223 Ok(PublicAccountSync::Apply(Box::new(public_update)))
1224 }
1225
1226 fn validate_account_proof(
1239 proof: AccountProof,
1240 proof_block_num: BlockNumber,
1241 account_id: AccountId,
1242 chain_tip_header: &BlockHeader,
1243 ) -> Result<AccountDetails, ClientError> {
1244 let target_block_num = chain_tip_header.block_num();
1245
1246 if proof_block_num != target_block_num {
1247 return Err(ClientError::ChainValidationError(format!(
1248 "get_account returned block {proof_block_num} but {target_block_num} was requested"
1249 )));
1250 }
1251
1252 let (witness, details) = proof.into_parts();
1253
1254 if witness.id() != account_id {
1256 return Err(ClientError::ChainValidationError(format!(
1257 "get_account returned account {} but {account_id} was requested",
1258 witness.id()
1259 )));
1260 }
1261
1262 let account_key = AccountIdKey::from(account_id).as_word();
1263 let state_commitment = witness.state_commitment();
1264 witness
1265 .into_proof()
1266 .verify_presence(&account_key, &state_commitment, &chain_tip_header.account_root())
1267 .map_err(|err| {
1268 ClientError::ChainValidationError(format!(
1269 "get_account witness for account {account_id} does not open under block \
1270 {target_block_num} account root: {err}"
1271 ))
1272 })?;
1273
1274 details.ok_or_else(|| {
1275 ClientError::ChainValidationError(format!(
1276 "get_account returned no details for public account {account_id}"
1277 ))
1278 })
1279 }
1280
1281 async fn build_patch_update(
1284 &self,
1285 account_id: AccountId,
1286 details: &AccountDetails,
1287 block_from: BlockNumber,
1288 block_to: BlockNumber,
1289 ) -> Result<PublicAccountUpdate, ClientError> {
1290 let value_slot_updates: Vec<(_, Word)> = details
1291 .storage_details
1292 .header
1293 .slots()
1294 .filter(|slot| slot.slot_type() == StorageSlotType::Value)
1295 .map(|slot| (slot.name().clone(), slot.value()))
1296 .collect();
1297
1298 let map_info = self
1301 .rpc_api
1302 .sync_storage_maps(block_from + 1, block_to, account_id)
1303 .await
1304 .map_err(ClientError::RpcError)?;
1305 let vault_info = self
1306 .rpc_api
1307 .sync_account_vault(block_from + 1, block_to, account_id)
1308 .await
1309 .map_err(ClientError::RpcError)?;
1310
1311 let patch = build_account_patch(
1312 &details.header,
1313 value_slot_updates,
1314 map_info.map_entries,
1315 vault_info.vault_patch,
1316 details.code.clone(),
1317 )
1318 .map_err(StoreError::AccountPatchError)?;
1319
1320 Ok(PublicAccountUpdate::Patch {
1321 new_header: details.header.clone(),
1322 patch,
1323 })
1324 }
1325
1326 async fn note_state_sync(
1344 &self,
1345 note_updates: &mut NoteUpdateTracker,
1346 notes: BTreeMap<NoteId, SyncedNote>,
1347 block_header: &BlockHeader,
1348 ) -> Result<NoteBlockRelevance, ClientError> {
1349 let mut relevance = NoteBlockRelevance::default();
1350
1351 for (_, mut note) in notes {
1352 if !self.note_observers.is_empty() {
1355 for obs in &self.note_observers {
1356 match obs.observe(¬e).await {
1357 Ok(true) => relevance.observer_requires_block = true,
1358 Ok(false) => {},
1359 Err(err) => {
1360 tracing::warn!(
1361 observer = obs.name(),
1362 error = ?err,
1363 "note observer failed; sync continues",
1364 );
1365 },
1366 }
1367 }
1368 }
1369
1370 let public_note = note.details.take().map(|details| {
1373 let state = UnverifiedNoteState {
1374 metadata: note.metadata,
1375 inclusion_proof: note.inclusion_proof.clone(),
1376 }
1377 .into();
1378 InputNoteRecord::new(details, note.attachments.clone(), None, state)
1379 });
1380
1381 let committed = note.into_committed_note();
1382
1383 match self.note_screener.on_note_received(committed, public_note).await? {
1384 NoteUpdateAction::Commit(committed_note) => {
1385 relevance.has_client_note |= note_updates
1389 .apply_committed_note_state_transitions(&committed_note, block_header)?;
1390 },
1391 NoteUpdateAction::Insert(public_note) => {
1392 relevance.has_client_note = true;
1393
1394 note_updates.apply_new_public_note(public_note, block_header)?;
1395 },
1396 NoteUpdateAction::Discard => {},
1397 }
1398 }
1399
1400 Ok(relevance)
1401 }
1402
1403 async fn nullifiers_state_sync(
1409 &self,
1410 note_updates: &mut NoteUpdateTracker,
1411 transaction_updates: &mut TransactionUpdateTracker,
1412 chain_tip: BlockNumber,
1413 current_block_num: BlockNumber,
1414 ) -> Result<(), ClientError> {
1415 let nullifiers_tags: Vec<u16> =
1421 note_updates.unspent_nullifiers().map(|nullifier| nullifier.prefix()).collect();
1422
1423 let mut new_nullifiers = self
1424 .rpc_api
1425 .sync_nullifiers(&nullifiers_tags, current_block_num + 1, chain_tip)
1426 .await?;
1427
1428 new_nullifiers.retain(|update| update.block_num <= chain_tip);
1431
1432 let consumptions: Vec<NoteConsumption> = new_nullifiers
1434 .into_iter()
1435 .map(|update| NoteConsumption {
1436 external_consumer: transaction_updates
1437 .external_nullifier_account(&update.nullifier),
1438 nullifier: update.nullifier,
1439 block_num: update.block_num,
1440 })
1441 .collect();
1442
1443 for consumption in consumptions {
1444 note_updates.apply_note_consumption(
1445 &consumption,
1446 transaction_updates.committed_transactions(),
1447 )?;
1448
1449 transaction_updates.apply_input_note_nullified(consumption.nullifier);
1453 }
1454
1455 Ok(())
1456 }
1457}
1458
1459pub struct ChainSyncData {
1468 pub(crate) block_from: BlockNumber,
1470 advance: Option<ChainAdvance>,
1473 superseded_states: Vec<Word>,
1475 pub(crate) note_updates: NoteUpdateTracker,
1478 transaction_updates: TransactionUpdateTracker,
1479 account_updates: AccountUpdates,
1480}
1481
1482struct ChainAdvance {
1484 chain_tip_header: BlockHeader,
1486 mmr_delta: MmrDelta,
1488 note_blocks_awaiting_screening: Vec<ResolvedSyncNotesBlock>,
1491 transactions: Vec<RpcTransactionRecord>,
1493 relevant_note_blocks: Vec<RelevantNoteBlock>,
1495 protocol_config: Option<ProtocolConfig>,
1497}
1498
1499pub(crate) fn block_num_from_forest(partial_mmr: &PartialMmr) -> Result<BlockNumber, ClientError> {
1504 Ok(u32::try_from(partial_mmr.forest().num_leaves().saturating_sub(1))
1505 .map_err(|_| ClientError::InvalidPartialMmrForest)?
1506 .into())
1507}
1508
1509fn group_txs_by_account_block(
1511 transaction_records: &[RpcTransactionRecord],
1512) -> BTreeMap<(AccountId, BlockNumber), Vec<&RpcTransactionRecord>> {
1513 let mut groups: BTreeMap<(AccountId, BlockNumber), Vec<&RpcTransactionRecord>> =
1514 BTreeMap::new();
1515 for record in transaction_records {
1516 let account_id = record.transaction_header.account_id();
1517 groups.entry((account_id, record.block_num)).or_default().push(record);
1518 }
1519 groups
1520}
1521
1522fn walk_execution_chain<'a>(
1528 txs: &'a [&'a RpcTransactionRecord],
1529) -> impl Iterator<Item = &'a RpcTransactionRecord> + 'a {
1530 let (self_loops, chained): (Vec<&RpcTransactionRecord>, Vec<&RpcTransactionRecord>) =
1531 txs.iter().copied().partition(|tx| {
1532 tx.transaction_header.initial_state_commitment()
1533 == tx.transaction_header.final_state_commitment()
1534 });
1535
1536 let final_states: BTreeSet<Word> = chained
1537 .iter()
1538 .map(|tx| tx.transaction_header.final_state_commitment())
1539 .collect();
1540
1541 let mut init_to_tx: BTreeMap<Word, &RpcTransactionRecord> = chained
1542 .iter()
1543 .map(|tx| (tx.transaction_header.initial_state_commitment(), *tx))
1544 .collect();
1545
1546 let start = chained
1547 .iter()
1548 .find(|tx| !final_states.contains(&tx.transaction_header.initial_state_commitment()))
1549 .copied();
1550
1551 assert!(start.is_some() || chained.is_empty(), "cannot walk cyclic execution chain");
1552
1553 let mut current =
1554 start.and_then(|tx| init_to_tx.remove(&tx.transaction_header.initial_state_commitment()));
1555 let mut self_loops_iter = self_loops.into_iter();
1556
1557 core::iter::from_fn(move || {
1558 if let Some(tx) = current {
1559 current = init_to_tx.remove(&tx.transaction_header.final_state_commitment());
1560 return Some(tx);
1561 }
1562 self_loops_iter.next()
1563 })
1564}
1565
1566fn derive_account_commitments(
1571 transaction_records: &[RpcTransactionRecord],
1572) -> Vec<(AccountId, Word)> {
1573 let mut latest_by_account: BTreeMap<AccountId, (BlockNumber, Word)> = BTreeMap::new();
1574
1575 for ((account_id, block_num), txs) in &group_txs_by_account_block(transaction_records) {
1576 let terminal_state = walk_execution_chain(txs)
1577 .last()
1578 .expect("account must have a final state")
1579 .transaction_header
1580 .final_state_commitment();
1581
1582 latest_by_account
1583 .entry(*account_id)
1584 .and_modify(|(existing_block, existing_state)| {
1585 if *block_num > *existing_block {
1586 *existing_block = *block_num;
1587 *existing_state = terminal_state;
1588 }
1589 })
1590 .or_insert((*block_num, terminal_state));
1591 }
1592
1593 latest_by_account
1594 .into_iter()
1595 .map(|(account_id, (_, state))| (account_id, state))
1596 .collect()
1597}
1598
1599fn compute_ordered_nullifiers(transaction_records: &[RpcTransactionRecord]) -> Vec<Nullifier> {
1606 let mut result = Vec::new();
1607
1608 for txs in group_txs_by_account_block(transaction_records).values() {
1609 for tx in walk_execution_chain(txs) {
1610 for commitment in tx.transaction_header.input_notes().iter() {
1611 result.push(commitment.nullifier());
1612 }
1613 }
1614 }
1615
1616 result
1617}
1618
1619#[cfg(all(test, feature = "testing"))]
1620mod tests {
1621 use alloc::collections::BTreeSet;
1622 use alloc::sync::Arc;
1623
1624 use async_trait::async_trait;
1625 use miden_protocol::account::Account;
1626 use miden_protocol::assembly::DefaultSourceManager;
1627 use miden_protocol::asset::{Asset, FungibleAsset};
1628 use miden_protocol::block::{BlockNumber, BlockSignatures};
1629 use miden_protocol::crypto::merkle::MerklePath;
1630 use miden_protocol::crypto::merkle::mmr::{Forest, InOrderIndex, PartialMmr};
1631 use miden_protocol::note::{
1632 Note,
1633 NoteAssets,
1634 NoteAttachment,
1635 NoteAttachments,
1636 NoteDetails,
1637 NoteHeader,
1638 NoteMetadata,
1639 NoteRecipient,
1640 NoteStorage,
1641 NoteTag,
1642 NoteType,
1643 PartialNoteMetadata,
1644 };
1645 use miden_protocol::testing::account_id::{
1646 ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET,
1647 ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET,
1648 ACCOUNT_ID_REGULAR_PRIVATE_ACCOUNT_UPDATABLE_CODE,
1649 ACCOUNT_ID_REGULAR_PUBLIC_ACCOUNT_IMMUTABLE_CODE,
1650 ACCOUNT_ID_SENDER,
1651 };
1652 use miden_protocol::transaction::{InputNotes, TransactionArgs, TransactionHeader};
1653 use miden_protocol::vm::AdviceMap;
1654 use miden_protocol::{EMPTY_WORD, Felt, Word, ZERO};
1655 use miden_standards::code_builder::CodeBuilder;
1656 use miden_standards::note::{NetworkAccountTarget, NoteExecutionHint};
1657 use miden_testing::{MockChainBuilder, MockTransactionInput};
1658
1659 use super::*;
1660 use crate::store::{OutputNoteRecord, OutputNoteState};
1661 use crate::test_utils::mock::MockRpcApi;
1662
1663 struct MockScreener;
1665
1666 #[async_trait(?Send)]
1667 impl OnNoteReceived for MockScreener {
1668 async fn on_note_received(
1669 &self,
1670 _committed_note: CommittedNote,
1671 _public_note: Option<InputNoteRecord>,
1672 ) -> Result<NoteUpdateAction, ClientError> {
1673 Ok(NoteUpdateAction::Discard)
1674 }
1675 }
1676
1677 struct AlwaysRelevantObserver;
1679
1680 #[async_trait(?Send)]
1681 impl NoteObserver for AlwaysRelevantObserver {
1682 fn name(&self) -> &'static str {
1683 "always-relevant"
1684 }
1685
1686 async fn observe(&self, _note: &SyncedNote) -> Result<bool, ClientError> {
1687 Ok(true)
1688 }
1689 }
1690
1691 fn genesis_validator_config(mock_rpc: &MockRpcApi) -> ValidatorConfig {
1693 mock_rpc.mock_chain.read().block_header(0).validator_config().clone()
1694 }
1695
1696 fn block_signatures(mock_rpc: &MockRpcApi, block_num: BlockNumber) -> BlockSignatures {
1698 mock_rpc
1699 .mock_chain
1700 .read()
1701 .proven_blocks()
1702 .iter()
1703 .find(|block| block.header().block_num() == block_num)
1704 .expect("the mock chain contains the block")
1705 .signatures()
1706 .clone()
1707 }
1708
1709 fn empty() -> StateSyncInput {
1710 StateSyncInput {
1711 accounts: vec![],
1712 note_tags: BTreeSet::new(),
1713 input_notes: vec![],
1714 output_notes: vec![],
1715 uncommitted_transactions: vec![],
1716 }
1717 }
1718
1719 fn word(n: u64) -> miden_protocol::Word {
1720 [
1721 Felt::new(n).expect("test value should fit into the base field"),
1722 ZERO,
1723 ZERO,
1724 ZERO,
1725 ]
1726 .into()
1727 }
1728
1729 fn header_with_account_root(header: &BlockHeader, account_root: Word) -> BlockHeader {
1730 BlockHeader::new(
1731 header.prev_block_commitment(),
1732 header.block_num(),
1733 header.chain_commitment(),
1734 account_root,
1735 header.nullifier_root(),
1736 header.note_root(),
1737 header.tx_commitment(),
1738 header.validator_config().clone(),
1739 header.fee_parameters().clone(),
1740 header.protocol_config_commitment(),
1741 header.next_protocol_config().cloned(),
1742 header.timestamp(),
1743 )
1744 }
1745
1746 #[tokio::test]
1747 async fn sync_public_accounts_ignores_older_node_snapshot() {
1748 let mut builder = MockChainBuilder::new();
1749 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1750 let rpc_api = MockRpcApi::new(builder.build().unwrap());
1751 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
1752 let validator_config = genesis_validator_config(&rpc_api);
1753 let state_sync =
1754 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1755
1756 let local_header =
1759 AccountHeader::new(account.id(), Felt::from(2u32), EMPTY_WORD, EMPTY_WORD, EMPTY_WORD);
1760 let current_public_accounts = vec![&local_header];
1761 let commitment_updates = vec![(account.id(), account.to_commitment())];
1762 let mut account_updates = AccountUpdates::default();
1763
1764 let superseded = state_sync
1765 .sync_public_accounts(
1766 &mut account_updates,
1767 &commitment_updates,
1768 ¤t_public_accounts,
1769 BlockNumber::GENESIS,
1770 &chain_tip_header,
1771 )
1772 .await
1773 .unwrap();
1774
1775 assert!(
1776 account_updates.updated_public_accounts().is_empty(),
1777 "public account sync should ignore node snapshots that are older than local"
1778 );
1779 assert!(
1780 superseded.is_empty(),
1781 "an older node snapshot must not supersede the local state"
1782 );
1783 }
1784
1785 #[tokio::test]
1786 async fn sync_public_accounts_marks_same_nonce_mismatch_as_superseded() {
1787 let mut builder = MockChainBuilder::new();
1788 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1789 let rpc_api = MockRpcApi::new(builder.build().unwrap());
1790 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
1791 let validator_config = genesis_validator_config(&rpc_api);
1792 let state_sync =
1793 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1794
1795 let local_header =
1798 AccountHeader::new(account.id(), account.nonce(), EMPTY_WORD, EMPTY_WORD, EMPTY_WORD);
1799 let current_public_accounts = vec![&local_header];
1800 let commitment_updates = vec![(account.id(), account.to_commitment())];
1801 let mut account_updates = AccountUpdates::default();
1802
1803 let superseded = state_sync
1804 .sync_public_accounts(
1805 &mut account_updates,
1806 &commitment_updates,
1807 ¤t_public_accounts,
1808 BlockNumber::GENESIS,
1809 &chain_tip_header,
1810 )
1811 .await
1812 .unwrap();
1813
1814 assert!(
1815 account_updates.updated_public_accounts().is_empty(),
1816 "a same-nonce fork must not overwrite the account while its tx is still pending"
1817 );
1818 assert_eq!(
1819 superseded,
1820 vec![local_header.to_commitment()],
1821 "the superseded local state should be reported so its transaction is discarded"
1822 );
1823 }
1824
1825 #[test]
1831 fn validate_transaction_records_range_rejects_out_of_range_blocks() {
1832 let account_id: AccountId = ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET.try_into().unwrap();
1833 let current = BlockNumber::from(5u32);
1834 let chain_tip = BlockNumber::from(10u32);
1835
1836 StateSync::validate_transaction_records_range(
1837 &[make_tx_record(account_id, 7)],
1838 current,
1839 chain_tip,
1840 )
1841 .unwrap();
1842
1843 let result = StateSync::validate_transaction_records_range(
1844 &[make_tx_record(account_id, 11)],
1845 current,
1846 chain_tip,
1847 );
1848 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
1849
1850 let result = StateSync::validate_transaction_records_range(
1851 &[make_tx_record(account_id, 5)],
1852 current,
1853 chain_tip,
1854 );
1855 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
1856 }
1857
1858 #[tokio::test]
1861 async fn verify_private_account_mismatch_ignores_forged_commitment() {
1862 let mut builder = MockChainBuilder::new();
1863 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1864 let rpc_api = MockRpcApi::new(builder.build().unwrap());
1865 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
1866 let on_chain_commitment = account.to_commitment();
1867 let validator_config = genesis_validator_config(&rpc_api);
1868 let state_sync =
1869 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1870
1871 let result = state_sync
1872 .verify_private_account_mismatch(account.id(), on_chain_commitment, &chain_tip_header)
1873 .await
1874 .unwrap();
1875
1876 assert!(
1877 result.is_none(),
1878 "an unproven commitment must not lock an account whose on-chain state matches local"
1879 );
1880 }
1881
1882 #[tokio::test]
1885 async fn verify_private_account_mismatch_reports_proven_divergence() {
1886 let mut builder = MockChainBuilder::new();
1887 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1888 let rpc_api = MockRpcApi::new(builder.build().unwrap());
1889 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
1890 let on_chain_commitment = account.to_commitment();
1891 let validator_config = genesis_validator_config(&rpc_api);
1892 let state_sync =
1893 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1894 let stale_local_commitment = word(0xdead_beef);
1895
1896 let result = state_sync
1897 .verify_private_account_mismatch(
1898 account.id(),
1899 stale_local_commitment,
1900 &chain_tip_header,
1901 )
1902 .await
1903 .unwrap();
1904
1905 assert_eq!(
1906 result,
1907 Some(on_chain_commitment),
1908 "a proven divergence should return the proven commitment to lock with"
1909 );
1910 }
1911
1912 #[tokio::test]
1915 async fn verify_private_account_mismatch_rejects_unverifiable_proof() {
1916 let mut builder = MockChainBuilder::new();
1917 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1918 let rpc_api = MockRpcApi::new(builder.build().unwrap());
1919 let real_header = rpc_api.mock_chain.read().latest_block_header();
1920 let validator_config = genesis_validator_config(&rpc_api);
1921 let state_sync =
1922 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1923
1924 let tampered_header = BlockHeader::new(
1927 real_header.prev_block_commitment(),
1928 real_header.block_num(),
1929 real_header.chain_commitment(),
1930 word(0xbad0_bad0),
1931 real_header.nullifier_root(),
1932 real_header.note_root(),
1933 real_header.tx_commitment(),
1934 real_header.validator_config().clone(),
1935 real_header.fee_parameters().clone(),
1936 real_header.protocol_config_commitment(),
1937 real_header.next_protocol_config().cloned(),
1938 real_header.timestamp(),
1939 );
1940
1941 let result = state_sync
1942 .verify_private_account_mismatch(
1943 account.id(),
1944 account.to_commitment(),
1945 &tampered_header,
1946 )
1947 .await;
1948 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
1949 }
1950
1951 #[tokio::test]
1954 async fn sync_public_accounts_pins_account_fetch_to_sync_target() {
1955 let mut builder = MockChainBuilder::new();
1956 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1957 let mut chain = builder.build().unwrap();
1958
1959 let sync_target_header = chain.latest_block_header();
1961 let tx = Box::pin(
1962 chain
1963 .build_transaction(MockTransactionInput::AccountId(account.id()))
1964 .build()
1965 .unwrap()
1966 .execute(),
1967 )
1968 .await
1969 .unwrap();
1970 let local_header = tx.final_account().clone();
1971 assert_ne!(local_header.to_commitment(), account.to_commitment());
1972 chain.add_pending_executed_transaction(&tx).unwrap();
1973
1974 let rpc_api = MockRpcApi::new(chain);
1975 rpc_api.prove_block();
1977 assert_eq!(
1978 rpc_api
1979 .mock_chain
1980 .read()
1981 .committed_account(account.id())
1982 .unwrap()
1983 .to_commitment(),
1984 local_header.to_commitment()
1985 );
1986 let validator_config = genesis_validator_config(&rpc_api);
1987 let state_sync =
1988 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1989
1990 let current_public_accounts = vec![&local_header];
1991 let commitment_updates = vec![(account.id(), account.to_commitment())];
1992 let mut account_updates = AccountUpdates::default();
1993
1994 let superseded = state_sync
1995 .sync_public_accounts(
1996 &mut account_updates,
1997 &commitment_updates,
1998 ¤t_public_accounts,
1999 BlockNumber::GENESIS,
2000 &sync_target_header,
2001 )
2002 .await
2003 .unwrap();
2004
2005 assert!(superseded.is_empty(), "the transaction must not be superseded");
2006 assert!(
2007 account_updates.updated_public_accounts().is_empty(),
2008 "the target state must not overwrite the local account"
2009 );
2010 }
2011
2012 async fn get_account_proof(
2014 rpc_api: &MockRpcApi,
2015 account_id: AccountId,
2016 ) -> (BlockNumber, AccountProof) {
2017 rpc_api
2018 .get_account(
2019 account_id,
2020 GetAccountRequest::new()
2021 .with_storage(StorageMapFetch::All)
2022 .with_vault(VaultFetch::Always),
2023 )
2024 .await
2025 .unwrap()
2026 }
2027
2028 #[tokio::test]
2030 async fn validate_account_proof_rejects_mismatched_account() {
2031 let mut builder = MockChainBuilder::new();
2032 let account_a = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2033 let account_b = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2034 let rpc_api = MockRpcApi::new(builder.build().unwrap());
2035 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
2036
2037 let (proof_block_num, proof) = get_account_proof(&rpc_api, account_b.id()).await;
2039 let result = StateSync::validate_account_proof(
2040 proof,
2041 proof_block_num,
2042 account_a.id(),
2043 &chain_tip_header,
2044 );
2045
2046 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2047 }
2048
2049 #[tokio::test]
2051 async fn validate_account_proof_rejects_wrong_account_root() {
2052 let mut builder = MockChainBuilder::new();
2053 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2054 let rpc_api = MockRpcApi::new(builder.build().unwrap());
2055 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
2056 let wrong_header = header_with_account_root(&chain_tip_header, word(999));
2057
2058 let (proof_block_num, proof) = get_account_proof(&rpc_api, account.id()).await;
2060 let result =
2061 StateSync::validate_account_proof(proof, proof_block_num, account.id(), &wrong_header);
2062
2063 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2064 }
2065
2066 #[tokio::test]
2068 async fn validate_account_proof_rejects_missing_details() {
2069 let mut builder = MockChainBuilder::new();
2070 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2071 let rpc_api = MockRpcApi::new(builder.build().unwrap());
2072 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
2073
2074 let (proof_block_num, proof) = get_account_proof(&rpc_api, account.id()).await;
2076 let (witness, _) = proof.into_parts();
2077 let proof = AccountProof::new(witness, None).unwrap();
2078 let result = StateSync::validate_account_proof(
2079 proof,
2080 proof_block_num,
2081 account.id(),
2082 &chain_tip_header,
2083 );
2084
2085 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2086 }
2087
2088 #[tokio::test]
2090 async fn validate_account_proof_rejects_wrong_block() {
2091 let mut builder = MockChainBuilder::new();
2092 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2093 let rpc_api = MockRpcApi::new(builder.build().unwrap());
2094 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
2095
2096 let (proof_block_num, proof) = get_account_proof(&rpc_api, account.id()).await;
2098 let result = StateSync::validate_account_proof(
2099 proof,
2100 proof_block_num + 1,
2101 account.id(),
2102 &chain_tip_header,
2103 );
2104
2105 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2106 }
2107
2108 mod compute_nullifiers_tests {
2112 use alloc::vec;
2113
2114 use miden_protocol::block::BlockNumber;
2115 use miden_protocol::note::Nullifier;
2116 use miden_protocol::transaction::{InputNoteCommitment, InputNotes, TransactionHeader};
2117
2118 use super::word;
2119 use crate::rpc::domain::transaction::TransactionRecord as RpcTransactionRecord;
2120
2121 fn make_rpc_tx(
2122 init_state: u64,
2123 final_state: u64,
2124 nullifier_vals: &[u64],
2125 block_number: u32,
2126 ) -> RpcTransactionRecord {
2127 let account_id = miden_protocol::account::AccountId::try_from(
2128 miden_protocol::testing::account_id::ACCOUNT_ID_REGULAR_PRIVATE_ACCOUNT_UPDATABLE_CODE,
2129 )
2130 .unwrap();
2131
2132 let input_notes = InputNotes::new_unchecked(
2133 nullifier_vals
2134 .iter()
2135 .map(|v| InputNoteCommitment::from(Nullifier::from_raw(word(*v))))
2136 .collect(),
2137 );
2138
2139 RpcTransactionRecord {
2140 block_num: BlockNumber::from(block_number),
2141 transaction_header: TransactionHeader::new(
2142 account_id,
2143 word(init_state),
2144 word(final_state),
2145 input_notes,
2146 vec![],
2147 )
2148 .unwrap(),
2149 output_notes: vec![],
2150 erased_output_notes: vec![],
2151 consumed_note_refs: vec![],
2152 }
2153 }
2154
2155 #[test]
2156 fn chains_rpc_transactions_by_state_commitment() {
2157 let tx_a = make_rpc_tx(1, 2, &[10], 5);
2161 let tx_b = make_rpc_tx(2, 3, &[20], 5);
2162 let tx_c = make_rpc_tx(3, 4, &[30], 5);
2163
2164 let result = super::super::compute_ordered_nullifiers(&[tx_c, tx_a, tx_b]);
2165
2166 assert_eq!(result[0], Nullifier::from_raw(word(10)));
2167 assert_eq!(result[1], Nullifier::from_raw(word(20)));
2168 assert_eq!(result[2], Nullifier::from_raw(word(30)));
2169 }
2170
2171 #[test]
2172 fn groups_independently_by_account_and_block() {
2173 let tx_a1 = make_rpc_tx(1, 2, &[10], 5);
2175 let tx_a2 = make_rpc_tx(2, 3, &[20], 5);
2176
2177 let tx_a3 = make_rpc_tx(3, 4, &[30], 6);
2179
2180 let account_b = miden_protocol::account::AccountId::try_from(
2182 miden_protocol::testing::account_id::ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET,
2183 )
2184 .unwrap();
2185
2186 let tx_b1 = RpcTransactionRecord {
2187 block_num: BlockNumber::from(5u32),
2188 transaction_header: TransactionHeader::new(
2189 account_b,
2190 word(100),
2191 word(200),
2192 InputNotes::new_unchecked(vec![InputNoteCommitment::from(
2193 Nullifier::from_raw(word(40)),
2194 )]),
2195 vec![],
2196 )
2197 .unwrap(),
2198 output_notes: vec![],
2199 erased_output_notes: vec![],
2200 consumed_note_refs: vec![],
2201 };
2202
2203 let result = super::super::compute_ordered_nullifiers(&[tx_a2, tx_b1, tx_a3, tx_a1]);
2204
2205 let pos = |val: u64| -> usize {
2208 result.iter().position(|n| *n == Nullifier::from_raw(word(val))).unwrap()
2209 };
2210
2211 assert!(pos(10) < pos(20)); assert!(result.contains(&Nullifier::from_raw(word(30)))); assert!(result.contains(&Nullifier::from_raw(word(40)))); }
2217
2218 #[test]
2219 fn multiple_nullifiers_per_transaction_are_consecutive() {
2220 let tx = make_rpc_tx(1, 2, &[10, 20, 30], 5);
2222
2223 let result = super::super::compute_ordered_nullifiers(&[tx]);
2224
2225 assert_eq!(result.len(), 3);
2226 assert!(result.contains(&Nullifier::from_raw(word(10))));
2227 assert!(result.contains(&Nullifier::from_raw(word(20))));
2228 assert!(result.contains(&Nullifier::from_raw(word(30))));
2229 }
2230
2231 #[test]
2232 fn empty_input_returns_empty_vec() {
2233 let result = super::super::compute_ordered_nullifiers(&[]);
2234 assert!(result.is_empty());
2235 }
2236 }
2237
2238 #[test]
2249 fn derive_account_commitments_walks_chains_per_account() {
2250 let make_tx = |account: AccountId, init_state: u64, final_state: u64, block_num: u32| {
2251 RpcTransactionRecord {
2252 block_num: BlockNumber::from(block_num),
2253 transaction_header: TransactionHeader::new(
2254 account,
2255 word(init_state),
2256 word(final_state),
2257 InputNotes::new_unchecked(vec![]),
2258 vec![],
2259 )
2260 .unwrap(),
2261 output_notes: vec![],
2262 erased_output_notes: vec![],
2263 consumed_note_refs: vec![],
2264 }
2265 };
2266
2267 let account_a: AccountId =
2268 ACCOUNT_ID_REGULAR_PRIVATE_ACCOUNT_UPDATABLE_CODE.try_into().unwrap();
2269 let account_b: AccountId = ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET.try_into().unwrap();
2270
2271 let tx_a_b5_1 = make_tx(account_a, 1, 2, 5);
2272 let tx_a_b5_2 = make_tx(account_a, 2, 3, 5);
2273 let tx_a_b6_1 = make_tx(account_a, 3, 4, 6);
2274 let tx_a_b6_2 = make_tx(account_a, 4, 5, 6);
2275 let tx_b_b6 = make_tx(account_b, 10, 20, 6);
2276
2277 let result = super::derive_account_commitments(&[
2279 tx_a_b6_1, tx_b_b6, tx_a_b5_2, tx_a_b6_2, tx_a_b5_1,
2280 ]);
2281
2282 assert_eq!(result.len(), 2, "one entry per account");
2283 assert!(
2284 result.contains(&(account_a, word(5))),
2285 "account A: must walk block 6's chain, not return block 5 or an intermediate",
2286 );
2287 assert!(
2288 result.contains(&(account_b, word(20))),
2289 "account B: must be resolved independently of account A",
2290 );
2291 }
2292
2293 struct CommitAllScreener;
2299
2300 #[async_trait(?Send)]
2301 impl OnNoteReceived for CommitAllScreener {
2302 async fn on_note_received(
2303 &self,
2304 committed_note: CommittedNote,
2305 _public_note: Option<InputNoteRecord>,
2306 ) -> Result<NoteUpdateAction, ClientError> {
2307 Ok(NoteUpdateAction::Commit(committed_note))
2308 }
2309 }
2310
2311 async fn build_chain_with_chained_consume_txs() -> (miden_testing::MockChain, Account, [Note; 3])
2315 {
2316 let sender_id: AccountId = ACCOUNT_ID_SENDER.try_into().unwrap();
2317 let faucet_id: AccountId = ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET.try_into().unwrap();
2318
2319 let mut builder = MockChainBuilder::new();
2320 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2321 let account_id = account.id();
2322
2323 let asset = Asset::from(FungibleAsset::new(faucet_id, 100u64).unwrap());
2324 let note1 = builder
2325 .add_p2id_note(sender_id, account_id, &[asset], NoteType::Public)
2326 .unwrap();
2327 let note2 = builder
2328 .add_p2id_note(sender_id, account_id, &[asset], NoteType::Public)
2329 .unwrap();
2330 let note3 = builder
2331 .add_p2id_note(sender_id, account_id, &[asset], NoteType::Public)
2332 .unwrap();
2333
2334 let mut chain = builder.build().unwrap();
2335 chain.prove_next_block().unwrap(); let mut current_account = account.clone();
2339 for note in [¬e1, ¬e2, ¬e3] {
2340 let tx = Box::pin(
2341 chain
2342 .build_transaction(MockTransactionInput::Account(current_account.clone()))
2343 .unauthenticated_input_note(note.clone())
2344 .build()
2345 .unwrap()
2346 .execute(),
2347 )
2348 .await
2349 .unwrap();
2350 current_account.apply_patch(tx.account_patch()).unwrap();
2351 chain.add_pending_executed_transaction(&tx).unwrap();
2352 }
2353
2354 chain.prove_next_block().unwrap(); (chain, account, [note1, note2, note3])
2356 }
2357
2358 #[tokio::test]
2361 async fn sync_state_sets_consumed_tx_order_for_chained_transactions() {
2362 use miden_protocol::note::NoteMetadata;
2363
2364 let (chain, account, [note1, note2, note3]) = build_chain_with_chained_consume_txs().await;
2365
2366 let mock_rpc = MockRpcApi::new(chain);
2367 let state_sync = StateSync::new(
2368 Arc::new(mock_rpc.clone()),
2369 Arc::new(CommitAllScreener),
2370 None,
2371 genesis_validator_config(&mock_rpc),
2372 );
2373
2374 let genesis_peaks =
2375 mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
2376 let mut partial_mmr = PartialMmr::from_peaks(genesis_peaks);
2377
2378 let input_notes: Vec<InputNoteRecord> = [¬e1, ¬e2, ¬e3]
2379 .into_iter()
2380 .map(|n| InputNoteRecord::from(n.clone()))
2381 .collect();
2382
2383 let note_tags: BTreeSet<NoteTag> =
2384 input_notes.iter().filter_map(|n| n.metadata().map(NoteMetadata::tag)).collect();
2385
2386 let account_id = account.id();
2387 let sync_input = StateSyncInput {
2388 accounts: vec![AccountHeader::from(&account)],
2389 note_tags,
2390 input_notes,
2391 output_notes: vec![],
2392 uncommitted_transactions: vec![],
2393 };
2394
2395 let update = state_sync.sync_state(&mut partial_mmr, sync_input).await.unwrap();
2396
2397 let updated_notes: Vec<_> = update.note_updates().updated_input_notes().collect();
2398
2399 let find_order = |details_commitment| -> Option<u32> {
2400 updated_notes
2401 .iter()
2402 .find(|n| n.inner().details_commitment() == details_commitment)
2403 .and_then(|n| n.consumed_tx_order())
2404 };
2405
2406 assert_eq!(find_order(note1.details_commitment()), Some(0), "note1 should have tx_order 0");
2407 assert_eq!(find_order(note2.details_commitment()), Some(1), "note2 should have tx_order 1");
2408 assert_eq!(find_order(note3.details_commitment()), Some(2), "note3 should have tx_order 2");
2409
2410 for note in &updated_notes {
2413 let record = note.inner();
2414 assert!(record.is_consumed(), "note should be in a consumed state");
2415 assert_eq!(
2416 record.consumer_account(),
2417 Some(account_id),
2418 "externally-consumed notes by a tracked account should have consumer_account set",
2419 );
2420 }
2421 }
2422
2423 #[tokio::test]
2424 async fn sync_state_across_multiple_iterations_with_same_mmr() {
2425 let mock_rpc = MockRpcApi::default();
2427 mock_rpc.advance_blocks(3);
2428 let chain_tip_1 = mock_rpc.get_chain_tip_block_num();
2429
2430 let state_sync = StateSync::new(
2431 Arc::new(mock_rpc.clone()),
2432 Arc::new(MockScreener),
2433 None,
2434 genesis_validator_config(&mock_rpc),
2435 );
2436
2437 let genesis_peaks =
2439 mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
2440 let mut partial_mmr = PartialMmr::from_peaks(genesis_peaks);
2441 assert_eq!(partial_mmr.forest().num_leaves(), 1);
2442
2443 let update = state_sync.sync_state(&mut partial_mmr, empty()).await.unwrap();
2445
2446 assert_eq!(update.block_num(), chain_tip_1);
2447 let forest_1 = partial_mmr.forest();
2448 assert_eq!(forest_1.num_leaves(), chain_tip_1.as_u32() as usize + 1);
2450
2451 mock_rpc.advance_blocks(2);
2453 let chain_tip_2 = mock_rpc.get_chain_tip_block_num();
2454
2455 let update = state_sync.sync_state(&mut partial_mmr, empty()).await.unwrap();
2456
2457 assert_eq!(update.block_num(), chain_tip_2);
2458 let forest_2 = partial_mmr.forest();
2459 assert!(forest_2 > forest_1);
2460 assert_eq!(forest_2.num_leaves(), chain_tip_2.as_u32() as usize + 1);
2461
2462 let update = state_sync.sync_state(&mut partial_mmr, empty()).await.unwrap();
2464
2465 assert_eq!(update.block_num(), chain_tip_2);
2466 assert_eq!(partial_mmr.forest(), forest_2);
2467 }
2468
2469 async fn build_chain_with_mint_notes(
2472 num_blocks: u64,
2473 ) -> (miden_testing::MockChain, BTreeSet<NoteTag>) {
2474 let mut builder = MockChainBuilder::new();
2475 let faucet = builder
2476 .add_existing_basic_faucet(
2477 miden_testing::Auth::BasicAuth {
2478 auth_scheme: miden_protocol::account::auth::AuthScheme::Falcon512Poseidon2,
2479 },
2480 "TST",
2481 10_000,
2482 None,
2483 )
2484 .unwrap();
2485 let _target = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2486 let mut chain = builder.build().unwrap();
2487
2488 let note_script = CodeBuilder::new()
2493 .compile_note_script("@note_script\npub proc main\n nop\nend")
2494 .unwrap();
2495 let note_recipient = NoteRecipient::new(
2496 Word::from([1u32, 2, 3, 4]),
2497 note_script,
2498 NoteStorage::new(vec![]).unwrap(),
2499 );
2500 let recipient = note_recipient.digest();
2501 let note_details = NoteDetails::new(NoteAssets::new(vec![]).unwrap(), note_recipient);
2504 let mut recipient_args = TransactionArgs::new(AdviceMap::default());
2505 recipient_args.add_output_note_recipient(¬e_details);
2506 let recipient_advice = recipient_args.advice_inputs().clone();
2507
2508 let tag = NoteTag::default();
2509 let mut faucet_account = faucet.clone();
2510 let mut note_tags = BTreeSet::new();
2511
2512 for i in 0..num_blocks {
2513 let amount = 100 + i;
2514 let source_manager = Arc::new(DefaultSourceManager::default());
2515 let mint_asset = FungibleAsset::new(faucet_account.id(), amount).unwrap();
2519 let asset_id_word = mint_asset.id().to_word();
2520 let asset_value_word = mint_asset.to_value_word();
2521 let tx_script_code = format!(
2522 "
2523 @transaction_script
2524 pub proc main
2525 push.{recipient}
2526 push.{note_type}
2527 push.{tag}
2528 push.{asset_value}
2529 push.{asset_id}
2530 call.::miden::standards::faucets::fungible::mint_and_send
2531 dropw dropw dropw dropw
2532 end
2533 ",
2534 recipient = recipient,
2535 note_type = NoteType::Private as u8,
2536 tag = u32::from(tag),
2537 asset_value = asset_value_word,
2538 asset_id = asset_id_word,
2539 );
2540 let tx_script = CodeBuilder::with_source_manager(source_manager.clone())
2541 .compile_tx_script(tx_script_code)
2542 .unwrap();
2543 let tx = Box::pin(
2544 chain
2545 .build_transaction(miden_testing::MockTransactionInput::Account(
2546 faucet_account.clone(),
2547 ))
2548 .extend_advice_inputs(recipient_advice.clone())
2549 .tx_script(tx_script)
2550 .with_source_manager(source_manager)
2551 .build()
2552 .unwrap()
2553 .execute(),
2554 )
2555 .await
2556 .unwrap();
2557
2558 for output_note in tx.output_notes().iter() {
2559 note_tags.insert(output_note.metadata().tag());
2560 }
2561
2562 faucet_account.apply_patch(tx.account_patch()).unwrap();
2563 chain.add_pending_executed_transaction(&tx).unwrap();
2564 chain.prove_next_block().unwrap();
2565 }
2566
2567 (chain, note_tags)
2568 }
2569
2570 #[tokio::test]
2573 async fn observer_relevance_persists_discarded_note_block() {
2574 let (chain, note_tags) = build_chain_with_mint_notes(2).await;
2575 let mock_rpc = MockRpcApi::new(chain);
2576 let chain_tip = mock_rpc.get_chain_tip_block_num();
2577
2578 let genesis_peaks =
2579 mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
2580 let mut partial_mmr = PartialMmr::from_peaks(genesis_peaks);
2581
2582 let validator_config = genesis_validator_config(&mock_rpc);
2583 let mut input = empty();
2584 let state_sync =
2585 StateSync::new(Arc::new(mock_rpc), Arc::new(MockScreener), None, validator_config)
2586 .with_note_observer(Arc::new(AlwaysRelevantObserver));
2587 input.note_tags = note_tags;
2588
2589 let update = state_sync.sync_state(&mut partial_mmr, input).await.unwrap();
2590 let observed_non_tip_block = BlockNumber::from(1u32);
2591
2592 assert!(
2593 update.partial_blockchain_updates().block_headers_to_store(chain_tip).any(
2594 |(header, is_relevant)| {
2595 header.block_num() == observed_non_tip_block && *is_relevant
2596 }
2597 ),
2598 "an observer-relevant block must be staged as relevant"
2599 );
2600 assert!(
2601 partial_mmr.is_tracked(observed_non_tip_block.as_usize()),
2602 "an observer-relevant block must remain tracked in the partial MMR"
2603 );
2604 }
2605
2606 #[tokio::test]
2616 async fn sync_state_tracks_note_blocks_in_mmr() {
2617 let (chain, note_tags) = build_chain_with_mint_notes(3).await;
2618 let mock_rpc = MockRpcApi::new(chain);
2619 let chain_tip = mock_rpc.get_chain_tip_block_num();
2620
2621 let note_blocks = mock_rpc
2623 .sync_notes(BlockNumber::from(0u32), chain_tip, ¬e_tags)
2624 .await
2625 .unwrap();
2626 assert!(
2627 note_blocks.len() >= 2,
2628 "expected notes in multiple blocks, got {}",
2629 note_blocks.len()
2630 );
2631
2632 let note_block_nums: BTreeSet<BlockNumber> =
2634 note_blocks.iter().map(|b| b.block_header.block_num()).collect();
2635
2636 let state_sync = StateSync::new(
2639 Arc::new(mock_rpc.clone()),
2640 Arc::new(MockScreener),
2641 None,
2642 genesis_validator_config(&mock_rpc),
2643 );
2644
2645 let genesis_peaks =
2646 mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
2647 let mut partial_mmr = PartialMmr::from_peaks(genesis_peaks);
2648
2649 let sync_data = state_sync
2650 .fetch_sync_data(BlockNumber::GENESIS, &[], &Arc::new(note_tags.clone()))
2651 .await
2652 .unwrap()
2653 .expect("should have progressed past genesis");
2654
2655 assert_eq!(sync_data.chain_tip_header.block_num(), chain_tip);
2657 assert!(!sync_data.note_blocks.is_empty(), "should have note blocks");
2658
2659 let _auth_nodes: Vec<(InOrderIndex, Word)> =
2661 partial_mmr.apply(sync_data.mmr_delta).map_err(StoreError::MmrError).unwrap();
2662 partial_mmr
2663 .add(sync_data.chain_tip_header.commitment(), false)
2664 .expect("chain tip should append to the partial MMR");
2665
2666 assert_eq!(partial_mmr.forest().num_leaves(), chain_tip.as_u32() as usize + 1);
2667
2668 for block in &sync_data.note_blocks {
2670 let bn = block.block_header.block_num();
2671 partial_mmr
2672 .track(bn.as_usize(), block.block_header.commitment(), &block.mmr_path)
2673 .map_err(StoreError::MmrError)
2674 .unwrap();
2675
2676 assert!(
2677 partial_mmr.is_tracked(bn.as_usize()),
2678 "block {bn} should be tracked after calling track()"
2679 );
2680 }
2681
2682 for &bn in ¬e_block_nums {
2684 assert!(
2685 partial_mmr.is_tracked(bn.as_usize()),
2686 "block {bn} with notes should be tracked in partial MMR"
2687 );
2688 }
2689 }
2690
2691 #[tokio::test]
2692 async fn sync_notes_with_content_fetches_inclusive_upper_bound_page() {
2693 let (chain, note_tags) = build_chain_with_mint_notes(10).await;
2694 let mock_rpc = MockRpcApi::new(chain);
2695
2696 let blocks = mock_rpc
2697 .sync_notes_with_content(
2698 4_u32.into(),
2699 10_u32.into(),
2700 ¬e_tags,
2701 NoteContentFetch::PublicDetailsAndAttachments,
2702 )
2703 .await
2704 .expect("sync notes should succeed");
2705
2706 assert_eq!(blocks.last().unwrap().block_header.block_num(), BlockNumber::from(10u32));
2707 assert!(
2708 blocks
2709 .iter()
2710 .any(|block| block.block_header.block_num() == BlockNumber::from(9u32))
2711 );
2712 }
2713
2714 #[tokio::test]
2720 async fn erased_notes_are_marked_as_consumed() {
2721 let sender_id: AccountId = ACCOUNT_ID_SENDER.try_into().unwrap();
2723 let partial_metadata = PartialNoteMetadata::new(sender_id, NoteType::Public);
2724 let metadata = NoteMetadata::new(partial_metadata, &NoteAttachments::empty());
2725 let script = CodeBuilder::new()
2726 .compile_note_script("@note_script\npub proc main\n nop\nend")
2727 .unwrap();
2728 let recipient = NoteRecipient::new(
2729 Word::from([1u32, 2, 3, 4]),
2730 script,
2731 NoteStorage::new(vec![]).unwrap(),
2732 );
2733 let output_note = OutputNoteRecord::new(
2734 recipient.digest(),
2735 NoteAssets::new(vec![]).unwrap(),
2736 metadata,
2737 OutputNoteState::ExpectedFull { recipient },
2738 BlockNumber::from(1u32),
2739 NoteAttachments::default(),
2740 );
2741 let note_id = output_note.id();
2742 let note_header = NoteHeader::new(output_note.details_commitment(), metadata);
2743
2744 let mut note_updates = NoteUpdateTracker::new(vec![], vec![output_note]);
2746
2747 let block_num = BlockNumber::from(3u32);
2749 note_updates
2750 .mark_erased_note_as_consumed(¬e_header, block_num)
2751 .expect("marking erased note should succeed");
2752
2753 let updated = note_updates
2754 .updated_output_notes()
2755 .find(|n| n.id() == note_id)
2756 .expect("output note should be in the update");
2757
2758 assert!(
2759 updated.inner().is_consumed(),
2760 "output note should be consumed after erasure detection, but state is: {}",
2761 updated.inner().state()
2762 );
2763 }
2764
2765 #[allow(clippy::too_many_lines)]
2781 #[ignore = "consumer derivation removed; see comment above"]
2782 #[tokio::test]
2783 async fn erased_notes_are_marked_as_consumed_by_network_account() {
2784 let mut builder = MockChainBuilder::new();
2787 let p2id_sender: AccountId = ACCOUNT_ID_SENDER.try_into().unwrap();
2788 let faucet_id: AccountId = ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET.try_into().unwrap();
2789 let sender_account =
2790 builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2791 let sender_id = sender_account.id();
2792
2793 let asset = Asset::from(FungibleAsset::new(faucet_id, 100u64).unwrap());
2794 let note = builder
2795 .add_p2id_note(p2id_sender, sender_id, &[asset], NoteType::Public)
2796 .unwrap();
2797
2798 let mut chain = builder.build().unwrap();
2799 chain.prove_next_block().unwrap();
2800
2801 let tx = Box::pin(
2802 chain
2803 .build_transaction(MockTransactionInput::Account(sender_account.clone()))
2804 .unauthenticated_input_note(note.clone())
2805 .build()
2806 .unwrap()
2807 .execute(),
2808 )
2809 .await
2810 .unwrap();
2811 chain.add_pending_executed_transaction(&tx).unwrap();
2812 chain.prove_next_block().unwrap();
2813
2814 let network_account_id: AccountId =
2816 ACCOUNT_ID_REGULAR_PUBLIC_ACCOUNT_IMMUTABLE_CODE.try_into().unwrap();
2817 let target =
2818 NetworkAccountTarget::new(network_account_id, NoteExecutionHint::Always).unwrap();
2819 let attachment: NoteAttachment = target.into();
2820 let attachments = NoteAttachments::new(vec![attachment]).unwrap();
2821 let partial_metadata = PartialNoteMetadata::new(sender_id, NoteType::Public);
2822 let metadata = NoteMetadata::new(partial_metadata, &attachments);
2823 let script = CodeBuilder::new()
2824 .compile_note_script("@note_script\npub proc main\n nop\nend")
2825 .unwrap();
2826 let recipient = NoteRecipient::new(
2827 Word::from([7u32, 8, 9, 10]),
2828 script,
2829 NoteStorage::new(vec![]).unwrap(),
2830 );
2831 let recipient_digest = recipient.digest();
2832 let assets = NoteAssets::new(vec![]).unwrap();
2833
2834 let output_note = OutputNoteRecord::new(
2837 recipient_digest,
2838 assets.clone(),
2839 metadata,
2840 OutputNoteState::ExpectedFull { recipient },
2841 BlockNumber::from(1u32),
2842 NoteAttachments::default(),
2843 );
2844 let erased_note_id = output_note.id();
2845 let erased_note_header = NoteHeader::new(output_note.details_commitment(), metadata);
2846
2847 let mock_rpc = MockRpcApi::new(chain);
2848 mock_rpc.mark_note_as_erased(erased_note_header);
2849
2850 let network_header =
2853 AccountHeader::new(network_account_id, ZERO, EMPTY_WORD, EMPTY_WORD, EMPTY_WORD);
2854
2855 let state_sync = StateSync::new(
2856 Arc::new(mock_rpc.clone()),
2857 Arc::new(MockScreener),
2858 None,
2859 genesis_validator_config(&mock_rpc),
2860 );
2861
2862 let genesis_peaks =
2863 mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
2864 let mut partial_mmr = PartialMmr::from_peaks(genesis_peaks);
2865
2866 let sync_input = StateSyncInput {
2867 accounts: vec![AccountHeader::from(&sender_account), network_header],
2868 note_tags: BTreeSet::new(),
2869 input_notes: vec![],
2870 output_notes: vec![output_note],
2871 uncommitted_transactions: vec![],
2872 };
2873
2874 let update = state_sync.sync_state(&mut partial_mmr, sync_input).await.unwrap();
2875
2876 let updated_output = update
2878 .note_updates()
2879 .updated_output_notes()
2880 .find(|n| n.id() == erased_note_id)
2881 .expect("output note should be in the update");
2882 assert!(
2883 updated_output.inner().is_consumed(),
2884 "output note should be consumed, got: {}",
2885 updated_output.inner().state()
2886 );
2887
2888 let input_note_update = update
2890 .note_updates()
2891 .updated_input_notes()
2892 .find(|n| n.id() == Some(erased_note_id))
2893 .expect("input note should be created from the erased output note");
2894
2895 let inner = input_note_update.inner();
2896 assert!(
2897 inner.is_consumed(),
2898 "input note should be in a consumed state, got: {}",
2899 inner.state()
2900 );
2901 assert_eq!(
2902 inner.consumer_account(),
2903 Some(network_account_id),
2904 "consumer should be the tracked network account"
2905 );
2906 }
2907
2908 #[tokio::test]
2911 async fn validate_chain_mmr_response_rejects_tampered_responses() {
2912 let mock_rpc = MockRpcApi::default();
2913 mock_rpc.advance_blocks(3);
2914 let chain_tip = mock_rpc.get_chain_tip_block_num();
2915 let current = BlockNumber::GENESIS;
2916 let validator_config = genesis_validator_config(&mock_rpc);
2917
2918 let header_of =
2919 |block_num: u32| mock_rpc.mock_chain.read().block_header(block_num as usize);
2920 let chain_mmr_response = || async {
2921 mock_rpc.sync_chain_mmr(current, SyncTarget::CommittedChainTip).await.unwrap()
2922 };
2923
2924 let response = chain_mmr_response().await;
2926 StateSync::validate_chain_mmr_response(&response, current, &validator_config).unwrap();
2927
2928 let mut response = chain_mmr_response().await;
2930 response.block_header = header_of(chain_tip.as_u32() - 1);
2931 let result = StateSync::validate_chain_mmr_response(&response, current, &validator_config);
2932 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2933
2934 let mut response = chain_mmr_response().await;
2936 response.block_from = current + 1;
2937 let result = StateSync::validate_chain_mmr_response(&response, current, &validator_config);
2938 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2939
2940 let mut response = chain_mmr_response().await;
2942 response.block_from = chain_tip;
2943 response.block_to = BlockNumber::GENESIS;
2944 response.block_header = header_of(0);
2945 let result =
2946 StateSync::validate_chain_mmr_response(&response, chain_tip, &validator_config);
2947 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2948 }
2949
2950 #[test]
2953 fn validate_note_blocks_range_rejects_out_of_range_blocks() {
2954 let mock_rpc = MockRpcApi::default();
2955 mock_rpc.advance_blocks(3);
2956 let chain_tip = mock_rpc.get_chain_tip_block_num();
2957 let current = BlockNumber::GENESIS;
2958
2959 StateSync::validate_note_blocks_range(&[], current, chain_tip).unwrap();
2961
2962 let genesis_note_block = ResolvedSyncNotesBlock {
2964 block_header: mock_rpc.mock_chain.read().block_header(0),
2965 mmr_path: MerklePath::new(Vec::new()),
2966 notes: BTreeMap::new(),
2967 };
2968 let result =
2969 StateSync::validate_note_blocks_range(&[genesis_note_block], current, chain_tip);
2970 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2971 }
2972
2973 #[tokio::test]
2976 async fn validate_chain_mmr_response_rejects_chain_tip_without_valid_signatures() {
2977 let mock_rpc = MockRpcApi::default();
2978 mock_rpc.advance_blocks(3);
2979 let chain_tip = mock_rpc.get_chain_tip_block_num();
2980 let current = BlockNumber::GENESIS;
2981
2982 let validator_config = genesis_validator_config(&mock_rpc);
2983 assert!(!validator_config.is_empty(), "the mock chain must commit a validator set");
2984
2985 let mut chain_mmr_info =
2987 mock_rpc.sync_chain_mmr(current, SyncTarget::CommittedChainTip).await.unwrap();
2988 StateSync::validate_chain_mmr_response(&chain_mmr_info, current, &validator_config)
2989 .unwrap();
2990
2991 let parent = BlockNumber::from(chain_tip.as_u32() - 1);
2992 chain_mmr_info.block_signatures = block_signatures(&mock_rpc, parent);
2994 let result =
2996 StateSync::validate_chain_mmr_response(&chain_mmr_info, current, &validator_config);
2997 assert!(
2998 matches!(result, Err(ClientError::ChainValidationError(_))),
2999 "signatures of another block must be rejected, got {result:?}"
3000 );
3001 }
3002
3003 #[test]
3006 fn advance_mmr_rejects_delta_inconsistent_with_chain_commitment() {
3007 let mock_rpc = MockRpcApi::default();
3008 mock_rpc.advance_blocks(3);
3009 let chain_tip = mock_rpc.get_chain_tip_block_num();
3010
3011 let chain_tip_header = mock_rpc.mock_chain.read().block_header(chain_tip.as_usize());
3012 let genesis_partial_mmr = || {
3013 let peaks = mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
3014 PartialMmr::from_peaks(peaks)
3015 };
3016
3017 let full_delta = mock_rpc
3019 .get_mmr()
3020 .get_delta(Forest::new(1).unwrap(), Forest::new(chain_tip.as_usize()).unwrap())
3021 .unwrap();
3022 StateSync::advance_mmr(
3023 full_delta,
3024 &chain_tip_header,
3025 &mut genesis_partial_mmr(),
3026 &mut PartialBlockchainUpdates::default(),
3027 )
3028 .unwrap();
3029
3030 let truncated_delta = mock_rpc
3032 .get_mmr()
3033 .get_delta(Forest::new(1).unwrap(), Forest::new(chain_tip.as_usize() - 1).unwrap())
3034 .unwrap();
3035 let result = StateSync::advance_mmr(
3036 truncated_delta,
3037 &chain_tip_header,
3038 &mut genesis_partial_mmr(),
3039 &mut PartialBlockchainUpdates::default(),
3040 );
3041 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
3042 }
3043
3044 fn make_tx_record(account_id: AccountId, block_num: u32) -> RpcTransactionRecord {
3046 RpcTransactionRecord {
3047 block_num: BlockNumber::from(block_num),
3048 transaction_header: TransactionHeader::new(
3049 account_id,
3050 word(1),
3051 word(2),
3052 InputNotes::new_unchecked(vec![]),
3053 vec![],
3054 )
3055 .unwrap(),
3056 output_notes: vec![],
3057 erased_output_notes: vec![],
3058 consumed_note_refs: vec![],
3059 }
3060 }
3061}