mod kad;
mod request_response;
mod swarm;
use crate::{driver::SwarmDriver, error::Result};
use core::fmt;
use custom_debug::Debug as CustomDebug;
#[cfg(feature = "local")]
use libp2p::mdns;
use libp2p::{
kad::{Record, RecordKey, K_VALUE},
request_response::ResponseChannel as PeerResponseChannel,
Multiaddr, PeerId,
};
use sn_evm::PaymentQuote;
#[cfg(feature = "open-metrics")]
use sn_protocol::CLOSE_GROUP_SIZE;
use sn_protocol::{
messages::{Query, Request, Response},
NetworkAddress, PrettyPrintRecordKey,
};
#[cfg(feature = "open-metrics")]
use std::collections::HashSet;
use std::{
collections::BTreeSet,
fmt::{Debug, Formatter},
};
use tokio::sync::oneshot;
#[derive(CustomDebug)]
pub(super) enum NodeEvent {
#[cfg(feature = "upnp")]
Upnp(libp2p::upnp::Event),
MsgReceived(libp2p::request_response::Event<Request, Response>),
Kademlia(libp2p::kad::Event),
#[cfg(feature = "local")]
Mdns(Box<mdns::Event>),
Identify(Box<libp2p::identify::Event>),
RelayClient(Box<libp2p::relay::client::Event>),
RelayServer(Box<libp2p::relay::Event>),
Void(void::Void),
}
#[cfg(feature = "upnp")]
impl From<libp2p::upnp::Event> for NodeEvent {
fn from(event: libp2p::upnp::Event) -> Self {
NodeEvent::Upnp(event)
}
}
impl From<libp2p::request_response::Event<Request, Response>> for NodeEvent {
fn from(event: libp2p::request_response::Event<Request, Response>) -> Self {
NodeEvent::MsgReceived(event)
}
}
impl From<libp2p::kad::Event> for NodeEvent {
fn from(event: libp2p::kad::Event) -> Self {
NodeEvent::Kademlia(event)
}
}
#[cfg(feature = "local")]
impl From<mdns::Event> for NodeEvent {
fn from(event: mdns::Event) -> Self {
NodeEvent::Mdns(Box::new(event))
}
}
impl From<libp2p::identify::Event> for NodeEvent {
fn from(event: libp2p::identify::Event) -> Self {
NodeEvent::Identify(Box::new(event))
}
}
impl From<libp2p::relay::client::Event> for NodeEvent {
fn from(event: libp2p::relay::client::Event) -> Self {
NodeEvent::RelayClient(Box::new(event))
}
}
impl From<libp2p::relay::Event> for NodeEvent {
fn from(event: libp2p::relay::Event) -> Self {
NodeEvent::RelayServer(Box::new(event))
}
}
impl From<void::Void> for NodeEvent {
fn from(event: void::Void) -> Self {
NodeEvent::Void(event)
}
}
#[derive(CustomDebug)]
pub enum MsgResponder {
FromSelf(Option<oneshot::Sender<Result<Response>>>),
FromPeer(PeerResponseChannel<Response>),
}
pub enum NetworkEvent {
QueryRequestReceived {
query: Query,
channel: MsgResponder,
},
ResponseReceived {
res: Response,
},
PeerAdded(PeerId, usize),
PeerRemoved(PeerId, usize),
PeerWithUnsupportedProtocol {
our_protocol: String,
their_protocol: String,
},
KeysToFetchForReplication(Vec<(PeerId, RecordKey)>),
NewListenAddr(Multiaddr),
UnverifiedRecord(Record),
TerminateNode { reason: TerminateNodeReason },
FailedToFetchHolders(BTreeSet<PeerId>),
QuoteVerification { quotes: Vec<(PeerId, PaymentQuote)> },
ChunkProofVerification {
peer_id: PeerId,
key_to_verify: NetworkAddress,
},
}
#[derive(Debug, Clone)]
pub enum TerminateNodeReason {
HardDiskWriteError,
UpnpGatewayNotFound,
}
impl Debug for NetworkEvent {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
match self {
NetworkEvent::QueryRequestReceived { query, .. } => {
write!(f, "NetworkEvent::QueryRequestReceived({query:?})")
}
NetworkEvent::ResponseReceived { res, .. } => {
write!(f, "NetworkEvent::ResponseReceived({res:?})")
}
NetworkEvent::PeerAdded(peer_id, connected_peers) => {
write!(f, "NetworkEvent::PeerAdded({peer_id:?}, {connected_peers})")
}
NetworkEvent::PeerRemoved(peer_id, connected_peers) => {
write!(
f,
"NetworkEvent::PeerRemoved({peer_id:?}, {connected_peers})"
)
}
NetworkEvent::PeerWithUnsupportedProtocol {
our_protocol,
their_protocol,
} => {
write!(f, "NetworkEvent::PeerWithUnsupportedProtocol({our_protocol:?}, {their_protocol:?})")
}
NetworkEvent::KeysToFetchForReplication(list) => {
let keys_len = list.len();
write!(f, "NetworkEvent::KeysForReplication({keys_len:?})")
}
NetworkEvent::NewListenAddr(addr) => {
write!(f, "NetworkEvent::NewListenAddr({addr:?})")
}
NetworkEvent::UnverifiedRecord(record) => {
let pretty_key = PrettyPrintRecordKey::from(&record.key);
write!(f, "NetworkEvent::UnverifiedRecord({pretty_key:?})")
}
NetworkEvent::TerminateNode { reason } => {
write!(f, "NetworkEvent::TerminateNode({reason:?})")
}
NetworkEvent::FailedToFetchHolders(bad_nodes) => {
write!(f, "NetworkEvent::FailedToFetchHolders({bad_nodes:?})")
}
NetworkEvent::QuoteVerification { quotes } => {
write!(
f,
"NetworkEvent::QuoteVerification({} quotes)",
quotes.len()
)
}
NetworkEvent::ChunkProofVerification {
peer_id,
key_to_verify: keys_to_verify,
} => {
write!(
f,
"NetworkEvent::ChunkProofVerification({peer_id:?} {keys_to_verify:?})"
)
}
}
}
}
impl SwarmDriver {
#[cfg(feature = "open-metrics")]
pub(crate) fn check_for_change_in_our_close_group(&mut self) {
let closest_k_peers = self.get_closest_k_value_local_peers();
let new_closest_peers: Vec<_> =
closest_k_peers.into_iter().take(CLOSE_GROUP_SIZE).collect();
let old = self.close_group.iter().cloned().collect::<HashSet<_>>();
let new_members: Vec<_> = new_closest_peers
.iter()
.filter(|p| !old.contains(p))
.collect();
if !new_members.is_empty() {
debug!("The close group has been updated. The new members are {new_members:?}");
debug!("New close group: {new_closest_peers:?}");
self.close_group = new_closest_peers.clone();
self.record_change_in_close_group(new_closest_peers);
}
}
pub(crate) fn update_on_peer_addition(&mut self, added_peer: PeerId) {
self.peers_in_rt = self.peers_in_rt.saturating_add(1);
let n_peers = self.peers_in_rt;
info!("New peer added to routing table: {added_peer:?}, now we have #{n_peers} connected peers");
#[cfg(feature = "loud")]
println!("New peer added to routing table: {added_peer:?}, now we have #{n_peers} connected peers");
self.log_kbuckets(&added_peer);
self.send_event(NetworkEvent::PeerAdded(added_peer, self.peers_in_rt));
#[cfg(feature = "open-metrics")]
if self.metrics_recorder.is_some() {
self.check_for_change_in_our_close_group();
}
#[cfg(feature = "open-metrics")]
if let Some(metrics_recorder) = &self.metrics_recorder {
metrics_recorder
.peers_in_routing_table
.set(self.peers_in_rt as i64);
}
}
pub(crate) fn update_on_peer_removal(&mut self, removed_peer: PeerId) {
self.peers_in_rt = self.peers_in_rt.saturating_sub(1);
let _result = self.swarm.disconnect_peer_id(removed_peer);
info!(
"Peer removed from routing table: {removed_peer:?}, now we have #{} connected peers",
self.peers_in_rt
);
self.log_kbuckets(&removed_peer);
self.send_event(NetworkEvent::PeerRemoved(removed_peer, self.peers_in_rt));
#[cfg(feature = "open-metrics")]
if self.metrics_recorder.is_some() {
self.check_for_change_in_our_close_group();
}
#[cfg(feature = "open-metrics")]
if let Some(metrics_recorder) = &self.metrics_recorder {
metrics_recorder
.peers_in_routing_table
.set(self.peers_in_rt as i64);
}
}
pub(crate) fn log_kbuckets(&mut self, peer: &PeerId) {
let distance = NetworkAddress::from_peer(self.self_peer_id)
.distance(&NetworkAddress::from_peer(*peer));
info!("Peer {peer:?} has a {:?} distance to us", distance.ilog2());
let mut kbucket_table_stats = vec![];
let mut index = 0;
let mut total_peers = 0;
let mut peers_in_non_full_buckets = 0;
let mut num_of_full_buckets = 0;
for kbucket in self.swarm.behaviour_mut().kademlia.kbuckets() {
let range = kbucket.range();
let num_entires = kbucket.num_entries();
if num_entires >= K_VALUE.get() {
num_of_full_buckets += 1;
} else {
peers_in_non_full_buckets += num_entires;
}
total_peers += num_entires;
if let Some(distance) = range.0.ilog2() {
kbucket_table_stats.push((index, num_entires, distance));
} else {
error!("bucket #{index:?} is ourself ???!!!");
}
index += 1;
}
let estimated_network_size =
Self::estimate_network_size(peers_in_non_full_buckets, num_of_full_buckets);
#[cfg(feature = "open-metrics")]
if let Some(metrics_recorder) = &self.metrics_recorder {
let _ = metrics_recorder
.estimated_network_size
.set(estimated_network_size as i64);
}
if total_peers % 10 == 0 && total_peers != self.peers_in_rt {
warn!(
"Total peers in routing table: {}, but kbucket table has {total_peers} peers",
self.peers_in_rt
);
}
info!("kBucketTable has {index:?} kbuckets {total_peers:?} peers, {kbucket_table_stats:?}, estimated network size: {estimated_network_size:?}");
#[cfg(feature = "loud")]
println!("Estimated network size: {estimated_network_size:?}");
}
fn estimate_network_size(
peers_in_non_full_buckets: usize,
num_of_full_buckets: usize,
) -> usize {
(peers_in_non_full_buckets + 1) * (2_usize.pow(num_of_full_buckets as u32))
}
}