//! `validator` subcommand: run a validator node. use crate::{ application::App, config::{NetworkConfig, NodeConfig}, types::{ self, BACKFILL_CHANNEL, BLOCKS_PER_EPOCH, BROADCAST_CHANNEL, Block, CERTIFICATE_CHANNEL, DKG_CHANNEL, DKG_PROBE_CHANNEL, DynamicProvider, FileSecretStore, IO_BUFFER_SIZE, LogReporter, MAILBOX_SIZE, MAX_MESSAGE_SIZE, MAX_PARTICIPANTS, MAX_SUPPORTED_MODE, MESSAGE_RATE, NAMESPACE, PAGE_CACHE_SIZE, PAGE_SIZE, Participants, QMDB_CHANNEL, RESOLVER_CHANNEL, REVEAL, Registrar, SHARING_MODE, Scheme, VOTE_CHANNEL, }, }; use clap::Args; use commonware_broadcast::buffered; use commonware_consensus::{ Reporters, marshal::{ self, core::Actor as MarshalActor, resolver::p2p as marshal_resolver, standard::Deferred, }, simplex::{ SkipBudget, config::{ForwardPolicy, SkipPolicy}, elector::RoundRobin, }, types::{Epoch, FixedEpocher, ViewDelta}, }; use commonware_cryptography::{ed25519, sha256::Sha256}; use commonware_glue::{ dkg::{ SecretStore as _, fence::Fence, orchestrator, probe, reshare, state_sync::{Config as StateSyncConfig, Plan as StateSyncPlan, StateSync}, }, stateful::{ Config as StatefulConfig, Stateful, SyncPlan, db::{DatabaseSet, p2p as qmdb_resolver}, }, }; use commonware_macros::boxed; use commonware_p2p::authenticated::{self, discovery}; use commonware_parallel::Sequential; use commonware_runtime::{Handle, Supervisor as _, buffer::paged::CacheRef, tokio}; use commonware_storage::{archive::prunable, translator::TwoCap}; use commonware_utils::{NZDuration, NZU64, NZUsize, sequence::Unit}; use std::{marker::PhantomData, path::PathBuf, time::Duration}; use tracing::error; /// Start a validator node. #[derive(Args)] pub struct Validator { /// Validator node directory containing config, genesis, secrets, and runtime storage. #[arg(long, default_value = "./data/validator-0")] pub node_dir: PathBuf, /// Run one-time peer state sync for a new late joiner. #[arg(long, default_value_t = false)] pub state_sync: bool, } /// Start every validator actor and run until one stops. #[boxed] pub async fn run(context: tokio::Context, args: Validator) { let node = NodeConfig::load(&args.node_dir).expect("failed to load node config"); let network = NetworkConfig::load(&args.node_dir).expect("failed to load network config"); network.validate().expect("invalid network config"); let genesis_info = types::read_genesis(&args.node_dir).expect("genesis is required"); let participants = Participants::new(&network).expect("invalid participants"); let local = node.public_key(); let partition_prefix = "validator"; let page_cache = CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE); let bootstrappers = network.bootstrappers(&local); let max_peers_per_set = authenticated::peer_set_limit(&network.participants, &local); let mut p2p_config = discovery::Config::local( node.signing_key.clone(), &[NAMESPACE, b"_P2P"].concat(), node.listen, node.dial, bootstrappers, max_peers_per_set, MAX_MESSAGE_SIZE, ); p2p_config.mailbox_size = MAILBOX_SIZE; let (mut p2p, oracle) = discovery::Network::new(context.child("network"), p2p_config); // Channel rates are enforced independently per peer. The network derives each shared inbound // mailbox capacity from the retained-peer bound and quota burst size. let vote_network = p2p.register(VOTE_CHANNEL, MESSAGE_RATE); let certificate_network = p2p.register(CERTIFICATE_CHANNEL, MESSAGE_RATE); let resolver_network = p2p.register(RESOLVER_CHANNEL, MESSAGE_RATE); let backfill_network = p2p.register(BACKFILL_CHANNEL, MESSAGE_RATE); let broadcast_network = p2p.register(BROADCAST_CHANNEL, MESSAGE_RATE); let qmdb_network = p2p.register(QMDB_CHANNEL, MESSAGE_RATE); let dkg_network = p2p.register(DKG_CHANNEL, MESSAGE_RATE); let dkg_probe_network = p2p.register(DKG_PROBE_CHANNEL, MESSAGE_RATE); let p2p_handle = p2p.start(); let provider = DynamicProvider::default(); let store = FileSecretStore::load(args.node_dir.join("secrets.json")) .expect("failed to load secret store"); let mut store_for_genesis = store.clone(); if let Some(share) = store_for_genesis.get_share(Epoch::zero()).await { provider.register( Epoch::zero(), Scheme::signer( NAMESPACE, genesis_info.output.players().clone(), genesis_info.output.public().clone(), share, ) .expect("epoch-0 share must match genesis"), ); } else { provider.register( Epoch::zero(), Scheme::verifier( NAMESPACE, genesis_info.output.players().clone(), genesis_info.output.public().clone(), ), ); } let resolver = marshal_resolver::init( context.child("marshal_resolver"), marshal_resolver::Config { public_key: local.clone(), peer_provider: oracle.clone(), blocker: oracle.clone(), mailbox_size: MAILBOX_SIZE, timeout: Duration::from_secs(2), fetch_retry_timeout: Duration::from_millis(100), priority_requests: false, priority_responses: false, }, backfill_network, ); let (broadcast_engine, buffer) = buffered::Engine::new( context.child("broadcast"), buffered::Config { public_key: local.clone(), mailbox_size: MAILBOX_SIZE, deque_size: 16, priority: false, codec_config: (), peer_provider: oracle.clone(), }, ); let broadcast_handle = broadcast_engine.start(broadcast_network); let finalizations_by_height = prunable::Archive::init( context.child("finalizations_by_height"), archive_config(partition_prefix, "finalizations", page_cache.clone(), ()), ) .await .expect("finalizations archive"); let finalized_blocks = prunable::Archive::init( context.child("finalized_blocks"), archive_config(partition_prefix, "blocks", page_cache.clone(), ()), ) .await .expect("blocks archive"); let genesis_target = as DatabaseSet>::initial_sync_targets(); let genesis = Block::genesis( network.participants[0].clone(), genesis_info.clone(), genesis_target, ); let (probe_actor, probe_mailbox) = probe::Actor::new(probe::Config { context: context.child("dkg_probe"), manager: oracle.clone(), bootstrap: probe::Bootstrap { epoch: Epoch::zero(), participants: genesis_info.participants(), directory: Unit, }, verifier: Scheme::certificate_verifier(NAMESPACE, *genesis_info.output.public().public()), genesis: genesis_info.clone(), strategy: Sequential, blocker: oracle.clone(), blocks_per_epoch: BLOCKS_PER_EPOCH, retry_timeout: NZDuration!(Duration::from_millis(500)), mailbox_size: MAILBOX_SIZE, block_codec_config: (), }); let probe_handle = probe_actor.start(dkg_probe_network); let stateful_startup = context.child("stateful_startup"); let mut plan = SyncPlan::init(&stateful_startup, partition_prefix).await; let should_state_sync = plan.should_state_sync(args.state_sync); let probe_artifact = if should_state_sync { let artifact = probe_mailbox.subscribe().await.expect("probe stopped"); provider.register( artifact.info.epoch, Scheme::verifier( NAMESPACE, artifact.info.output.players().clone(), artifact.info.output.public().clone(), ), ); plan = plan.with_floor(artifact.floor.clone()); Some(artifact) } else { None }; let (marshal_actor, marshal, floor) = MarshalActor::init( context.child("marshal"), finalizations_by_height, finalized_blocks, marshal::Config { provider: provider.clone(), epocher: FixedEpocher::new(BLOCKS_PER_EPOCH), start: plan.marshal_start(genesis.clone()), partition_prefix: partition_prefix.to_string(), mailbox_size: MAILBOX_SIZE, view_retention: ViewDelta::new(10), prunable_items_per_section: NZU64!(10), page_cache: page_cache.clone(), replay_buffer: types::IO_BUFFER_SIZE, key_write_buffer: types::IO_BUFFER_SIZE, value_write_buffer: types::IO_BUFFER_SIZE, block_codec_config: (), max_repair: NZUsize!(10), max_pending_acks: NZUsize!(1), strategy: Sequential, }, ) .await; let (qmdb_actor, qmdb_sync_resolver) = qmdb_resolver::Actor::new( context.child("qmdb_resolver"), qmdb_resolver::Config { peer_provider: oracle.clone(), blocker: oracle.clone(), database: None, mailbox_size: MAILBOX_SIZE, me: Some(local.clone()), timeout: Duration::from_secs(2), fetch_retry_timeout: Duration::from_millis(100), max_serve_ops: NZU64!(16), priority_requests: false, priority_responses: false, }, ); let qmdb_handle = qmdb_actor.start(qmdb_network); let fence_epoch = probe_artifact .as_ref() .map_or_else(Epoch::zero, |artifact| artifact.info.epoch); let state_sync = probe_artifact.map(|artifact| { let floor = plan .floor() .cloned() .expect("state sync startup must have floor"); StateSync { info: artifact.info, floor, } }); let state_sync = StateSyncPlan::init( context.child("dkg_state_sync_plan"), StateSyncConfig { partition_prefix: partition_prefix.to_string(), max_participants: MAX_PARTICIPANTS, max_supported_mode: MAX_SUPPORTED_MODE, }, state_sync, ) .await; let (fence, gate) = Fence::new(fence_epoch); let (reshare_actor, reshare_mailbox) = reshare::Actor::new( context.child("reshare"), reshare::Config { signer: node.signing_key, manager: oracle.clone(), blocker: oracle.clone(), participants_provider: participants, secret_store: store, strategy: Sequential, registrar: Registrar::new(provider.clone()), marshal: marshal.clone(), state_sync: state_sync.clone(), fence, namespace: NAMESPACE, sharing_mode: SHARING_MODE, reveal: REVEAL, mailbox_size: MAILBOX_SIZE, partition_prefix: format!("{partition_prefix}-reshare"), max_participants: MAX_PARTICIPANTS, blocks_per_epoch: BLOCKS_PER_EPOCH, batch_verifier: PhantomData::, }, ); let reshare_handle = reshare_actor.start(dkg_network); let (stateful_actor, stateful_mailbox) = Stateful::init( context.child("stateful"), StatefulConfig { application: App::new(genesis.clone()), db_config: types::db_config(partition_prefix, page_cache.clone()), provider: (), marshal: (marshal.clone(), floor), mailbox_size: MAILBOX_SIZE, plan, resolvers: qmdb_sync_resolver, sync_config: types::sync_config(), prune_config: None, }, ); // The reshare wrapper drives the payload for the stateful application. let deferred = Deferred::new( context.child("deferred"), reshare::Application::new( stateful_mailbox.clone(), reshare_mailbox.clone(), BLOCKS_PER_EPOCH, ), marshal.clone(), FixedEpocher::new(BLOCKS_PER_EPOCH), ); let (orchestrator_actor, orchestrator_mailbox) = orchestrator::Actor::new( context.child("orchestrator"), orchestrator::Config { oracle: oracle.clone(), manager: oracle.clone(), provider: provider.clone(), marshal: marshal.clone(), application: deferred, strategy: Sequential, simplex: orchestrator::SimplexConfig { elector: RoundRobin::::default(), mailbox_size: NZUsize!(3), replay_buffer: IO_BUFFER_SIZE, write_buffer: IO_BUFFER_SIZE, page_cache_page_size: PAGE_SIZE, page_cache_pages: PAGE_CACHE_SIZE, leader_timeout: Duration::from_secs(1), certification_timeout: Duration::from_secs(2), timeout_retry: Duration::from_millis(500), fetch_timeout: Duration::from_secs(2), view_retention: ViewDelta::new(10), skip: SkipPolicy::Enabled { timeout: Duration::from_secs(5), budget: SkipBudget::Participants, }, forward: ForwardPolicy::Disabled, track_historical_votes: false, }, gate, state_sync, blocks_per_epoch: BLOCKS_PER_EPOCH, muxer_size: 128, mailbox_size: MAILBOX_SIZE, partition_prefix: format!("{partition_prefix}-orchestrator"), }, ); let orchestrator_handle = orchestrator_actor.start(vote_network, certificate_network, resolver_network); let reporters = Reporters::from(( stateful_mailbox.clone(), Reporters::from(( orchestrator_mailbox, Reporters::from((reshare_mailbox, LogReporter)), )), )); let marshal_handle = marshal_actor.start(reporters, buffer, resolver); probe_mailbox.attach(marshal.clone()); let stateful_handle = stateful_actor.start(); if let Err(err) = Handle::select([ p2p_handle, broadcast_handle, probe_handle, qmdb_handle, reshare_handle, orchestrator_handle, marshal_handle, stateful_handle, ]) .await { error!(?err, "validator task failed"); } } fn archive_config( prefix: &str, name: &str, page_cache: CacheRef, codec_config: C, ) -> prunable::Config { prunable::Config { translator: TwoCap, metadata_partition: format!("{prefix}-{name}-metadata"), key_partition: format!("{prefix}-{name}-key"), key_page_cache: page_cache, value_partition: format!("{prefix}-{name}-value"), compression: None, codec_config, items_per_section: NZU64!(10), key_write_buffer: IO_BUFFER_SIZE, value_write_buffer: IO_BUFFER_SIZE, replay_buffer: IO_BUFFER_SIZE, } } #[cfg(test)] mod tests { use super::*; use futures::{FutureExt as _, future::pending}; use std::sync::{ Arc, atomic::{AtomicUsize, Ordering}, }; struct CountDrop(Arc); impl Drop for CountDrop { fn drop(&mut self) { self.0.fetch_add(1, Ordering::Relaxed); } } fn pending_handle(dropped: Arc) -> Handle<()> { let count_drop = CountDrop(dropped); Handle::from_future(async move { let _count_drop = count_drop; pending().await }) } #[test] fn successful_actor_completion_stops_validator() { let dropped = Arc::new(AtomicUsize::new(0)); // Model a clean actor exit alongside siblings that would otherwise run forever. let actors = [ Handle::ready(Ok(())), pending_handle(dropped.clone()), pending_handle(dropped.clone()), pending_handle(dropped.clone()), pending_handle(dropped.clone()), pending_handle(dropped.clone()), pending_handle(dropped.clone()), pending_handle(dropped.clone()), ]; // Supervision must complete and abort every pending sibling. assert!(matches!( Handle::select(actors).now_or_never(), Some(Ok(())) )); assert_eq!(dropped.load(Ordering::Relaxed), 7); } }