use std::{
collections::{HashMap, HashSet},
sync::Arc,
};
use anyhow::Context as _;
pub use network::RpcConfig;
use zksync_concurrency::{ctx, error::Wrap as _, limiter, net, scope, time};
use zksync_consensus_bft as bft;
use zksync_consensus_engine::EngineManager;
use zksync_consensus_network as network;
use zksync_consensus_roles::{node, validator};
use zksync_protobuf::kB;
#[cfg(test)]
mod tests;
#[derive(Debug)]
pub struct Config {
pub build_version: Option<semver::Version>,
pub server_addr: std::net::SocketAddr,
pub public_addr: net::Host,
pub max_payload_size: usize,
pub max_tx_size: usize,
pub view_timeout: time::Duration,
pub gossip_dynamic_inbound_limit: usize,
pub gossip_static_inbound: HashSet<node::PublicKey>,
pub gossip_static_outbound: HashMap<node::PublicKey, net::Host>,
pub rpc: RpcConfig,
pub node_key: node::SecretKey,
pub validator_key: Option<validator::SecretKey>,
pub debug_page: Option<network::debug_page::Config>,
}
impl Config {
pub(crate) fn gossip(&self) -> network::GossipConfig {
network::GossipConfig {
key: self.node_key.clone(),
dynamic_inbound_limit: self.gossip_dynamic_inbound_limit,
static_inbound: self.gossip_static_inbound.clone(),
static_outbound: self.gossip_static_outbound.clone(),
}
}
}
#[derive(Debug)]
pub struct Executor {
pub config: Config,
pub engine_manager: Arc<EngineManager>,
}
impl Executor {
pub async fn run(self, ctx: &ctx::Ctx) -> anyhow::Result<()> {
let res = scope::run!(ctx, |ctx, s| async {
tracing::trace!("Starting executor.");
if self.engine_manager.head().next() < self.engine_manager.first_block() {
s.spawn(async {
self.spawn_components(ctx, self.network_config(), None, None)
.await
.wrap("Components for pregenesis stopped")
});
}
let mut cur_epoch = self
.engine_manager
.wait_until_epoch_schedule_populated(ctx)
.await?;
loop {
let schedule = self
.engine_manager
.wait_for_validator_schedule(ctx, cur_epoch)
.await?
.schedule;
let epoch = ctx::NoCopy(cur_epoch);
s.spawn(async {
let epoch = epoch.into();
self.spawn_components(ctx, self.network_config(), Some(epoch), Some(schedule))
.await
.wrap(format!("Components for epoch {} stopped", epoch))
});
cur_epoch = cur_epoch.next();
}
})
.await;
match res {
Ok(()) | Err(ctx::Error::Canceled(_)) => Ok(()),
Err(ctx::Error::Internal(err)) => Err(err),
}
}
async fn spawn_components(
&self,
ctx: &ctx::Ctx,
network_config: network::Config,
epoch: Option<validator::EpochNumber>,
validator_schedule: Option<validator::Schedule>,
) -> ctx::Result<()> {
if let Some(epoch_number) = epoch {
tracing::trace!("Spawning components for epoch {}", epoch_number);
} else {
tracing::trace!("Spawning components for pregenesis");
}
let (consensus_send, consensus_recv) = bft::create_input_channel();
let (network_send, network_recv) = ctx::channel::unbounded();
scope::run!(ctx, |ctx, s| async {
tracing::trace!("Starting network component.");
let (net, runner) = network::Network::new(
network_config,
self.engine_manager.clone(),
epoch,
consensus_send,
network_recv,
)?;
net.register_metrics();
s.spawn(async { runner.run(ctx, true).await.context("Network stopped") });
if let Some(cfg) = self.config.debug_page.clone() {
s.spawn(async {
network::debug_page::Server::new(cfg, net)
.run(ctx, true)
.await
.context("Debug page server stopped")
});
}
if epoch.is_some()
&& validator_schedule.is_some()
&& self.config.validator_key.is_some()
&& validator_schedule
.as_ref()
.unwrap()
.contains(&self.config.validator_key.clone().unwrap().public())
{
tracing::trace!(
"This node is an active validator for epoch {}. Starting bft component.",
epoch.unwrap()
);
s.spawn(async {
bft::Config::new(
self.config.validator_key.clone().unwrap(),
self.config.max_payload_size,
self.config.view_timeout,
self.engine_manager.clone(),
epoch.unwrap(),
)?
.run(ctx, network_send, consensus_recv)
.await
.context("Consensus stopped")
});
} else {
tracing::trace!(
"Running the node in non-validator mode for epoch {}.",
epoch.unwrap()
);
}
Ok(())
})
.await?;
Ok(())
}
fn network_config(&self) -> network::Config {
network::Config {
build_version: self.config.build_version.clone(),
server_addr: net::tcp::ListenerAddr::new(self.config.server_addr),
public_addr: self.config.public_addr.clone(),
gossip: self.config.gossip(),
validator_key: self.config.validator_key.clone(),
ping_timeout: Some(time::Duration::seconds(10)),
max_block_size: self.config.max_payload_size.saturating_add(kB),
max_tx_size: self.config.max_tx_size,
max_block_queue_size: 20,
tcp_accept_rate: limiter::Rate {
burst: 10,
refresh: time::Duration::milliseconds(100),
},
rpc: self.config.rpc.clone(),
}
}
}