1use std::collections::{BTreeMap, BTreeSet, HashSet};
7use std::num::NonZeroUsize;
8use std::path::{Path, PathBuf};
9use std::sync::Arc;
10
11use miden_node_proto::domain::batch::BatchInputs;
12use miden_node_utils::clap::StorageOptions;
13use miden_node_utils::formatting::format_array;
14use miden_node_utils::tracing::miden_instrument;
15use miden_protocol::Word;
16use miden_protocol::account::AccountId;
17use miden_protocol::block::account_tree::AccountWitness;
18use miden_protocol::block::nullifier_tree::{NullifierTree, NullifierWitness};
19use miden_protocol::block::{BlockHeader, BlockInputs, BlockNumber, Blockchain};
20use miden_protocol::crypto::merkle::mmr::{MmrProof, PartialMmr};
21use miden_protocol::crypto::merkle::smt::{LargeSmt, SmtStorage};
22use miden_protocol::note::{NoteId, NoteScript, Nullifier};
23use miden_protocol::transaction::PartialBlockchain;
24use tokio::sync::{Mutex, RwLock, watch};
25use tracing::{Instrument, Span};
26
27use crate::account_state_forest::{AccountStateForest, AccountStateForestBackend};
28use crate::accounts::AccountTreeWithHistory;
29use crate::blocks::BlockStore;
30use crate::db::{Db, NoteRecord, NullifierInfo};
31use crate::errors::{
32 DatabaseError,
33 GetBatchInputsError,
34 GetBlockHeaderError,
35 GetBlockInputsError,
36 StateInitializationError,
37};
38use crate::proven_tip::ProvenTipWriter;
39use crate::{COMPONENT, DataDirectory, DatabaseOptions};
40
41const BLOCK_CACHE_CAPACITY: NonZeroUsize = NonZeroUsize::new(512).unwrap();
43
44const PROOF_CACHE_CAPACITY: NonZeroUsize = NonZeroUsize::new(512).unwrap();
46
47mod loader;
48use loader::{
49 ACCOUNT_STATE_FOREST_STORAGE_DIR,
50 ACCOUNT_TREE_STORAGE_DIR,
51 AccountForestLoader,
52 NULLIFIER_TREE_STORAGE_DIR,
53 TreeStorage,
54 TreeStorageLoader,
55 load_mmr,
56 verify_account_state_forest_consistency,
57 verify_tree_consistency,
58};
59
60mod replica;
61pub use replica::{BlockCache, BlockNotification, ProofCache, ProofNotification};
62
63mod account;
64
65mod apply_block;
66mod apply_proof;
67mod bootstrap;
68mod disk_monitor;
69mod sync_state;
70
71#[derive(Debug, Clone, Copy)]
76pub enum Finality {
77 Committed,
79 Proven,
81}
82
83#[derive(Debug, Default)]
87pub struct TransactionInputs {
88 pub account_commitment: Word,
89 pub nullifiers: Vec<NullifierInfo>,
90 pub found_unauthenticated_notes: HashSet<Word>,
91 pub new_account_id_prefix_is_unique: Option<bool>,
92}
93
94type BlockInputWitnesses = (
95 BlockNumber,
96 BTreeMap<AccountId, AccountWitness>,
97 BTreeMap<Nullifier, NullifierWitness>,
98 PartialMmr,
99);
100
101struct InnerState<S>
103where
104 S: SmtStorage,
105{
106 nullifier_tree: NullifierTree<LargeSmt<S>>,
107 blockchain: Blockchain,
108 account_tree: AccountTreeWithHistory<S>,
109}
110
111impl<S: SmtStorage> InnerState<S> {
112 fn latest_block_num(&self) -> BlockNumber {
114 self.blockchain
115 .chain_tip()
116 .expect("chain should always have at least the genesis block")
117 }
118}
119
120pub struct State {
125 data_directory: PathBuf,
127
128 db: Arc<Db>,
131
132 block_store: Arc<BlockStore>,
134
135 inner: RwLock<InnerState<TreeStorage>>,
139
140 forest: RwLock<AccountStateForest<AccountStateForestBackend>>,
142
143 writer: Mutex<()>,
146
147 proven_tip: ProvenTipWriter,
149
150 committed_tip_tx: watch::Sender<BlockNumber>,
153
154 pub(crate) block_cache: BlockCache,
157
158 pub(crate) proof_cache: ProofCache,
161}
162
163impl State {
164 #[miden_instrument(
172 target = COMPONENT,
173 skip_all,
174 )]
175 pub async fn load(
176 data_path: &Path,
177 storage_options: StorageOptions,
178 ) -> Result<Self, StateInitializationError> {
179 Self::load_with_database_options(data_path, storage_options, DatabaseOptions::default())
180 .await
181 }
182
183 #[miden_instrument(
188 target = COMPONENT,
189 skip_all,
190 )]
191 pub async fn load_with_database_options(
192 data_path: &Path,
193 storage_options: StorageOptions,
194 database_options: DatabaseOptions,
195 ) -> Result<Self, StateInitializationError> {
196 let data_directory = DataDirectory::load(data_path.to_path_buf())
197 .map_err(StateInitializationError::DataDirectoryLoadError)?;
198
199 let block_store = Arc::new(
200 BlockStore::load(data_directory.block_store_dir())
201 .map_err(StateInitializationError::BlockStoreLoadError)?,
202 );
203
204 let database_filepath = data_directory.database_path();
205 let mut db = Db::load_with_pool_size(
206 database_filepath.clone(),
207 database_options.connection_pool_size,
208 )
209 .await
210 .map_err(StateInitializationError::DatabaseLoadError)?;
211
212 let blockchain = load_mmr(&mut db).await?;
213 let latest_block_num = blockchain.chain_tip().unwrap_or(BlockNumber::GENESIS);
214
215 #[cfg(feature = "rocksdb")]
216 let (account_storage_config, nullifier_storage_config, forest_storage_config) = (
217 storage_options.account_tree.into(),
218 storage_options.nullifier_tree.into(),
219 storage_options.account_state_forest.into(),
220 );
221 #[cfg(not(feature = "rocksdb"))]
222 let (account_storage_config, nullifier_storage_config, forest_storage_config) = {
223 let _ = &storage_options;
224 ((), (), ())
225 };
226 let account_storage =
227 TreeStorage::create(data_path, &account_storage_config, ACCOUNT_TREE_STORAGE_DIR)?;
228 let account_tree = account_storage.load_account_tree(&mut db).await?;
229
230 let nullifier_storage =
231 TreeStorage::create(data_path, &nullifier_storage_config, NULLIFIER_TREE_STORAGE_DIR)?;
232 let nullifier_tree = nullifier_storage.load_nullifier_tree(&mut db).await?;
233
234 verify_tree_consistency(account_tree.root(), nullifier_tree.root(), &mut db).await?;
238
239 let account_tree = AccountTreeWithHistory::new(account_tree, latest_block_num);
240
241 let forest_backend = AccountStateForestBackend::create(
242 data_path,
243 &forest_storage_config,
244 ACCOUNT_STATE_FOREST_STORAGE_DIR,
245 )?;
246 let forest = forest_backend.load_account_state_forest(&mut db, latest_block_num).await?;
247 verify_account_state_forest_consistency(&forest, &mut db).await?;
248
249 let inner = RwLock::new(InnerState { nullifier_tree, blockchain, account_tree });
250
251 let forest = RwLock::new(forest);
252 let writer = Mutex::new(());
253 let db = Arc::new(db);
254
255 let proven_tip_init = block_store
257 .load_proven_tip()
258 .map_err(StateInitializationError::ProvenTipLoadError)?;
259 let (proven_tip, _rx) = ProvenTipWriter::new(proven_tip_init);
260
261 let (committed_tip_tx, _rx) = watch::channel(latest_block_num);
263
264 Ok(Self {
265 data_directory: data_path.to_path_buf(),
266 db,
267 block_store,
268 inner,
269 forest,
270 writer,
271 proven_tip,
272 committed_tip_tx,
273 block_cache: BlockCache::new(BLOCK_CACHE_CAPACITY),
274 proof_cache: ProofCache::new(PROOF_CACHE_CAPACITY),
275 })
276 }
277
278 pub fn subscribe_committed_tip(&self) -> watch::Receiver<BlockNumber> {
280 self.committed_tip_tx.subscribe()
281 }
282
283 pub async fn load_proving_inputs(
285 &self,
286 block_num: BlockNumber,
287 ) -> std::io::Result<Option<Vec<u8>>> {
288 self.block_store.load_proving_inputs(block_num).await
289 }
290
291 pub fn subscribe_proven_tip(&self) -> watch::Receiver<BlockNumber> {
293 self.proven_tip.subscribe()
294 }
295
296 fn with_inner_read_blocking<R>(&self, f: impl FnOnce(&InnerState<TreeStorage>) -> R) -> R {
305 let span = Span::current();
306 tokio::task::block_in_place(|| {
307 span.in_scope(|| {
308 let inner = self.inner.blocking_read();
309 f(&inner)
310 })
311 })
312 }
313
314 fn with_inner_write_blocking<R>(&self, f: impl FnOnce(&mut InnerState<TreeStorage>) -> R) -> R {
318 let span = Span::current();
319 tokio::task::block_in_place(|| {
320 span.in_scope(|| {
321 let mut inner = self.inner.blocking_write();
322 f(&mut inner)
323 })
324 })
325 }
326
327 fn with_forest_read_blocking<R>(
333 &self,
334 f: impl FnOnce(&AccountStateForest<AccountStateForestBackend>) -> R,
335 ) -> R {
336 let span = Span::current();
337 tokio::task::block_in_place(|| {
338 span.in_scope(|| {
339 let forest = self.forest.blocking_read();
340 f(&forest)
341 })
342 })
343 }
344
345 fn with_forest_write_blocking<R>(
349 &self,
350 f: impl FnOnce(&mut AccountStateForest<AccountStateForestBackend>) -> R,
351 ) -> R {
352 let span = Span::current();
353 tokio::task::block_in_place(|| {
354 span.in_scope(|| {
355 let mut forest = self.forest.blocking_write();
356 f(&mut forest)
357 })
358 })
359 }
360
361 #[miden_instrument(
369 level = "debug",
370 target = COMPONENT,
371 skip_all,
372 err,
373 )]
374 pub async fn get_block_header(
375 &self,
376 block_num: Option<BlockNumber>,
377 include_mmr_proof: bool,
378 ) -> Result<(Option<BlockHeader>, Option<MmrProof>), GetBlockHeaderError> {
379 let block_header = self.db.select_block_header_by_block_num(block_num).await?;
380 if let Some(header) = block_header {
381 let mmr_proof = if include_mmr_proof {
382 let inner = self.inner.read().await;
383 let mmr_proof = inner.blockchain.open(header.block_num())?;
384 Some(mmr_proof)
385 } else {
386 None
387 };
388 Ok((Some(header), mmr_proof))
389 } else {
390 Ok((None, None))
391 }
392 }
393
394 pub async fn get_notes_by_id(
399 &self,
400 note_ids: Vec<NoteId>,
401 ) -> Result<Vec<NoteRecord>, DatabaseError> {
402 self.db.select_notes_by_id(note_ids).await
403 }
404
405 pub async fn get_batch_inputs(
424 &self,
425 tx_reference_blocks: BTreeSet<BlockNumber>,
426 unauthenticated_note_commitments: BTreeSet<Word>,
427 ) -> Result<BatchInputs, GetBatchInputsError> {
428 if tx_reference_blocks.is_empty() {
429 return Err(GetBatchInputsError::TransactionBlockReferencesEmpty);
430 }
431
432 let note_proofs = self
436 .db
437 .select_note_inclusion_proofs(unauthenticated_note_commitments)
438 .await
439 .map_err(GetBatchInputsError::SelectNoteInclusionProofError)?;
440
441 let note_blocks = note_proofs.values().map(|proof| proof.location().block_num());
443
444 let mut blocks: BTreeSet<BlockNumber> = tx_reference_blocks;
448 blocks.extend(note_blocks);
449
450 let (batch_reference_block, partial_mmr) = {
453 let inner_state = self.inner.read().await;
454
455 let latest_block_num = inner_state.latest_block_num();
456
457 let highest_block_num =
458 *blocks.last().expect("we should have checked for empty block references");
459 if highest_block_num > latest_block_num {
460 return Err(GetBatchInputsError::UnknownTransactionBlockReference {
461 highest_block_num,
462 latest_block_num,
463 });
464 }
465
466 blocks.remove(&latest_block_num);
470
471 let partial_mmr = inner_state
479 .blockchain
480 .partial_mmr_from_blocks(&blocks, latest_block_num)
481 .expect("latest block num should exist and all blocks in set should be < than latest block");
482
483 (latest_block_num, partial_mmr)
484 };
485
486 let mut headers = self
489 .db
490 .select_block_headers(blocks.into_iter().chain(std::iter::once(batch_reference_block)))
491 .await
492 .map_err(GetBatchInputsError::SelectBlockHeaderError)?;
493
494 let header_index = headers
496 .iter()
497 .enumerate()
498 .find_map(|(index, header)| {
499 (header.block_num() == batch_reference_block).then_some(index)
500 })
501 .expect("DB should have returned the header of the batch reference block");
502
503 let batch_reference_block_header = headers.swap_remove(header_index);
505
506 let partial_block_chain = PartialBlockchain::new_unchecked(partial_mmr, headers)
515 .expect("partial mmr and block headers should be consistent");
516
517 Ok(BatchInputs {
518 batch_reference_block_header,
519 note_proofs,
520 partial_block_chain,
521 })
522 }
523
524 pub async fn get_block_inputs(
526 &self,
527 account_ids: Vec<AccountId>,
528 nullifiers: Vec<Nullifier>,
529 unauthenticated_note_commitments: BTreeSet<Word>,
530 reference_blocks: BTreeSet<BlockNumber>,
531 ) -> Result<BlockInputs, GetBlockInputsError> {
532 let unauthenticated_note_proofs = self
536 .db
537 .select_note_inclusion_proofs(unauthenticated_note_commitments)
538 .await
539 .map_err(GetBlockInputsError::SelectNoteInclusionProofError)?;
540
541 let note_proof_reference_blocks =
543 unauthenticated_note_proofs.values().map(|proof| proof.location().block_num());
544
545 let mut blocks = reference_blocks;
547 blocks.extend(note_proof_reference_blocks);
548
549 let (latest_block_number, account_witnesses, nullifier_witnesses, partial_mmr) =
550 self.get_block_inputs_witnesses(&mut blocks, &account_ids, &nullifiers)?;
551
552 let mut headers = self
555 .db
556 .select_block_headers(blocks.into_iter().chain(std::iter::once(latest_block_number)))
557 .await
558 .map_err(GetBlockInputsError::SelectBlockHeaderError)?;
559
560 let latest_block_header_index = headers
563 .iter()
564 .enumerate()
565 .find_map(|(index, header)| {
566 (header.block_num() == latest_block_number).then_some(index)
567 })
568 .expect("DB should have returned the header of the latest block header");
569
570 let latest_block_header = headers.swap_remove(latest_block_header_index);
572
573 let partial_block_chain = PartialBlockchain::new_unchecked(partial_mmr, headers)
582 .expect("partial mmr and block headers should be consistent");
583
584 Ok(BlockInputs::new(
585 latest_block_header,
586 partial_block_chain,
587 account_witnesses,
588 nullifier_witnesses,
589 unauthenticated_note_proofs,
590 ))
591 }
592
593 fn get_block_inputs_witnesses(
600 &self,
601 blocks: &mut BTreeSet<BlockNumber>,
602 account_ids: &[AccountId],
603 nullifiers: &[Nullifier],
604 ) -> Result<BlockInputWitnesses, GetBlockInputsError> {
605 self.with_inner_read_blocking(|inner| {
606 let latest_block_number = inner.latest_block_num();
607
608 let highest_block_number = blocks.last().copied().unwrap_or(latest_block_number);
610 if highest_block_number > latest_block_number {
611 return Err(GetBlockInputsError::UnknownBatchBlockReference {
612 highest_block_number,
613 latest_block_number,
614 });
615 }
616
617 blocks.remove(&latest_block_number);
620
621 let partial_mmr =
632 inner.blockchain.partial_mmr_from_blocks(blocks, latest_block_number).expect(
633 "latest block num should exist and all blocks in set should be < than latest block",
634 );
635
636 let account_witnesses = account_ids
638 .iter()
639 .copied()
640 .map(|account_id| (account_id, inner.account_tree.open_latest(account_id)))
641 .collect::<BTreeMap<AccountId, AccountWitness>>();
642
643 let nullifier_witnesses: BTreeMap<Nullifier, NullifierWitness> = nullifiers
646 .iter()
647 .copied()
648 .map(|nullifier| (nullifier, inner.nullifier_tree.open(&nullifier)))
649 .collect();
650
651 Ok((latest_block_number, account_witnesses, nullifier_witnesses, partial_mmr))
652 })
653 }
654
655 #[miden_instrument(
657 target = COMPONENT,
658 skip_all,
659 fields(
660 account.id=%account_id,
661 nullifiers = %format_array(nullifiers),
662 ),
663 )]
664 pub async fn get_transaction_inputs(
665 &self,
666 account_id: AccountId,
667 nullifiers: &[Nullifier],
668 unauthenticated_note_commitments: Vec<Word>,
669 ) -> Result<TransactionInputs, DatabaseError> {
670 let tree_inputs = self.with_inner_read_blocking(|inner| {
671 let account_commitment = inner.account_tree.get_latest_commitment(account_id);
672
673 let new_account_id_prefix_is_unique = if account_commitment.is_empty() {
674 Some(!inner.account_tree.contains_account_id_prefix_in_latest(account_id.prefix()))
675 } else {
676 None
677 };
678
679 if let Some(false) = new_account_id_prefix_is_unique {
681 return Err(TransactionInputs {
682 new_account_id_prefix_is_unique,
683 ..Default::default()
684 });
685 }
686
687 let nullifiers = nullifiers
688 .iter()
689 .map(|nullifier| NullifierInfo {
690 nullifier: *nullifier,
691 block_num: inner.nullifier_tree.get_block_num(nullifier).unwrap_or_default(),
692 })
693 .collect();
694
695 Ok((account_commitment, nullifiers, new_account_id_prefix_is_unique))
696 });
697 let (account_commitment, nullifiers, new_account_id_prefix_is_unique) = match tree_inputs {
698 Ok(inputs) => inputs,
699 Err(inputs) => return Ok(inputs),
700 };
701
702 let found_unauthenticated_notes = self
703 .db
704 .select_existing_note_commitments(unauthenticated_note_commitments)
705 .await?;
706
707 Ok(TransactionInputs {
708 account_commitment,
709 nullifiers,
710 found_unauthenticated_notes,
711 new_account_id_prefix_is_unique,
712 })
713 }
714
715 pub async fn filter_network_accounts(
717 &self,
718 account_ids: &[AccountId],
719 ) -> Result<HashSet<AccountId>, DatabaseError> {
720 self.db.select_network_accounts_subset(account_ids.to_vec()).await
721 }
722
723 pub async fn chain_tip(&self, finality: Finality) -> BlockNumber {
729 match finality {
730 Finality::Committed => self
731 .inner
732 .read()
733 .instrument(tracing::info_span!("acquire_inner"))
734 .await
735 .latest_block_num(),
736 Finality::Proven => self.proven_tip.read(),
737 }
738 }
739
740 pub async fn load_block(
743 &self,
744 block_num: BlockNumber,
745 ) -> Result<Option<Vec<u8>>, DatabaseError> {
746 if block_num > self.chain_tip(Finality::Committed).await {
747 return Ok(None);
748 }
749 if let Some(block) = self.block_cache.get(block_num) {
750 return Ok(Some(block.block_bytes().to_vec()));
751 }
752 self.block_store.load_block(block_num).await.map_err(Into::into)
753 }
754
755 pub async fn load_proof(
758 &self,
759 block_num: BlockNumber,
760 ) -> Result<Option<Vec<u8>>, DatabaseError> {
761 if block_num > self.chain_tip(Finality::Proven).await {
762 return Ok(None);
763 }
764 if let Some(proof) = self.proof_cache.get(block_num) {
765 return Ok(Some(proof.proof_bytes().to_vec()));
766 }
767 self.block_store.load_proof(block_num).await.map_err(Into::into)
768 }
769
770 pub async fn get_note_script_by_root(
772 &self,
773 root: Word,
774 ) -> Result<Option<NoteScript>, DatabaseError> {
775 self.db.select_note_script_by_root(root).await
776 }
777}