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