Skip to main content

miden_node_store/state/
lifecycle.rs

1//! Store lifecycle: loading the state, starting its write worker, and stopping the store.
2
3use std::num::NonZeroUsize;
4use std::path::Path;
5use std::sync::Arc;
6use std::sync::atomic::AtomicUsize;
7
8use arc_swap::ArcSwap;
9use miden_node_tracing::spawn::spawn_blocking_in_current_span;
10use miden_node_tracing::{ErrorReport, miden_instrument};
11use miden_node_utils::clap::StorageOptions;
12use miden_node_utils::shutdown::CancellationToken;
13use tokio::sync::{mpsc, watch};
14use tokio::task::JoinHandle;
15use tracing::Instrument;
16
17use crate::account_state_forest::AccountStateForestBackend;
18use crate::accounts::AccountTreeWithHistory;
19use crate::blocks::BlockStore;
20use crate::db::Db;
21use crate::errors::StateInitializationError;
22use crate::proven_tip::ProvenTipWriter;
23use crate::state::loader::{
24    ACCOUNT_STATE_FOREST_STORAGE_DIR,
25    ACCOUNT_TREE_STORAGE_DIR,
26    AccountForestLoader,
27    NULLIFIER_TREE_STORAGE_DIR,
28    TreeStorage,
29    TreeStorageLoader,
30    load_mmr,
31    verify_account_state_forest_consistency,
32    verify_tree_consistency,
33};
34use crate::state::writer::{WriteRequest, WriteWorker, WriterTask};
35use crate::state::{
36    BlockCache,
37    BlockWriter,
38    ProofCache,
39    ProofWriter,
40    SnapshotGuard,
41    State,
42    StateSnapshot,
43};
44use crate::{COMPONENT, DataDirectory, DatabaseOptions};
45
46/// Awaits a spawned load task, forwarding its result.
47///
48/// The load tasks are never aborted, so a join error is a panic from the task; it is resumed on
49/// the caller so panics keep propagating as panics.
50async fn join_load_task<T>(
51    handle: JoinHandle<Result<T, StateInitializationError>>,
52) -> Result<T, StateInitializationError> {
53    match handle.await {
54        Ok(result) => result,
55        Err(err) => std::panic::resume_unwind(err.into_panic()),
56    }
57}
58
59/// Number of recent committed blocks held in the in-memory cache for replica subscriptions.
60const BLOCK_CACHE_CAPACITY: NonZeroUsize = NonZeroUsize::new(512).unwrap();
61
62/// Number of recent block proofs held in the in-memory cache for replica subscriptions.
63const PROOF_CACHE_CAPACITY: NonZeroUsize = NonZeroUsize::new(512).unwrap();
64
65// LOADED STATE
66// ================================================================================================
67
68/// A loaded store state whose write worker has not been started yet.
69///
70/// Returned by [`State::load`]. [`Self::start`] spawns the writer and yields the read-only
71/// [`State`] together with the write capabilities; since this is the only way to obtain a
72/// [`BlockWriter`], [`BlockWriter::apply_block`] can always make progress.
73#[must_use = "call `start` to spawn the write worker and obtain the state"]
74pub struct LoadedState {
75    state: State,
76    writer: WriteWorker,
77    write_tx: mpsc::Sender<WriteRequest>,
78}
79
80impl LoadedState {
81    /// Spawns the write worker onto the current runtime and returns the read-only state together
82    /// with the write capabilities and the writer task's handle.
83    ///
84    /// [`Arc<State>`] is the read-only view shared with every component that queries or
85    /// subscribes; it exposes no mutating methods. [`BlockWriter`] and [`ProofWriter`] are the
86    /// only handles able to mutate the store — hand them to the single task driving each write
87    /// path (block production or sync, and proof scheduling or sync respectively). The
88    /// capabilities expose no read access: tasks that both read and write receive the
89    /// [`Arc<State>`] alongside their capability.
90    ///
91    /// The writer exits once `shutdown` is cancelled or the [`BlockWriter`] (holding the only
92    /// request sender) is dropped — an in-flight block write always completes first. Awaiting the
93    /// returned handle after either event guarantees the writer has released the tree storage it
94    /// owns; a join error carries a writer panic.
95    ///
96    /// Callers without a token to cancel should pass [`CancellationToken::new`] and stop the store
97    /// via [`BlockWriter::stop`] instead of dropping and joining by hand.
98    pub fn start(
99        self,
100        shutdown: CancellationToken,
101    ) -> (Arc<State>, BlockWriter, ProofWriter, WriterTask) {
102        let writer_task = tokio::spawn(self.writer.run(shutdown));
103        let state = Arc::new(self.state);
104        let block_writer = BlockWriter {
105            block_store: Arc::clone(&state.block_store),
106            write_tx: self.write_tx,
107        };
108        let proof_writer = ProofWriter { state: Arc::clone(&state) };
109        (state, block_writer, proof_writer, WriterTask(writer_task))
110    }
111}
112
113// LOAD
114// ================================================================================================
115
116impl State {
117    /// Loads the state from the data directory.
118    ///
119    /// The loaded state owns all store data structures and exposes subscription methods for
120    /// sequencer and replica tasks. Call [`LoadedState::start`] on the result to spawn the block
121    /// writer and obtain the usable [`State`].
122    #[miden_instrument(
123        target = COMPONENT,
124    )]
125    pub async fn load(
126        data_path: &Path,
127        storage_options: StorageOptions,
128    ) -> Result<LoadedState, StateInitializationError> {
129        Self::load_with_database_options(data_path, storage_options, DatabaseOptions::default())
130            .await
131    }
132
133    /// Loads the state from the data directory using explicit database options.
134    ///
135    /// The loaded state owns all store data structures and exposes subscription methods for
136    /// sequencer and replica tasks. Call [`LoadedState::start`] on the result to spawn the block
137    /// writer and obtain the usable [`State`].
138    #[miden_instrument(
139        target = COMPONENT,
140    )]
141    pub async fn load_with_database_options(
142        data_path: &Path,
143        storage_options: StorageOptions,
144        database_options: DatabaseOptions,
145    ) -> Result<LoadedState, StateInitializationError> {
146        let data_directory = DataDirectory::load(data_path.to_path_buf())
147            .map_err(StateInitializationError::DataDirectoryLoadError)?;
148
149        let block_store = Arc::new(
150            BlockStore::load(data_directory.block_store_dir())
151                .map_err(StateInitializationError::BlockStoreLoadError)?,
152        );
153
154        let database_filepath = data_directory.database_path();
155        let db = Arc::new(
156            Db::load_with_pool_size(
157                database_filepath.clone(),
158                database_options.connection_pool_size,
159            )
160            .await
161            .map_err(StateInitializationError::DatabaseLoadError)?,
162        );
163
164        // The chain tip drives forest loading and the account tree history below; `load_mmr`'s
165        // consistency check also pins the chain MMR to this header.
166        let latest_block_num = db
167            .select_block_header_by_block_num(None)
168            .await?
169            .ok_or(StateInitializationError::GenesisBlockMissing)?
170            .block_num();
171
172        let apply_block_thread_priority = storage_options.apply_block_thread_priority;
173
174        #[cfg(feature = "rocksdb")]
175        let (account_storage_config, nullifier_storage_config, forest_storage_config) = (
176            storage_options.account_tree.into(),
177            storage_options.nullifier_tree.into(),
178            storage_options.account_state_forest.into(),
179        );
180        #[cfg(not(feature = "rocksdb"))]
181        let (account_storage_config, nullifier_storage_config, forest_storage_config) =
182            ((), (), ());
183
184        // The four structures live in independent storages and the database pool supports
185        // concurrent readers, so open and load them concurrently. Each branch is a spawned task
186        // because loading has long synchronous sections (RocksDB opens, MMR hashing, SMT top
187        // reconstruction) that would serialize if polled from a single task. Spawning is eager, so
188        // all four run from this point; the join below only collects their results.
189        let mmr_task = tokio::spawn(
190            {
191                let db = Arc::clone(&db);
192                async move { load_mmr(&db).await }
193            }
194            .in_current_span(),
195        );
196        let account_tree_task = tokio::spawn(
197            {
198                let (db, path) = (Arc::clone(&db), data_path.to_path_buf());
199                async move {
200                    join_load_task(spawn_blocking_in_current_span(move || {
201                        TreeStorage::create(
202                            &path,
203                            &account_storage_config,
204                            ACCOUNT_TREE_STORAGE_DIR,
205                        )
206                    }))
207                    .await?
208                    .load_account_tree(&db)
209                    .await
210                }
211            }
212            .in_current_span(),
213        );
214        let nullifier_tree_task = tokio::spawn(
215            {
216                let (db, path) = (Arc::clone(&db), data_path.to_path_buf());
217                async move {
218                    join_load_task(spawn_blocking_in_current_span(move || {
219                        TreeStorage::create(
220                            &path,
221                            &nullifier_storage_config,
222                            NULLIFIER_TREE_STORAGE_DIR,
223                        )
224                    }))
225                    .await?
226                    .load_nullifier_tree(&db)
227                    .await
228                }
229            }
230            .in_current_span(),
231        );
232        let forest_task = tokio::spawn(
233            {
234                let (db, path) = (Arc::clone(&db), data_path.to_path_buf());
235                async move {
236                    let forest = join_load_task(spawn_blocking_in_current_span(move || {
237                        AccountStateForestBackend::create(
238                            &path,
239                            &forest_storage_config,
240                            ACCOUNT_STATE_FOREST_STORAGE_DIR,
241                        )
242                    }))
243                    .await?
244                    .load_account_state_forest(&db, latest_block_num)
245                    .await?;
246                    verify_account_state_forest_consistency(&forest, &db).await?;
247                    Ok(forest)
248                }
249            }
250            .in_current_span(),
251        );
252        let (blockchain, account_tree, nullifier_tree, forest) = tokio::try_join!(
253            join_load_task(mmr_task),
254            join_load_task(account_tree_task),
255            join_load_task(nullifier_tree_task),
256            join_load_task(forest_task),
257        )?;
258
259        // Verify that tree roots match the expected roots from the database. This catches any
260        // divergence between persistent storage and the database caused by corruption or incomplete
261        // shutdown.
262        verify_tree_consistency(account_tree.root(), nullifier_tree.root(), &db).await?;
263
264        let account_tree = AccountTreeWithHistory::new(account_tree, latest_block_num);
265
266        // Initialize the proven tip from the block store.
267        let proven_tip_init = block_store
268            .load_proven_tip()
269            .map_err(StateInitializationError::ProvenTipLoadError)?;
270        let (proven_tip, _rx) = ProvenTipWriter::new(proven_tip_init);
271
272        // Committed-tip watch: fires after each successful apply_block.
273        let (committed_tip_tx, _rx) = watch::channel(latest_block_num);
274        let committed_tip_tx = Arc::new(committed_tip_tx);
275
276        let block_cache = BlockCache::new(BLOCK_CACHE_CAPACITY);
277        let proof_cache = ProofCache::new(PROOF_CACHE_CAPACITY);
278
279        // Shared counter of live snapshot generations, for observability.
280        let snapshots_live = Arc::new(AtomicUsize::new(0));
281
282        // Create the initial snapshot from reader views of the just-loaded trees.
283        let initial_snapshot = Arc::new(StateSnapshot::new(
284            nullifier_tree
285                .reader()
286                .map_err(|e| StateInitializationError::NullifierTreeIoError(e.as_report()))?,
287            blockchain.clone(),
288            account_tree.reader(),
289            forest
290                .reader()
291                .map_err(|e| StateInitializationError::AccountStateForestIoError(e.as_report()))?,
292            SnapshotGuard::new(Arc::clone(&snapshots_live), latest_block_num),
293        ));
294        let latest_snapshot = Arc::new(ArcSwap::from(initial_snapshot));
295
296        // Assemble the write worker. It owns the writable trees and processes write requests
297        // serially, publishing a new snapshot after each committed block. The caller runs it; it
298        // exits when the shutdown token is cancelled or the `BlockWriter` (holding the only request
299        // sender) is dropped.
300        let (write_tx, write_rx) = mpsc::channel(1);
301        let block_writer = WriteWorker::new(
302            Arc::clone(&db),
303            Arc::clone(&block_store),
304            Arc::clone(&latest_snapshot),
305            Arc::clone(&committed_tip_tx),
306            block_cache.clone(),
307            write_rx,
308            nullifier_tree,
309            account_tree,
310            blockchain,
311            forest,
312            snapshots_live,
313            apply_block_thread_priority,
314        );
315        let state = Self {
316            data_directory: data_path.to_path_buf(),
317            db,
318            block_store,
319            latest_snapshot,
320            proven_tip,
321            committed_tip_tx,
322            block_cache,
323            proof_cache,
324        };
325
326        Ok(LoadedState { state, writer: block_writer, write_tx })
327    }
328
329    /// Loads the state with default options and starts its write worker, detaching the worker
330    /// task.
331    ///
332    /// Test-only helper for tests in sibling crates that don't manage the writer's lifecycle:
333    /// the detached writer exits once the returned [`BlockWriter`] is dropped. Hidden from public
334    /// docs and not part of the stable API.
335    ///
336    /// # Panics
337    ///
338    /// Panics if the state fails to load.
339    #[doc(hidden)]
340    pub async fn for_tests(data_path: &Path) -> (Arc<Self>, BlockWriter, ProofWriter) {
341        let (state, block_writer, proof_writer, _writer_task) =
342            Self::load(data_path, StorageOptions::default())
343                .await
344                .expect("state should load")
345                .start(CancellationToken::new());
346        (state, block_writer, proof_writer)
347    }
348}