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;
#[derive(Args)]
pub struct Validator {
#[arg(long, default_value = "./data/validator-0")]
pub node_dir: PathBuf,
#[arg(long, default_value_t = false)]
pub state_sync: bool,
}
#[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);
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 =
<types::Database<tokio::Context> as DatabaseSet<tokio::Context>>::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::<ed25519::Batch>,
},
);
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,
},
);
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::<Sha256>::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<C>(
prefix: &str,
name: &str,
page_cache: CacheRef,
codec_config: C,
) -> prunable::Config<TwoCap, C> {
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<AtomicUsize>);
impl Drop for CountDrop {
fn drop(&mut self) {
self.0.fetch_add(1, Ordering::Relaxed);
}
}
fn pending_handle(dropped: Arc<AtomicUsize>) -> 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));
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()),
];
assert!(matches!(
Handle::select(actors).now_or_never(),
Some(Ok(()))
));
assert_eq!(dropped.load(Ordering::Relaxed), 7);
}
}