miden_node_store/state/writer/mod.rs
1//! Serialized block-write path for the store state.
2//!
3//! A single [`WriteWorker`] task owns the mutable trees and processes incoming [`WriteRequest`]s
4//! one at a time via an mpsc channel. After each successful commit it publishes a new
5//! [`StateSnapshot`] snapshot via an [`ArcSwap`], making the updated trees immediately visible to
6//! wait-free readers.
7//!
8//! The [`BlockWriter`] and [`ProofWriter`] capabilities defined here are the only handles able
9//! to feed this worker and to commit proofs; the submodules hold their entry points.
10
11mod apply_block;
12mod apply_proof;
13
14mod worker;
15use std::future::Future;
16use std::pin::Pin;
17use std::sync::Arc;
18use std::task::{Context, Poll};
19
20use miden_node_tracing::{ErrorReport, miden_instrument};
21use miden_protocol::block::SignedBlock;
22use miden_protocol::protocol_config::ProtocolConfig;
23use tokio::sync::{mpsc, oneshot};
24pub(in crate::state) use worker::WriteWorker;
25
26use crate::COMPONENT;
27use crate::blocks::BlockStore;
28use crate::errors::ApplyBlockError;
29
30// WRITE CAPABILITIES
31// ================================================================================================
32
33/// The store's block-write capability.
34///
35/// Only handle able to apply blocks; obtained exactly once from [`LoadedState::start`](crate::state::LoadedState::start) and
36/// deliberately not cloneable, so granting it to a single task (the block builder in sequencer
37/// mode, the block sync loop in full-node mode) statically prevents every other component from
38/// writing blocks.
39///
40/// Exposes no read access: holders that also need to query the store receive the
41/// [`Arc<State>`](crate::state::State) returned alongside this capability by
42/// [`LoadedState::start`](crate::state::LoadedState::start).
43pub struct BlockWriter {
44 /// The block store, used to persist proving inputs alongside applied blocks.
45 pub(super) block_store: Arc<BlockStore>,
46 /// Sender for block-write requests to the [`WriteWorker`] task. Never cloned out of this
47 /// struct: the writer exits once it is dropped.
48 pub(super) write_tx: mpsc::Sender<WriteRequest>,
49}
50
51/// The store's proof-write capability.
52///
53/// Only handle able to commit block proofs and advance the proven tip; obtained exactly once from
54/// [`LoadedState::start`](crate::state::LoadedState::start) and deliberately not cloneable, so granting it to a single task (the
55/// proof scheduler in sequencer mode, the proof sync loop in full-node mode) statically prevents
56/// every other component from writing proofs.
57///
58/// Exposes no read access: the held state is only used internally to commit proofs and advance
59/// the proven tip. Holders that also need to query the store receive the
60/// [`Arc<State>`](crate::state::State) returned alongside this capability by
61/// [`LoadedState::start`](crate::state::LoadedState::start).
62pub struct ProofWriter {
63 pub(super) state: Arc<crate::state::State>,
64}
65
66// WRITER TASK
67// ================================================================================================
68
69/// Handle of the store's write worker task, returned by [`LoadedState::start`](crate::state::LoadedState::start).
70///
71/// Awaiting it resolves once the writer has exited and released the tree storage it owns; a join
72/// error carries a writer panic. The newtype ensures [`BlockWriter::stop`] can only be given the
73/// store's own writer task, and deliberately does not expose [`tokio::task::JoinHandle::abort`]:
74/// aborting the writer mid-write could leave the trees lagging the committed database state,
75/// voiding the guarantee that an in-flight block write always completes.
76#[must_use = "await the writer task to observe its exit, or pass it to `BlockWriter::stop`"]
77pub struct WriterTask(pub(super) tokio::task::JoinHandle<()>);
78
79impl Future for WriterTask {
80 type Output = Result<(), tokio::task::JoinError>;
81
82 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
83 Pin::new(&mut self.0).poll(cx)
84 }
85}
86
87// WRITE REQUEST
88// ================================================================================================
89
90/// A request to apply a block, paired with a one-shot channel for the result.
91pub(super) struct WriteRequest {
92 signed_block: SignedBlock,
93 protocol_config: Option<ProtocolConfig>,
94 result_tx: oneshot::Sender<Result<(), ApplyBlockError>>,
95 /// Span of the `apply_block` caller. The worker runs the write under it, keeping the write path
96 /// in the caller's trace across the channel hop.
97 span: miden_node_tracing::Span,
98}
99
100impl BlockWriter {
101 /// Stops the store, waiting until the write worker has released the tree storage it owns.
102 ///
103 /// Consumes the capability — closing the write channel the write worker listens on — and then
104 /// joins the writer task returned by [`LoadedState::start`](crate::state::LoadedState::start). The drop must precede the join
105 /// or the write worker never observes the closed channel; doing both here keeps that ordering
106 /// out of caller hands. Read-only [`State`](crate::state::State) references may outlive the
107 /// stop.
108 ///
109 /// Callers that need the storage released deterministically must use this method instead of
110 /// dropping: the node's `recover` command stops the store before the process exits, and the
111 /// stress-test's store seeding stops it so the same data directory can be re-loaded (or its
112 /// temporary directory deleted) immediately afterwards. The running node does not use this
113 /// method — its writer exits via the shutdown token passed to
114 /// [`State::load`](crate::state::State::load) and is joined through the node's task set.
115 ///
116 /// # Panics
117 ///
118 /// Panics if the writer task panicked.
119 pub async fn stop(self, writer_task: WriterTask) {
120 drop(self);
121 writer_task.await.expect("write worker task should not panic");
122 }
123
124 /// Apply changes of a new block to the DB and in-memory data structures.
125 ///
126 /// Supply the active configuration if its commitment is not yet stored. The configuration
127 /// must match the block header. New configurations are committed with the block.
128 /// If the configuration is omitted and its commitment is not stored, this method returns
129 /// `DatabaseError::ProtocolConfigNotFound` through `ApplyBlockError`.
130 ///
131 /// Blocks are forwarded to the store's write worker task, which processes them one at a
132 /// time.
133 /// Readers are unaffected while a block is being applied: they keep reading from the previous
134 /// in-memory snapshot until the writer atomically publishes the new one.
135 #[miden_instrument(
136 target = COMPONENT,
137 err,
138 )]
139 pub async fn apply_block(
140 &mut self,
141 signed_block: SignedBlock,
142 protocol_config: Option<ProtocolConfig>,
143 ) -> Result<(), ApplyBlockError> {
144 let (result_tx, result_rx) = oneshot::channel();
145 self.write_tx
146 .send(WriteRequest {
147 signed_block,
148 protocol_config,
149 result_tx,
150 span: miden_node_tracing::Span::current(),
151 })
152 .await
153 .map_err(|e| ApplyBlockError::WriterTaskSendFailed(e.as_report()))?;
154 result_rx.await?
155 }
156}