Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
40 commits
Select commit Hold shift + click to select a range
7195581
Initial apply_block lock-free refactor
sergerad Jul 15, 2026
5c7dbeb
Add snapshot count observability
sergerad Jul 15, 2026
b7699fc
Merge branch 'next' of github.com:0xMiden/miden-node into sergerad-lo…
sergerad Jul 20, 2026
46b8df6
Fix ui tests
sergerad Jul 20, 2026
bb3cf85
Lint
sergerad Jul 20, 2026
891d702
Warn about long lived snapshots
sergerad Jul 20, 2026
8c2bc02
Warn on snapshot count
sergerad Jul 20, 2026
8373d9b
Scape all requests
sergerad Jul 20, 2026
3a27431
Lint
sergerad Jul 20, 2026
8fa1502
Merge branch 'next' of github.com:0xMiden/miden-node into sergerad-lo…
sergerad Jul 22, 2026
c21a950
Replace result with control flow
sergerad Jul 22, 2026
0b562a8
Compartmentalize apply_block
sergerad Jul 22, 2026
ddf4f15
snapshots_live fn
sergerad Jul 22, 2026
6556067
Shutdown writer
sergerad Jul 22, 2026
2ad2bbf
Lint
sergerad Jul 22, 2026
09b11f8
Merge branch 'next' of github.com:0xMiden/miden-node into sergerad-lo…
sergerad Jul 22, 2026
dc1e08a
Fix comments
sergerad Jul 22, 2026
d7e6c52
Lint
sergerad Jul 23, 2026
d14362e
Add branch comment
sergerad Jul 23, 2026
d6bf7bd
Fix shutdown
sergerad Jul 23, 2026
220f911
Shutdown integration without token
sergerad Jul 23, 2026
8a531c2
integrate cancellation token
sergerad Jul 23, 2026
3b648ff
Refactor writer task management
sergerad Jul 23, 2026
aea62ce
Split module
sergerad Jul 23, 2026
7f8ebc5
WriterTask
sergerad Jul 23, 2026
c528f19
Add benchmarks
sergerad Jul 23, 2026
3ceabd1
Revert "Add benchmarks"
sergerad Jul 24, 2026
0a7f05a
Merge branch 'next' of github.com:0xMiden/miden-node into sergerad-lo…
sergerad Jul 24, 2026
927a777
Merge branch 'next' of github.com:0xMiden/miden-node into sergerad-lo…
sergerad Jul 28, 2026
974aa12
Llint
sergerad Jul 28, 2026
29765f0
Merge branch 'next' of github.com:0xMiden/miden-node into sergerad-lo…
sergerad Jul 28, 2026
5b6f022
RM deadcode
sergerad Jul 29, 2026
4369f8d
State::for_tests
sergerad Jul 30, 2026
858ee24
BlockWriter and WriteWorker
sergerad Jul 30, 2026
9b40e4b
Merge branch 'next' of github.com:0xMiden/miden-node into sergerad-lo…
sergerad Jul 30, 2026
d49626f
Merge branch 'next' of github.com:0xMiden/miden-node into sergerad-lo…
sergerad Jul 30, 2026
082ed21
StateView (#2415)
sergerad Jul 31, 2026
d9b0f04
Add ScopedBlockNum
sergerad Jul 31, 2026
db5fe1e
Lint
sergerad Jul 31, 2026
f12d4ce
Fix worker module position
sergerad Jul 31, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ miden-crypto = { version = "0.28" }

# External dependencies
anyhow = { version = "1.0" }
arc-swap = { version = "1.7" }
assert_matches = { version = "1.5" }
axum = { version = "0.8" }
backon = { version = "1.6" }
Expand Down
40 changes: 30 additions & 10 deletions bin/node/src/commands/modes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ use miden_node_proto::clients::{
ValidatorClient,
};
use miden_node_rpc::{PreAuthSubmission, Rpc, RpcMode, SequencerInternal, ValidatorClients};
use miden_node_store::State;
use miden_node_store::{BlockWriter, ProofWriter, State, WriterTask};
use miden_node_utils::clap::{GrpcOptionsInternal, duration_to_human_readable_string};
use miden_node_utils::formatting::format_endpoint;
use miden_node_utils::shutdown::CancellationToken;
Expand Down Expand Up @@ -57,11 +57,14 @@ impl SequencerCommand {
let runtime = self.runtime.runtime_config(&self.store);
self.block_producer.validate()?;
let network_tx_auth = self.runtime.rpc.network_tx_auth()?;
let state = load_state(&runtime).await?;
let (state, block_writer, proof_writer, writer_task) =
load_state(&runtime, shutdown.clone()).await?;
let _disk_monitor = state.spawn_disk_monitor(shutdown.clone());

let sequencer = Sequencer {
store: Arc::clone(&state),
state: Arc::clone(&state),
block_writer,
proof_writer,
validator_urls: self.external_services.validator_urls.clone(),
validator_timeout: self.external_services.validator_timeout,
batch_prover_url: self.block_producer.batch.prover_url,
Expand All @@ -75,24 +78,25 @@ impl SequencerCommand {
batch_workers: self.block_producer.batch.workers,
}
.spawn(shutdown.clone())
.await
.context("failed to spawn sequencer")?;
let block_producer = sequencer.api();

let rpc = Rpc {
listener: bind_rpc(runtime.rpc_listen).await?,
store: state,
state,
mode: RpcMode::sequencer(
block_producer.clone(),
self.external_services.validator_clients()?,
),
sync_writers: None,
ntx_builder: Some(self.external_services.ntx_builder_client()?),
grpc_options: runtime.external_grpc_options,
network_tx_auth,
};
let mut tasks = Tasks::new();
tasks.spawn("sequencer", sequencer.wait());
tasks.spawn("RPC server", rpc.serve(shutdown.clone()));
tasks.spawn("store block writer", join_store_writer(writer_task));
if let Some(internal_listen) = self.internal {
let sequencer_internal = SequencerInternal {
listener: bind_rpc(internal_listen).await?,
Expand Down Expand Up @@ -240,19 +244,22 @@ impl FullNodeCommand {
let source_rpc = self.sync.source_rpc_client()?;
let pre_auth = self.pre_auth_submission()?;
let network_tx_auth = self.runtime.rpc.network_tx_auth()?;
let state = load_state(&runtime).await?;
let (state, block_writer, proof_writer, writer_task) =
load_state(&runtime, shutdown.clone()).await?;
let _disk_monitor = state.spawn_disk_monitor(shutdown.clone());

let rpc = Rpc {
listener: bind_rpc(runtime.rpc_listen).await?,
store: state,
state,
mode: RpcMode::full_node(source_rpc, self.sync.readiness_threshold, pre_auth),
sync_writers: Some((block_writer, proof_writer)),
ntx_builder: None,
grpc_options: runtime.external_grpc_options,
network_tx_auth,
};
let mut tasks = Tasks::new();
tasks.spawn("RPC server", rpc.serve(shutdown.clone()));
tasks.spawn("store block writer", join_store_writer(writer_task));

tasks.join_next_or_cancelled(shutdown).await
}
Expand Down Expand Up @@ -326,16 +333,29 @@ impl SyncOptions {
}
}

async fn load_state(runtime: &RuntimeConfig) -> anyhow::Result<Arc<State>> {
let state = State::load_with_database_options(
async fn load_state(
runtime: &RuntimeConfig,
shutdown: CancellationToken,
) -> anyhow::Result<(Arc<State>, BlockWriter, ProofWriter, WriterTask)> {
let loaded = State::load_with_database_options(
&runtime.data_directory,
runtime.storage_options.clone(),
runtime.database_options,
shutdown,
)
.await
.context("failed to load state")?;

Ok(Arc::new(state))
Ok(loaded.start())
}

/// Supervises the store's write worker task.
///
/// On shutdown the task-drain loop waits for the writer to finish any in-flight block write and
/// close its storage; an early exit or panic surfaces through the task set like any other task
/// failure.
async fn join_store_writer(writer_task: WriterTask) -> anyhow::Result<()> {
writer_task.await.map_err(anyhow::Error::from)
}

async fn bind_rpc(listen: SocketAddr) -> anyhow::Result<TcpListener> {
Expand Down
33 changes: 22 additions & 11 deletions bin/node/src/commands/recover.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,8 @@ use std::time::Duration;
use anyhow::Context;
use miden_node_proto::clients::{Builder, ValidatorClient};
use miden_node_proto::generated::validator::BlockSubscriptionRequest;
use miden_node_store::State;
use miden_node_store::state::Finality;
use miden_node_store::{BlockWriter, State, WriterTask};
use miden_node_utils::shutdown::CancellationToken;
use miden_protocol::block::{BlockNumber, SignedBlock};
use miden_protocol::utils::serde::Deserializable;
use tokio_stream::StreamExt;
Expand Down Expand Up @@ -36,21 +36,28 @@ pub struct RecoverCommand {

impl RecoverCommand {
pub async fn handle(self) -> anyhow::Result<()> {
let state = self.load_state().await?;
let (state, block_writer, writer_task) = self.load_state().await?;
let validator = self.validator_client()?;
recover_from_validator(&state, validator).await
let result = recover_from_validator(&state, &block_writer, validator).await;
// Wait for the writer to drain and release the backing storage before the process exits.
block_writer.stop(writer_task).await;
result
}

async fn load_state(&self) -> anyhow::Result<Arc<State>> {
let state = State::load_with_database_options(
async fn load_state(&self) -> anyhow::Result<(Arc<State>, BlockWriter, WriterTask)> {
// Recovery is not wired into the node's shutdown token; the writer exits once the
// `BlockWriter` (holding the only write handle) is dropped after recovery completes.
let loaded = State::load_with_database_options(
&self.data_directory,
self.store.storage.clone().into(),
self.store.sqlite.database_options(),
CancellationToken::new(),
)
.await
.context("failed to load state")?;

Ok(Arc::new(state))
let (state, block_writer, _proof_writer, writer_task) = loaded.start();
Ok((state, block_writer, writer_task))
}

fn validator_client(&self) -> anyhow::Result<ValidatorClient> {
Expand All @@ -66,7 +73,8 @@ impl RecoverCommand {

/// Streams blocks from the validator into the local store until the chain tip is reached.
async fn recover_from_validator(
state: &Arc<State>,
state: &State,
block_writer: &BlockWriter,
mut validator: ValidatorClient,
) -> anyhow::Result<()> {
// Capture the validator's chain tip as the recovery target. The validator's block stream
Expand All @@ -81,7 +89,7 @@ async fn recover_from_validator(
.chain_tip,
);

let local_tip = state.chain_tip(Finality::Committed).await;
let local_tip = state.committed_tip();
if local_tip >= validator_tip {
info!(
target: LOG_TARGET,
Expand Down Expand Up @@ -111,7 +119,10 @@ async fn recover_from_validator(
let block = SignedBlock::read_from_bytes(&event.block)
.context("failed to deserialize block from validator")?;
let block_num = block.header().block_num();
state.apply_block(block).await.context("failed to apply recovered block")?;
block_writer
.apply_block(block)
.await
.context("failed to apply recovered block")?;
info!(target: LOG_TARGET, block_number = %block_num.as_u32(), "Applied recovered block");

// Stop once we reach the tip captured at the start of recovery.
Expand All @@ -121,7 +132,7 @@ async fn recover_from_validator(
}

// The stream can end before reaching the tip if the validator restarts or drops the connection.
let final_tip = state.chain_tip(Finality::Committed).await;
let final_tip = state.committed_tip();
anyhow::ensure!(
final_tip >= validator_tip,
"validator block stream ended at block {} before reaching the chain tip {}",
Expand Down
Loading
Loading