1use alloc::boxed::Box;
2use alloc::collections::{BTreeMap, BTreeSet};
3use alloc::sync::Arc;
4use alloc::vec::Vec;
5
6use miden_protocol::account::{
7 Account,
8 AccountCode,
9 AccountId,
10 PartialAccount,
11 StorageMapKey,
12 StorageMapWitness,
13 StorageSlot,
14 StorageSlotContent,
15 StorageSlotName,
16};
17use miden_protocol::asset::{AssetId, AssetVault, AssetWitness};
18use miden_protocol::block::{BlockHeader, BlockNumber};
19use miden_protocol::crypto::merkle::MerklePath;
20use miden_protocol::crypto::merkle::mmr::{InOrderIndex, MmrPeaks, PartialMmr};
21use miden_protocol::note::{NoteScript, NoteScriptRoot};
22use miden_protocol::protocol_config::ProtocolConfig;
23use miden_protocol::transaction::{AccountInputs, PartialBlockchain};
24use miden_protocol::vm::FutureMaybeSend;
25use miden_protocol::{Word, ZERO};
26use miden_tx::{
27 DataStore,
28 DataStoreError,
29 LoadedMastForest,
30 MastForestStore,
31 TransactionMastStore,
32};
33
34use super::{AccountStorageFilter, PartialBlockchainFilter, Store};
35use crate::rpc::domain::account::{
36 AccountStorageRequirements,
37 GetAccountRequest,
38 StorageMapEntries,
39 StorageMapFetch,
40 VaultFetch,
41};
42use crate::rpc::{AccountStateAt, NodeRpcClient};
43use crate::store::StoreError;
44use crate::transaction::{ChainAnchor, ChainAnchorError, fetch_public_account_inputs};
45
46mod cache;
47use cache::DataStoreCache;
48
49pub struct ClientDataStore {
54 store: alloc::sync::Arc<dyn Store>,
56 cache: DataStoreCache,
58 rpc_api: Arc<dyn NodeRpcClient>,
60 anchor: Option<Box<ChainAnchor>>,
64}
65
66impl ClientDataStore {
67 pub fn new(store: alloc::sync::Arc<dyn Store>, rpc_api: Arc<dyn NodeRpcClient>) -> Self {
68 Self {
69 store,
70 cache: DataStoreCache::new(),
71 rpc_api,
72 anchor: None,
73 }
74 }
75
76 #[must_use]
84 pub fn with_chain_anchor(mut self, anchor: ChainAnchor) -> Self {
85 self.anchor = Some(Box::new(anchor));
86 self
87 }
88
89 #[must_use]
97 pub fn with_execution_input_cache(mut self) -> Self {
98 self.cache.enable_execution_input_cache();
99 self
100 }
101
102 pub fn mast_store(&self) -> Arc<TransactionMastStore> {
103 self.cache.mast_store.clone()
104 }
105
106 pub fn register_foreign_account_inputs(
109 &self,
110 foreign_accounts: impl IntoIterator<Item = AccountInputs>,
111 ) {
112 self.cache.replace_foreign_account_inputs(foreign_accounts);
113 }
114
115 pub(crate) fn register_block_numbers(
117 &self,
118 block_numbers: impl IntoIterator<Item = BlockNumber>,
119 ) {
120 self.cache.replace_block_numbers(block_numbers);
121 }
122
123 pub fn register_note_scripts(&self, note_scripts: impl IntoIterator<Item = NoteScript>) {
129 self.cache.insert_note_scripts(note_scripts);
130 }
131
132 async fn get_local_storage_map_witness(
139 &self,
140 account_id: AccountId,
141 map_root: Word,
142 map_key: StorageMapKey,
143 ) -> Result<Option<StorageMapWitness>, DataStoreError> {
144 match self
145 .store
146 .get_account_storage(account_id, AccountStorageFilter::Root(map_root))
147 .await
148 {
149 Ok(account_storage) => {
150 match account_storage.slots().first().map(StorageSlot::content) {
151 Some(StorageSlotContent::Map(map)) => Ok(Some(map.open(&map_key))),
152 Some(StorageSlotContent::Value(value)) => Err(DataStoreError::other(format!(
153 "found StorageSlotContent::Value with {value} as its value."
154 ))),
155 _ => Ok(None),
156 }
157 },
158 Err(err) => {
159 tracing::debug!(
160 %account_id,
161 %err,
162 "storage map not found locally, will try remote fetch"
163 );
164 Ok(None)
165 },
166 }
167 }
168
169 async fn fetch_and_cache_foreign_account(
172 &self,
173 account_id: AccountId,
174 account_state_at: AccountStateAt,
175 ) -> Result<AccountInputs, DataStoreError> {
176 let account_inputs = fetch_public_account_inputs(
177 &self.store,
178 &self.rpc_api,
179 account_id,
180 AccountStorageRequirements::default(),
181 account_state_at,
182 )
183 .await
184 .map_err(|err| {
185 DataStoreError::other_with_source("failed to fetch foreign account inputs", err)
186 })?;
187
188 self.cache.mast_store.load_account_code(account_inputs.code());
189 self.cache.insert_foreign_account_inputs(account_inputs.clone());
190
191 Ok(account_inputs)
192 }
193
194 async fn fetch_and_cache_storage_map_witness(
199 &self,
200 account_id: AccountId,
201 map_root: Word,
202 slot_name: StorageSlotName,
203 map_key: StorageMapKey,
204 known_code: AccountCode,
205 ) -> Result<StorageMapWitness, DataStoreError> {
206 let account_state_at =
207 self.cache.ref_block().map_or(AccountStateAt::ChainTip, AccountStateAt::Block);
208
209 let storage_requirements = AccountStorageRequirements::new([(slot_name, &[map_key])]);
210 let (_, account_proof): (BlockNumber, _) = self
211 .rpc_api
212 .get_account(
213 account_id,
214 GetAccountRequest::new()
215 .with_storage(StorageMapFetch::Slots(storage_requirements))
216 .with_known_code(Some(known_code))
217 .at(account_state_at),
218 )
219 .await
220 .map_err(|err| {
221 DataStoreError::other_with_source("failed to fetch storage map via RPC", err)
222 })?;
223
224 let (_, account_details) = account_proof.into_parts();
225 let details = account_details.ok_or_else(|| {
226 DataStoreError::other(format!(
227 "RPC returned no account details for account {account_id}"
228 ))
229 })?;
230
231 let map_detail =
232 details.storage_details.map_details.into_iter().next().ok_or_else(|| {
233 DataStoreError::other(format!(
234 "RPC returned no storage map details for account {account_id}"
235 ))
236 })?;
237
238 let StorageMapEntries::PartialMap { partial_smt, .. } = map_detail.entries else {
239 return Err(DataStoreError::other(
240 "expected a partial storage map in response to a specific-key request",
241 ));
242 };
243
244 let map_detail_root = partial_smt.root();
247 if map_detail_root != map_root {
248 return Err(DataStoreError::other(format!(
249 "storage map fetched for account {account_id} verifies against root \
250 {map_detail_root} but the executor requires root {map_root}"
251 )));
252 }
253
254 let proof = partial_smt.open(&map_key.hash().as_word()).map_err(|err| {
255 DataStoreError::other_with_source("failed to open the requested storage map key", err)
256 })?;
257
258 let witness = StorageMapWitness::new(proof, [map_key]).map_err(|err| {
259 DataStoreError::other_with_source("failed to create storage map witness", err)
260 })?;
261 self.cache.insert_storage_map_witness(map_root, map_key, witness.clone());
262 Ok(witness)
263 }
264
265 async fn fetch_vault_via_rpc(
270 &self,
271 account_id: AccountId,
272 vault_root: Word,
273 ) -> Result<AssetVault, DataStoreError> {
274 let account_state_at =
275 self.cache.ref_block().map_or(AccountStateAt::ChainTip, AccountStateAt::Block);
276
277 let known_code = self
279 .cache
280 .with_foreign_account_inputs(account_id, |inputs| inputs.code().clone());
281
282 let (block_num, mut account_proof) = self
283 .rpc_api
284 .get_account(
285 account_id,
286 GetAccountRequest::new()
287 .at(account_state_at)
288 .with_known_code(known_code)
289 .with_vault(VaultFetch::Always),
290 )
291 .await
292 .map_err(|err| {
293 DataStoreError::other_with_source("failed to fetch account vault via RPC", err)
294 })?;
295
296 let details = account_proof.details_mut().ok_or_else(|| {
297 DataStoreError::other(format!(
298 "RPC returned no account details for account {account_id}"
299 ))
300 })?;
301
302 self.rpc_api
303 .resolve_oversize_vault(account_id, block_num, details)
304 .await
305 .map_err(|err| {
306 DataStoreError::other_with_source("failed to resolve oversize vault via RPC", err)
307 })?;
308
309 let vault = AssetVault::new(&details.vault_details.assets).map_err(|err| {
310 DataStoreError::other_with_source("failed to build the fetched vault", err)
311 })?;
312
313 if vault.root() != vault_root {
314 return Err(DataStoreError::other(format!(
315 "vault fetched for account {account_id} has root {} but the executor requires \
316 root {vault_root}",
317 vault.root()
318 )));
319 }
320
321 Ok(vault)
322 }
323}
324
325impl DataStore for ClientDataStore {
326 async fn get_transaction_inputs(
327 &self,
328 account_id: AccountId,
329 mut block_refs: BTreeSet<BlockNumber>,
330 ) -> Result<(PartialAccount, BlockHeader, ProtocolConfig, PartialBlockchain), DataStoreError>
331 {
332 let ref_block = *block_refs.last().ok_or(DataStoreError::other("block set is empty"))?;
334
335 for block_num in self.cache.block_numbers() {
336 if block_num > ref_block {
337 return Err(DataStoreError::other(format!(
338 "requested block {block_num} is after transaction reference block {ref_block}"
339 )));
340 }
341 block_refs.insert(block_num);
342 }
343
344 self.cache.set_ref_block(ref_block);
346
347 let partial_account =
348 if let Some(partial_account) = self.cache.get_partial_account(account_id) {
349 partial_account
350 } else {
351 let partial_account_record = self
352 .store
353 .get_minimal_partial_account(account_id)
354 .await?
355 .ok_or(DataStoreError::AccountNotFound(account_id))?;
356
357 let partial_account: PartialAccount = if partial_account_record.nonce() == ZERO {
362 let full_record = self
363 .store
364 .get_account(account_id)
365 .await?
366 .ok_or(DataStoreError::AccountNotFound(account_id))?;
367 let account: Account = full_record
368 .try_into()
369 .map_err(|_| DataStoreError::AccountNotFound(account_id))?;
370 PartialAccount::from(&account)
371 } else {
372 partial_account_record
373 .try_into()
374 .map_err(|_| DataStoreError::AccountNotFound(account_id))?
375 };
376
377 self.cache.insert_partial_account(&partial_account);
378 partial_account
379 };
380
381 let (block_header, partial_blockchain) = if let Some(anchor) = &self.anchor {
382 if ref_block != anchor.block_num() {
386 return Err(DataStoreError::other_with_source(
387 "anchored data store cannot serve the requested reference block",
388 ChainAnchorError::ReferenceBlockMismatch {
389 requested: ref_block,
390 anchor: anchor.block_num(),
391 },
392 ));
393 }
394
395 for block_num in block_refs.iter().filter(|block_num| **block_num != ref_block) {
396 if !anchor.partial_blockchain().contains_block(*block_num) {
397 return Err(DataStoreError::other_with_source(
398 "anchored data store cannot serve an untracked block",
399 ChainAnchorError::BlockNotTracked { block_num: *block_num },
400 ));
401 }
402 }
403
404 (anchor.header().clone(), anchor.partial_blockchain().clone())
405 } else if let Some((block_header, partial_blockchain)) =
406 self.cache.get_blockchain(&block_refs)
407 {
408 (block_header, partial_blockchain)
409 } else {
410 let cache_key = block_refs.clone();
413 block_refs.remove(&ref_block);
414
415 let current_peaks = self.store.get_current_blockchain_peaks().await?;
416
417 let (block_header, _had_notes) = self
418 .store
419 .get_block_header_by_num(ref_block)
420 .await?
421 .ok_or(DataStoreError::BlockNotFound(ref_block))?;
422
423 let (partial_mmr, block_headers) = build_partial_mmr_and_headers_with_fallback(
427 &self.store,
428 &self.rpc_api,
429 current_peaks,
430 &block_refs,
431 )
432 .await?;
433
434 let partial_blockchain =
435 PartialBlockchain::new(partial_mmr, block_headers).map_err(|err| {
436 DataStoreError::other_with_source(
437 "error creating PartialBlockchain from internal data",
438 err,
439 )
440 })?;
441
442 self.cache.insert_blockchain(cache_key, &block_header, &partial_blockchain);
443 (block_header, partial_blockchain)
444 };
445
446 let protocol_config = crate::protocol_config::load_protocol_config(
447 self.store.as_ref(),
448 block_header.protocol_config_commitment(),
449 )
450 .await?;
451 Ok((partial_account, block_header, protocol_config, partial_blockchain))
452 }
453
454 async fn get_vault_asset_witnesses(
460 &self,
461 account_id: AccountId,
462 vault_root: Word,
463 asset_ids: BTreeSet<AssetId>,
464 ) -> Result<Vec<AssetWitness>, DataStoreError> {
465 if let Some(witnesses) = self.cache.get_vault_asset_witnesses(vault_root, &asset_ids) {
466 return Ok(witnesses);
467 }
468
469 let asset_witnesses = match self
470 .store
471 .get_vault_asset_witnesses(account_id, vault_root, asset_ids.clone())
472 .await
473 {
474 Ok(witnesses) => witnesses,
475 Err(err) => {
476 tracing::debug!(
477 %account_id,
478 requested_root = %vault_root,
479 %err,
480 "local store cannot serve the requested vault root, will fetch it via RPC"
481 );
482 let vault = self.fetch_vault_via_rpc(account_id, vault_root).await?;
483 asset_ids.iter().copied().map(|asset_id| vault.open(asset_id)).collect()
484 },
485 };
486
487 self.cache
488 .insert_vault_asset_witnesses(vault_root, &asset_ids, &asset_witnesses);
489 Ok(asset_witnesses)
490 }
491
492 async fn get_storage_map_witness(
496 &self,
497 account_id: AccountId,
498 map_root: Word,
499 map_key: StorageMapKey,
500 ) -> Result<StorageMapWitness, DataStoreError> {
501 if let Some(witness) = self.cache.get_storage_map_witness(map_root, map_key) {
503 return Ok(witness);
504 }
505
506 if let Some(witness) =
508 self.get_local_storage_map_witness(account_id, map_root, map_key).await?
509 {
510 return Ok(witness);
511 }
512
513 let resolution = if let Some(resolution) =
516 self.cache.with_foreign_account_inputs(account_id, |inputs| {
517 resolve_witness_from_inputs(inputs, map_root, map_key)
518 }) {
519 resolution?
520 } else {
521 let account_state_at = self
522 .cache
523 .ref_block()
524 .map(AccountStateAt::Block)
525 .expect("reference block should be set");
526 let inputs = self.fetch_and_cache_foreign_account(account_id, account_state_at).await?;
527 resolve_witness_from_inputs(&inputs, map_root, map_key)?
528 };
529
530 match resolution {
531 WitnessResolution::Witness(witness) => Ok(witness),
532 WitnessResolution::FetchParams(slot_name, known_code) => {
533 self.fetch_and_cache_storage_map_witness(
534 account_id, map_root, slot_name, map_key, known_code,
535 )
536 .await
537 },
538 }
539 }
540
541 async fn get_foreign_account_inputs(
544 &self,
545 foreign_account_id: AccountId,
546 ref_block: BlockNumber,
547 ) -> Result<AccountInputs, DataStoreError> {
548 if let Some(inputs) = self.cache.get_foreign_account_inputs(foreign_account_id) {
550 return Ok(inputs);
551 }
552
553 self.fetch_and_cache_foreign_account(foreign_account_id, AccountStateAt::Block(ref_block))
554 .await
555 }
556
557 fn get_note_script(
560 &self,
561 script_root: NoteScriptRoot,
562 ) -> impl FutureMaybeSend<Result<Option<NoteScript>, DataStoreError>> {
563 let registered_script = self.cache.get_note_script(script_root.into());
564 let store = self.store.clone();
565 let rpc_api = self.rpc_api.clone();
566
567 async move {
568 if let Some(note_script) = registered_script {
570 return Ok(Some(note_script));
571 }
572
573 match store.get_note_script(script_root.into()).await {
575 Ok(note_script) => return Ok(Some(note_script)),
576 Err(StoreError::NoteScriptNotFound(_)) => {},
577 Err(err) => {
578 return Err(DataStoreError::other_with_source(
579 format!("failed to get note script {script_root} from store"),
580 err,
581 ));
582 },
583 }
584
585 let Some(note_script) =
587 rpc_api.get_note_script_by_root(script_root.into()).await.map_err(|err| {
588 DataStoreError::other_with_source("failed to fetch note script via RPC", err)
589 })?
590 else {
591 return Ok(None);
592 };
593
594 if let Err(err) = store.upsert_note_scripts(core::slice::from_ref(¬e_script)).await {
596 tracing::warn!(
597 %err,
598 "Failed to persist fetched note script to store"
599 );
600 }
601
602 Ok(Some(note_script))
603 }
604 }
605}
606
607impl MastForestStore for ClientDataStore {
611 fn get(&self, procedure_hash: &Word) -> Option<LoadedMastForest> {
612 self.cache.mast_store.get(procedure_hash)
613 }
614}
615
616enum WitnessResolution {
622 Witness(StorageMapWitness),
623 FetchParams(StorageSlotName, AccountCode),
626}
627
628fn resolve_witness_from_inputs(
632 inputs: &AccountInputs,
633 map_root: Word,
634 map_key: StorageMapKey,
635) -> Result<WitnessResolution, DataStoreError> {
636 if let Some(partial_map) = inputs.storage().maps().find(|m| m.root() == map_root)
637 && let Ok(witness) = partial_map.open(&map_key)
638 {
639 return Ok(WitnessResolution::Witness(witness));
640 }
641
642 let account_id = inputs.id();
643 let slot_name = inputs
644 .storage()
645 .header()
646 .slots()
647 .find(|slot| slot.slot_type().is_map() && slot.value() == map_root)
648 .map(|slot| slot.name().clone())
649 .ok_or_else(|| {
650 DataStoreError::other(format!(
651 "did not find map slot with root {map_root} for foreign account {account_id}"
652 ))
653 })?;
654
655 Ok(WitnessResolution::FetchParams(slot_name, inputs.code().clone()))
656}
657
658pub(crate) async fn build_partial_mmr_with_paths(
664 store: &alloc::sync::Arc<dyn Store>,
665 rpc_api: &Arc<dyn NodeRpcClient>,
666 peaks: MmrPeaks,
667 authenticated_blocks: &[BlockHeader],
668) -> Result<PartialMmr, DataStoreError> {
669 let mut partial_mmr: PartialMmr = PartialMmr::from_peaks(peaks);
670
671 let block_nums: Vec<BlockNumber> =
672 authenticated_blocks.iter().map(BlockHeader::block_num).collect();
673
674 let authentication_paths =
675 get_authentication_path_for_blocks(store, &block_nums, partial_mmr.forest().num_leaves())
676 .await?;
677
678 for (header, local_path) in authenticated_blocks.iter().zip(authentication_paths.iter()) {
679 if partial_mmr
680 .track(header.block_num().as_usize(), header.commitment(), local_path)
681 .is_ok()
682 {
683 continue;
684 }
685
686 fetch_and_track_block_header(rpc_api, &mut partial_mmr, header.block_num(), Some(header))
687 .await?;
688 }
689
690 Ok(partial_mmr)
691}
692
693pub(crate) async fn build_partial_mmr_and_headers_with_fallback(
698 store: &alloc::sync::Arc<dyn Store>,
699 rpc_api: &Arc<dyn NodeRpcClient>,
700 peaks: MmrPeaks,
701 block_numbers: &BTreeSet<BlockNumber>,
702) -> Result<(PartialMmr, Vec<BlockHeader>), DataStoreError> {
703 let mut headers: BTreeMap<BlockNumber, BlockHeader> = store
704 .get_block_headers(block_numbers)
705 .await?
706 .into_iter()
707 .map(|(header, _has_notes)| (header.block_num(), header))
708 .collect();
709
710 let local_headers: Vec<BlockHeader> = headers.values().cloned().collect();
711 let mut partial_mmr =
712 build_partial_mmr_with_paths(store, rpc_api, peaks, &local_headers).await?;
713
714 for &block_num in block_numbers {
715 if headers.contains_key(&block_num) {
716 continue;
717 }
718
719 let header =
720 fetch_and_track_block_header(rpc_api, &mut partial_mmr, block_num, None).await?;
721 headers.insert(block_num, header);
722 }
723
724 Ok((partial_mmr, headers.into_values().collect()))
725}
726
727async fn fetch_and_track_block_header(
731 rpc_api: &Arc<dyn NodeRpcClient>,
732 partial_mmr: &mut PartialMmr,
733 block_num: BlockNumber,
734 expected_header: Option<&BlockHeader>,
735) -> Result<BlockHeader, DataStoreError> {
736 let (header, proof) = rpc_api.get_block_header_with_proof(block_num).await.map_err(|err| {
737 DataStoreError::other_with_source(
738 format!("failed to fetch block header and MMR proof for block {block_num}"),
739 err,
740 )
741 })?;
742
743 if header.block_num() != block_num {
744 return Err(DataStoreError::other(format!(
745 "node returned block header {} for requested block {block_num}",
746 header.block_num()
747 )));
748 }
749 if expected_header.is_some_and(|expected| header != *expected) {
750 return Err(DataStoreError::other(format!(
751 "node returned a different header for block {block_num}"
752 )));
753 }
754 if proof.leaf() != header.commitment() {
755 return Err(DataStoreError::other(format!(
756 "node returned an invalid MMR proof for block {block_num}"
757 )));
758 }
759
760 let proof = proof.with_forest(partial_mmr.forest()).map_err(|err| {
761 DataStoreError::other(format!("failed to adjust MMR proof for block {block_num}: {err}"))
762 })?;
763 partial_mmr
764 .track(block_num.as_usize(), header.commitment(), proof.merkle_path())
765 .map_err(|err| DataStoreError::other(format!("error constructing MMR: {err}")))?;
766
767 Ok(header)
768}
769
770async fn get_authentication_path_for_blocks(
776 store: &alloc::sync::Arc<dyn Store>,
777 block_nums: &[BlockNumber],
778 forest: usize,
779) -> Result<Vec<MerklePath>, StoreError> {
780 let mut node_indices = BTreeSet::new();
781
782 for block_num in block_nums {
784 let path_depth = mmr_merkle_path_len(block_num.as_usize(), forest);
785
786 let mut idx = InOrderIndex::from_leaf_pos(block_num.as_usize());
787
788 for _ in 0..path_depth {
789 node_indices.insert(idx.sibling());
790 idx = idx.parent();
791 }
792 }
793
794 let node_indices: Vec<InOrderIndex> = node_indices.into_iter().collect();
795
796 let filter = PartialBlockchainFilter::List(node_indices);
797 let mmr_nodes = store.get_partial_blockchain_nodes(filter).await?;
798
799 let mut authentication_paths = vec![];
801 for block_num in block_nums {
802 let mut merkle_nodes = vec![];
803 let mut idx = InOrderIndex::from_leaf_pos(block_num.as_usize());
804
805 while let Some(node) = mmr_nodes.get(&idx.sibling()) {
806 merkle_nodes.push(*node);
807 idx = idx.parent();
808 }
809 let path = MerklePath::new(merkle_nodes);
810 authentication_paths.push(path);
811 }
812
813 Ok(authentication_paths)
814}
815
816fn mmr_merkle_path_len(leaf_index: usize, forest: usize) -> usize {
819 let before: usize = forest & leaf_index;
820 let after = forest ^ before;
821
822 after.ilog2() as usize
823}