use std::sync::{atomic::AtomicUsize, Arc};
use fetch::RequestItem;
use tracing::Instrument;
pub(crate) use validator_addrs::*;
use zksync_concurrency::{ctx, scope, sync};
use zksync_consensus_engine::{EngineManager, Transaction};
use zksync_consensus_roles::{node, validator};
use crate::{gossip::ValidatorAddrsWatch, io, pool::PoolWatch, Config, MeteredStreamStats};
mod fetch;
mod handshake;
pub mod loadtest;
mod runner;
#[cfg(test)]
mod testonly;
#[cfg(test)]
mod tests;
mod validator_addrs;
#[derive(Debug)]
pub(crate) struct Connection {
pub(crate) key: node::PublicKey,
pub(crate) build_version: Option<semver::Version>,
pub(crate) stats: Arc<MeteredStreamStats>,
}
pub(crate) struct Network {
pub(crate) epoch_number: Option<validator::EpochNumber>,
pub(crate) cfg: Config,
pub(crate) inbound: PoolWatch<node::PublicKey, Arc<Connection>>,
pub(crate) outbound: PoolWatch<node::PublicKey, Arc<Connection>>,
pub(crate) validator_addrs: ValidatorAddrsWatch,
pub(crate) engine_manager: Arc<EngineManager>,
pub(crate) consensus_sender: sync::prunable_mpsc::Sender<io::ConsensusReq>,
pub(crate) fetch_queue: fetch::Queue,
pub(crate) tx_pool: sync::broadcast::Sender<Transaction>,
pub(crate) push_validator_addrs_calls: AtomicUsize,
}
impl Network {
pub(crate) fn new(
cfg: Config,
engine_manager: Arc<EngineManager>,
epoch_number: Option<validator::EpochNumber>,
consensus_sender: sync::prunable_mpsc::Sender<io::ConsensusReq>,
) -> Arc<Self> {
Arc::new(Self {
epoch_number,
consensus_sender,
inbound: PoolWatch::new(
cfg.gossip.static_inbound.clone(),
cfg.gossip.dynamic_inbound_limit,
),
outbound: PoolWatch::new(cfg.gossip.static_outbound.keys().cloned().collect(), 0),
validator_addrs: ValidatorAddrsWatch::default(),
cfg,
fetch_queue: fetch::Queue::default(),
tx_pool: engine_manager.tx_pool_sender(),
push_validator_addrs_calls: 0.into(),
engine_manager,
})
}
pub(crate) fn genesis_hash(&self) -> validator::GenesisHash {
self.engine_manager.genesis_hash()
}
pub(crate) fn first_block(&self) -> validator::BlockNumber {
self.engine_manager.first_block()
}
pub(crate) fn validator_schedule(&self) -> anyhow::Result<Option<validator::Schedule>> {
match self.epoch_number {
None => Ok(None),
Some(epoch_number) => match self.engine_manager.validator_schedule(epoch_number) {
Some(vs) => Ok(Some(vs.schedule)),
None => anyhow::bail!(
"validator schedule not found for epoch {} in network component",
epoch_number
),
},
}
}
pub(crate) async fn run_block_fetcher(&self, ctx: &ctx::Ctx) {
let sem = sync::Semaphore::new(self.cfg.max_block_queue_size);
let _: ctx::OrCanceled<()> = scope::run!(ctx, |ctx, s| async {
let mut next = self.engine_manager.queued().next();
loop {
let permit = sync::acquire(ctx, &sem).await?;
let number = ctx::NoCopy(next);
next = next + 1;
s.spawn(
async {
let _permit = permit;
let number = number.into();
let _: ctx::OrCanceled<()> = scope::run!(ctx, |ctx, s| async {
s.spawn_bg(
self.fetch_queue
.request(ctx, RequestItem::Block(number))
.instrument(tracing::trace_span!("fetch_block_request")),
);
self.engine_manager.wait_until_queued(ctx, number).await?;
Err(ctx::Canceled)
})
.instrument(tracing::trace_span!("wait_for_block_to_queue"))
.await;
self.engine_manager.wait_until_persisted(ctx, number).await
}
.instrument(tracing::trace_span!("fetch_block_from_peer", l2_block = %next)),
);
}
})
.await;
}
}