Skip to main content

miden_client/store/data_store/
mod.rs

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
49// DATA STORE
50// ================================================================================================
51
52/// Wrapper structure that implements [`DataStore`] over any [`Store`].
53pub struct ClientDataStore {
54    /// Local database containing information about the accounts managed by this client.
55    store: alloc::sync::Arc<dyn Store>,
56    /// In-memory state served to the executor for the duration of the execution session.
57    cache: DataStoreCache,
58    /// RPC client used to lazy-load foreign account data on cache miss.
59    rpc_api: Arc<dyn NodeRpcClient>,
60    /// When set, chain data (reference block header and partial blockchain) is served from this
61    /// anchor instead of being rebuilt at the store's sync height. Boxed to keep the data store
62    /// small: it is held inline by every execution future.
63    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    /// Serves chain data from the provided [`ChainAnchor`] instead of rebuilding it at the store's
77    /// sync height, pinning execution to the anchor's reference block.
78    ///
79    /// The store's account data is still used as-is: only the reference block header and the
80    /// partial blockchain come from the anchor. The anchor must track all blocks required by the
81    /// request and all creation blocks of authenticated input notes. Otherwise,
82    /// `get_transaction_inputs` fails.
83    #[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    /// Enables memoization of `get_transaction_inputs` and `get_vault_asset_witnesses` for the
90    /// lifetime of this data store.
91    ///
92    /// This is only correct when the account state served to the executor does not change between
93    /// executions, as is the case while the [`crate::note::NoteScreener`] runs trial executions
94    /// against the same accounts and reference block. It must stay disabled for real transaction
95    /// execution, where account state evolves between executions.
96    #[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    /// Stores the provided foreign account inputs so they can be served to the executor upon
107    /// request.
108    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    /// Stores the blocks that the current transaction must be able to authenticate.
116    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    /// Registers note scripts so they can be served to the executor upon request.
124    ///
125    /// Scripts accumulate across calls (they are not cleared) so that a data store reused for
126    /// several executions — e.g. by [`crate::transaction::BatchBuilder`] — keeps serving the
127    /// scripts registered for earlier transactions.
128    pub fn register_note_scripts(&self, note_scripts: impl IntoIterator<Item = NoteScript>) {
129        self.cache.insert_note_scripts(note_scripts);
130    }
131
132    /// Attempts to resolve a storage map witness from the local store.
133    ///
134    /// This covers any account present in the store (local or foreign) as well as any foreign
135    /// account previously cached in `foreign_account_inputs`.
136    ///
137    /// Returns `Ok(None)` when the map is not found locally.
138    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    /// Lazily fetches a foreign account's inputs from the network, loads its code into the MAST
170    /// store, and caches the result in [`Self::foreign_account_inputs`].
171    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    /// Fetches a storage map witness for a specific key from the network via RPC and caches it.
195    ///
196    /// Anchored at the transaction reference block: a chain-tip query would return proofs for a
197    /// newer map root whenever the account changed after the caller's last sync.
198    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        // Reject a wrong-root response here rather than as an opaque merkle error inside the VM.
245        // The whole tree shares one root, so this covers every opening taken from it.
246        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    /// Fetches an account's full vault via RPC — anchored at the transaction reference block — and
266    /// verifies it against the vault root the executor requires. Fallback for vault reads the local
267    /// store cannot serve, typically foreign accounts whose [`AccountInputs`] carry only their
268    /// vault root.
269    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        // The cached foreign inputs hold the code, letting the node omit it from the response.
278        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        // Last block is used as reference (it does not need to be authenticated manually)
333        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        // Cache the reference block so lazy-loading methods can use it
345        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                // New accounts (nonce == 0) need full storage maps as advice inputs for the kernel
358                // to validate during account creation. For these, fetch the full account and
359                // convert to PartialAccount (which includes full storage for new accounts).
360                // Existing accounts use the minimal partial record directly.
361                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            // Anchored execution: serve the pinned chain data. The executor-derived reference block
383            // must match the anchor. The anchor's partial blockchain must track every other block
384            // in the set.
385            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            // The full set identifies the served blockchain, so keep it as the cache key before the
411            // reference block is removed from it below.
412            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            // TODO: the client stores only the peaks of the MMR at the current sync height, so we
424            // are not actually following the block_ref here. If the block_ref !=
425            // current_sync_height, this would return an invalid partial blockchain.
426            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    /// Retrieves witnesses for the requested assets from the local store, falling back to a single
455    /// RPC vault fetch when the store cannot serve the requested root.
456    ///
457    /// Assets absent from the vault are served too: the store returns an emptiness proof for them,
458    /// which the executor needs when an asset is being added to the vault.
459    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    /// Retrieves the [`StorageMapWitness`] requested from the store. Alternatively fetching it from
493    /// the RPC if not available locally. Witnesses fetched via RPC are cached in memory so that
494    /// repeated accesses to the same map entry within a transaction avoid additional RPC calls.
495    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        // Check the in-memory witness cache first.
502        if let Some(witness) = self.cache.get_storage_map_witness(map_root, map_key) {
503            return Ok(witness);
504        }
505
506        // Try the local store.
507        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        // Resolve against the cached account inputs (without cloning them), fetching and caching
514        // the account first if it isn't cached yet.
515        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    /// Returns the [`AccountInputs`] for the given foreign account from the cache or alternatively
542    /// fetching them from the RPC if not available locally.
543    async fn get_foreign_account_inputs(
544        &self,
545        foreign_account_id: AccountId,
546        ref_block: BlockNumber,
547    ) -> Result<AccountInputs, DataStoreError> {
548        // Fast path: check the cache first.
549        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    /// Returns the [`NoteScript`] for the given script root from the registered session scripts,
558    /// the store, or alternatively fetching it from the RPC if not available locally.
559    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            // Fastest path: scripts registered for the in-flight transaction request.
569            if let Some(note_script) = registered_script {
570                return Ok(Some(note_script));
571            }
572
573            // Fast path: check the local store first.
574            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            // Store miss, fetch from the network via RPC.
586            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            // Persist for future lookups.
595            if let Err(err) = store.upsert_note_scripts(core::slice::from_ref(&note_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
607// MAST FOREST STORE
608// ================================================================================================
609
610impl MastForestStore for ClientDataStore {
611    fn get(&self, procedure_hash: &Word) -> Option<LoadedMastForest> {
612        self.cache.mast_store.get(procedure_hash)
613    }
614}
615
616// HELPER FUNCTIONS
617// ================================================================================================
618
619/// Outcome of resolving a storage map witness against an account's inputs: either the witness
620/// itself, or the parameters needed to fetch it via RPC.
621enum WitnessResolution {
622    Witness(StorageMapWitness),
623    /// The [`AccountCode`] is not needed to build the witness: it is only sent along with the RPC
624    /// request so the node can omit the account code from its response.
625    FetchParams(StorageSlotName, AccountCode),
626}
627
628/// Tries to open the witness from the inputs' partial storage maps (this can miss if the account's
629/// storage is too big); on a miss, resolves the slot name and account code needed to fetch the
630/// witness via RPC.
631fn 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
658/// Builds a [`PartialMmr`] from the given peaks and a list of blocks that should be authenticated
659/// against them.
660///
661/// `authenticated_blocks` must not contain the block whose forest matches `peaks`. For that block
662/// the kernel extends the MMR itself, so an authentication path is not needed.
663pub(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
693/// Builds a [`PartialMmr`] and returns the authenticated block headers.
694///
695/// Headers in the store use local MMR nodes when possible. If a local path is incomplete, the node
696/// returns the header and its proof in one call. A missing header always uses this combined call.
697pub(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
727/// Fetches a block header with its MMR proof and tracks the authenticated header.
728///
729/// If `expected_header` is set, the fetched header must match it.
730async 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
770/// Retrieves all Partial Blockchain nodes required for authenticating the set of blocks, and then
771/// constructs the path for each of them.
772///
773/// This function assumes `block_nums` doesn't contain values above or equal to `forest`. If there
774/// are any such values, the function will panic when calling `mmr_merkle_path_len()`.
775async 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    // Calculate all needed nodes indices for generating the paths
783    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    // Construct authentication paths
800    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
816/// Calculates the merkle path length for an MMR of a specific forest and a leaf index `leaf_index`
817/// is a 0-indexed leaf number and `forest` is the total amount of leaves in the MMR at this point.
818fn 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}