use crate::{
behaviour::{TenzroBehaviour, TenzroBehaviourEvent},
block_sync_proto::{BlockSyncRequest, BlockSyncResponse},
config::NetworkConfig,
consensus_direct_proto::{
ConsensusDirectError, ConsensusDirectRequest, ConsensusDirectResponse,
MAX_INBOUND_STREAMS_PER_PEER,
},
error::{NetworkError, Result},
gossip::{MessageDeduplicator, MessageValidation, validate_gossip_message},
message::{ConsensusMessage, NetworkMessage, MessagePayload},
metrics::NetworkMetrics,
peer_manager::{PeerManager, ManagedPeer},
};
use async_trait::async_trait;
use futures::StreamExt;
use libp2p::{
autonat, dcutr,
gossipsub::{self, IdentTopic, TopicHash},
identify,
kad::{self, QueryResult},
ping, relay,
request_response::{self, InboundRequestId, OutboundRequestId, ResponseChannel},
swarm::SwarmEvent,
Multiaddr, PeerId, Swarm,
};
use parking_lot::Mutex;
use prometheus_client::registry::Registry;
use std::collections::{HashMap, HashSet};
use std::net::IpAddr;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::{mpsc, oneshot};
use tokio::time::{interval, MissedTickBehavior};
use tenzro_types::network::PeerStatus;
fn load_or_generate_keypair(data_dir: &Option<PathBuf>) -> Result<libp2p::identity::Keypair> {
let Some(dir) = data_dir else {
tracing::warn!("No data_dir configured — generating ephemeral keypair (peer ID will change on restart)");
return Ok(libp2p::identity::Keypair::generate_ed25519());
};
let key_path = dir.join("p2p_key");
if key_path.exists() {
match std::fs::read(&key_path) {
Ok(bytes) => {
match libp2p::identity::Keypair::from_protobuf_encoding(&bytes) {
Ok(keypair) => {
tracing::info!("Loaded persistent keypair from {}", key_path.display());
return Ok(keypair);
}
Err(e) => {
tracing::warn!("Failed to decode keypair from {}: {} — generating new one", key_path.display(), e);
}
}
}
Err(e) => {
tracing::warn!("Failed to read keypair file {}: {} — generating new one", key_path.display(), e);
}
}
}
let keypair = libp2p::identity::Keypair::generate_ed25519();
if let Some(parent) = key_path.parent()
&& let Err(e) = std::fs::create_dir_all(parent)
{
tracing::warn!("Failed to create directory {}: {} — keypair will be ephemeral", parent.display(), e);
return Ok(keypair);
}
match keypair.to_protobuf_encoding() {
Ok(bytes) => {
match std::fs::write(&key_path, &bytes) {
Ok(()) => {
tracing::info!("Generated and saved new keypair to {}", key_path.display());
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
if let Err(e) = std::fs::set_permissions(&key_path, std::fs::Permissions::from_mode(0o600)) {
tracing::warn!("Failed to set keypair file permissions: {}", e);
}
}
}
Err(e) => {
tracing::warn!("Failed to write keypair to {}: {} — keypair will be ephemeral", key_path.display(), e);
}
}
}
Err(e) => {
tracing::warn!("Failed to encode keypair: {} — keypair will be ephemeral", e);
}
}
Ok(keypair)
}
fn is_globally_routable(addr: &Multiaddr) -> bool {
use libp2p::multiaddr::Protocol;
for proto in addr.iter() {
match proto {
Protocol::Ip4(ip) => {
let octets = ip.octets();
let is_docker_or_corp = octets[0] == 172 && (octets[1] & 0xf0) == 16;
let is_consumer_nat = octets[0] == 192 && octets[1] == 168;
let is_cgn = octets[0] == 100 && (octets[1] & 0xc0) == 64;
if ip.is_loopback()
|| ip.is_link_local()
|| ip.is_unspecified()
|| ip.is_broadcast()
|| ip.is_documentation()
|| is_docker_or_corp
|| is_consumer_nat
|| is_cgn
{
return false;
}
return true;
}
Protocol::Ip6(ip) => {
let seg0 = ip.segments()[0];
if ip.is_loopback()
|| ip.is_unspecified()
|| ip.is_multicast()
|| (seg0 & 0xffc0) == 0xfe80 || (seg0 & 0xfe00) == 0xfc00 {
return false;
}
return true;
}
_ => {}
}
}
false }
fn extract_ip(addr: &Multiaddr) -> Option<IpAddr> {
use libp2p::multiaddr::Protocol;
for proto in addr.iter() {
match proto {
Protocol::Ip4(ip) => return Some(IpAddr::V4(ip)),
Protocol::Ip6(ip) => return Some(IpAddr::V6(ip)),
_ => continue,
}
}
None
}
fn extract_port(addr: &Multiaddr) -> Option<u16> {
use libp2p::multiaddr::Protocol;
for proto in addr.iter() {
match proto {
Protocol::Tcp(p) | Protocol::Udp(p) => return Some(p),
_ => continue,
}
}
None
}
fn is_observed_port_one_of_ours(observed: &Multiaddr, our_listen_addrs: &[Multiaddr]) -> bool {
let Some(observed_port) = extract_port(observed) else {
return false;
};
our_listen_addrs
.iter()
.filter_map(extract_port)
.any(|p| p == observed_port)
}
#[derive(Debug)]
pub struct InboundBlockSync {
pub peer: PeerId,
pub request_id: InboundRequestId,
pub request: BlockSyncRequest,
}
#[derive(Debug)]
pub struct OutboundBlockSyncResult {
pub peer: PeerId,
pub request_id: OutboundRequestId,
pub result: std::result::Result<BlockSyncResponse, BlockSyncOutboundError>,
}
#[derive(Debug, Clone)]
pub enum PeerEvent {
Connected(PeerId),
Disconnected(PeerId),
}
#[derive(Debug, thiserror::Error)]
pub enum BlockSyncOutboundError {
#[error("dial failure")]
DialFailure,
#[error("request timed out")]
Timeout,
#[error("connection closed")]
ConnectionClosed,
#[error("remote does not speak the block-sync protocol")]
UnsupportedProtocols,
#[error("io error: {0}")]
Io(String),
}
impl From<request_response::OutboundFailure> for BlockSyncOutboundError {
fn from(e: request_response::OutboundFailure) -> Self {
match e {
request_response::OutboundFailure::DialFailure => Self::DialFailure,
request_response::OutboundFailure::Timeout => Self::Timeout,
request_response::OutboundFailure::ConnectionClosed => Self::ConnectionClosed,
request_response::OutboundFailure::UnsupportedProtocols => Self::UnsupportedProtocols,
request_response::OutboundFailure::Io(io) => Self::Io(io.to_string()),
}
}
}
#[async_trait]
pub trait NetworkService: Send + Sync {
async fn broadcast(&self, topic: &str, message: NetworkMessage) -> Result<()>;
async fn send_to(&self, peer_id: PeerId, message: NetworkMessage) -> Result<()>;
async fn subscribe(&self, topic: &str) -> Result<mpsc::UnboundedReceiver<NetworkMessage>>;
async fn connected_peers(&self) -> Result<Vec<PeerId>>;
async fn peer_info(&self, peer_id: &PeerId) -> Result<Option<ManagedPeer>>;
async fn ban_peer(&self, peer_id: &PeerId) -> Result<()>;
async fn unban_peer(&self, peer_id: &PeerId) -> Result<()>;
async fn local_peer_id(&self) -> Result<PeerId>;
async fn dial(&self, addr: Multiaddr) -> Result<()>;
async fn set_validator_registry(&self, registry: std::sync::Arc<dyn crate::peer_manager::ValidatorRegistry>) -> Result<()>;
async fn listen_addresses(&self) -> Result<Vec<Multiaddr>>;
async fn broadcast_to_validators(&self, message: ConsensusMessage) -> Result<usize>;
async fn connected_validator_count(&self) -> Result<usize>;
async fn subscribe_consensus_direct(
&self,
) -> Result<mpsc::UnboundedReceiver<ConsensusMessage>>;
async fn set_mpc_did_resolver(
&self,
resolver: std::sync::Arc<dyn crate::mpc_relay::MpcDidResolver>,
) -> Result<()>;
async fn send_mpc_relay_message(
&self,
message: crate::mpc_relay::MpcRelayRequest,
) -> Result<()>;
async fn subscribe_mpc_relay(
&self,
) -> Result<mpsc::UnboundedReceiver<crate::mpc_relay::MpcRelayRequest>>;
}
#[allow(clippy::large_enum_variant)]
enum NetworkCommand {
Broadcast {
topic: String,
message: NetworkMessage,
response: oneshot::Sender<Result<()>>,
},
Subscribe {
topic: String,
response: oneshot::Sender<Result<mpsc::UnboundedReceiver<NetworkMessage>>>,
},
ConnectedPeers {
response: oneshot::Sender<Result<Vec<PeerId>>>,
},
PeerInfo {
peer_id: PeerId,
response: oneshot::Sender<Result<Option<ManagedPeer>>>,
},
BanPeer {
peer_id: PeerId,
response: oneshot::Sender<Result<()>>,
},
UnbanPeer {
peer_id: PeerId,
response: oneshot::Sender<Result<()>>,
},
LocalPeerId {
response: oneshot::Sender<Result<PeerId>>,
},
Dial {
addr: Multiaddr,
response: oneshot::Sender<Result<()>>,
},
SetValidatorRegistry {
registry: std::sync::Arc<dyn crate::peer_manager::ValidatorRegistry>,
response: oneshot::Sender<Result<()>>,
},
MeshPeerCount {
topic: String,
response: oneshot::Sender<Result<usize>>,
},
ListenAddresses {
response: oneshot::Sender<Result<Vec<Multiaddr>>>,
},
AdmittedMeshPeers {
topic: String,
response: oneshot::Sender<Result<usize>>,
},
SendBlockSyncRequest {
peer: PeerId,
request: BlockSyncRequest,
response: oneshot::Sender<Result<OutboundRequestId>>,
},
SendBlockSyncResponse {
request_id: InboundRequestId,
response_payload: BlockSyncResponse,
response: oneshot::Sender<Result<()>>,
},
SubscribeBlockSyncRequests {
response: oneshot::Sender<Result<mpsc::UnboundedReceiver<InboundBlockSync>>>,
},
SubscribeBlockSyncResults {
response: oneshot::Sender<Result<mpsc::UnboundedReceiver<OutboundBlockSyncResult>>>,
},
SubscribePeerEvents {
response: oneshot::Sender<Result<mpsc::UnboundedReceiver<PeerEvent>>>,
},
BroadcastToValidators {
message: ConsensusMessage,
response: oneshot::Sender<Result<usize>>,
},
ConnectedValidatorCount {
response: oneshot::Sender<Result<usize>>,
},
SubscribeConsensusDirect {
response: oneshot::Sender<Result<mpsc::UnboundedReceiver<ConsensusMessage>>>,
},
SetMpcDidResolver {
resolver: std::sync::Arc<dyn crate::mpc_relay::MpcDidResolver>,
response: oneshot::Sender<Result<()>>,
},
SendMpcRelayMessage {
message: crate::mpc_relay::MpcRelayRequest,
response: oneshot::Sender<Result<()>>,
},
SubscribeMpcRelay {
response: oneshot::Sender<
Result<mpsc::UnboundedReceiver<crate::mpc_relay::MpcRelayRequest>>,
>,
},
Shutdown {
response: oneshot::Sender<Result<()>>,
},
}
pub struct TenzroNetworkService {
command_tx: mpsc::UnboundedSender<NetworkCommand>,
metrics: Arc<NetworkMetrics>,
metrics_registry: Arc<Mutex<Registry>>,
}
impl TenzroNetworkService {
pub async fn new(config: NetworkConfig) -> Result<Self> {
let mut registry = Registry::default();
let metrics = NetworkMetrics::register(&mut registry);
Self::new_with_registry(config, Arc::new(Mutex::new(registry)), metrics).await
}
pub async fn new_with_registry(
config: NetworkConfig,
metrics_registry: Arc<Mutex<Registry>>,
metrics: Arc<NetworkMetrics>,
) -> Result<Self> {
config.validate()?;
let (command_tx, command_rx) = mpsc::unbounded_channel();
let loop_metrics = metrics.clone();
tokio::spawn(async move {
if let Err(e) = run_event_loop(config, command_rx, loop_metrics).await {
tracing::error!("Network event loop error: {}", e);
}
});
Ok(Self {
command_tx,
metrics,
metrics_registry,
})
}
pub fn metrics(&self) -> Arc<NetworkMetrics> {
self.metrics.clone()
}
pub fn metrics_registry(&self) -> Arc<Mutex<Registry>> {
self.metrics_registry.clone()
}
pub async fn shutdown(&self) -> Result<()> {
self.send_command(|response| NetworkCommand::Shutdown { response }).await
}
pub async fn mesh_peer_count(&self, topic: &str) -> Result<usize> {
let topic_owned = topic.to_string();
self.send_command(move |response| NetworkCommand::MeshPeerCount {
topic: topic_owned,
response,
})
.await
}
pub async fn wait_for_mesh(
&self,
topic: &str,
min_peers: usize,
timeout: std::time::Duration,
) -> Result<usize> {
let deadline = tokio::time::Instant::now() + timeout;
#[allow(unused_assignments)]
let mut last_seen = 0usize;
loop {
match self.mesh_peer_count(topic).await {
Ok(count) => {
last_seen = count;
if count >= min_peers {
tracing::info!(
topic = topic,
count = count,
min_peers = min_peers,
"Gossipsub mesh ready"
);
return Ok(count);
}
}
Err(e) => {
return Err(e);
}
}
if tokio::time::Instant::now() >= deadline {
tracing::warn!(
topic = topic,
count = last_seen,
min_peers = min_peers,
"wait_for_mesh timed out — proceeding with degraded mesh"
);
return Ok(last_seen);
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
}
pub async fn admitted_mesh_peer_count(&self, topic: &str) -> Result<usize> {
let topic_owned = topic.to_string();
self.send_command(move |response| NetworkCommand::AdmittedMeshPeers {
topic: topic_owned,
response,
})
.await
}
pub async fn wait_for_admitted_mesh(
&self,
topic: &str,
min_admitted: usize,
timeout: std::time::Duration,
) -> Result<usize> {
let deadline = tokio::time::Instant::now() + timeout;
#[allow(unused_assignments)]
let mut last_seen = 0usize;
loop {
match self.admitted_mesh_peer_count(topic).await {
Ok(count) => {
last_seen = count;
if count >= min_admitted {
tracing::info!(
topic = topic,
admitted = count,
min_admitted = min_admitted,
"Admitted mesh ready — first publish safe"
);
return Ok(count);
}
}
Err(e) => return Err(e),
}
if tokio::time::Instant::now() >= deadline {
tracing::warn!(
topic = topic,
admitted = last_seen,
min_admitted = min_admitted,
"wait_for_admitted_mesh timed out — first publish may be silently dropped \
by receivers' validator-only topic gate"
);
return Ok(last_seen);
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
}
pub async fn wait_for_connected_validators(
&self,
min_validators: usize,
timeout: std::time::Duration,
) -> Result<usize> {
let deadline = tokio::time::Instant::now() + timeout;
#[allow(unused_assignments)]
let mut last_seen = 0usize;
loop {
match self.connected_validator_count().await {
Ok(count) => {
last_seen = count;
if count >= min_validators {
tracing::info!(
connected_validators = count,
min_validators = min_validators,
"Connected validators ready — first consensus-direct broadcast safe"
);
return Ok(count);
}
}
Err(e) => return Err(e),
}
if tokio::time::Instant::now() >= deadline {
tracing::warn!(
connected_validators = last_seen,
min_validators = min_validators,
"wait_for_connected_validators timed out — first \
consensus-direct broadcast may dispatch to zero peers"
);
return Ok(last_seen);
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
}
pub async fn request_blocks(
&self,
peer: PeerId,
start: tenzro_types::primitives::BlockHeight,
count: u32,
) -> Result<OutboundRequestId> {
let request = BlockSyncRequest::GetBlockRange { start, count };
self.send_command(move |response| NetworkCommand::SendBlockSyncRequest {
peer,
request,
response,
})
.await
}
pub async fn request_tip_info(&self, peer: PeerId) -> Result<OutboundRequestId> {
self.send_command(move |response| NetworkCommand::SendBlockSyncRequest {
peer,
request: BlockSyncRequest::GetTipInfo,
response,
})
.await
}
pub async fn request_block_by_hash(
&self,
peer: PeerId,
hash: tenzro_types::primitives::Hash,
) -> Result<OutboundRequestId> {
self.send_command(move |response| NetworkCommand::SendBlockSyncRequest {
peer,
request: BlockSyncRequest::GetBlockByHash { hash },
response,
})
.await
}
pub async fn subscribe_block_sync_requests(
&self,
) -> Result<mpsc::UnboundedReceiver<InboundBlockSync>> {
self.send_command(|response| NetworkCommand::SubscribeBlockSyncRequests { response })
.await
}
pub async fn subscribe_block_sync_results(
&self,
) -> Result<mpsc::UnboundedReceiver<OutboundBlockSyncResult>> {
self.send_command(|response| NetworkCommand::SubscribeBlockSyncResults { response })
.await
}
pub async fn subscribe_peer_events(
&self,
) -> Result<mpsc::UnboundedReceiver<PeerEvent>> {
self.send_command(|response| NetworkCommand::SubscribePeerEvents { response })
.await
}
pub async fn send_block_sync_response(
&self,
request_id: InboundRequestId,
response_payload: BlockSyncResponse,
) -> Result<()> {
self.send_command(move |response| NetworkCommand::SendBlockSyncResponse {
request_id,
response_payload,
response,
})
.await
}
async fn send_command<F, T>(&self, f: F) -> Result<T>
where
F: FnOnce(oneshot::Sender<Result<T>>) -> NetworkCommand,
{
let (tx, rx) = oneshot::channel();
let command = f(tx);
self.command_tx
.send(command)
.map_err(|_| NetworkError::ChannelSend)?;
rx.await.map_err(|_| NetworkError::ChannelReceive)?
}
}
#[async_trait]
impl NetworkService for TenzroNetworkService {
async fn broadcast(&self, topic: &str, message: NetworkMessage) -> Result<()> {
self.send_command(|response| NetworkCommand::Broadcast {
topic: topic.to_string(),
message,
response,
})
.await
}
async fn send_to(&self, peer_id: PeerId, message: NetworkMessage) -> Result<()> {
let inner_bytes = message
.to_bytes()
.map_err(NetworkError::Serialization)?
.to_vec();
let peer_multihash = peer_id.to_bytes();
let mut payload = Vec::with_capacity(peer_multihash.len() + inner_bytes.len());
payload.extend_from_slice(&peer_multihash);
payload.extend_from_slice(&inner_bytes);
let direct_message = NetworkMessage::new(MessagePayload::Custom {
topic: "tenzro/direct".to_string(),
data: payload,
});
self.send_command(|response| NetworkCommand::Broadcast {
topic: "tenzro/direct".to_string(),
message: direct_message,
response,
})
.await
}
async fn subscribe(&self, topic: &str) -> Result<mpsc::UnboundedReceiver<NetworkMessage>> {
self.send_command(|response| NetworkCommand::Subscribe {
topic: topic.to_string(),
response,
})
.await
}
async fn connected_peers(&self) -> Result<Vec<PeerId>> {
self.send_command(|response| NetworkCommand::ConnectedPeers { response })
.await
}
async fn peer_info(&self, peer_id: &PeerId) -> Result<Option<ManagedPeer>> {
self.send_command(|response| NetworkCommand::PeerInfo {
peer_id: *peer_id,
response,
})
.await
}
async fn ban_peer(&self, peer_id: &PeerId) -> Result<()> {
self.send_command(|response| NetworkCommand::BanPeer {
peer_id: *peer_id,
response,
})
.await
}
async fn unban_peer(&self, peer_id: &PeerId) -> Result<()> {
self.send_command(|response| NetworkCommand::UnbanPeer {
peer_id: *peer_id,
response,
})
.await
}
async fn local_peer_id(&self) -> Result<PeerId> {
self.send_command(|response| NetworkCommand::LocalPeerId { response })
.await
}
async fn dial(&self, addr: Multiaddr) -> Result<()> {
self.send_command(|response| NetworkCommand::Dial { addr, response })
.await
}
async fn set_validator_registry(&self, registry: std::sync::Arc<dyn crate::peer_manager::ValidatorRegistry>) -> Result<()> {
self.send_command(|response| NetworkCommand::SetValidatorRegistry { registry, response })
.await
}
async fn listen_addresses(&self) -> Result<Vec<Multiaddr>> {
self.send_command(|response| NetworkCommand::ListenAddresses { response })
.await
}
async fn broadcast_to_validators(&self, message: ConsensusMessage) -> Result<usize> {
self.send_command(move |response| NetworkCommand::BroadcastToValidators {
message,
response,
})
.await
}
async fn subscribe_consensus_direct(
&self,
) -> Result<mpsc::UnboundedReceiver<ConsensusMessage>> {
self.send_command(|response| NetworkCommand::SubscribeConsensusDirect { response })
.await
}
async fn connected_validator_count(&self) -> Result<usize> {
self.send_command(|response| NetworkCommand::ConnectedValidatorCount { response })
.await
}
async fn set_mpc_did_resolver(
&self,
resolver: std::sync::Arc<dyn crate::mpc_relay::MpcDidResolver>,
) -> Result<()> {
self.send_command(|response| NetworkCommand::SetMpcDidResolver { resolver, response })
.await
}
async fn send_mpc_relay_message(
&self,
message: crate::mpc_relay::MpcRelayRequest,
) -> Result<()> {
self.send_command(move |response| NetworkCommand::SendMpcRelayMessage {
message,
response,
})
.await
}
async fn subscribe_mpc_relay(
&self,
) -> Result<mpsc::UnboundedReceiver<crate::mpc_relay::MpcRelayRequest>> {
self.send_command(|response| NetworkCommand::SubscribeMpcRelay { response })
.await
}
}
struct EventLoopState {
swarm: Swarm<TenzroBehaviour>,
peer_manager: PeerManager,
subscribers: HashMap<TopicHash, Vec<mpsc::UnboundedSender<NetworkMessage>>>,
deduplicator: MessageDeduplicator,
metrics: Arc<NetworkMetrics>,
listen_addresses: Vec<Multiaddr>,
pending_inbound_block_sync: HashMap<InboundRequestId, ResponseChannel<BlockSyncResponse>>,
block_sync_request_subscriber: Option<mpsc::UnboundedSender<InboundBlockSync>>,
block_sync_request_subscriber_ever_attached: bool,
block_sync_result_subscriber: Option<mpsc::UnboundedSender<OutboundBlockSyncResult>>,
peer_event_subscriber: Option<mpsc::UnboundedSender<PeerEvent>>,
consensus_direct_subscriber: Option<mpsc::UnboundedSender<ConsensusMessage>>,
consensus_direct_subscriber_ever_attached: bool,
consensus_direct_inbound_inflight: HashMap<PeerId, usize>,
mpc_relay_subscriber: Option<mpsc::UnboundedSender<crate::mpc_relay::MpcRelayRequest>>,
mpc_relay_subscriber_ever_attached: bool,
mpc_relay_inbound_inflight: HashMap<PeerId, usize>,
mpc_did_resolver: Option<std::sync::Arc<dyn crate::mpc_relay::MpcDidResolver>>,
observed_addrs: HashMap<Multiaddr, HashSet<PeerId>>,
advertised_external_addrs: HashSet<Multiaddr>,
bootstrap_peers: Vec<(PeerId, Multiaddr)>,
}
const OBSERVED_ADDR_CONFIRMATION_THRESHOLD: usize = 1;
async fn run_event_loop(
config: NetworkConfig,
mut command_rx: mpsc::UnboundedReceiver<NetworkCommand>,
metrics: Arc<NetworkMetrics>,
) -> Result<()> {
let local_key = load_or_generate_keypair(&config.data_dir)?;
let local_peer_id = PeerId::from(local_key.public());
tracing::info!("Local peer ID: {}", local_peer_id);
let enable_relay = config.enable_relay;
let enable_hole_punching = config.enable_hole_punching;
let protocol_version = config.protocol_version.clone();
let user_agent = config.user_agent.clone();
let idle_timeout = config.connection_idle_timeout;
let mut swarm = libp2p::SwarmBuilder::with_existing_identity(local_key.clone())
.with_tokio()
.with_tcp(
libp2p::tcp::Config::default().nodelay(true),
libp2p::tls::Config::new,
libp2p::yamux::Config::default,
)
.map_err(|e| NetworkError::Transport(format!("TCP/TLS upgrade failed: {}", e)))?
.with_quic()
.with_dns()
.map_err(|e| NetworkError::Transport(format!("DNS transport failed: {}", e)))?
.with_relay_client(libp2p::tls::Config::new, libp2p::yamux::Config::default)
.map_err(|e| NetworkError::Transport(format!("Relay client transport failed: {}", e)))?
.with_behaviour(|key, relay_client| {
TenzroBehaviour::new(
local_peer_id,
key,
protocol_version,
user_agent,
enable_relay,
enable_hole_punching,
Some(relay_client),
)
.map_err(|e| -> Box<dyn std::error::Error + Send + Sync> {
e.to_string().into()
})
})
.map_err(|e| NetworkError::Transport(format!("Behaviour construction failed: {}", e)))?
.with_swarm_config(|cfg| cfg.with_idle_connection_timeout(idle_timeout))
.build();
for addr in &config.listen_addresses {
swarm
.listen_on(addr.clone())
.map_err(|e| NetworkError::Transport(format!("Failed to listen on {}: {}", addr, e)))?;
tracing::info!("Listening on {}", addr);
}
for addr in &config.external_addresses {
if !is_globally_routable(addr) {
tracing::warn!(
"Refusing to advertise non-routable external address {} — \
check NetworkConfig::external_addresses",
addr
);
continue;
}
swarm.add_external_address(addr.clone());
tracing::info!("Advertising external address {}", addr);
}
let preconfigured_external: HashSet<Multiaddr> = config
.external_addresses
.iter()
.filter(|a| is_globally_routable(a))
.cloned()
.collect();
let peer_manager = PeerManager::new(
(config.max_inbound_peers + config.max_outbound_peers) as usize,
);
for addr in &config.boot_nodes {
for proto in addr.iter() {
if let libp2p::multiaddr::Protocol::P2p(peer_id) = proto {
peer_manager.add_protected_peer(peer_id);
tracing::info!(
peer = %peer_id,
"Registered boot node as protected peer"
);
}
}
}
for topic_str in &config.gossip_topics {
let topic = IdentTopic::new(topic_str.as_str());
if let Err(e) = swarm.behaviour_mut().subscribe(&topic) {
tracing::warn!("Failed to subscribe to topic {}: {}", topic_str, e);
} else {
tracing::info!("Subscribed to topic: {}", topic_str);
}
}
let bootstrap_peers: Vec<(PeerId, Multiaddr)> = config
.boot_nodes
.iter()
.filter_map(|addr| {
addr.iter()
.find_map(|proto| {
if let libp2p::multiaddr::Protocol::P2p(peer_id) = proto {
Some(peer_id)
} else {
None
}
})
.map(|peer_id| (peer_id, addr.clone()))
})
.collect();
if config.enable_dht && !bootstrap_peers.is_empty() {
crate::discovery::bootstrap_dht(
&mut swarm.behaviour_mut().kademlia,
bootstrap_peers.clone(),
);
}
for addr in &config.boot_nodes {
if let Err(e) = swarm.dial(addr.clone()) {
tracing::warn!("Failed to dial boot node {}: {}", addr, e);
} else {
tracing::info!("Dialing boot node: {}", addr);
}
}
let mut state = EventLoopState {
swarm,
peer_manager,
subscribers: HashMap::new(),
deduplicator: MessageDeduplicator::default(),
metrics,
listen_addresses: Vec::new(),
pending_inbound_block_sync: HashMap::new(),
block_sync_request_subscriber: None,
block_sync_request_subscriber_ever_attached: false,
block_sync_result_subscriber: None,
peer_event_subscriber: None,
consensus_direct_subscriber: None,
consensus_direct_subscriber_ever_attached: false,
consensus_direct_inbound_inflight: HashMap::new(),
mpc_relay_subscriber: None,
mpc_relay_subscriber_ever_attached: false,
mpc_relay_inbound_inflight: HashMap::new(),
mpc_did_resolver: None,
observed_addrs: HashMap::new(),
advertised_external_addrs: preconfigured_external,
bootstrap_peers,
};
let mut cleanup_interval = interval(Duration::from_secs(60));
cleanup_interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
let dht_enabled = config.enable_dht;
let mut kademlia_bootstrap_interval = interval(Duration::from_secs(60));
kademlia_bootstrap_interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
'event_loop: loop {
tokio::select! {
Some(command) = command_rx.recv() => {
if matches!(command, NetworkCommand::Shutdown { .. }) {
handle_command(&mut state, command).await;
tracing::info!("Network event loop shutting down gracefully");
break 'event_loop;
}
handle_command(&mut state, command).await;
}
event = state.swarm.select_next_some() => {
handle_swarm_event(&mut state, event).await;
}
_ = cleanup_interval.tick() => {
state.peer_manager.cleanup_expired_bans();
state.peer_manager.cleanup_stale_peers(Duration::from_secs(86400));
let stats = state.peer_manager.stats();
state.metrics.peers_connected.set(stats.connected as i64);
state.metrics.peers_banned.set(stats.banned as i64);
tracing::debug!(
"Peer stats: {} total, {} connected, {} banned",
stats.total,
stats.connected,
stats.banned
);
for (peer_id, addr) in &state.bootstrap_peers {
if !state.swarm.is_connected(peer_id) {
match state.swarm.dial(addr.clone()) {
Ok(()) => tracing::info!(
%peer_id,
%addr,
"Re-dialing configured bootstrap peer (not connected)"
),
Err(e) => tracing::debug!(
%peer_id,
%addr,
error = %e,
"Bootstrap re-dial skipped"
),
}
}
}
}
_ = kademlia_bootstrap_interval.tick(), if dht_enabled => {
match state.swarm.behaviour_mut().kademlia.bootstrap() {
Ok(query_id) => {
tracing::debug!(
?query_id,
"Periodic Kademlia bootstrap initiated"
);
}
Err(e) => {
tracing::warn!(
error = %e,
"Periodic Kademlia bootstrap failed to start \
(no known peers in routing table?)"
);
}
}
}
}
}
Ok(())
}
fn handle_block_sync_event(
state: &mut EventLoopState,
event: request_response::Event<BlockSyncRequest, BlockSyncResponse>,
) {
use request_response::{Event as RrEvent, Message};
match event {
RrEvent::Message {
peer,
message: Message::Request {
request_id,
request,
channel,
},
..
} => {
let Some(tx) = state.block_sync_request_subscriber.as_ref() else {
if state.block_sync_request_subscriber_ever_attached {
tracing::warn!(
%peer,
"Inbound block-sync request received but subscriber was dropped — \
replying with Storage error"
);
} else {
tracing::debug!(
%peer,
"Inbound block-sync request received during bootstrap window \
(subscriber not yet attached) — replying with Storage error"
);
}
let err_resp = BlockSyncResponse::Error(
crate::block_sync_proto::BlockSyncError::Storage(
"no block-sync subscriber attached".to_string(),
),
);
let _ = state
.swarm
.behaviour_mut()
.block_sync
.send_response(channel, err_resp);
return;
};
state
.pending_inbound_block_sync
.insert(request_id, channel);
let inbound = InboundBlockSync {
peer,
request_id,
request,
};
if tx.send(inbound).is_err() {
tracing::warn!(
%peer,
"Block-sync request subscriber dropped — discarding inbound request"
);
state.block_sync_request_subscriber = None;
if let Some(channel) = state.pending_inbound_block_sync.remove(&request_id) {
let err_resp = BlockSyncResponse::Error(
crate::block_sync_proto::BlockSyncError::Storage(
"subscriber dropped".to_string(),
),
);
let _ = state
.swarm
.behaviour_mut()
.block_sync
.send_response(channel, err_resp);
}
}
}
RrEvent::Message {
peer,
message: Message::Response {
request_id,
response,
},
..
} => {
if let Some(tx) = state.block_sync_result_subscriber.as_ref() {
let item = OutboundBlockSyncResult {
peer,
request_id,
result: Ok(response),
};
if tx.send(item).is_err() {
tracing::warn!("Block-sync result subscriber dropped");
state.block_sync_result_subscriber = None;
}
} else {
tracing::warn!(
%peer,
%request_id,
"Block-sync response received but no result subscriber attached"
);
}
}
RrEvent::OutboundFailure {
peer,
request_id,
error,
..
} => {
tracing::warn!(%peer, %request_id, %error, "Outbound block-sync failure");
if let Some(tx) = state.block_sync_result_subscriber.as_ref() {
let item = OutboundBlockSyncResult {
peer,
request_id,
result: Err(error.into()),
};
if tx.send(item).is_err() {
state.block_sync_result_subscriber = None;
}
}
}
RrEvent::InboundFailure {
peer,
request_id,
error,
..
} => {
tracing::debug!(%peer, %request_id, %error, "Inbound block-sync failure");
state.pending_inbound_block_sync.remove(&request_id);
}
RrEvent::ResponseSent { peer, request_id, .. } => {
tracing::trace!(%peer, %request_id, "Block-sync response flushed to wire");
}
}
}
fn handle_consensus_direct_event(
state: &mut EventLoopState,
event: request_response::Event<ConsensusDirectRequest, ConsensusDirectResponse>,
) {
use request_response::{Event as RrEvent, Message};
match event {
RrEvent::Message {
peer,
message: Message::Request {
request_id: _,
request,
channel,
},
..
} => {
let inflight = state
.consensus_direct_inbound_inflight
.entry(peer)
.or_insert(0);
if *inflight >= MAX_INBOUND_STREAMS_PER_PEER {
tracing::warn!(
%peer,
limit = MAX_INBOUND_STREAMS_PER_PEER,
"consensus-direct: rejecting overflow inbound stream"
);
let _ = state
.swarm
.behaviour_mut()
.consensus_direct
.send_response(
channel,
ConsensusDirectResponse::Error(ConsensusDirectError::ServerBusy {
limit: MAX_INBOUND_STREAMS_PER_PEER,
}),
);
return;
}
let Some(tx) = state.consensus_direct_subscriber.as_ref() else {
if state.consensus_direct_subscriber_ever_attached {
tracing::warn!(
%peer,
"consensus-direct: subscriber was dropped — replying NoSubscriber"
);
} else {
tracing::debug!(
%peer,
"consensus-direct: inbound during bootstrap window \
(subscriber not yet attached) — replying NoSubscriber"
);
}
let _ = state
.swarm
.behaviour_mut()
.consensus_direct
.send_response(
channel,
ConsensusDirectResponse::Error(ConsensusDirectError::NoSubscriber),
);
return;
};
let ConsensusDirectRequest::Message(consensus_msg) = request;
let send_result = tx.send(consensus_msg);
if send_result.is_err() {
tracing::warn!(
%peer,
"consensus-direct: subscriber dropped — replying NoSubscriber"
);
state.consensus_direct_subscriber = None;
let _ = state
.swarm
.behaviour_mut()
.consensus_direct
.send_response(
channel,
ConsensusDirectResponse::Error(ConsensusDirectError::NoSubscriber),
);
return;
}
*inflight += 1;
let _ = state
.swarm
.behaviour_mut()
.consensus_direct
.send_response(channel, ConsensusDirectResponse::Ack);
}
RrEvent::Message {
peer,
message: Message::Response {
request_id,
response,
},
..
} => {
match response {
ConsensusDirectResponse::Ack => {
tracing::trace!(%peer, %request_id, "consensus-direct ack");
}
ConsensusDirectResponse::Error(err) => {
tracing::warn!(
%peer,
%request_id,
%err,
"consensus-direct: peer returned error response"
);
}
}
}
RrEvent::OutboundFailure {
peer,
request_id,
error,
..
} => {
tracing::warn!(
%peer,
%request_id,
%error,
"consensus-direct outbound failure"
);
}
RrEvent::InboundFailure {
peer,
request_id,
error,
..
} => {
tracing::debug!(
%peer,
%request_id,
%error,
"consensus-direct inbound failure"
);
if let Some(count) = state.consensus_direct_inbound_inflight.get_mut(&peer)
&& *count > 0
{
*count -= 1;
}
}
RrEvent::ResponseSent { peer, request_id, .. } => {
tracing::trace!(
%peer,
%request_id,
"consensus-direct response flushed to wire"
);
if let Some(count) = state.consensus_direct_inbound_inflight.get_mut(&peer)
&& *count > 0
{
*count -= 1;
}
}
}
}
fn handle_mpc_relay_event(
state: &mut EventLoopState,
event: request_response::Event<
crate::mpc_relay::MpcRelayRequest,
crate::mpc_relay::MpcRelayResponse,
>,
) {
use crate::mpc_relay::{MpcRelayError, MpcRelayResponse, MAX_INBOUND_STREAMS_PER_PEER};
use request_response::{Event as RrEvent, Message};
match event {
RrEvent::Message {
peer,
message: Message::Request {
request_id: _,
request,
channel,
},
..
} => {
let inflight = state
.mpc_relay_inbound_inflight
.entry(peer)
.or_insert(0);
if *inflight >= MAX_INBOUND_STREAMS_PER_PEER {
tracing::warn!(
%peer,
limit = MAX_INBOUND_STREAMS_PER_PEER,
"mpc-relay: rejecting overflow inbound stream"
);
let _ = state
.swarm
.behaviour_mut()
.mpc_relay
.send_response(
channel,
MpcRelayResponse::Error(MpcRelayError::ServerBusy {
limit: MAX_INBOUND_STREAMS_PER_PEER,
}),
);
return;
}
let claimed_did = request.from_did.clone();
let audit_ok = match state.mpc_did_resolver.as_ref() {
None => {
tracing::warn!(
%peer,
claimed_did = %claimed_did,
"mpc-relay: rejecting inbound — no DID resolver installed"
);
false
}
Some(resolver) => match resolver.did_for_peer_id(&peer) {
Some(actual_did) if actual_did == claimed_did => true,
Some(actual_did) => {
tracing::warn!(
%peer,
claimed_did = %claimed_did,
actual_did = %actual_did,
"mpc-relay: from_did does not match peer's bound DID — rejecting"
);
false
}
None => {
tracing::warn!(
%peer,
claimed_did = %claimed_did,
"mpc-relay: peer has no bound DID — rejecting"
);
false
}
},
};
if !audit_ok {
let _ = state
.swarm
.behaviour_mut()
.mpc_relay
.send_response(
channel,
MpcRelayResponse::Error(MpcRelayError::UnknownSender),
);
return;
}
let Some(tx) = state.mpc_relay_subscriber.as_ref() else {
if state.mpc_relay_subscriber_ever_attached {
tracing::warn!(
%peer,
"mpc-relay: subscriber was dropped — replying NoSubscriber"
);
} else {
tracing::debug!(
%peer,
"mpc-relay: inbound during bootstrap window \
(subscriber not yet attached) — replying NoSubscriber"
);
}
let _ = state
.swarm
.behaviour_mut()
.mpc_relay
.send_response(
channel,
MpcRelayResponse::Error(MpcRelayError::NoSubscriber),
);
return;
};
let send_result = tx.send(request);
if send_result.is_err() {
tracing::warn!(
%peer,
"mpc-relay: subscriber dropped — replying NoSubscriber"
);
state.mpc_relay_subscriber = None;
let _ = state
.swarm
.behaviour_mut()
.mpc_relay
.send_response(
channel,
MpcRelayResponse::Error(MpcRelayError::NoSubscriber),
);
return;
}
*inflight += 1;
let _ = state
.swarm
.behaviour_mut()
.mpc_relay
.send_response(channel, MpcRelayResponse::Ack);
}
RrEvent::Message {
peer,
message: Message::Response {
request_id,
response,
},
..
} => match response {
MpcRelayResponse::Ack => {
tracing::trace!(%peer, %request_id, "mpc-relay ack");
}
MpcRelayResponse::Error(err) => {
tracing::warn!(
%peer,
%request_id,
%err,
"mpc-relay: peer returned error response"
);
}
},
RrEvent::OutboundFailure {
peer,
request_id,
error,
..
} => {
tracing::warn!(
%peer,
%request_id,
%error,
"mpc-relay outbound failure"
);
}
RrEvent::InboundFailure {
peer,
request_id,
error,
..
} => {
tracing::debug!(
%peer,
%request_id,
%error,
"mpc-relay inbound failure"
);
if let Some(count) = state.mpc_relay_inbound_inflight.get_mut(&peer)
&& *count > 0
{
*count -= 1;
}
}
RrEvent::ResponseSent { peer, request_id, .. } => {
tracing::trace!(
%peer,
%request_id,
"mpc-relay response flushed to wire"
);
if let Some(count) = state.mpc_relay_inbound_inflight.get_mut(&peer)
&& *count > 0
{
*count -= 1;
}
}
}
}
async fn handle_swarm_event(
state: &mut EventLoopState,
event: SwarmEvent<TenzroBehaviourEvent>,
) {
match event {
SwarmEvent::Behaviour(behaviour_event) => match behaviour_event {
TenzroBehaviourEvent::Gossipsub(gossipsub::Event::Message {
propagation_source,
message_id,
message,
}) => {
tracing::debug!(
"Received message {} from peer {}",
message_id,
propagation_source
);
if state.peer_manager.is_banned(&propagation_source) {
tracing::warn!("Ignoring message from banned peer {}", propagation_source);
state.metrics.gossip_rejected_invalid.inc();
return;
}
if !state.peer_manager.check_rate_limit(&propagation_source) {
tracing::warn!("Rate limited message from peer {}", propagation_source);
state.metrics.gossip_rejected_invalid.inc();
return;
}
if state.deduplicator.is_duplicate(&message.data) {
tracing::trace!("Dropping duplicate message from peer {}", propagation_source);
state.metrics.gossip_rejected_duplicate.inc();
return;
}
let topic_str = message.topic.to_string();
if !state.peer_manager.authorize_peer_for_topic(&propagation_source, &topic_str) {
state.metrics.gossip_rejected_validator_only.inc();
return;
}
match validate_gossip_message(&message.topic, &message.data) {
MessageValidation::Accept => {}
MessageValidation::Reject => {
tracing::warn!("Message validation rejected from peer {}", propagation_source);
state.peer_manager.decrease_reputation(&propagation_source, 5);
state.metrics.gossip_rejected_invalid.inc();
return;
}
MessageValidation::Ignore => {
state.metrics.gossip_rejected_invalid.inc();
return;
}
}
match NetworkMessage::from_bytes(&message.data) {
Ok(net_msg) => {
state
.peer_manager
.increase_reputation(&propagation_source, 1);
state.metrics.gossip_accepted.inc();
if let Some(subs) = state.subscribers.get_mut(&message.topic) {
subs.retain(|tx| tx.send(net_msg.clone()).is_ok());
}
}
Err(e) => {
tracing::warn!("Failed to parse network message: {}", e);
state
.peer_manager
.decrease_reputation(&propagation_source, 5);
state.metrics.gossip_rejected_invalid.inc();
}
}
}
TenzroBehaviourEvent::Gossipsub(gossipsub::Event::Subscribed { peer_id, topic }) => {
tracing::debug!("Peer {} subscribed to topic {:?}", peer_id, topic);
}
TenzroBehaviourEvent::Gossipsub(gossipsub::Event::Unsubscribed { peer_id, topic }) => {
tracing::debug!("Peer {} unsubscribed from topic {:?}", peer_id, topic);
}
TenzroBehaviourEvent::Identify(identify::Event::Received { peer_id, info, connection_id: _ }) => {
tracing::info!(
"Identified peer {}: protocol={}, agent={}",
peer_id,
info.protocol_version,
info.agent_version
);
state
.peer_manager
.try_register_validator_on_identify(&peer_id, &info.protocol_version);
state
.peer_manager
.update_protocol_version(&peer_id, info.protocol_version);
let already_connected = state.swarm.is_connected(&peer_id);
for addr in &info.listen_addrs {
if is_globally_routable(addr) {
state
.swarm
.behaviour_mut()
.kademlia
.add_address(&peer_id, addr.clone());
state.swarm.add_peer_address(peer_id, addr.clone());
if !already_connected {
match state.swarm.dial(addr.clone()) {
Ok(()) => tracing::info!(
%peer_id,
%addr,
"Dialing Kademlia-discovered peer (Identify)"
),
Err(e) => tracing::debug!(
%peer_id,
%addr,
error = %e,
"Dial-on-discovery skipped"
),
}
}
} else {
tracing::debug!(
"Skipping non-routable DHT address from {}: {}",
peer_id,
addr
);
}
}
let obs = info.observed_addr.clone();
if !is_globally_routable(&obs) {
tracing::trace!(
%peer_id,
observed = %obs,
"Ignoring non-routable observed address from peer"
);
} else if !is_observed_port_one_of_ours(&obs, &state.listen_addresses) {
tracing::debug!(
%peer_id,
observed = %obs,
"Ignoring observed address — port not in our listen-port set \
(NAT-translated source port, not a reachable listen address)"
);
} else if !state.advertised_external_addrs.contains(&obs) {
let reporters = state.observed_addrs.entry(obs.clone()).or_default();
if reporters.insert(peer_id) {
tracing::debug!(
%peer_id,
observed = %obs,
count = reporters.len(),
threshold = OBSERVED_ADDR_CONFIRMATION_THRESHOLD,
"Observed-address report received via Identify"
);
}
if reporters.len() >= OBSERVED_ADDR_CONFIRMATION_THRESHOLD {
tracing::info!(
address = %obs,
reporters = reporters.len(),
"Promoting observed address to advertised external address \
(N distinct peers agree — permissionless NAT discovery)"
);
state.swarm.add_external_address(obs.clone());
state.advertised_external_addrs.insert(obs.clone());
state.observed_addrs.remove(&obs);
}
}
}
TenzroBehaviourEvent::Kademlia(kad::Event::OutboundQueryProgressed { result, .. }) => {
match result {
QueryResult::GetProviders(Ok(kad::GetProvidersOk::FoundProviders { providers, .. })) => {
tracing::debug!("Found {} provider(s)", providers.len());
for peer in providers {
tracing::debug!("Found provider: {}", peer);
}
}
QueryResult::GetProviders(Ok(kad::GetProvidersOk::FinishedWithNoAdditionalRecord { closest_peers })) => {
tracing::debug!("GetProviders query finished with {} closest peers", closest_peers.len());
}
QueryResult::Bootstrap(Ok(_)) => {
tracing::info!("DHT bootstrap completed");
}
_ => {}
}
}
TenzroBehaviourEvent::Ping(ping::Event { peer, result, .. }) => {
match result {
Ok(duration) => {
tracing::trace!("Ping to {} successful: {:?}", peer, duration);
state.peer_manager.increase_reputation(&peer, 1);
}
Err(e) => {
tracing::trace!("Ping to {} failed: {}", peer, e);
}
}
}
TenzroBehaviourEvent::BlockSync(rr_event) => {
handle_block_sync_event(state, rr_event);
}
TenzroBehaviourEvent::ConsensusDirect(rr_event) => {
handle_consensus_direct_event(state, rr_event);
}
TenzroBehaviourEvent::MpcRelay(rr_event) => {
handle_mpc_relay_event(state, rr_event);
}
TenzroBehaviourEvent::AutonatClient(autonat::v2::client::Event {
tested_addr,
bytes_sent,
server,
result,
}) => {
match result {
Ok(()) => tracing::info!(
%server,
address = %tested_addr,
bytes_sent,
"AutoNAT probe succeeded — address reachable"
),
Err(e) => tracing::debug!(
%server,
address = %tested_addr,
bytes_sent,
error = %e,
"AutoNAT probe failed — address not reachable from server"
),
}
}
TenzroBehaviourEvent::AutonatServer(_) => {
}
TenzroBehaviourEvent::RelayClient(relay_event) => match relay_event {
relay::client::Event::ReservationReqAccepted {
relay_peer_id,
renewal,
limit,
} => {
tracing::info!(
%relay_peer_id,
renewal,
?limit,
"Circuit-Relay v2 reservation accepted — this node is now \
reachable via /p2p/<relay>/p2p-circuit/p2p/<self>"
);
}
relay::client::Event::OutboundCircuitEstablished {
relay_peer_id,
limit,
} => {
tracing::info!(
%relay_peer_id,
?limit,
"Outbound circuit established via relay"
);
}
relay::client::Event::InboundCircuitEstablished {
src_peer_id,
limit,
} => {
tracing::info!(
%src_peer_id,
?limit,
"Inbound circuit established via relay"
);
}
},
TenzroBehaviourEvent::Relay(_) => {
}
TenzroBehaviourEvent::Dcutr(dcutr::Event {
remote_peer_id,
result,
}) => match result {
Ok(conn_id) => tracing::info!(
%remote_peer_id,
?conn_id,
"DCUtR hole-punch succeeded — direct connection upgraded from relayed"
),
Err(e) => tracing::debug!(
%remote_peer_id,
error = %e,
"DCUtR hole-punch failed — will continue using relayed connection"
),
},
_ => {}
},
SwarmEvent::ConnectionEstablished {
peer_id,
endpoint,
num_established,
..
} => {
if state.peer_manager.is_banned(&peer_id) {
tracing::info!("Auto-unbanning reconnecting peer {}", peer_id);
state.peer_manager.unban_peer(&peer_id);
}
tracing::info!(
"Connection established with {} (endpoint: {:?}, num_established: {})",
peer_id,
endpoint,
num_established
);
state.metrics.connections_established.inc();
match &endpoint {
libp2p::core::ConnectedPoint::Listener { .. } => {
state.metrics.connections_inbound_total.inc();
}
libp2p::core::ConnectedPoint::Dialer { .. } => {
state.metrics.connections_outbound_total.inc();
}
}
state.peer_manager.add_peer(peer_id);
state
.peer_manager
.update_status(&peer_id, PeerStatus::Connected);
let remote_addr = endpoint.get_remote_address().clone();
if state.peer_manager.update_endpoint(&peer_id, remote_addr.clone()) {
state.metrics.peer_address_migrations_total.inc();
tracing::info!(
%peer_id,
new_addr = %remote_addr,
"Peer address migration observed — counter incremented"
);
}
if num_established.get() == 1 {
if let Some(tx) = state.peer_event_subscriber.as_ref() {
if tx.send(PeerEvent::Connected(peer_id)).is_err() {
tracing::warn!(
"Peer-event subscriber dropped while sending Connected({}); detaching",
peer_id
);
state.peer_event_subscriber = None;
}
}
}
}
SwarmEvent::ConnectionClosed {
peer_id,
cause,
num_established,
..
} => {
tracing::info!(
"Connection closed with {} (cause: {:?}, remaining: {})",
peer_id,
cause,
num_established
);
state.metrics.connections_established.dec();
if num_established == 0 {
state
.peer_manager
.update_status(&peer_id, PeerStatus::Disconnected);
if let Some(tx) = state.peer_event_subscriber.as_ref() {
if tx.send(PeerEvent::Disconnected(peer_id)).is_err() {
tracing::warn!(
"Peer-event subscriber dropped while sending Disconnected({}); detaching",
peer_id
);
state.peer_event_subscriber = None;
}
}
}
}
SwarmEvent::NewListenAddr { address, .. } => {
tracing::info!("Listening on {}", address);
if !state.listen_addresses.contains(&address) {
state.listen_addresses.push(address);
}
}
SwarmEvent::ExpiredListenAddr { address, .. } => {
tracing::info!("Listen address expired: {}", address);
state.listen_addresses.retain(|a| a != &address);
}
SwarmEvent::IncomingConnection { send_back_addr, local_addr, .. } => {
tracing::debug!("Incoming connection from {} to {}", send_back_addr, local_addr);
if let Some(ip) = extract_ip(&send_back_addr)
&& !state.peer_manager.check_dial_rate_limit(ip)
{
tracing::warn!("Dial rate-limit exceeded for IP {}", ip);
state.metrics.dials_rejected_per_ip.inc();
}
}
SwarmEvent::IncomingConnectionError { send_back_addr, error, .. } => {
tracing::warn!("Incoming connection error from {}: {}", send_back_addr, error);
}
SwarmEvent::OutgoingConnectionError { peer_id, error, .. } => {
if let Some(peer_id) = peer_id {
tracing::warn!("Outgoing connection error with {}: {}", peer_id, error);
state.peer_manager.record_failed_connection(&peer_id);
} else {
tracing::warn!("Outgoing connection error (no peer ID): {}", error);
}
}
SwarmEvent::NewExternalAddrCandidate { address } => {
tracing::debug!(%address, "New external address candidate");
}
SwarmEvent::ExternalAddrConfirmed { address } => {
tracing::info!(
%address,
"External address confirmed (AutoNAT probe-back succeeded)"
);
state.advertised_external_addrs.insert(address.clone());
state.observed_addrs.remove(&address);
}
SwarmEvent::ExternalAddrExpired { address } => {
tracing::info!(
%address,
"External address expired — stopping advertisement"
);
state.advertised_external_addrs.remove(&address);
}
_ => {}
}
}
async fn handle_command(state: &mut EventLoopState, command: NetworkCommand) {
match command {
NetworkCommand::Broadcast {
topic,
message,
response,
} => {
let result = (|| {
let topic_obj = IdentTopic::new(topic);
let bytes = message.to_bytes().map_err(|e| {
NetworkError::Serialization(e)
})?;
state
.swarm
.behaviour_mut()
.publish(&topic_obj, bytes.to_vec())
.map_err(|e| NetworkError::PublishError(e.to_string()))?;
Ok(())
})();
if result.is_ok() {
state.metrics.gossip_published.inc();
}
let _ = response.send(result);
}
NetworkCommand::Subscribe { topic, response } => {
let result = (|| {
let topic_obj = IdentTopic::new(&topic);
let topic_hash = topic_obj.hash();
state
.swarm
.behaviour_mut()
.subscribe(&topic_obj)
.map_err(|e| NetworkError::SubscriptionError(e.to_string()))?;
let (tx, rx) = mpsc::unbounded_channel();
state
.subscribers
.entry(topic_hash)
.or_default()
.push(tx);
Ok(rx)
})();
let _ = response.send(result);
}
NetworkCommand::ConnectedPeers { response } => {
let peers: Vec<PeerId> = state.swarm.connected_peers().cloned().collect();
let _ = response.send(Ok(peers));
}
NetworkCommand::PeerInfo { peer_id, response } => {
let info = state.peer_manager.get_peer(&peer_id);
let _ = response.send(Ok(info));
}
NetworkCommand::BanPeer { peer_id, response } => {
state.peer_manager.ban_peer(&peer_id);
state.swarm.behaviour_mut().block_peer(peer_id);
tracing::info!("Peer {} banned and blocked at libp2p layer", peer_id);
let _ = response.send(Ok(()));
}
NetworkCommand::UnbanPeer { peer_id, response } => {
state.peer_manager.unban_peer(&peer_id);
state.swarm.behaviour_mut().unblock_peer(peer_id);
tracing::info!("Peer {} unbanned and unblocked at libp2p layer", peer_id);
let _ = response.send(Ok(()));
}
NetworkCommand::LocalPeerId { response } => {
let peer_id = *state.swarm.local_peer_id();
let _ = response.send(Ok(peer_id));
}
NetworkCommand::Dial { addr, response } => {
let result = state
.swarm
.dial(addr)
.map_err(|e| NetworkError::Connection(e.to_string()));
let _ = response.send(result);
}
NetworkCommand::SetValidatorRegistry { registry, response } => {
state.peer_manager.set_validator_registry(registry);
tracing::info!("Validator registry installed in peer manager");
let _ = response.send(Ok(()));
}
NetworkCommand::MeshPeerCount { topic, response } => {
let topic_hash = IdentTopic::new(topic).hash();
let count = state
.swarm
.behaviour()
.mesh_peers(&topic_hash)
.len();
let _ = response.send(Ok(count));
}
NetworkCommand::ListenAddresses { response } => {
let _ = response.send(Ok(state.listen_addresses.clone()));
}
NetworkCommand::AdmittedMeshPeers { topic, response } => {
let topic_hash = IdentTopic::new(topic).hash();
let mesh: Vec<PeerId> = state
.swarm
.behaviour()
.mesh_peers(&topic_hash)
.into_iter()
.copied()
.collect();
let admitted = match state.peer_manager.validator_registry() {
None => {
mesh.len()
}
Some(registry) => {
let validator_set = registry.validator_peer_ids();
mesh.iter().filter(|p| validator_set.contains(*p)).count()
}
};
let _ = response.send(Ok(admitted));
}
NetworkCommand::SendBlockSyncRequest { peer, request, response } => {
let request_id = state
.swarm
.behaviour_mut()
.block_sync
.send_request(&peer, request);
let _ = response.send(Ok(request_id));
}
NetworkCommand::SendBlockSyncResponse {
request_id,
response_payload,
response,
} => {
let result = match state.pending_inbound_block_sync.remove(&request_id) {
Some(channel) => state
.swarm
.behaviour_mut()
.block_sync
.send_response(channel, response_payload)
.map_err(|_| {
NetworkError::ChannelSend
}),
None => Err(NetworkError::PeerNotFound(format!(
"no parked inbound block-sync request for id {}",
request_id
))),
};
let _ = response.send(result);
}
NetworkCommand::SubscribeBlockSyncRequests { response } => {
let (tx, rx) = mpsc::unbounded_channel();
state.block_sync_request_subscriber = Some(tx);
state.block_sync_request_subscriber_ever_attached = true;
let _ = response.send(Ok(rx));
}
NetworkCommand::SubscribeBlockSyncResults { response } => {
let (tx, rx) = mpsc::unbounded_channel();
state.block_sync_result_subscriber = Some(tx);
let _ = response.send(Ok(rx));
}
NetworkCommand::SubscribePeerEvents { response } => {
let (tx, rx) = mpsc::unbounded_channel();
state.peer_event_subscriber = Some(tx);
let _ = response.send(Ok(rx));
}
NetworkCommand::BroadcastToValidators { message, response } => {
let result = match state.peer_manager.validator_registry() {
None => Err(NetworkError::InvalidConfig(
"consensus-direct broadcast requires an installed ValidatorRegistry"
.to_string(),
)),
Some(registry) => {
let validator_set = registry.validator_peer_ids();
let local = *state.swarm.local_peer_id();
let request = ConsensusDirectRequest::Message(message);
let mut dispatched = 0usize;
for peer in validator_set.iter() {
if *peer == local {
continue;
}
let _request_id = state
.swarm
.behaviour_mut()
.consensus_direct
.send_request(peer, request.clone());
dispatched += 1;
}
Ok(dispatched)
}
};
let _ = response.send(result);
}
NetworkCommand::SubscribeConsensusDirect { response } => {
let (tx, rx) = mpsc::unbounded_channel();
state.consensus_direct_subscriber = Some(tx);
state.consensus_direct_subscriber_ever_attached = true;
let _ = response.send(Ok(rx));
}
NetworkCommand::ConnectedValidatorCount { response } => {
let count = match state.peer_manager.validator_registry() {
None => 0,
Some(registry) => {
let validator_set = registry.validator_peer_ids();
let local = *state.swarm.local_peer_id();
state
.swarm
.connected_peers()
.filter(|p| **p != local && validator_set.contains(*p))
.count()
}
};
let _ = response.send(Ok(count));
}
NetworkCommand::SetMpcDidResolver { resolver, response } => {
state.mpc_did_resolver = Some(resolver);
let _ = response.send(Ok(()));
}
NetworkCommand::SendMpcRelayMessage { message, response } => {
let result = match state.mpc_did_resolver.as_ref() {
None => Err(NetworkError::InvalidConfig(
"mpc-relay send requires an installed MpcDidResolver".to_string(),
)),
Some(resolver) => match resolver.peer_id_for_did(&message.to_did) {
None => Err(NetworkError::PeerNotFound(message.to_did.clone())),
Some(peer) => {
let _req_id = state
.swarm
.behaviour_mut()
.mpc_relay
.send_request(&peer, message);
Ok(())
}
},
};
let _ = response.send(result);
}
NetworkCommand::SubscribeMpcRelay { response } => {
let (tx, rx) = mpsc::unbounded_channel();
state.mpc_relay_subscriber = Some(tx);
state.mpc_relay_subscriber_ever_attached = true;
let _ = response.send(Ok(rx));
}
NetworkCommand::Shutdown { response } => {
let _ = response.send(Ok(()));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn ma(s: &str) -> Multiaddr {
s.parse().expect("valid multiaddr")
}
#[test]
fn extract_port_pulls_tcp_port() {
assert_eq!(extract_port(&ma("/ip4/35.184.63.8/tcp/9000")), Some(9000));
}
#[test]
fn extract_port_pulls_udp_port_for_quic() {
assert_eq!(
extract_port(&ma("/ip4/35.184.63.8/udp/9000/quic-v1")),
Some(9000)
);
}
#[test]
fn extract_port_handles_p2p_suffix() {
let m = ma("/ip4/10.0.0.5/tcp/9000/p2p/12D3KooWGgjoKhKXBvN6jFWqn5sJE6KD38bXtNaoE6itbRUjGBxK");
assert_eq!(extract_port(&m), Some(9000));
}
#[test]
fn observed_port_accepted_when_matches_listen() {
let listen = vec![ma("/ip4/0.0.0.0/tcp/9000")];
let observed = ma("/ip4/35.184.63.8/tcp/9000");
assert!(is_observed_port_one_of_ours(&observed, &listen));
}
#[test]
fn observed_port_rejected_when_ephemeral_source_port() {
let listen = vec![ma("/ip4/0.0.0.0/tcp/9000")];
let observed = ma("/ip4/35.184.63.8/tcp/39692");
assert!(!is_observed_port_one_of_ours(&observed, &listen));
}
#[test]
fn observed_port_accepted_against_quic_listener() {
let listen = vec![ma("/ip4/0.0.0.0/udp/9000/quic-v1")];
let observed = ma("/ip4/35.184.63.8/tcp/9000");
assert!(is_observed_port_one_of_ours(&observed, &listen));
}
#[test]
fn observed_port_rejected_when_no_listeners() {
let listen: Vec<Multiaddr> = vec![];
let observed = ma("/ip4/35.184.63.8/tcp/9000");
assert!(!is_observed_port_one_of_ours(&observed, &listen));
}
#[test]
fn observed_port_rejected_when_no_port_in_observation() {
let listen = vec![ma("/ip4/0.0.0.0/tcp/9000")];
let observed = ma("/ip4/35.184.63.8");
assert!(!is_observed_port_one_of_ours(&observed, &listen));
}
}