use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Once};
use arc_swap::ArcSwap;
use miden_node_utils::ErrorReport;
use miden_node_utils::shutdown::CancellationToken;
use miden_node_utils::tracing::{miden_instrument, miden_span_record};
use miden_protocol::Word;
use miden_protocol::account::AccountUpdateDetails;
use miden_protocol::block::account_tree::AccountMutationSet;
use miden_protocol::block::nullifier_tree::{NullifierMutationSet, NullifierTree};
use miden_protocol::block::{BlockBody, BlockHeader, BlockNumber, Blockchain, SignedBlock};
use miden_protocol::crypto::merkle::smt::LargeSmt;
use miden_protocol::note::{NoteDetails, Nullifier};
use miden_protocol::transaction::OutputNote;
use miden_protocol::utils::serde::Serializable;
use rayon::ThreadPool;
use thread_priority::{ThreadPriority, set_current_thread_priority};
use tokio::sync::{mpsc, watch};
use super::WriteRequest;
use crate::account_state_forest::{
AccountStateForest,
AccountStateForestBackend,
PreparedAccountStateForestBlockUpdate,
};
use crate::accounts::AccountTreeWithHistory;
use crate::blocks::BlockStore;
use crate::db::{Db, NoteRecord};
use crate::errors::{ApplyBlockError, InvalidBlockError};
use crate::state::block_lifecycle::{BlockLifecycle, lifecycle_events_enabled};
use crate::state::loader::TreeStorage;
use crate::state::view::{
PublishedGenerations,
SNAPSHOTS_LIVE_WARN_THRESHOLD,
SnapshotGuard,
StateSnapshot,
};
use crate::state::{BlockCache, BlockNotification};
use crate::{COMPONENT, HistoricalError, LOG_TARGET};
pub(in crate::state) struct WriteWorker {
db: Arc<Db>,
block_store: Arc<BlockStore>,
latest_snapshot: Arc<ArcSwap<StateSnapshot>>,
committed_tip_tx: Arc<watch::Sender<BlockNumber>>,
block_cache: BlockCache,
rx: mpsc::Receiver<WriteRequest>,
nullifier_tree: NullifierTree<LargeSmt<TreeStorage>>,
account_tree: AccountTreeWithHistory<TreeStorage>,
blockchain: Blockchain,
forest: AccountStateForest<AccountStateForestBackend>,
snapshots_live: Arc<AtomicUsize>,
published_generations: PublishedGenerations,
apply_pool: Arc<ThreadPool>,
}
struct PreparedBlockUpdate {
notes: Vec<(NoteRecord, Option<Nullifier>)>,
nullifier_tree_update: NullifierMutationSet,
account_tree_update: AccountMutationSet,
account_forest_update: PreparedAccountStateForestBlockUpdate<AccountStateForestBackend>,
}
impl WriteWorker {
#[expect(clippy::too_many_arguments)]
pub(in crate::state) fn new(
db: Arc<Db>,
block_store: Arc<BlockStore>,
latest_snapshot: Arc<ArcSwap<StateSnapshot>>,
committed_tip_tx: Arc<watch::Sender<BlockNumber>>,
block_cache: BlockCache,
rx: mpsc::Receiver<WriteRequest>,
nullifier_tree: NullifierTree<LargeSmt<TreeStorage>>,
account_tree: AccountTreeWithHistory<TreeStorage>,
blockchain: Blockchain,
forest: AccountStateForest<AccountStateForestBackend>,
snapshots_live: Arc<AtomicUsize>,
apply_block_thread_priority: bool,
) -> Self {
let mut published_generations = PublishedGenerations::new();
let initial_snapshot = latest_snapshot.load_full();
published_generations.record(initial_snapshot.latest_block_num(), &initial_snapshot);
let mut pool_builder =
rayon::ThreadPoolBuilder::new().thread_name(|index| format!("apply_block_{index}"));
if apply_block_thread_priority {
pool_builder = pool_builder.start_handler(|_| raise_thread_priority());
}
let apply_pool =
Arc::new(pool_builder.build().expect("apply_block thread pool should build"));
Self {
db,
block_store,
latest_snapshot,
committed_tip_tx,
block_cache,
rx,
nullifier_tree,
account_tree,
blockchain,
forest,
snapshots_live,
published_generations,
apply_pool,
}
}
pub async fn run(mut self, shutdown: CancellationToken) {
loop {
let req = tokio::select! {
biased;
() = shutdown.cancelled() => break,
req = self.rx.recv() => match req {
Some(req) => req,
None => break,
},
};
let result = self.write_block(req.signed_block).await;
let _ = req.result_tx.send(result);
}
}
#[miden_instrument(
target = COMPONENT,
err,
)]
async fn write_block(&mut self, signed_block: SignedBlock) -> Result<(), ApplyBlockError> {
let header = signed_block.header();
let body = signed_block.body();
let block_num = header.block_num();
let block_commitment = header.commitment();
let num_transactions = body.transactions().as_slice().len();
miden_span_record!(
block.number = %block_num,
block.commitment = %block_commitment,
block.transactions.count = num_transactions,
);
self.validate_block_header(header).await?;
let block_lifecycle =
lifecycle_events_enabled().then(|| BlockLifecycle::from_block_body(block_num, body));
let unresolved_note_nullifiers = block_lifecycle
.as_ref()
.map_or_else(Vec::new, BlockLifecycle::unresolved_note_nullifiers);
let (prepared, signed_block_bytes) = self.prepare_block_update(&signed_block)?;
let PreparedBlockUpdate {
notes,
nullifier_tree_update,
account_tree_update,
account_forest_update,
} = prepared;
let precomputed_public_states = account_forest_update.account_states.clone();
self.block_store.save_block(block_num, &signed_block_bytes).await?;
let prune_tip = self.published_generations.prune_tip(block_num);
let resolved_note_ids = self
.db
.apply_block(
signed_block,
notes,
precomputed_public_states,
unresolved_note_nullifiers,
prune_tip,
)
.await
.map_err(|err| ApplyBlockError::DbUpdateTaskFailed(err.as_report()))?;
let snapshot = self.apply_prepared_mutations(
block_num,
block_commitment,
nullifier_tree_update,
account_tree_update,
account_forest_update,
);
self.published_generations.record(block_num, &snapshot);
self.latest_snapshot.swap(snapshot).mark_superseded();
let snapshots_live = self.check_live_snapshots(block_num);
miden_span_record!(snapshots.live = snapshots_live);
self.block_cache
.push(block_num, BlockNotification::new(block_num, signed_block_bytes))
.expect("block cache receives sequential block numbers");
self.committed_tip_tx.send_replace(block_num);
if let Some(block_lifecycle) = block_lifecycle {
block_lifecycle.emit(&resolved_note_ids);
}
tracing::debug!(target: LOG_TARGET, "Block applied");
Ok(())
}
fn check_live_snapshots(&self, block_num: BlockNumber) -> u64 {
let snapshots_live = self.snapshots_live.load(Ordering::Relaxed) as u64;
if snapshots_live > SNAPSHOTS_LIVE_WARN_THRESHOLD {
tracing::warn!(
target: COMPONENT,
block_num = block_num.as_u32(),
snapshots.live = snapshots_live,
"too many live state snapshots; slow readers are pinning old generations",
);
}
snapshots_live
}
fn prepare_block_update(
&self,
signed_block: &SignedBlock,
) -> Result<(PreparedBlockUpdate, Vec<u8>), ApplyBlockError> {
run_on_pool(&self.apply_pool, || {
let header = signed_block.header();
let body = signed_block.body();
let tx_commitment = body.transactions().commitment();
if header.tx_commitment() != tx_commitment {
return Err(InvalidBlockError::InvalidBlockTxCommitment {
expected: tx_commitment,
actual: header.tx_commitment(),
}
.into());
}
let notes = Self::build_note_records(header, body)?;
let (nullifier_tree_update, account_tree_update) =
self.compute_tree_mutations(header, body)?;
let account_patches =
body.updated_accounts().iter().filter_map(|update| match update.details() {
AccountUpdateDetails::Public(patch) => Some(patch.clone()),
AccountUpdateDetails::Private => None,
});
let account_forest_update = self
.forest
.compute_block_update_mutations(header.block_num(), account_patches)
.map_err(ApplyBlockError::AccountStateForestPreparation)?;
let prepared = PreparedBlockUpdate {
notes,
nullifier_tree_update,
account_tree_update,
account_forest_update,
};
Ok((prepared, signed_block.to_bytes()))
})
}
fn apply_prepared_mutations(
&mut self,
block_num: BlockNumber,
block_commitment: Word,
nullifier_tree_update: NullifierMutationSet,
account_tree_update: AccountMutationSet,
account_forest_update: PreparedAccountStateForestBlockUpdate<AccountStateForestBackend>,
) -> Arc<StateSnapshot> {
let apply_pool = Arc::clone(&self.apply_pool);
run_on_pool(&apply_pool, || {
self.nullifier_tree
.apply_mutations(nullifier_tree_update)
.unwrap_or_else(|error| {
panic!("nullifier tree update failed after database commit: {error}")
});
self.account_tree.apply_mutations(account_tree_update).unwrap_or_else(|error| {
panic!("account tree update failed after database commit: {error}")
});
self.blockchain.push(block_commitment);
self.forest
.apply_precomputed_block_update(block_num, account_forest_update)
.unwrap_or_else(|error| {
panic!("account-state forest update failed after database commit: {error}")
});
Arc::new(StateSnapshot::new(
self.nullifier_tree
.reader()
.expect("nullifier tree snapshot creation should not fail"),
self.blockchain.clone(),
self.account_tree.reader(),
self.forest.reader().expect("forest snapshot creation should not fail"),
SnapshotGuard::new(Arc::clone(&self.snapshots_live), block_num),
))
})
}
#[miden_instrument(
target = COMPONENT,
err,
)]
async fn validate_block_header(&self, header: &BlockHeader) -> Result<(), ApplyBlockError> {
let block_num = header.block_num();
let prev_block = self
.db
.select_block_header_by_block_num(None)
.await?
.ok_or(ApplyBlockError::DbBlockHeaderEmpty)?;
let expected_block_num = prev_block.block_num().child();
if block_num != expected_block_num {
return Err(InvalidBlockError::NewBlockInvalidBlockNum {
expected: expected_block_num,
submitted: block_num,
}
.into());
}
if header.prev_block_commitment() != prev_block.commitment() {
return Err(InvalidBlockError::NewBlockInvalidPrevCommitment.into());
}
Ok(())
}
#[miden_instrument(
target = COMPONENT,
err,
)]
fn compute_tree_mutations(
&self,
header: &BlockHeader,
body: &BlockBody,
) -> Result<(NullifierMutationSet, AccountMutationSet), ApplyBlockError> {
let block_num = header.block_num();
let duplicate_nullifiers: Vec<_> = body
.created_nullifiers()
.iter()
.filter(|&nullifier| self.nullifier_tree.get_block_num(nullifier).is_some())
.copied()
.collect();
if !duplicate_nullifiers.is_empty() {
return Err(InvalidBlockError::DuplicatedNullifiers(duplicate_nullifiers).into());
}
let peaks = self.blockchain.peaks();
if peaks.hash_peaks() != header.chain_commitment() {
return Err(InvalidBlockError::NewBlockInvalidChainCommitment.into());
}
let nullifier_tree_update = self
.nullifier_tree
.compute_mutations(
body.created_nullifiers().iter().map(|nullifier| (*nullifier, block_num)),
)
.map_err(InvalidBlockError::NewBlockNullifierAlreadySpent)?;
if nullifier_tree_update.as_mutation_set().root() != header.nullifier_root() {
return Err(InvalidBlockError::NewBlockInvalidNullifierRoot.into());
}
let account_tree_update = self
.account_tree
.compute_mutations(
body.updated_accounts()
.iter()
.map(|update| (update.account_id(), update.final_state_commitment())),
)
.map_err(|e| match e {
HistoricalError::AccountTreeError(err) => {
InvalidBlockError::NewBlockDuplicateAccountIdPrefix(err)
},
HistoricalError::MerkleError(_) => {
panic!("Unexpected MerkleError during account tree mutation computation")
},
})?;
if account_tree_update.as_mutation_set().root() != header.account_root() {
return Err(InvalidBlockError::NewBlockInvalidAccountRoot.into());
}
Ok((nullifier_tree_update, account_tree_update))
}
#[miden_instrument(
target = COMPONENT,
err,
)]
fn build_note_records(
header: &BlockHeader,
body: &BlockBody,
) -> Result<Vec<(NoteRecord, Option<Nullifier>)>, ApplyBlockError> {
let block_num = header.block_num();
let note_tree = body.compute_block_note_tree();
if note_tree.root() != header.note_root() {
return Err(InvalidBlockError::NewBlockInvalidNoteRoot.into());
}
let notes = body
.output_notes()
.map(|(note_index, note)| {
let (details, attachments, nullifier) = match note {
OutputNote::Public(public) => (
Some(NoteDetails::from(public.as_note())),
public.as_note().attachments().clone(),
Some(public.as_note().nullifier()),
),
OutputNote::Private(private) => (None, private.attachments().clone(), None),
};
let inclusion_path = note_tree.open(note_index);
let note_record = NoteRecord {
block_num,
note_index,
note_id: note.id().as_word(),
metadata: *note.metadata(),
details,
attachments,
inclusion_path,
};
Ok((note_record, nullifier))
})
.collect::<Result<Vec<_>, InvalidBlockError>>()?;
Ok(notes)
}
}
fn run_on_pool<T: Send>(pool: &ThreadPool, op: impl FnOnce() -> T + Send) -> T {
let span = tracing::Span::current();
tokio::task::block_in_place(|| pool.install(|| span.in_scope(op)))
}
fn raise_thread_priority() {
static WARN_ONCE: Once = Once::new();
if let Err(error) = set_current_thread_priority(ThreadPriority::Max) {
WARN_ONCE.call_once(|| {
tracing::warn!(
target: COMPONENT,
?error,
"failed to raise apply-block thread priority; continuing at normal priority",
);
});
}
}