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, AccountWitness};
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
45pub(crate) const 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<(AccountWitness, 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 let mut account_witnesses = Vec::with_capacity(synced_accounts.len());
1140 for ((account_id, local_header), (witness, synced_account)) in
1141 diverging_accounts.iter().zip(synced_accounts)
1142 {
1143 account_witnesses.push((*account_id, witness));
1144
1145 match synced_account {
1146 PublicAccountSync::Apply(public_update) => {
1147 account_updates.extend(AccountUpdates::new(vec![*public_update], Vec::new()));
1148 },
1149 PublicAccountSync::Superseded => {
1150 superseded_states.push(local_header.to_commitment());
1151 },
1152 PublicAccountSync::Ignore => {},
1153 }
1154 }
1155
1156 account_updates.extend(AccountUpdates::default().with_account_witnesses(account_witnesses));
1157
1158 Ok(superseded_states)
1159 }
1160
1161 async fn sync_public_account(
1172 &self,
1173 account_id: AccountId,
1174 local_header: &AccountHeader,
1175 block_from: BlockNumber,
1176 chain_tip_header: &BlockHeader,
1177 ) -> Result<(AccountWitness, PublicAccountSync), ClientError> {
1178 let target_block_num = chain_tip_header.block_num();
1179
1180 let (proof_block_num, proof) = self
1183 .rpc_api
1184 .get_account(
1185 account_id,
1186 GetAccountRequest::new()
1187 .at(AccountStateAt::Block(target_block_num))
1188 .with_storage(StorageMapFetch::All)
1189 .with_vault(VaultFetch::Always),
1190 )
1191 .await
1192 .map_err(ClientError::RpcError)?;
1193
1194 let (witness, details) =
1195 Self::validate_account_proof(proof, proof_block_num, account_id, chain_tip_header)?;
1196
1197 match details
1198 .header
1199 .nonce()
1200 .as_canonical_u64()
1201 .cmp(&local_header.nonce().as_canonical_u64())
1202 {
1203 Ordering::Less => return Ok((witness, PublicAccountSync::Ignore)),
1206 Ordering::Equal => return Ok((witness, PublicAccountSync::Superseded)),
1208 Ordering::Greater => {},
1210 }
1211
1212 let vault_oversized = details.vault_details.too_many_assets;
1213 let any_map_oversized = details
1214 .storage_details
1215 .map_details
1216 .iter()
1217 .any(AccountStorageMapDetails::is_limit_exceeded);
1218
1219 let public_update = if vault_oversized || any_map_oversized {
1223 self.build_patch_update(account_id, local_header, &details, block_from, proof_block_num)
1225 .await?
1226 } else {
1227 let account = Account::try_from(&details).map_err(ClientError::RpcError)?;
1229 PublicAccountUpdate::Full(account)
1230 };
1231
1232 Ok((witness, PublicAccountSync::Apply(Box::new(public_update))))
1233 }
1234
1235 fn validate_account_proof(
1248 proof: AccountProof,
1249 proof_block_num: BlockNumber,
1250 account_id: AccountId,
1251 chain_tip_header: &BlockHeader,
1252 ) -> Result<(AccountWitness, AccountDetails), ClientError> {
1253 let target_block_num = chain_tip_header.block_num();
1254
1255 if proof_block_num != target_block_num {
1256 return Err(ClientError::ChainValidationError(format!(
1257 "get_account returned block {proof_block_num} but {target_block_num} was requested"
1258 )));
1259 }
1260
1261 let (witness, details) = proof.into_parts();
1262
1263 validate_account_witness(&witness, account_id, chain_tip_header)?;
1264
1265 let details = details.ok_or_else(|| {
1266 ClientError::ChainValidationError(format!(
1267 "get_account returned no details for public account {account_id}"
1268 ))
1269 })?;
1270
1271 Ok((witness, details))
1272 }
1273
1274 async fn build_patch_update(
1292 &self,
1293 account_id: AccountId,
1294 local_header: &AccountHeader,
1295 details: &AccountDetails,
1296 block_from: BlockNumber,
1297 block_to: BlockNumber,
1298 ) -> Result<PublicAccountUpdate, ClientError> {
1299 let value_slot_updates: Vec<(_, Word)> = details
1300 .storage_details
1301 .header
1302 .slots()
1303 .filter(|slot| slot.slot_type() == StorageSlotType::Value)
1304 .map(|slot| (slot.name().clone(), slot.value()))
1305 .collect();
1306
1307 let map_info = self
1310 .rpc_api
1311 .sync_storage_maps(block_from + 1, block_to, account_id)
1312 .await
1313 .map_err(ClientError::RpcError)?;
1314 let vault_info = self
1315 .rpc_api
1316 .sync_account_vault(block_from + 1, block_to, account_id)
1317 .await
1318 .map_err(ClientError::RpcError)?;
1319
1320 let patch = build_account_patch(
1321 &details.header,
1322 value_slot_updates,
1323 map_info.map_entries,
1324 vault_info.vault_patch,
1325 details.code.clone(),
1326 local_header.code_commitment(),
1327 )
1328 .map_err(StoreError::AccountPatchError)?;
1329
1330 Ok(PublicAccountUpdate::Patch {
1331 new_header: details.header.clone(),
1332 patch,
1333 })
1334 }
1335
1336 async fn note_state_sync(
1354 &self,
1355 note_updates: &mut NoteUpdateTracker,
1356 notes: BTreeMap<NoteId, SyncedNote>,
1357 block_header: &BlockHeader,
1358 ) -> Result<NoteBlockRelevance, ClientError> {
1359 let mut relevance = NoteBlockRelevance::default();
1360
1361 for (_, mut note) in notes {
1362 if !self.note_observers.is_empty() {
1365 for obs in &self.note_observers {
1366 match obs.observe(¬e).await {
1367 Ok(true) => relevance.observer_requires_block = true,
1368 Ok(false) => {},
1369 Err(err) => {
1370 tracing::warn!(
1371 observer = obs.name(),
1372 error = ?err,
1373 "note observer failed; sync continues",
1374 );
1375 },
1376 }
1377 }
1378 }
1379
1380 let public_note = note.details.take().map(|details| {
1383 let state = UnverifiedNoteState {
1384 metadata: note.metadata,
1385 inclusion_proof: note.inclusion_proof.clone(),
1386 }
1387 .into();
1388 InputNoteRecord::new(details, note.attachments.clone(), None, state)
1389 });
1390
1391 let committed = note.into_committed_note();
1392
1393 match self.note_screener.on_note_received(committed, public_note).await? {
1394 NoteUpdateAction::Commit(committed_note) => {
1395 relevance.has_client_note |= note_updates
1399 .apply_committed_note_state_transitions(&committed_note, block_header)?;
1400 },
1401 NoteUpdateAction::Insert(public_note) => {
1402 relevance.has_client_note = true;
1403
1404 note_updates.apply_new_public_note(public_note, block_header)?;
1405 },
1406 NoteUpdateAction::Discard => {},
1407 }
1408 }
1409
1410 Ok(relevance)
1411 }
1412
1413 async fn nullifiers_state_sync(
1419 &self,
1420 note_updates: &mut NoteUpdateTracker,
1421 transaction_updates: &mut TransactionUpdateTracker,
1422 chain_tip: BlockNumber,
1423 current_block_num: BlockNumber,
1424 ) -> Result<(), ClientError> {
1425 let nullifiers_tags: Vec<u16> =
1431 note_updates.unspent_nullifiers().map(|nullifier| nullifier.prefix()).collect();
1432
1433 let mut new_nullifiers = self
1434 .rpc_api
1435 .sync_nullifiers(&nullifiers_tags, current_block_num + 1, chain_tip)
1436 .await?;
1437
1438 new_nullifiers.retain(|update| update.block_num <= chain_tip);
1441
1442 let consumptions: Vec<NoteConsumption> = new_nullifiers
1444 .into_iter()
1445 .map(|update| NoteConsumption {
1446 external_consumer: transaction_updates
1447 .external_nullifier_account(&update.nullifier),
1448 nullifier: update.nullifier,
1449 block_num: update.block_num,
1450 })
1451 .collect();
1452
1453 for consumption in consumptions {
1454 note_updates.apply_note_consumption(
1455 &consumption,
1456 transaction_updates.committed_transactions(),
1457 )?;
1458
1459 transaction_updates.apply_input_note_nullified(consumption.nullifier);
1463 }
1464
1465 Ok(())
1466 }
1467}
1468
1469pub struct ChainSyncData {
1478 pub(crate) block_from: BlockNumber,
1480 advance: Option<ChainAdvance>,
1483 superseded_states: Vec<Word>,
1485 pub(crate) note_updates: NoteUpdateTracker,
1488 transaction_updates: TransactionUpdateTracker,
1489 pub(crate) account_updates: AccountUpdates,
1492}
1493
1494impl ChainSyncData {
1495 pub(crate) fn chain_tip_header(&self) -> Option<&BlockHeader> {
1498 self.advance.as_ref().map(|advance| &advance.chain_tip_header)
1499 }
1500}
1501
1502struct ChainAdvance {
1504 chain_tip_header: BlockHeader,
1506 mmr_delta: MmrDelta,
1508 note_blocks_awaiting_screening: Vec<ResolvedSyncNotesBlock>,
1511 transactions: Vec<RpcTransactionRecord>,
1513 relevant_note_blocks: Vec<RelevantNoteBlock>,
1515 protocol_config: Option<ProtocolConfig>,
1517}
1518
1519pub(crate) fn validate_account_witness(
1532 witness: &AccountWitness,
1533 account_id: AccountId,
1534 chain_tip_header: &BlockHeader,
1535) -> Result<(), ClientError> {
1536 if witness.id() != account_id {
1538 return Err(ClientError::ChainValidationError(format!(
1539 "get_account returned account {} but {account_id} was requested",
1540 witness.id()
1541 )));
1542 }
1543
1544 let account_key = AccountIdKey::from(account_id).as_word();
1545 let state_commitment = witness.state_commitment();
1546 witness
1547 .clone()
1548 .into_proof()
1549 .verify_presence(&account_key, &state_commitment, &chain_tip_header.account_root())
1550 .map_err(|err| {
1551 ClientError::ChainValidationError(format!(
1552 "get_account witness for account {account_id} does not open under block {} \
1553 account root: {err}",
1554 chain_tip_header.block_num()
1555 ))
1556 })
1557}
1558
1559pub(crate) fn block_num_from_forest(partial_mmr: &PartialMmr) -> Result<BlockNumber, ClientError> {
1561 Ok(u32::try_from(partial_mmr.forest().num_leaves().saturating_sub(1))
1562 .map_err(|_| ClientError::InvalidPartialMmrForest)?
1563 .into())
1564}
1565
1566fn group_txs_by_account_block(
1568 transaction_records: &[RpcTransactionRecord],
1569) -> BTreeMap<(AccountId, BlockNumber), Vec<&RpcTransactionRecord>> {
1570 let mut groups: BTreeMap<(AccountId, BlockNumber), Vec<&RpcTransactionRecord>> =
1571 BTreeMap::new();
1572 for record in transaction_records {
1573 let account_id = record.transaction_header.account_id();
1574 groups.entry((account_id, record.block_num)).or_default().push(record);
1575 }
1576 groups
1577}
1578
1579fn walk_execution_chain<'a>(
1585 txs: &'a [&'a RpcTransactionRecord],
1586) -> impl Iterator<Item = &'a RpcTransactionRecord> + 'a {
1587 let (self_loops, chained): (Vec<&RpcTransactionRecord>, Vec<&RpcTransactionRecord>) =
1588 txs.iter().copied().partition(|tx| {
1589 tx.transaction_header.initial_state_commitment()
1590 == tx.transaction_header.final_state_commitment()
1591 });
1592
1593 let final_states: BTreeSet<Word> = chained
1594 .iter()
1595 .map(|tx| tx.transaction_header.final_state_commitment())
1596 .collect();
1597
1598 let mut init_to_tx: BTreeMap<Word, &RpcTransactionRecord> = chained
1599 .iter()
1600 .map(|tx| (tx.transaction_header.initial_state_commitment(), *tx))
1601 .collect();
1602
1603 let start = chained
1604 .iter()
1605 .find(|tx| !final_states.contains(&tx.transaction_header.initial_state_commitment()))
1606 .copied();
1607
1608 assert!(start.is_some() || chained.is_empty(), "cannot walk cyclic execution chain");
1609
1610 let mut current =
1611 start.and_then(|tx| init_to_tx.remove(&tx.transaction_header.initial_state_commitment()));
1612 let mut self_loops_iter = self_loops.into_iter();
1613
1614 core::iter::from_fn(move || {
1615 if let Some(tx) = current {
1616 current = init_to_tx.remove(&tx.transaction_header.final_state_commitment());
1617 return Some(tx);
1618 }
1619 self_loops_iter.next()
1620 })
1621}
1622
1623fn derive_account_commitments(
1628 transaction_records: &[RpcTransactionRecord],
1629) -> Vec<(AccountId, Word)> {
1630 let mut latest_by_account: BTreeMap<AccountId, (BlockNumber, Word)> = BTreeMap::new();
1631
1632 for ((account_id, block_num), txs) in &group_txs_by_account_block(transaction_records) {
1633 let terminal_state = walk_execution_chain(txs)
1634 .last()
1635 .expect("account must have a final state")
1636 .transaction_header
1637 .final_state_commitment();
1638
1639 latest_by_account
1640 .entry(*account_id)
1641 .and_modify(|(existing_block, existing_state)| {
1642 if *block_num > *existing_block {
1643 *existing_block = *block_num;
1644 *existing_state = terminal_state;
1645 }
1646 })
1647 .or_insert((*block_num, terminal_state));
1648 }
1649
1650 latest_by_account
1651 .into_iter()
1652 .map(|(account_id, (_, state))| (account_id, state))
1653 .collect()
1654}
1655
1656fn compute_ordered_nullifiers(transaction_records: &[RpcTransactionRecord]) -> Vec<Nullifier> {
1663 let mut result = Vec::new();
1664
1665 for txs in group_txs_by_account_block(transaction_records).values() {
1666 for tx in walk_execution_chain(txs) {
1667 for commitment in tx.transaction_header.input_notes().iter() {
1668 result.push(commitment.nullifier());
1669 }
1670 }
1671 }
1672
1673 result
1674}
1675
1676#[cfg(all(test, feature = "testing"))]
1677mod tests {
1678 use alloc::collections::BTreeSet;
1679 use alloc::sync::Arc;
1680
1681 use async_trait::async_trait;
1682 use miden_protocol::account::Account;
1683 use miden_protocol::assembly::DefaultSourceManager;
1684 use miden_protocol::asset::{Asset, FungibleAsset};
1685 use miden_protocol::block::{BlockNumber, BlockSignatures};
1686 use miden_protocol::crypto::merkle::MerklePath;
1687 use miden_protocol::crypto::merkle::mmr::{Forest, InOrderIndex, PartialMmr};
1688 use miden_protocol::note::{
1689 Note,
1690 NoteAssets,
1691 NoteAttachment,
1692 NoteAttachments,
1693 NoteDetails,
1694 NoteHeader,
1695 NoteMetadata,
1696 NoteRecipient,
1697 NoteStorage,
1698 NoteTag,
1699 NoteType,
1700 PartialNoteMetadata,
1701 };
1702 use miden_protocol::testing::account_id::{
1703 ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET,
1704 ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET,
1705 ACCOUNT_ID_REGULAR_PRIVATE_ACCOUNT_UPDATABLE_CODE,
1706 ACCOUNT_ID_REGULAR_PUBLIC_ACCOUNT_IMMUTABLE_CODE,
1707 ACCOUNT_ID_SENDER,
1708 };
1709 use miden_protocol::transaction::{InputNotes, TransactionArgs, TransactionHeader};
1710 use miden_protocol::vm::AdviceMap;
1711 use miden_protocol::{EMPTY_WORD, Felt, Word, ZERO};
1712 use miden_standards::code_builder::CodeBuilder;
1713 use miden_standards::note::{NetworkAccountTarget, NoteExecutionHint};
1714 use miden_testing::{MockChainBuilder, MockTransactionInput};
1715
1716 use super::*;
1717 use crate::store::{OutputNoteRecord, OutputNoteState};
1718 use crate::test_utils::mock::MockRpcApi;
1719
1720 struct MockScreener;
1722
1723 #[async_trait(?Send)]
1724 impl OnNoteReceived for MockScreener {
1725 async fn on_note_received(
1726 &self,
1727 _committed_note: CommittedNote,
1728 _public_note: Option<InputNoteRecord>,
1729 ) -> Result<NoteUpdateAction, ClientError> {
1730 Ok(NoteUpdateAction::Discard)
1731 }
1732 }
1733
1734 struct AlwaysRelevantObserver;
1736
1737 #[async_trait(?Send)]
1738 impl NoteObserver for AlwaysRelevantObserver {
1739 fn name(&self) -> &'static str {
1740 "always-relevant"
1741 }
1742
1743 async fn observe(&self, _note: &SyncedNote) -> Result<bool, ClientError> {
1744 Ok(true)
1745 }
1746 }
1747
1748 fn genesis_validator_config(mock_rpc: &MockRpcApi) -> ValidatorConfig {
1750 mock_rpc.mock_chain.read().block_header(0).validator_config().clone()
1751 }
1752
1753 fn block_signatures(mock_rpc: &MockRpcApi, block_num: BlockNumber) -> BlockSignatures {
1755 mock_rpc
1756 .mock_chain
1757 .read()
1758 .proven_blocks()
1759 .iter()
1760 .find(|block| block.header().block_num() == block_num)
1761 .expect("the mock chain contains the block")
1762 .signatures()
1763 .clone()
1764 }
1765
1766 fn empty() -> StateSyncInput {
1767 StateSyncInput {
1768 accounts: vec![],
1769 note_tags: BTreeSet::new(),
1770 input_notes: vec![],
1771 output_notes: vec![],
1772 uncommitted_transactions: vec![],
1773 }
1774 }
1775
1776 fn word(n: u64) -> miden_protocol::Word {
1777 [
1778 Felt::new(n).expect("test value should fit into the base field"),
1779 ZERO,
1780 ZERO,
1781 ZERO,
1782 ]
1783 .into()
1784 }
1785
1786 fn header_with_account_root(header: &BlockHeader, account_root: Word) -> BlockHeader {
1787 BlockHeader::new(
1788 header.prev_block_commitment(),
1789 header.block_num(),
1790 header.chain_commitment(),
1791 account_root,
1792 header.nullifier_root(),
1793 header.note_root(),
1794 header.tx_commitment(),
1795 header.validator_config().clone(),
1796 header.fee_parameters().clone(),
1797 header.protocol_config_commitment(),
1798 header.next_protocol_config().cloned(),
1799 header.timestamp(),
1800 )
1801 }
1802
1803 #[tokio::test]
1804 async fn sync_public_accounts_ignores_older_node_snapshot() {
1805 let mut builder = MockChainBuilder::new();
1806 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1807 let rpc_api = MockRpcApi::new(builder.build().unwrap());
1808 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
1809 let validator_config = genesis_validator_config(&rpc_api);
1810 let state_sync =
1811 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1812
1813 let local_header =
1816 AccountHeader::new(account.id(), Felt::from(2u32), EMPTY_WORD, EMPTY_WORD, EMPTY_WORD);
1817 let current_public_accounts = vec![&local_header];
1818 let commitment_updates = vec![(account.id(), account.to_commitment())];
1819 let mut account_updates = AccountUpdates::default();
1820
1821 let superseded = state_sync
1822 .sync_public_accounts(
1823 &mut account_updates,
1824 &commitment_updates,
1825 ¤t_public_accounts,
1826 BlockNumber::GENESIS,
1827 &chain_tip_header,
1828 )
1829 .await
1830 .unwrap();
1831
1832 assert!(
1833 account_updates.updated_public_accounts().is_empty(),
1834 "public account sync should ignore node snapshots that are older than local"
1835 );
1836 assert!(
1837 superseded.is_empty(),
1838 "an older node snapshot must not supersede the local state"
1839 );
1840 }
1841
1842 #[tokio::test]
1843 async fn sync_public_accounts_marks_same_nonce_mismatch_as_superseded() {
1844 let mut builder = MockChainBuilder::new();
1845 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1846 let rpc_api = MockRpcApi::new(builder.build().unwrap());
1847 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
1848 let validator_config = genesis_validator_config(&rpc_api);
1849 let state_sync =
1850 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1851
1852 let local_header =
1855 AccountHeader::new(account.id(), account.nonce(), EMPTY_WORD, EMPTY_WORD, EMPTY_WORD);
1856 let current_public_accounts = vec![&local_header];
1857 let commitment_updates = vec![(account.id(), account.to_commitment())];
1858 let mut account_updates = AccountUpdates::default();
1859
1860 let superseded = state_sync
1861 .sync_public_accounts(
1862 &mut account_updates,
1863 &commitment_updates,
1864 ¤t_public_accounts,
1865 BlockNumber::GENESIS,
1866 &chain_tip_header,
1867 )
1868 .await
1869 .unwrap();
1870
1871 assert!(
1872 account_updates.updated_public_accounts().is_empty(),
1873 "a same-nonce fork must not overwrite the account while its tx is still pending"
1874 );
1875 assert_eq!(
1876 superseded,
1877 vec![local_header.to_commitment()],
1878 "the superseded local state should be reported so its transaction is discarded"
1879 );
1880 }
1881
1882 #[test]
1888 fn validate_transaction_records_range_rejects_out_of_range_blocks() {
1889 let account_id: AccountId = ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET.try_into().unwrap();
1890 let current = BlockNumber::from(5u32);
1891 let chain_tip = BlockNumber::from(10u32);
1892
1893 StateSync::validate_transaction_records_range(
1894 &[make_tx_record(account_id, 7)],
1895 current,
1896 chain_tip,
1897 )
1898 .unwrap();
1899
1900 let result = StateSync::validate_transaction_records_range(
1901 &[make_tx_record(account_id, 11)],
1902 current,
1903 chain_tip,
1904 );
1905 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
1906
1907 let result = StateSync::validate_transaction_records_range(
1908 &[make_tx_record(account_id, 5)],
1909 current,
1910 chain_tip,
1911 );
1912 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
1913 }
1914
1915 #[tokio::test]
1918 async fn verify_private_account_mismatch_ignores_forged_commitment() {
1919 let mut builder = MockChainBuilder::new();
1920 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1921 let rpc_api = MockRpcApi::new(builder.build().unwrap());
1922 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
1923 let on_chain_commitment = account.to_commitment();
1924 let validator_config = genesis_validator_config(&rpc_api);
1925 let state_sync =
1926 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1927
1928 let result = state_sync
1929 .verify_private_account_mismatch(account.id(), on_chain_commitment, &chain_tip_header)
1930 .await
1931 .unwrap();
1932
1933 assert!(
1934 result.is_none(),
1935 "an unproven commitment must not lock an account whose on-chain state matches local"
1936 );
1937 }
1938
1939 #[tokio::test]
1942 async fn verify_private_account_mismatch_reports_proven_divergence() {
1943 let mut builder = MockChainBuilder::new();
1944 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1945 let rpc_api = MockRpcApi::new(builder.build().unwrap());
1946 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
1947 let on_chain_commitment = account.to_commitment();
1948 let validator_config = genesis_validator_config(&rpc_api);
1949 let state_sync =
1950 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1951 let stale_local_commitment = word(0xdead_beef);
1952
1953 let result = state_sync
1954 .verify_private_account_mismatch(
1955 account.id(),
1956 stale_local_commitment,
1957 &chain_tip_header,
1958 )
1959 .await
1960 .unwrap();
1961
1962 assert_eq!(
1963 result,
1964 Some(on_chain_commitment),
1965 "a proven divergence should return the proven commitment to lock with"
1966 );
1967 }
1968
1969 #[tokio::test]
1972 async fn verify_private_account_mismatch_rejects_unverifiable_proof() {
1973 let mut builder = MockChainBuilder::new();
1974 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
1975 let rpc_api = MockRpcApi::new(builder.build().unwrap());
1976 let real_header = rpc_api.mock_chain.read().latest_block_header();
1977 let validator_config = genesis_validator_config(&rpc_api);
1978 let state_sync =
1979 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
1980
1981 let tampered_header = BlockHeader::new(
1984 real_header.prev_block_commitment(),
1985 real_header.block_num(),
1986 real_header.chain_commitment(),
1987 word(0xbad0_bad0),
1988 real_header.nullifier_root(),
1989 real_header.note_root(),
1990 real_header.tx_commitment(),
1991 real_header.validator_config().clone(),
1992 real_header.fee_parameters().clone(),
1993 real_header.protocol_config_commitment(),
1994 real_header.next_protocol_config().cloned(),
1995 real_header.timestamp(),
1996 );
1997
1998 let result = state_sync
1999 .verify_private_account_mismatch(
2000 account.id(),
2001 account.to_commitment(),
2002 &tampered_header,
2003 )
2004 .await;
2005 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2006 }
2007
2008 #[tokio::test]
2011 async fn sync_public_accounts_pins_account_fetch_to_sync_target() {
2012 let mut builder = MockChainBuilder::new();
2013 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2014 let mut chain = builder.build().unwrap();
2015
2016 let sync_target_header = chain.latest_block_header();
2018 let tx = Box::pin(
2019 chain
2020 .build_transaction(MockTransactionInput::AccountId(account.id()))
2021 .build()
2022 .unwrap()
2023 .execute(),
2024 )
2025 .await
2026 .unwrap();
2027 let local_header = tx.final_account().clone();
2028 assert_ne!(local_header.to_commitment(), account.to_commitment());
2029 chain.add_pending_executed_transaction(&tx).unwrap();
2030
2031 let rpc_api = MockRpcApi::new(chain);
2032 rpc_api.prove_block();
2034 assert_eq!(
2035 rpc_api
2036 .mock_chain
2037 .read()
2038 .committed_account(account.id())
2039 .unwrap()
2040 .to_commitment(),
2041 local_header.to_commitment()
2042 );
2043 let validator_config = genesis_validator_config(&rpc_api);
2044 let state_sync =
2045 StateSync::new(Arc::new(rpc_api), Arc::new(MockScreener), None, validator_config);
2046
2047 let current_public_accounts = vec![&local_header];
2048 let commitment_updates = vec![(account.id(), account.to_commitment())];
2049 let mut account_updates = AccountUpdates::default();
2050
2051 let superseded = state_sync
2052 .sync_public_accounts(
2053 &mut account_updates,
2054 &commitment_updates,
2055 ¤t_public_accounts,
2056 BlockNumber::GENESIS,
2057 &sync_target_header,
2058 )
2059 .await
2060 .unwrap();
2061
2062 assert!(superseded.is_empty(), "the transaction must not be superseded");
2063 assert!(
2064 account_updates.updated_public_accounts().is_empty(),
2065 "the target state must not overwrite the local account"
2066 );
2067 }
2068
2069 async fn get_account_proof(
2071 rpc_api: &MockRpcApi,
2072 account_id: AccountId,
2073 ) -> (BlockNumber, AccountProof) {
2074 rpc_api
2075 .get_account(
2076 account_id,
2077 GetAccountRequest::new()
2078 .with_storage(StorageMapFetch::All)
2079 .with_vault(VaultFetch::Always),
2080 )
2081 .await
2082 .unwrap()
2083 }
2084
2085 #[tokio::test]
2087 async fn validate_account_proof_rejects_mismatched_account() {
2088 let mut builder = MockChainBuilder::new();
2089 let account_a = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2090 let account_b = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2091 let rpc_api = MockRpcApi::new(builder.build().unwrap());
2092 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
2093
2094 let (proof_block_num, proof) = get_account_proof(&rpc_api, account_b.id()).await;
2096 let result = StateSync::validate_account_proof(
2097 proof,
2098 proof_block_num,
2099 account_a.id(),
2100 &chain_tip_header,
2101 );
2102
2103 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2104 }
2105
2106 #[tokio::test]
2108 async fn validate_account_proof_rejects_wrong_account_root() {
2109 let mut builder = MockChainBuilder::new();
2110 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2111 let rpc_api = MockRpcApi::new(builder.build().unwrap());
2112 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
2113 let wrong_header = header_with_account_root(&chain_tip_header, word(999));
2114
2115 let (proof_block_num, proof) = get_account_proof(&rpc_api, account.id()).await;
2117 let result =
2118 StateSync::validate_account_proof(proof, proof_block_num, account.id(), &wrong_header);
2119
2120 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2121 }
2122
2123 #[tokio::test]
2125 async fn validate_account_proof_rejects_missing_details() {
2126 let mut builder = MockChainBuilder::new();
2127 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2128 let rpc_api = MockRpcApi::new(builder.build().unwrap());
2129 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
2130
2131 let (proof_block_num, proof) = get_account_proof(&rpc_api, account.id()).await;
2133 let (witness, _) = proof.into_parts();
2134 let proof = AccountProof::new(witness, None).unwrap();
2135 let result = StateSync::validate_account_proof(
2136 proof,
2137 proof_block_num,
2138 account.id(),
2139 &chain_tip_header,
2140 );
2141
2142 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2143 }
2144
2145 #[tokio::test]
2147 async fn validate_account_proof_rejects_wrong_block() {
2148 let mut builder = MockChainBuilder::new();
2149 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2150 let rpc_api = MockRpcApi::new(builder.build().unwrap());
2151 let chain_tip_header = rpc_api.mock_chain.read().latest_block_header();
2152
2153 let (proof_block_num, proof) = get_account_proof(&rpc_api, account.id()).await;
2155 let result = StateSync::validate_account_proof(
2156 proof,
2157 proof_block_num + 1,
2158 account.id(),
2159 &chain_tip_header,
2160 );
2161
2162 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2163 }
2164
2165 mod compute_nullifiers_tests {
2169 use alloc::vec;
2170
2171 use miden_protocol::block::BlockNumber;
2172 use miden_protocol::note::Nullifier;
2173 use miden_protocol::transaction::{InputNoteCommitment, InputNotes, TransactionHeader};
2174
2175 use super::word;
2176 use crate::rpc::domain::transaction::TransactionRecord as RpcTransactionRecord;
2177
2178 fn make_rpc_tx(
2179 init_state: u64,
2180 final_state: u64,
2181 nullifier_vals: &[u64],
2182 block_number: u32,
2183 ) -> RpcTransactionRecord {
2184 let account_id = miden_protocol::account::AccountId::try_from(
2185 miden_protocol::testing::account_id::ACCOUNT_ID_REGULAR_PRIVATE_ACCOUNT_UPDATABLE_CODE,
2186 )
2187 .unwrap();
2188
2189 let input_notes = InputNotes::new_unchecked(
2190 nullifier_vals
2191 .iter()
2192 .map(|v| InputNoteCommitment::from(Nullifier::from_raw(word(*v))))
2193 .collect(),
2194 );
2195
2196 RpcTransactionRecord {
2197 block_num: BlockNumber::from(block_number),
2198 transaction_header: TransactionHeader::new(
2199 account_id,
2200 word(init_state),
2201 word(final_state),
2202 input_notes,
2203 vec![],
2204 )
2205 .unwrap(),
2206 output_notes: vec![],
2207 erased_output_notes: vec![],
2208 consumed_note_refs: vec![],
2209 }
2210 }
2211
2212 #[test]
2213 fn chains_rpc_transactions_by_state_commitment() {
2214 let tx_a = make_rpc_tx(1, 2, &[10], 5);
2218 let tx_b = make_rpc_tx(2, 3, &[20], 5);
2219 let tx_c = make_rpc_tx(3, 4, &[30], 5);
2220
2221 let result = super::super::compute_ordered_nullifiers(&[tx_c, tx_a, tx_b]);
2222
2223 assert_eq!(result[0], Nullifier::from_raw(word(10)));
2224 assert_eq!(result[1], Nullifier::from_raw(word(20)));
2225 assert_eq!(result[2], Nullifier::from_raw(word(30)));
2226 }
2227
2228 #[test]
2229 fn groups_independently_by_account_and_block() {
2230 let tx_a1 = make_rpc_tx(1, 2, &[10], 5);
2232 let tx_a2 = make_rpc_tx(2, 3, &[20], 5);
2233
2234 let tx_a3 = make_rpc_tx(3, 4, &[30], 6);
2236
2237 let account_b = miden_protocol::account::AccountId::try_from(
2239 miden_protocol::testing::account_id::ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET,
2240 )
2241 .unwrap();
2242
2243 let tx_b1 = RpcTransactionRecord {
2244 block_num: BlockNumber::from(5u32),
2245 transaction_header: TransactionHeader::new(
2246 account_b,
2247 word(100),
2248 word(200),
2249 InputNotes::new_unchecked(vec![InputNoteCommitment::from(
2250 Nullifier::from_raw(word(40)),
2251 )]),
2252 vec![],
2253 )
2254 .unwrap(),
2255 output_notes: vec![],
2256 erased_output_notes: vec![],
2257 consumed_note_refs: vec![],
2258 };
2259
2260 let result = super::super::compute_ordered_nullifiers(&[tx_a2, tx_b1, tx_a3, tx_a1]);
2261
2262 let pos = |val: u64| -> usize {
2265 result.iter().position(|n| *n == Nullifier::from_raw(word(val))).unwrap()
2266 };
2267
2268 assert!(pos(10) < pos(20)); assert!(result.contains(&Nullifier::from_raw(word(30)))); assert!(result.contains(&Nullifier::from_raw(word(40)))); }
2274
2275 #[test]
2276 fn multiple_nullifiers_per_transaction_are_consecutive() {
2277 let tx = make_rpc_tx(1, 2, &[10, 20, 30], 5);
2279
2280 let result = super::super::compute_ordered_nullifiers(&[tx]);
2281
2282 assert_eq!(result.len(), 3);
2283 assert!(result.contains(&Nullifier::from_raw(word(10))));
2284 assert!(result.contains(&Nullifier::from_raw(word(20))));
2285 assert!(result.contains(&Nullifier::from_raw(word(30))));
2286 }
2287
2288 #[test]
2289 fn empty_input_returns_empty_vec() {
2290 let result = super::super::compute_ordered_nullifiers(&[]);
2291 assert!(result.is_empty());
2292 }
2293 }
2294
2295 #[test]
2306 fn derive_account_commitments_walks_chains_per_account() {
2307 let make_tx = |account: AccountId, init_state: u64, final_state: u64, block_num: u32| {
2308 RpcTransactionRecord {
2309 block_num: BlockNumber::from(block_num),
2310 transaction_header: TransactionHeader::new(
2311 account,
2312 word(init_state),
2313 word(final_state),
2314 InputNotes::new_unchecked(vec![]),
2315 vec![],
2316 )
2317 .unwrap(),
2318 output_notes: vec![],
2319 erased_output_notes: vec![],
2320 consumed_note_refs: vec![],
2321 }
2322 };
2323
2324 let account_a: AccountId =
2325 ACCOUNT_ID_REGULAR_PRIVATE_ACCOUNT_UPDATABLE_CODE.try_into().unwrap();
2326 let account_b: AccountId = ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET.try_into().unwrap();
2327
2328 let tx_a_b5_1 = make_tx(account_a, 1, 2, 5);
2329 let tx_a_b5_2 = make_tx(account_a, 2, 3, 5);
2330 let tx_a_b6_1 = make_tx(account_a, 3, 4, 6);
2331 let tx_a_b6_2 = make_tx(account_a, 4, 5, 6);
2332 let tx_b_b6 = make_tx(account_b, 10, 20, 6);
2333
2334 let result = super::derive_account_commitments(&[
2336 tx_a_b6_1, tx_b_b6, tx_a_b5_2, tx_a_b6_2, tx_a_b5_1,
2337 ]);
2338
2339 assert_eq!(result.len(), 2, "one entry per account");
2340 assert!(
2341 result.contains(&(account_a, word(5))),
2342 "account A: must walk block 6's chain, not return block 5 or an intermediate",
2343 );
2344 assert!(
2345 result.contains(&(account_b, word(20))),
2346 "account B: must be resolved independently of account A",
2347 );
2348 }
2349
2350 struct CommitAllScreener;
2356
2357 #[async_trait(?Send)]
2358 impl OnNoteReceived for CommitAllScreener {
2359 async fn on_note_received(
2360 &self,
2361 committed_note: CommittedNote,
2362 _public_note: Option<InputNoteRecord>,
2363 ) -> Result<NoteUpdateAction, ClientError> {
2364 Ok(NoteUpdateAction::Commit(committed_note))
2365 }
2366 }
2367
2368 async fn build_chain_with_chained_consume_txs() -> (miden_testing::MockChain, Account, [Note; 3])
2372 {
2373 let sender_id: AccountId = ACCOUNT_ID_SENDER.try_into().unwrap();
2374 let faucet_id: AccountId = ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET.try_into().unwrap();
2375
2376 let mut builder = MockChainBuilder::new();
2377 let account = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2378 let account_id = account.id();
2379
2380 let asset = Asset::from(FungibleAsset::new(faucet_id, 100u64).unwrap());
2381 let note1 = builder
2382 .add_p2id_note(sender_id, account_id, &[asset], NoteType::Public)
2383 .unwrap();
2384 let note2 = builder
2385 .add_p2id_note(sender_id, account_id, &[asset], NoteType::Public)
2386 .unwrap();
2387 let note3 = builder
2388 .add_p2id_note(sender_id, account_id, &[asset], NoteType::Public)
2389 .unwrap();
2390
2391 let mut chain = builder.build().unwrap();
2392 chain.prove_next_block().unwrap(); let mut current_account = account.clone();
2396 for note in [¬e1, ¬e2, ¬e3] {
2397 let tx = Box::pin(
2398 chain
2399 .build_transaction(MockTransactionInput::Account(current_account.clone()))
2400 .unauthenticated_input_note(note.clone())
2401 .build()
2402 .unwrap()
2403 .execute(),
2404 )
2405 .await
2406 .unwrap();
2407 current_account.apply_patch(tx.account_patch()).unwrap();
2408 chain.add_pending_executed_transaction(&tx).unwrap();
2409 }
2410
2411 chain.prove_next_block().unwrap(); (chain, account, [note1, note2, note3])
2413 }
2414
2415 #[tokio::test]
2418 async fn sync_state_sets_consumed_tx_order_for_chained_transactions() {
2419 use miden_protocol::note::NoteMetadata;
2420
2421 let (chain, account, [note1, note2, note3]) = build_chain_with_chained_consume_txs().await;
2422
2423 let mock_rpc = MockRpcApi::new(chain);
2424 let state_sync = StateSync::new(
2425 Arc::new(mock_rpc.clone()),
2426 Arc::new(CommitAllScreener),
2427 None,
2428 genesis_validator_config(&mock_rpc),
2429 );
2430
2431 let genesis_peaks =
2432 mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
2433 let mut partial_mmr = PartialMmr::from_peaks(genesis_peaks);
2434
2435 let input_notes: Vec<InputNoteRecord> = [¬e1, ¬e2, ¬e3]
2436 .into_iter()
2437 .map(|n| InputNoteRecord::from(n.clone()))
2438 .collect();
2439
2440 let note_tags: BTreeSet<NoteTag> =
2441 input_notes.iter().filter_map(|n| n.metadata().map(NoteMetadata::tag)).collect();
2442
2443 let account_id = account.id();
2444 let sync_input = StateSyncInput {
2445 accounts: vec![AccountHeader::from(&account)],
2446 note_tags,
2447 input_notes,
2448 output_notes: vec![],
2449 uncommitted_transactions: vec![],
2450 };
2451
2452 let update = state_sync.sync_state(&mut partial_mmr, sync_input).await.unwrap();
2453
2454 let updated_notes: Vec<_> = update.note_updates().updated_input_notes().collect();
2455
2456 let find_order = |details_commitment| -> Option<u32> {
2457 updated_notes
2458 .iter()
2459 .find(|n| n.inner().details_commitment() == details_commitment)
2460 .and_then(|n| n.consumed_tx_order())
2461 };
2462
2463 assert_eq!(find_order(note1.details_commitment()), Some(0), "note1 should have tx_order 0");
2464 assert_eq!(find_order(note2.details_commitment()), Some(1), "note2 should have tx_order 1");
2465 assert_eq!(find_order(note3.details_commitment()), Some(2), "note3 should have tx_order 2");
2466
2467 for note in &updated_notes {
2470 let record = note.inner();
2471 assert!(record.is_consumed(), "note should be in a consumed state");
2472 assert_eq!(
2473 record.consumer_account(),
2474 Some(account_id),
2475 "externally-consumed notes by a tracked account should have consumer_account set",
2476 );
2477 }
2478 }
2479
2480 #[tokio::test]
2481 async fn sync_state_across_multiple_iterations_with_same_mmr() {
2482 let mock_rpc = MockRpcApi::default();
2484 mock_rpc.advance_blocks(3);
2485 let chain_tip_1 = mock_rpc.get_chain_tip_block_num();
2486
2487 let state_sync = StateSync::new(
2488 Arc::new(mock_rpc.clone()),
2489 Arc::new(MockScreener),
2490 None,
2491 genesis_validator_config(&mock_rpc),
2492 );
2493
2494 let genesis_peaks =
2496 mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
2497 let mut partial_mmr = PartialMmr::from_peaks(genesis_peaks);
2498 assert_eq!(partial_mmr.forest().num_leaves(), 1);
2499
2500 let update = state_sync.sync_state(&mut partial_mmr, empty()).await.unwrap();
2502
2503 assert_eq!(update.block_num(), chain_tip_1);
2504 let forest_1 = partial_mmr.forest();
2505 assert_eq!(forest_1.num_leaves(), chain_tip_1.as_u32() as usize + 1);
2507
2508 mock_rpc.advance_blocks(2);
2510 let chain_tip_2 = mock_rpc.get_chain_tip_block_num();
2511
2512 let update = state_sync.sync_state(&mut partial_mmr, empty()).await.unwrap();
2513
2514 assert_eq!(update.block_num(), chain_tip_2);
2515 let forest_2 = partial_mmr.forest();
2516 assert!(forest_2 > forest_1);
2517 assert_eq!(forest_2.num_leaves(), chain_tip_2.as_u32() as usize + 1);
2518
2519 let update = state_sync.sync_state(&mut partial_mmr, empty()).await.unwrap();
2521
2522 assert_eq!(update.block_num(), chain_tip_2);
2523 assert_eq!(partial_mmr.forest(), forest_2);
2524 }
2525
2526 async fn build_chain_with_mint_notes(
2529 num_blocks: u64,
2530 ) -> (miden_testing::MockChain, BTreeSet<NoteTag>) {
2531 let mut builder = MockChainBuilder::new();
2532 let faucet = builder
2533 .add_existing_basic_faucet(
2534 miden_testing::Auth::BasicAuth {
2535 auth_scheme: miden_protocol::account::auth::AuthScheme::Falcon512Poseidon2,
2536 },
2537 "TST",
2538 10_000,
2539 None,
2540 )
2541 .unwrap();
2542 let _target = builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2543 let mut chain = builder.build().unwrap();
2544
2545 let note_script = CodeBuilder::new()
2550 .compile_note_script("@note_script\npub proc main\n nop\nend")
2551 .unwrap();
2552 let note_recipient = NoteRecipient::new(
2553 Word::from([1u32, 2, 3, 4]),
2554 note_script,
2555 NoteStorage::new(vec![]).unwrap(),
2556 );
2557 let recipient = note_recipient.digest();
2558 let note_details = NoteDetails::new(NoteAssets::new(vec![]).unwrap(), note_recipient);
2561 let mut recipient_args = TransactionArgs::new(AdviceMap::default());
2562 recipient_args.add_output_note_recipient(¬e_details);
2563 let recipient_advice = recipient_args.advice_inputs().clone();
2564
2565 let tag = NoteTag::default();
2566 let mut faucet_account = faucet.clone();
2567 let mut note_tags = BTreeSet::new();
2568
2569 for i in 0..num_blocks {
2570 let amount = 100 + i;
2571 let source_manager = Arc::new(DefaultSourceManager::default());
2572 let mint_asset = FungibleAsset::new(faucet_account.id(), amount).unwrap();
2576 let asset_id_word = mint_asset.id().to_word();
2577 let asset_value_word = mint_asset.to_value_word();
2578 let tx_script_code = format!(
2579 "
2580 @transaction_script
2581 pub proc main
2582 push.{recipient}
2583 push.{note_type}
2584 push.{tag}
2585 push.{asset_value}
2586 push.{asset_id}
2587 call.::miden::standards::faucets::fungible::mint_and_send
2588 dropw dropw dropw dropw
2589 end
2590 ",
2591 recipient = recipient,
2592 note_type = NoteType::Private as u8,
2593 tag = u32::from(tag),
2594 asset_value = asset_value_word,
2595 asset_id = asset_id_word,
2596 );
2597 let tx_script = CodeBuilder::with_source_manager(source_manager.clone())
2598 .compile_tx_script(tx_script_code)
2599 .unwrap();
2600 let tx = Box::pin(
2601 chain
2602 .build_transaction(miden_testing::MockTransactionInput::Account(
2603 faucet_account.clone(),
2604 ))
2605 .extend_advice_inputs(recipient_advice.clone())
2606 .tx_script(tx_script)
2607 .with_source_manager(source_manager)
2608 .build()
2609 .unwrap()
2610 .execute(),
2611 )
2612 .await
2613 .unwrap();
2614
2615 for output_note in tx.output_notes().iter() {
2616 note_tags.insert(output_note.metadata().tag());
2617 }
2618
2619 faucet_account.apply_patch(tx.account_patch()).unwrap();
2620 chain.add_pending_executed_transaction(&tx).unwrap();
2621 chain.prove_next_block().unwrap();
2622 }
2623
2624 (chain, note_tags)
2625 }
2626
2627 #[tokio::test]
2630 async fn observer_relevance_persists_discarded_note_block() {
2631 let (chain, note_tags) = build_chain_with_mint_notes(2).await;
2632 let mock_rpc = MockRpcApi::new(chain);
2633 let chain_tip = mock_rpc.get_chain_tip_block_num();
2634
2635 let genesis_peaks =
2636 mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
2637 let mut partial_mmr = PartialMmr::from_peaks(genesis_peaks);
2638
2639 let validator_config = genesis_validator_config(&mock_rpc);
2640 let mut input = empty();
2641 let state_sync =
2642 StateSync::new(Arc::new(mock_rpc), Arc::new(MockScreener), None, validator_config)
2643 .with_note_observer(Arc::new(AlwaysRelevantObserver));
2644 input.note_tags = note_tags;
2645
2646 let update = state_sync.sync_state(&mut partial_mmr, input).await.unwrap();
2647 let observed_non_tip_block = BlockNumber::from(1u32);
2648
2649 assert!(
2650 update.partial_blockchain_updates().block_headers_to_store(chain_tip).any(
2651 |(header, is_relevant)| {
2652 header.block_num() == observed_non_tip_block && *is_relevant
2653 }
2654 ),
2655 "an observer-relevant block must be staged as relevant"
2656 );
2657 assert!(
2658 partial_mmr.is_tracked(observed_non_tip_block.as_usize()),
2659 "an observer-relevant block must remain tracked in the partial MMR"
2660 );
2661 }
2662
2663 #[tokio::test]
2673 async fn sync_state_tracks_note_blocks_in_mmr() {
2674 let (chain, note_tags) = build_chain_with_mint_notes(3).await;
2675 let mock_rpc = MockRpcApi::new(chain);
2676 let chain_tip = mock_rpc.get_chain_tip_block_num();
2677
2678 let note_blocks = mock_rpc
2680 .sync_notes(BlockNumber::from(0u32), chain_tip, ¬e_tags)
2681 .await
2682 .unwrap();
2683 assert!(
2684 note_blocks.len() >= 2,
2685 "expected notes in multiple blocks, got {}",
2686 note_blocks.len()
2687 );
2688
2689 let note_block_nums: BTreeSet<BlockNumber> =
2691 note_blocks.iter().map(|b| b.block_header.block_num()).collect();
2692
2693 let state_sync = StateSync::new(
2696 Arc::new(mock_rpc.clone()),
2697 Arc::new(MockScreener),
2698 None,
2699 genesis_validator_config(&mock_rpc),
2700 );
2701
2702 let genesis_peaks =
2703 mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
2704 let mut partial_mmr = PartialMmr::from_peaks(genesis_peaks);
2705
2706 let sync_data = state_sync
2707 .fetch_sync_data(BlockNumber::GENESIS, &[], &Arc::new(note_tags.clone()))
2708 .await
2709 .unwrap()
2710 .expect("should have progressed past genesis");
2711
2712 assert_eq!(sync_data.chain_tip_header.block_num(), chain_tip);
2714 assert!(!sync_data.note_blocks.is_empty(), "should have note blocks");
2715
2716 let _auth_nodes: Vec<(InOrderIndex, Word)> =
2718 partial_mmr.apply(sync_data.mmr_delta).map_err(StoreError::MmrError).unwrap();
2719 partial_mmr
2720 .add(sync_data.chain_tip_header.commitment(), false)
2721 .expect("chain tip should append to the partial MMR");
2722
2723 assert_eq!(partial_mmr.forest().num_leaves(), chain_tip.as_u32() as usize + 1);
2724
2725 for block in &sync_data.note_blocks {
2727 let bn = block.block_header.block_num();
2728 partial_mmr
2729 .track(bn.as_usize(), block.block_header.commitment(), &block.mmr_path)
2730 .map_err(StoreError::MmrError)
2731 .unwrap();
2732
2733 assert!(
2734 partial_mmr.is_tracked(bn.as_usize()),
2735 "block {bn} should be tracked after calling track()"
2736 );
2737 }
2738
2739 for &bn in ¬e_block_nums {
2741 assert!(
2742 partial_mmr.is_tracked(bn.as_usize()),
2743 "block {bn} with notes should be tracked in partial MMR"
2744 );
2745 }
2746 }
2747
2748 #[tokio::test]
2749 async fn sync_notes_with_content_fetches_inclusive_upper_bound_page() {
2750 let (chain, note_tags) = build_chain_with_mint_notes(10).await;
2751 let mock_rpc = MockRpcApi::new(chain);
2752
2753 let blocks = mock_rpc
2754 .sync_notes_with_content(
2755 4_u32.into(),
2756 10_u32.into(),
2757 ¬e_tags,
2758 NoteContentFetch::PublicDetailsAndAttachments,
2759 )
2760 .await
2761 .expect("sync notes should succeed");
2762
2763 assert_eq!(blocks.last().unwrap().block_header.block_num(), BlockNumber::from(10u32));
2764 assert!(
2765 blocks
2766 .iter()
2767 .any(|block| block.block_header.block_num() == BlockNumber::from(9u32))
2768 );
2769 }
2770
2771 #[tokio::test]
2777 async fn erased_notes_are_marked_as_consumed() {
2778 let sender_id: AccountId = ACCOUNT_ID_SENDER.try_into().unwrap();
2780 let partial_metadata = PartialNoteMetadata::new(sender_id, NoteType::Public);
2781 let metadata = NoteMetadata::new(partial_metadata, &NoteAttachments::empty());
2782 let script = CodeBuilder::new()
2783 .compile_note_script("@note_script\npub proc main\n nop\nend")
2784 .unwrap();
2785 let recipient = NoteRecipient::new(
2786 Word::from([1u32, 2, 3, 4]),
2787 script,
2788 NoteStorage::new(vec![]).unwrap(),
2789 );
2790 let output_note = OutputNoteRecord::new(
2791 recipient.digest(),
2792 NoteAssets::new(vec![]).unwrap(),
2793 metadata,
2794 OutputNoteState::ExpectedFull { recipient },
2795 BlockNumber::from(1u32),
2796 NoteAttachments::default(),
2797 );
2798 let note_id = output_note.id();
2799 let note_header = NoteHeader::new(output_note.details_commitment(), metadata);
2800
2801 let mut note_updates = NoteUpdateTracker::new(vec![], vec![output_note]);
2803
2804 let block_num = BlockNumber::from(3u32);
2806 note_updates
2807 .mark_erased_note_as_consumed(¬e_header, block_num)
2808 .expect("marking erased note should succeed");
2809
2810 let updated = note_updates
2811 .updated_output_notes()
2812 .find(|n| n.id() == note_id)
2813 .expect("output note should be in the update");
2814
2815 assert!(
2816 updated.inner().is_consumed(),
2817 "output note should be consumed after erasure detection, but state is: {}",
2818 updated.inner().state()
2819 );
2820 }
2821
2822 #[allow(clippy::too_many_lines)]
2838 #[ignore = "consumer derivation removed; see comment above"]
2839 #[tokio::test]
2840 async fn erased_notes_are_marked_as_consumed_by_network_account() {
2841 let mut builder = MockChainBuilder::new();
2844 let p2id_sender: AccountId = ACCOUNT_ID_SENDER.try_into().unwrap();
2845 let faucet_id: AccountId = ACCOUNT_ID_PRIVATE_FUNGIBLE_FAUCET.try_into().unwrap();
2846 let sender_account =
2847 builder.add_existing_mock_account(miden_testing::Auth::IncrNonce).unwrap();
2848 let sender_id = sender_account.id();
2849
2850 let asset = Asset::from(FungibleAsset::new(faucet_id, 100u64).unwrap());
2851 let note = builder
2852 .add_p2id_note(p2id_sender, sender_id, &[asset], NoteType::Public)
2853 .unwrap();
2854
2855 let mut chain = builder.build().unwrap();
2856 chain.prove_next_block().unwrap();
2857
2858 let tx = Box::pin(
2859 chain
2860 .build_transaction(MockTransactionInput::Account(sender_account.clone()))
2861 .unauthenticated_input_note(note.clone())
2862 .build()
2863 .unwrap()
2864 .execute(),
2865 )
2866 .await
2867 .unwrap();
2868 chain.add_pending_executed_transaction(&tx).unwrap();
2869 chain.prove_next_block().unwrap();
2870
2871 let network_account_id: AccountId =
2873 ACCOUNT_ID_REGULAR_PUBLIC_ACCOUNT_IMMUTABLE_CODE.try_into().unwrap();
2874 let target =
2875 NetworkAccountTarget::new(network_account_id, NoteExecutionHint::Always).unwrap();
2876 let attachment: NoteAttachment = target.into();
2877 let attachments = NoteAttachments::new(vec![attachment]).unwrap();
2878 let partial_metadata = PartialNoteMetadata::new(sender_id, NoteType::Public);
2879 let metadata = NoteMetadata::new(partial_metadata, &attachments);
2880 let script = CodeBuilder::new()
2881 .compile_note_script("@note_script\npub proc main\n nop\nend")
2882 .unwrap();
2883 let recipient = NoteRecipient::new(
2884 Word::from([7u32, 8, 9, 10]),
2885 script,
2886 NoteStorage::new(vec![]).unwrap(),
2887 );
2888 let recipient_digest = recipient.digest();
2889 let assets = NoteAssets::new(vec![]).unwrap();
2890
2891 let output_note = OutputNoteRecord::new(
2894 recipient_digest,
2895 assets.clone(),
2896 metadata,
2897 OutputNoteState::ExpectedFull { recipient },
2898 BlockNumber::from(1u32),
2899 NoteAttachments::default(),
2900 );
2901 let erased_note_id = output_note.id();
2902 let erased_note_header = NoteHeader::new(output_note.details_commitment(), metadata);
2903
2904 let mock_rpc = MockRpcApi::new(chain);
2905 mock_rpc.mark_note_as_erased(erased_note_header);
2906
2907 let network_header =
2910 AccountHeader::new(network_account_id, ZERO, EMPTY_WORD, EMPTY_WORD, EMPTY_WORD);
2911
2912 let state_sync = StateSync::new(
2913 Arc::new(mock_rpc.clone()),
2914 Arc::new(MockScreener),
2915 None,
2916 genesis_validator_config(&mock_rpc),
2917 );
2918
2919 let genesis_peaks =
2920 mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
2921 let mut partial_mmr = PartialMmr::from_peaks(genesis_peaks);
2922
2923 let sync_input = StateSyncInput {
2924 accounts: vec![AccountHeader::from(&sender_account), network_header],
2925 note_tags: BTreeSet::new(),
2926 input_notes: vec![],
2927 output_notes: vec![output_note],
2928 uncommitted_transactions: vec![],
2929 };
2930
2931 let update = state_sync.sync_state(&mut partial_mmr, sync_input).await.unwrap();
2932
2933 let updated_output = update
2935 .note_updates()
2936 .updated_output_notes()
2937 .find(|n| n.id() == erased_note_id)
2938 .expect("output note should be in the update");
2939 assert!(
2940 updated_output.inner().is_consumed(),
2941 "output note should be consumed, got: {}",
2942 updated_output.inner().state()
2943 );
2944
2945 let input_note_update = update
2947 .note_updates()
2948 .updated_input_notes()
2949 .find(|n| n.id() == Some(erased_note_id))
2950 .expect("input note should be created from the erased output note");
2951
2952 let inner = input_note_update.inner();
2953 assert!(
2954 inner.is_consumed(),
2955 "input note should be in a consumed state, got: {}",
2956 inner.state()
2957 );
2958 assert_eq!(
2959 inner.consumer_account(),
2960 Some(network_account_id),
2961 "consumer should be the tracked network account"
2962 );
2963 }
2964
2965 #[tokio::test]
2968 async fn validate_chain_mmr_response_rejects_tampered_responses() {
2969 let mock_rpc = MockRpcApi::default();
2970 mock_rpc.advance_blocks(3);
2971 let chain_tip = mock_rpc.get_chain_tip_block_num();
2972 let current = BlockNumber::GENESIS;
2973 let validator_config = genesis_validator_config(&mock_rpc);
2974
2975 let header_of =
2976 |block_num: u32| mock_rpc.mock_chain.read().block_header(block_num as usize);
2977 let chain_mmr_response = || async {
2978 mock_rpc.sync_chain_mmr(current, SyncTarget::CommittedChainTip).await.unwrap()
2979 };
2980
2981 let response = chain_mmr_response().await;
2983 StateSync::validate_chain_mmr_response(&response, current, &validator_config).unwrap();
2984
2985 let mut response = chain_mmr_response().await;
2987 response.block_header = header_of(chain_tip.as_u32() - 1);
2988 let result = StateSync::validate_chain_mmr_response(&response, current, &validator_config);
2989 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2990
2991 let mut response = chain_mmr_response().await;
2993 response.block_from = current + 1;
2994 let result = StateSync::validate_chain_mmr_response(&response, current, &validator_config);
2995 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
2996
2997 let mut response = chain_mmr_response().await;
2999 response.block_from = chain_tip;
3000 response.block_to = BlockNumber::GENESIS;
3001 response.block_header = header_of(0);
3002 let result =
3003 StateSync::validate_chain_mmr_response(&response, chain_tip, &validator_config);
3004 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
3005 }
3006
3007 #[test]
3010 fn validate_note_blocks_range_rejects_out_of_range_blocks() {
3011 let mock_rpc = MockRpcApi::default();
3012 mock_rpc.advance_blocks(3);
3013 let chain_tip = mock_rpc.get_chain_tip_block_num();
3014 let current = BlockNumber::GENESIS;
3015
3016 StateSync::validate_note_blocks_range(&[], current, chain_tip).unwrap();
3018
3019 let genesis_note_block = ResolvedSyncNotesBlock {
3021 block_header: mock_rpc.mock_chain.read().block_header(0),
3022 mmr_path: MerklePath::new(Vec::new()),
3023 notes: BTreeMap::new(),
3024 };
3025 let result =
3026 StateSync::validate_note_blocks_range(&[genesis_note_block], current, chain_tip);
3027 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
3028 }
3029
3030 #[tokio::test]
3033 async fn validate_chain_mmr_response_rejects_chain_tip_without_valid_signatures() {
3034 let mock_rpc = MockRpcApi::default();
3035 mock_rpc.advance_blocks(3);
3036 let chain_tip = mock_rpc.get_chain_tip_block_num();
3037 let current = BlockNumber::GENESIS;
3038
3039 let validator_config = genesis_validator_config(&mock_rpc);
3040 assert!(!validator_config.is_empty(), "the mock chain must commit a validator set");
3041
3042 let mut chain_mmr_info =
3044 mock_rpc.sync_chain_mmr(current, SyncTarget::CommittedChainTip).await.unwrap();
3045 StateSync::validate_chain_mmr_response(&chain_mmr_info, current, &validator_config)
3046 .unwrap();
3047
3048 let parent = BlockNumber::from(chain_tip.as_u32() - 1);
3049 chain_mmr_info.block_signatures = block_signatures(&mock_rpc, parent);
3051 let result =
3053 StateSync::validate_chain_mmr_response(&chain_mmr_info, current, &validator_config);
3054 assert!(
3055 matches!(result, Err(ClientError::ChainValidationError(_))),
3056 "signatures of another block must be rejected, got {result:?}"
3057 );
3058 }
3059
3060 #[test]
3063 fn advance_mmr_rejects_delta_inconsistent_with_chain_commitment() {
3064 let mock_rpc = MockRpcApi::default();
3065 mock_rpc.advance_blocks(3);
3066 let chain_tip = mock_rpc.get_chain_tip_block_num();
3067
3068 let chain_tip_header = mock_rpc.mock_chain.read().block_header(chain_tip.as_usize());
3069 let genesis_partial_mmr = || {
3070 let peaks = mock_rpc.get_mmr().peaks_at(Forest::new(1).expect("valid forest")).unwrap();
3071 PartialMmr::from_peaks(peaks)
3072 };
3073
3074 let full_delta = mock_rpc
3076 .get_mmr()
3077 .get_delta(Forest::new(1).unwrap(), Forest::new(chain_tip.as_usize()).unwrap())
3078 .unwrap();
3079 StateSync::advance_mmr(
3080 full_delta,
3081 &chain_tip_header,
3082 &mut genesis_partial_mmr(),
3083 &mut PartialBlockchainUpdates::default(),
3084 )
3085 .unwrap();
3086
3087 let truncated_delta = mock_rpc
3089 .get_mmr()
3090 .get_delta(Forest::new(1).unwrap(), Forest::new(chain_tip.as_usize() - 1).unwrap())
3091 .unwrap();
3092 let result = StateSync::advance_mmr(
3093 truncated_delta,
3094 &chain_tip_header,
3095 &mut genesis_partial_mmr(),
3096 &mut PartialBlockchainUpdates::default(),
3097 );
3098 assert!(matches!(result, Err(ClientError::ChainValidationError(_))));
3099 }
3100
3101 fn make_tx_record(account_id: AccountId, block_num: u32) -> RpcTransactionRecord {
3103 RpcTransactionRecord {
3104 block_num: BlockNumber::from(block_num),
3105 transaction_header: TransactionHeader::new(
3106 account_id,
3107 word(1),
3108 word(2),
3109 InputNotes::new_unchecked(vec![]),
3110 vec![],
3111 )
3112 .unwrap(),
3113 output_notes: vec![],
3114 erased_output_notes: vec![],
3115 consumed_note_refs: vec![],
3116 }
3117 }
3118}