use std::sync::atomic::Ordering;
use miden_node_proto::{DecodeMessageExt, generated as grpc};
use miden_node_tracing::spawn::spawn_blocking_in_current_span;
use miden_node_tracing::{ErrorReport, Instrument, info_span, miden_instrument};
use miden_protocol::Word;
use miden_protocol::block::{BlockHeader, BlockNumber, ProposedBlock};
use miden_protocol::crypto::dsa::ecdsa_k256_keccak::{PublicKey, Signature};
use miden_protocol::protocol_config::ProtocolConfig;
use miden_protocol::transaction::{TransactionHeader, TransactionId};
use super::{StatusResultExt, ValidatorService};
use crate::COMPONENT;
#[tonic::async_trait]
impl grpc::server::miden_validator_v1_validator_service::SignBlock for ValidatorService {
type Input = grpc::miden::validator::v1::SignBlockRequest;
type Output = (Signature, Word, PublicKey);
#[miden_instrument(
target = COMPONENT,
err,
)]
fn decode(request: grpc::miden::validator::v1::SignBlockRequest) -> tonic::Result<Self::Input> {
Ok(request)
}
#[miden_instrument(
target = COMPONENT,
err,
)]
fn encode(
output: Self::Output,
) -> tonic::Result<grpc::miden::validator::v1::SignBlockResponse> {
let (signature, block_commitment, public_key) = output;
Ok(grpc::miden::validator::v1::SignBlockResponse {
signature: Some(signature.into()),
block_commitment: Some(block_commitment.into()),
public_key: Some((&public_key).into()),
})
}
async fn handle(
&self,
request: Self::Input,
_metadata: &tonic::metadata::MetadataMap,
_extensions: &tonic::codegen::http::Extensions,
) -> tonic::Result<Self::Output> {
let _guard = self.serve_lock.try_read().map_err(|_| {
tonic::Status::resource_exhausted("validator is busy streaming a backup")
})?;
let _permit = self
.sign_block_semaphore
.acquire()
.instrument(info_span!("acquire_permit"))
.await
.or_internal("sign_block semaphore closed")?;
let (proposed_block, protocol_config, protocol_config_commitment) =
spawn_blocking_in_current_span(move || {
let request = request
.decode_and_build_unchecked()
.map_err(miden_node_proto::errors::ConversionError::into_status)?;
let protocol_config = request.protocol_config;
let protocol_config_commitment = request.block_header.protocol_config_commitment();
let proposed_block = ProposedBlock::new_at(
request.block_inputs,
request.tx_batches.into_vec(),
request.block_header.timestamp(),
)
.map(|block| {
block
.with_next_validator_config(request.block_header.validator_config().clone())
.with_next_protocol_config(
request.block_header.next_protocol_config().cloned(),
)
})
.or_invalid_argument("Failed to build proposed block")?;
Ok::<_, tonic::Status>((
proposed_block,
protocol_config,
protocol_config_commitment,
))
})
.await
.or_internal("Block decoding task failed")??;
let protocol_config = self
.resolve_protocol_config(protocol_config_commitment, protocol_config)
.await?;
let block_num = proposed_block.block_num();
let previous_backup = self
.block_store
.load_block(block_num)
.await
.or_internal("Failed to load previous block backup")?;
let chain_tip = self
.db
.load_chain_tip()
.await
.or_internal("Failed to load chain tip")?
.ok_or_else(|| tonic::Status::internal("Chain tip not found in database"))?;
let block_transactions: Vec<TransactionId> =
proposed_block.transactions().map(TransactionHeader::id).collect();
let chain_tip_num = chain_tip.block_num();
let (signature, header) = self
.validate_block(proposed_block, chain_tip)
.await
.or_invalid_argument("Failed to validate block")?;
let block_commitment = header.commitment();
let new_block_num = header.block_num().as_u32();
let is_replacement = header.block_num() == chain_tip_num;
self.persist_signed_block(
header,
protocol_config,
block_transactions,
is_replacement,
previous_backup,
)
.await?;
self.committed_tip.send_replace(BlockNumber::from(new_block_num));
if !is_replacement {
self.signed_blocks_count.fetch_add(1, Ordering::Relaxed);
}
Ok((signature, block_commitment, self.signer.public_key()))
}
}
impl ValidatorService {
async fn resolve_protocol_config(
&self,
commitment: Word,
supplied: Option<ProtocolConfig>,
) -> tonic::Result<ProtocolConfig> {
let stored = self
.db
.load_protocol_config(commitment)
.await
.or_internal("Failed to load protocol config")?;
match (supplied, stored) {
(Some(supplied), Some(stored)) if supplied != stored => Err(tonic::Status::internal(
format!("Stored protocol config {commitment} differs from the supplied config"),
)),
(Some(supplied), _) => Ok(supplied),
(None, Some(stored)) => Ok(stored),
(None, None) => Err(tonic::Status::invalid_argument(format!(
"Protocol config {commitment} is not stored"
))),
}
}
async fn persist_signed_block(
&self,
header: BlockHeader,
protocol_config: ProtocolConfig,
transactions: Vec<TransactionId>,
is_replacement: bool,
previous_backup: Option<Vec<u8>>,
) -> tonic::Result<()> {
let block_num = header.block_num();
let persisted = if is_replacement {
self.db.replace_signed_block(header, protocol_config, transactions).await
} else {
self.db.insert_signed_block(header, protocol_config, transactions).await
};
let Err(err) = persisted else {
return Ok(());
};
if let Some(previous_backup) = previous_backup {
self.block_store
.save_block(block_num, &previous_backup)
.await
.map_err(|restore_err| {
tonic::Status::internal(format!(
"Failed to persist block header: {}; failed to restore block backup: {restore_err}",
err.as_report()
))
})?;
}
Err(tonic::Status::internal(format!(
"Failed to persist block header: {}",
err.as_report()
)))
}
}