Skip to main content

ab_networking/
node_runner.rs

1#![expect(
2    clippy::rest_pattern_accessible_field,
3    reason = "Too verbose otherwise"
4)]
5
6use crate::behavior::persistent_parameters::{
7    KnownPeersRegistry, PeerAddressRemovedEvent, append_p2p_suffix, remove_p2p_suffix,
8};
9use crate::behavior::{Behavior, Event};
10use crate::constructor::DummyRecordStore;
11use crate::constructor::temporary_bans::TemporaryBans;
12use crate::protocols::request_response::request_response_factory::{
13    Event as RequestResponseEvent, IfDisconnected,
14};
15use crate::shared::{Command, CreatedSubscription, PeerDiscovered, Shared};
16use crate::utils::{SubspaceMetrics, is_global_address_or_dns, strip_peer_id};
17use async_lock::Mutex as AsyncMutex;
18use bytes::Bytes;
19use event_listener_primitives::HandlerId;
20use futures::channel::mpsc;
21use futures::future::Fuse;
22use futures::{FutureExt, StreamExt};
23use libp2p::autonat::{Event as AutonatEvent, NatStatus, OutboundProbeEvent};
24use libp2p::core::ConnectedPoint;
25use libp2p::gossipsub::{Event as GossipsubEvent, TopicHash};
26use libp2p::identify::Event as IdentifyEvent;
27use libp2p::kad::{
28    Behaviour as Kademlia, BootstrapOk, Event as KademliaEvent, GetClosestPeersError,
29    GetClosestPeersOk, GetProvidersError, GetProvidersOk, GetRecordError, GetRecordOk,
30    InboundRequest, KBucketKey, PeerRecord, ProgressStep, PutRecordOk, QueryId, QueryResult,
31    Quorum, Record, RecordKey,
32};
33use libp2p::metrics::{Metrics, Recorder};
34use libp2p::multiaddr::Protocol;
35use libp2p::swarm::dial_opts::DialOpts;
36use libp2p::swarm::{DialError, SwarmEvent};
37use libp2p::{Multiaddr, PeerId, Swarm, TransportError};
38use nohash_hasher::IntMap;
39use parking_lot::Mutex;
40use std::collections::hash_map::Entry;
41use std::collections::{HashMap, HashSet};
42use std::net::IpAddr;
43use std::pin::Pin;
44use std::sync::atomic::Ordering;
45use std::sync::{Arc, Weak};
46use std::time::Duration;
47use std::{fmt, slice};
48use tokio::sync::OwnedSemaphorePermit;
49use tokio::task::yield_now;
50use tokio::time::Sleep;
51use tracing::{debug, error, trace, warn};
52
53enum QueryResultSender {
54    Value {
55        sender: mpsc::UnboundedSender<PeerRecord>,
56        // Just holding onto permit while data structure is not dropped
57        _permit: OwnedSemaphorePermit,
58    },
59    ClosestPeers {
60        sender: mpsc::UnboundedSender<PeerId>,
61        // Just holding onto permit while data structure is not dropped
62        _permit: Option<OwnedSemaphorePermit>,
63    },
64    Providers {
65        key: RecordKey,
66        sender: mpsc::UnboundedSender<PeerId>,
67        // Just holding onto permit while data structure is not dropped
68        _permit: Option<OwnedSemaphorePermit>,
69    },
70    PutValue {
71        sender: mpsc::UnboundedSender<()>,
72        // Just holding onto permit while data structure is not dropped
73        _permit: OwnedSemaphorePermit,
74    },
75    Bootstrap {
76        sender: mpsc::UnboundedSender<()>,
77    },
78}
79
80#[derive(Debug, Default)]
81enum BootstrapCommandState {
82    #[default]
83    NotStarted,
84    InProgress(mpsc::UnboundedReceiver<()>),
85    Finished,
86}
87
88/// Runner for the Node.
89#[must_use = "Node does not function properly unless its runner is driven forward"]
90pub struct NodeRunner {
91    /// Should non-global addresses be added to the DHT?
92    allow_non_global_addresses_in_dht: bool,
93    /// Whether node is listening on some addresses
94    is_listening: bool,
95    command_receiver: mpsc::Receiver<Command>,
96    swarm: Swarm<Behavior>,
97    shared_weak: Weak<Shared>,
98    /// How frequently should random queries be done using Kademlia DHT to populate routing table.
99    next_random_query_interval: Duration,
100    query_id_receivers: HashMap<QueryId, QueryResultSender>,
101    /// Global subscription counter, is assigned to every (logical) subscription and is used for
102    /// unsubscribing.
103    next_subscription_id: usize,
104    /// Topic subscription senders for logical subscriptions (multiple logical subscriptions can be
105    /// present for the same physical subscription).
106    topic_subscription_senders: HashMap<TopicHash, IntMap<usize, mpsc::UnboundedSender<Bytes>>>,
107    random_query_timeout: Pin<Box<Fuse<Sleep>>>,
108    /// Defines an interval between periodical tasks.
109    periodical_tasks_interval: Pin<Box<Fuse<Sleep>>>,
110    /// Manages the networking parameters like known peers and addresses
111    known_peers_registry: Box<dyn KnownPeersRegistry>,
112    connected_servers: HashSet<PeerId>,
113    /// Defines set of peers with a permanent connection (and reconnection if necessary).
114    reserved_peers: HashMap<PeerId, Multiaddr>,
115    /// Temporarily banned peers.
116    temporary_bans: Arc<Mutex<TemporaryBans>>,
117    /// Libp2p Prometheus metrics.
118    libp2p_metrics: Option<Metrics>,
119    /// Subspace Prometheus metrics.
120    metrics: Option<SubspaceMetrics>,
121    /// Mapping from specific peer to ip addresses
122    peer_ip_addresses: HashMap<PeerId, HashSet<IpAddr>>,
123    /// Defines protocol version for the network peers. Affects network partition.
124    protocol_version: String,
125    /// Addresses to bootstrap Kademlia network
126    bootstrap_addresses: Vec<Multiaddr>,
127    /// Ensures a single bootstrap on run() invocation.
128    bootstrap_command_state: Arc<AsyncMutex<BootstrapCommandState>>,
129    /// Receives an event on peer address removal from the persistent storage.
130    removed_addresses_rx: mpsc::UnboundedReceiver<PeerAddressRemovedEvent>,
131    /// Optional storage for the [`HandlerId`] of the address removal task.
132    /// We keep to stop the task along with the rest of the networking.
133    _address_removal_task_handler_id: Option<HandlerId>,
134}
135
136impl fmt::Debug for NodeRunner {
137    #[inline]
138    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
139        f.debug_struct("NodeRunner").finish_non_exhaustive()
140    }
141}
142
143// Helper struct for NodeRunner configuration (clippy requirement).
144pub(crate) struct NodeRunnerConfig {
145    pub(crate) allow_non_global_addresses_in_dht: bool,
146    /// Whether node is listening on some addresses
147    pub(crate) is_listening: bool,
148    pub(crate) command_receiver: mpsc::Receiver<Command>,
149    pub(crate) swarm: Swarm<Behavior>,
150    pub(crate) shared_weak: Weak<Shared>,
151    pub(crate) next_random_query_interval: Duration,
152    pub(crate) known_peers_registry: Box<dyn KnownPeersRegistry>,
153    pub(crate) reserved_peers: HashMap<PeerId, Multiaddr>,
154    pub(crate) temporary_bans: Arc<Mutex<TemporaryBans>>,
155    pub(crate) libp2p_metrics: Option<Metrics>,
156    pub(crate) metrics: Option<SubspaceMetrics>,
157    pub(crate) protocol_version: String,
158    pub(crate) bootstrap_addresses: Vec<Multiaddr>,
159}
160
161impl NodeRunner {
162    pub(crate) fn new(
163        NodeRunnerConfig {
164            allow_non_global_addresses_in_dht,
165            is_listening,
166            command_receiver,
167            swarm,
168            shared_weak,
169            next_random_query_interval,
170            mut known_peers_registry,
171            reserved_peers,
172            temporary_bans,
173            libp2p_metrics,
174            metrics,
175            protocol_version,
176            bootstrap_addresses,
177        }: NodeRunnerConfig,
178    ) -> Self {
179        // Setup the address removal events exchange between persistent params storage and Kademlia.
180        let (removed_addresses_tx, removed_addresses_rx) = mpsc::unbounded();
181        let mut address_removal_task_handler_id = None;
182        if let Some(handler_id) = known_peers_registry.on_unreachable_address({
183            Arc::new(move |event| {
184                if let Err(error) = removed_addresses_tx.unbounded_send(event.clone()) {
185                    debug!(?error, ?event, "Cannot send PeerAddressRemovedEvent");
186                }
187            })
188        }) {
189            address_removal_task_handler_id.replace(handler_id);
190        }
191
192        Self {
193            allow_non_global_addresses_in_dht,
194            is_listening,
195            command_receiver,
196            swarm,
197            shared_weak,
198            next_random_query_interval,
199            query_id_receivers: HashMap::default(),
200            next_subscription_id: 0,
201            topic_subscription_senders: HashMap::default(),
202            // We'll make the first query right away and continue at the interval.
203            random_query_timeout: Box::pin(tokio::time::sleep(Duration::from_secs(0)).fuse()),
204            // We'll make the first dial right away and continue at the interval.
205            periodical_tasks_interval: Box::pin(tokio::time::sleep(Duration::from_secs(0)).fuse()),
206            known_peers_registry,
207            connected_servers: HashSet::new(),
208            reserved_peers,
209            temporary_bans,
210            libp2p_metrics,
211            metrics,
212            peer_ip_addresses: HashMap::new(),
213            protocol_version,
214            bootstrap_addresses,
215            bootstrap_command_state: Arc::new(AsyncMutex::new(BootstrapCommandState::default())),
216            removed_addresses_rx,
217            _address_removal_task_handler_id: address_removal_task_handler_id,
218        }
219    }
220
221    /// Drives the main networking future forward.
222    pub async fn run(&mut self) {
223        if self.is_listening {
224            // Wait for listen addresses, otherwise we will get ephemeral addresses in external
225            // address candidates that we do not want
226            loop {
227                if self.swarm.listeners().next().is_some() {
228                    break;
229                }
230
231                if let Some(swarm_event) = self.swarm.next().await {
232                    self.register_event_metrics(&swarm_event);
233                    self.handle_swarm_event(swarm_event).await;
234                } else {
235                    break;
236                }
237            }
238        }
239
240        self.bootstrap().await;
241
242        loop {
243            futures::select! {
244                () = &mut self.random_query_timeout => {
245                    self.handle_random_query_interval();
246                    // Increase interval 2x, but to at most 1 minute.
247                    self.random_query_timeout =
248                        Box::pin(tokio::time::sleep(self.next_random_query_interval).fuse());
249                    self.next_random_query_interval =
250                        (self.next_random_query_interval * 2).min(Duration::from_mins(1));
251                },
252                swarm_event = self.swarm.next() => {
253                    if let Some(swarm_event) = swarm_event {
254                        self.register_event_metrics(&swarm_event);
255                        self.handle_swarm_event(swarm_event).await;
256                    } else {
257                        break;
258                    }
259                },
260                command = self.command_receiver.next() => {
261                    if let Some(command) = command {
262                        self.handle_command(command);
263                    } else {
264                        break;
265                    }
266                },
267                () = self.known_peers_registry.run().fuse() => {
268                    trace!("Network parameters registry runner exited");
269                },
270                () = &mut self.periodical_tasks_interval => {
271                    self.handle_periodical_tasks();
272
273                    self.periodical_tasks_interval =
274                        Box::pin(tokio::time::sleep(Duration::from_secs(5)).fuse());
275                },
276                event = self.removed_addresses_rx.select_next_some() => {
277                    self.handle_removed_address_event(event);
278                },
279            }
280
281            // Allow to exit from busy loop during graceful shutdown
282            yield_now().await;
283        }
284    }
285
286    /// Bootstraps Kademlia network
287    async fn bootstrap(&mut self) {
288        // Add bootstrap nodes first to make sure there is space for them in k-buckets
289        for (peer_id, address) in strip_peer_id(self.bootstrap_addresses.clone()) {
290            self.swarm
291                .behaviour_mut()
292                .kademlia
293                .add_address(&peer_id, address);
294        }
295
296        let known_peers = self.known_peers_registry.all_known_peers().await;
297
298        if !known_peers.is_empty() {
299            for (peer_id, addresses) in known_peers {
300                for address in addresses.clone() {
301                    let address = match address.with_p2p(peer_id) {
302                        Ok(address) => address,
303                        Err(address) => {
304                            warn!(%peer_id, %address, "Failed to add peer ID to known peer address");
305                            break;
306                        }
307                    };
308                    self.swarm
309                        .behaviour_mut()
310                        .kademlia
311                        .add_address(&peer_id, address);
312                }
313
314                if let Err(error) = self
315                    .swarm
316                    .dial(DialOpts::peer_id(peer_id).addresses(addresses).build())
317                {
318                    warn!(%peer_id, %error, "Failed to dial peer during bootstrapping");
319                }
320            }
321
322            // Do bootstrap asynchronously
323            self.handle_command(Command::Bootstrap {
324                result_sender: None,
325            });
326            return;
327        }
328
329        let bootstrap_command_state = Arc::clone(&self.bootstrap_command_state);
330        let mut bootstrap_command_state = bootstrap_command_state.lock().await;
331        let bootstrap_command_receiver = match &mut *bootstrap_command_state {
332            BootstrapCommandState::NotStarted => {
333                debug!("Bootstrap started.");
334
335                let (bootstrap_command_sender, bootstrap_command_receiver) = mpsc::unbounded();
336
337                self.handle_command(Command::Bootstrap {
338                    result_sender: Some(bootstrap_command_sender),
339                });
340
341                *bootstrap_command_state =
342                    BootstrapCommandState::InProgress(bootstrap_command_receiver);
343                match &mut *bootstrap_command_state {
344                    BootstrapCommandState::InProgress(bootstrap_command_receiver) => {
345                        bootstrap_command_receiver
346                    }
347                    _ => {
348                        unreachable!("Was just set to that exact value");
349                    }
350                }
351            }
352            BootstrapCommandState::InProgress(bootstrap_command_receiver) => {
353                bootstrap_command_receiver
354            }
355            BootstrapCommandState::Finished => {
356                return;
357            }
358        };
359
360        let mut bootstrap_step = 0usize;
361        loop {
362            futures::select! {
363                swarm_event = self.swarm.next() => {
364                    if let Some(swarm_event) = swarm_event {
365                        self.register_event_metrics(&swarm_event);
366                        self.handle_swarm_event(swarm_event).await;
367                    } else {
368                        break;
369                    }
370                },
371                result = bootstrap_command_receiver.next() => {
372                    if result.is_some() {
373                        debug!(%bootstrap_step, "Kademlia bootstrapping...");
374                        bootstrap_step += 1;
375                    } else {
376                        break;
377                    }
378                }
379            }
380        }
381
382        debug!("Bootstrap finished.");
383        *bootstrap_command_state = BootstrapCommandState::Finished;
384    }
385
386    /// Handles periodical tasks.
387    fn handle_periodical_tasks(&mut self) {
388        // Log current connections.
389        let network_info = self.swarm.network_info();
390        let connections = network_info.connection_counters();
391
392        debug!(?connections, "Current connections and limits.");
393
394        // Renew known external addresses.
395        let mut external_addresses = self.swarm.external_addresses().cloned().collect::<Vec<_>>();
396
397        if let Some(shared) = self.shared_weak.upgrade() {
398            debug!(?external_addresses, "Renew external addresses.");
399            let mut addresses = shared.external_addresses.lock();
400            addresses.clear();
401            addresses.append(&mut external_addresses);
402        }
403
404        self.log_kademlia_stats();
405    }
406
407    fn handle_random_query_interval(&mut self) {
408        let random_peer_id = PeerId::random();
409
410        trace!("Starting random Kademlia query for {}", random_peer_id);
411
412        self.swarm
413            .behaviour_mut()
414            .kademlia
415            .get_closest_peers(random_peer_id);
416    }
417
418    fn handle_removed_address_event(&mut self, event: PeerAddressRemovedEvent) {
419        trace!(?event, "Peer address removed event.");
420
421        let bootstrap_node_ids = strip_peer_id(self.bootstrap_addresses.clone())
422            .into_iter()
423            .map(|(peer_id, _)| peer_id)
424            .collect::<Vec<_>>();
425
426        if bootstrap_node_ids.contains(&event.peer_id) {
427            debug!(
428                ?event,
429                ?bootstrap_node_ids,
430                "Skipped removing bootstrap node from Kademlia buckets."
431            );
432
433            return;
434        }
435
436        // Remove both versions of the address
437        self.swarm.behaviour_mut().kademlia.remove_address(
438            &event.peer_id,
439            &append_p2p_suffix(event.peer_id, event.address.clone()),
440        );
441
442        self.swarm
443            .behaviour_mut()
444            .kademlia
445            .remove_address(&event.peer_id, &remove_p2p_suffix(event.address));
446    }
447
448    fn handle_remove_listeners(&mut self, removed_listeners: &[Multiaddr]) {
449        let Some(shared) = self.shared_weak.upgrade() else {
450            return;
451        };
452
453        // Remove both versions of the address
454        let peer_id = shared.id;
455        shared.listeners.lock().retain(|old_listener| {
456            !removed_listeners.contains(&append_p2p_suffix(peer_id, old_listener.clone()))
457                && !removed_listeners.contains(&remove_p2p_suffix(old_listener.clone()))
458        });
459    }
460
461    async fn handle_swarm_event(&mut self, swarm_event: SwarmEvent<Event>) {
462        #[expect(clippy::ref_patterns, reason = "Much less awkward this way")]
463        match swarm_event {
464            SwarmEvent::Behaviour(Event::Identify(event)) => {
465                self.handle_identify_event(*event);
466            }
467            SwarmEvent::Behaviour(Event::Kademlia(event)) => {
468                self.handle_kademlia_event(event);
469            }
470            SwarmEvent::Behaviour(Event::Gossipsub(event)) => {
471                self.handle_gossipsub_event(event);
472            }
473            SwarmEvent::Behaviour(Event::RequestResponse(event)) => {
474                self.handle_request_response_event(event);
475            }
476            SwarmEvent::Behaviour(Event::Autonat(event)) => {
477                self.handle_autonat_event(event);
478            }
479            ref event @ SwarmEvent::NewListenAddr { ref address, .. } => {
480                trace!(?event, "New local listener  event.");
481
482                let Some(shared) = self.shared_weak.upgrade() else {
483                    return;
484                };
485                shared.listeners.lock().push(address.clone());
486                shared.handlers.new_listener.call_simple(address);
487            }
488            ref event @ SwarmEvent::ListenerClosed { ref addresses, .. } => {
489                trace!(?event, "Local listener closed event.");
490                self.handle_remove_listeners(addresses);
491            }
492            ref event @ SwarmEvent::ExpiredListenAddr { ref address, .. } => {
493                trace!(?event, "Local listener expired event.");
494                self.handle_remove_listeners(slice::from_ref(address));
495            }
496            SwarmEvent::ConnectionEstablished {
497                peer_id,
498                endpoint,
499                num_established,
500                ..
501            } => {
502                // Save known addresses that were successfully dialed.
503                if let ConnectedPoint::Dialer { address, .. } = &endpoint {
504                    // filter non-global addresses when non-globals addresses are disabled
505                    if self.allow_non_global_addresses_in_dht || is_global_address_or_dns(address) {
506                        self.known_peers_registry
507                            .add_known_peer(peer_id, vec![address.clone()])
508                            .await;
509                    }
510                }
511
512                let Some(shared) = self.shared_weak.upgrade() else {
513                    return;
514                };
515
516                let is_reserved_peer = self.reserved_peers.contains_key(&peer_id);
517                debug!(
518                    %peer_id,
519                    %is_reserved_peer,
520                    ?endpoint,
521                    %num_established,
522                    "Connection established"
523                );
524
525                let maybe_remote_ip =
526                    endpoint
527                        .get_remote_address()
528                        .iter()
529                        .find_map(|protocol| match protocol {
530                            Protocol::Ip4(ip) => Some(IpAddr::V4(ip)),
531                            Protocol::Ip6(ip) => Some(IpAddr::V6(ip)),
532                            _ => None,
533                        });
534                if let Some(ip) = maybe_remote_ip {
535                    self.peer_ip_addresses
536                        .entry(peer_id)
537                        .and_modify(|ips| {
538                            ips.insert(ip);
539                        })
540                        .or_insert(HashSet::from([ip]));
541                }
542
543                let num_established_peer_connections = shared
544                    .num_established_peer_connections
545                    .fetch_add(1, Ordering::SeqCst)
546                    + 1;
547
548                shared
549                    .handlers
550                    .num_established_peer_connections_change
551                    .call_simple(&num_established_peer_connections);
552
553                // A new connection
554                if num_established.get() == 1 {
555                    shared.handlers.connected_peer.call_simple(&peer_id);
556                }
557
558                if let Some(metrics) = self.metrics.as_ref() {
559                    metrics.inc_established_connections();
560                }
561            }
562            SwarmEvent::ConnectionClosed {
563                peer_id,
564                num_established,
565                cause,
566                ..
567            } => {
568                let Some(shared) = self.shared_weak.upgrade() else {
569                    return;
570                };
571
572                debug!(
573                    %peer_id,
574                    ?cause,
575                    %num_established,
576                    "Connection closed with peer"
577                );
578
579                if num_established == 0 {
580                    self.peer_ip_addresses.remove(&peer_id);
581                    self.connected_servers.remove(&peer_id);
582                }
583                let num_established_peer_connections = shared
584                    .num_established_peer_connections
585                    .fetch_sub(1, Ordering::SeqCst)
586                    - 1;
587
588                shared
589                    .handlers
590                    .num_established_peer_connections_change
591                    .call_simple(&num_established_peer_connections);
592
593                // No more connections
594                if num_established == 0 {
595                    shared.handlers.disconnected_peer.call_simple(&peer_id);
596                }
597
598                if let Some(metrics) = self.metrics.as_ref() {
599                    metrics.dec_established_connections();
600                }
601            }
602            SwarmEvent::OutgoingConnectionError { peer_id, error, .. } => {
603                if let Some(peer_id) = &peer_id {
604                    let should_ban_temporarily =
605                        self.should_temporary_ban_on_dial_error(peer_id, &error);
606
607                    trace!(%should_ban_temporarily, "Temporary bans conditions.");
608
609                    if should_ban_temporarily {
610                        self.temporary_bans.lock().create_or_extend(peer_id);
611                        debug!(%peer_id, ?error, "Peer was temporarily banned.");
612                    }
613                }
614
615                debug!(
616                    ?peer_id,
617                    ?error,
618                    "SwarmEvent::OutgoingConnectionError for peer."
619                );
620
621                match error {
622                    DialError::Transport(ref addresses) => {
623                        for (addr, _) in addresses {
624                            trace!(?error, ?peer_id, %addr, "SwarmEvent::OutgoingConnectionError (DialError::Transport) for peer.");
625                            if let Some(peer_id) = peer_id {
626                                self.known_peers_registry
627                                    .remove_known_peer_addresses(peer_id, vec![addr.clone()])
628                                    .await;
629                            }
630                        }
631                    }
632                    DialError::WrongPeerId { obtained, .. } => {
633                        trace!(?error, ?peer_id, obtained_peer_id=?obtained, "SwarmEvent::WrongPeerId (DialError::WrongPeerId) for peer.");
634
635                        if let Some(ref peer_id) = peer_id {
636                            let kademlia = &mut self.swarm.behaviour_mut().kademlia;
637                            let _: Option<_> = kademlia.remove_peer(peer_id);
638                        }
639                    }
640                    _ => {
641                        trace!(?error, ?peer_id, "SwarmEvent::OutgoingConnectionError");
642                    }
643                }
644            }
645            SwarmEvent::NewExternalAddrCandidate { address } => {
646                trace!(%address, "External address candidate");
647            }
648            SwarmEvent::ExternalAddrConfirmed { address } => {
649                debug!(%address, "Confirmed external address");
650
651                let connected_peers = self.swarm.connected_peers().copied().collect::<Vec<_>>();
652                self.swarm.behaviour_mut().identify.push(connected_peers);
653            }
654            SwarmEvent::ExternalAddrExpired { address } => {
655                debug!(%address, "External address expired");
656
657                let connected_peers = self.swarm.connected_peers().copied().collect::<Vec<_>>();
658                self.swarm.behaviour_mut().identify.push(connected_peers);
659            }
660            other => {
661                trace!("Other swarm event: {:?}", other);
662            }
663        }
664    }
665
666    fn should_temporary_ban_on_dial_error(&self, peer_id: &PeerId, error: &DialError) -> bool {
667        // TODO: Replace with banning of addresses rather peer IDs if this helps
668        if true {
669            return false;
670        }
671
672        // Ban temporarily only peers without active connections.
673        if self.swarm.is_connected(peer_id) {
674            return false;
675        }
676
677        match &error {
678            DialError::Transport(addresses) => {
679                for (_, error) in addresses {
680                    match error {
681                        TransportError::MultiaddrNotSupported(_) => {
682                            return true;
683                        }
684                        TransportError::Other(_) => {
685                            // Ignore "temporary ban" errors
686                            if self.temporary_bans.lock().is_banned(peer_id) {
687                                return false;
688                            }
689                        }
690                    }
691                }
692                // Other errors that are not related to temporary bans
693                true
694            }
695            DialError::LocalPeerId { .. } => {
696                // We don't ban ourselves
697                debug!("Local peer dial attempt detected.");
698
699                false
700            }
701            DialError::NoAddresses => {
702                // Let's wait until we get addresses
703                true
704            }
705            DialError::DialPeerConditionFalse(_) => {
706                // These are local conditions, we don't need to ban remote peers
707                false
708            }
709            DialError::Aborted => {
710                // Seems like a transient event
711                false
712            }
713            DialError::WrongPeerId { .. } => {
714                // It's likely that peer was restarted with different identity
715                false
716            }
717            DialError::Denied { .. } => {
718                // We exceeded the connection limits or we hit a black listed peer
719                false
720            }
721        }
722    }
723
724    fn handle_identify_event(&mut self, event: IdentifyEvent) {
725        let local_peer_id = *self.swarm.local_peer_id();
726
727        if let IdentifyEvent::Received {
728            peer_id, mut info, ..
729        } = event
730        {
731            debug!(?peer_id, protocols = ?info.protocols, "IdentifyEvent::Received");
732
733            // Check for network partition
734            if info.protocol_version != self.protocol_version {
735                debug!(
736                    %local_peer_id,
737                    %peer_id,
738                    local_protocol_version = %self.protocol_version,
739                    peer_protocol_version = %info.protocol_version,
740                    "Peer has different protocol version, banning temporarily",
741                );
742
743                self.temporary_bans.lock().create_or_extend(&peer_id);
744                // Forget about this peer until they upgrade
745                let _: Result<(), ()> = self.swarm.disconnect_peer_id(peer_id);
746                self.swarm.behaviour_mut().kademlia.remove_peer(&peer_id);
747                self.known_peers_registry
748                    .remove_all_known_peer_addresses(peer_id);
749
750                return;
751            }
752
753            // Remove temporary ban if there was any
754            self.temporary_bans.lock().remove(&peer_id);
755
756            if info.listen_addrs.len() > 30 {
757                debug!(
758                    %local_peer_id,
759                    %peer_id,
760                    "Node has reported more than 30 addresses; it is identified by {} and {}",
761                    info.protocol_version, info.agent_version
762                );
763                info.listen_addrs.truncate(30);
764            }
765
766            let kademlia = &mut self.swarm.behaviour_mut().kademlia;
767            let full_kademlia_support = kademlia
768                .protocol_names()
769                .iter()
770                .all(|local_protocol| info.protocols.contains(local_protocol));
771
772            if full_kademlia_support {
773                let received_addresses = info
774                    .listen_addrs
775                    .into_iter()
776                    .filter(|address| {
777                        if self.allow_non_global_addresses_in_dht
778                            || is_global_address_or_dns(address)
779                        {
780                            true
781                        } else {
782                            trace!(
783                                %local_peer_id,
784                                %peer_id,
785                                %address,
786                                "Ignoring self-reported non-global address",
787                            );
788
789                            false
790                        }
791                    })
792                    .collect::<Vec<_>>();
793                let received_address_strings = received_addresses
794                    .iter()
795                    .map(ToString::to_string)
796                    .collect::<Vec<_>>();
797                let old_addresses = kademlia
798                    .kbucket(peer_id)
799                    .and_then(|peers| {
800                        let key = peer_id.into();
801                        peers.iter().find_map(|peer| {
802                            (peer.node.key == &key).then_some(
803                                peer.node
804                                    .value
805                                    .iter()
806                                    .filter(|existing_address| {
807                                        let existing_address = existing_address.to_string();
808
809                                        !received_address_strings.iter().any(|received_address| {
810                                            received_address.starts_with(&existing_address)
811                                                || existing_address.starts_with(received_address)
812                                        })
813                                    })
814                                    .cloned()
815                                    .collect::<Vec<_>>(),
816                            )
817                        })
818                    })
819                    .unwrap_or_default();
820
821                for address in received_addresses {
822                    debug!(
823                        %local_peer_id,
824                        %peer_id,
825                        %address,
826                        protocol_names = ?kademlia.protocol_names(),
827                        "Adding self-reported address to Kademlia DHT",
828                    );
829
830                    kademlia.add_address(&peer_id, address);
831                }
832
833                for old_address in old_addresses {
834                    trace!(
835                        %local_peer_id,
836                        %peer_id,
837                        %old_address,
838                        "Removing old self-reported address from Kademlia DHT",
839                    );
840
841                    kademlia.remove_address(&peer_id, &old_address);
842                }
843
844                self.connected_servers.insert(peer_id);
845            } else {
846                debug!(
847                    %local_peer_id,
848                    %peer_id,
849                    peer_protocols = ?info.protocols,
850                    protocol_names = ?kademlia.protocol_names(),
851                    "Peer doesn't support our Kademlia DHT protocol",
852                );
853
854                kademlia.remove_peer(&peer_id);
855                self.connected_servers.remove(&peer_id);
856            }
857        }
858    }
859
860    fn handle_kademlia_event(&mut self, event: KademliaEvent) {
861        trace!("Kademlia event: {:?}", event);
862
863        match event {
864            KademliaEvent::InboundRequest {
865                request: InboundRequest::AddProvider { record },
866            } => {
867                debug!("Unexpected AddProvider request received: {:?}", record);
868            }
869            KademliaEvent::UnroutablePeer { peer } => {
870                debug!(%peer, "Unroutable peer detected");
871
872                self.swarm.behaviour_mut().kademlia.remove_peer(&peer);
873
874                if let Some(shared) = self.shared_weak.upgrade() {
875                    shared
876                        .handlers
877                        .peer_discovered
878                        .call_simple(&PeerDiscovered::UnroutablePeer { peer_id: peer });
879                }
880            }
881            KademliaEvent::RoutablePeer { peer, address } => {
882                debug!(?address, "Routable peer detected: {:?}", peer);
883
884                if let Some(shared) = self.shared_weak.upgrade() {
885                    shared
886                        .handlers
887                        .peer_discovered
888                        .call_simple(&PeerDiscovered::RoutablePeer {
889                            peer_id: peer,
890                            address,
891                        });
892                }
893            }
894            KademliaEvent::PendingRoutablePeer { peer, address } => {
895                debug!(?address, "Pending routable peer detected: {:?}", peer);
896
897                if let Some(shared) = self.shared_weak.upgrade() {
898                    shared
899                        .handlers
900                        .peer_discovered
901                        .call_simple(&PeerDiscovered::RoutablePeer {
902                            peer_id: peer,
903                            address,
904                        });
905                }
906            }
907            KademliaEvent::OutboundQueryProgressed {
908                step: ProgressStep { last, .. },
909                id,
910                result: QueryResult::GetClosestPeers(result),
911                ..
912            } => {
913                let mut cancelled = false;
914                if let Some(QueryResultSender::ClosestPeers { sender, .. }) =
915                    self.query_id_receivers.get(&id)
916                {
917                    match result {
918                        Ok(GetClosestPeersOk { key, peers }) => {
919                            trace!(
920                                "Get closest peers query for {} yielded {} results",
921                                hex::encode(key),
922                                peers.len(),
923                            );
924
925                            if peers.is_empty()
926                                // Connected peers collection is not empty.
927                                && self.swarm.connected_peers().next().is_some()
928                            {
929                                debug!("Random Kademlia query has yielded empty list of peers");
930                            }
931
932                            for peer in peers {
933                                cancelled = Self::unbounded_send_and_cancel_on_error(
934                                    &mut self.swarm.behaviour_mut().kademlia,
935                                    sender,
936                                    peer.peer_id,
937                                    "GetClosestPeersOk",
938                                    id,
939                                ) || cancelled;
940                            }
941                        }
942                        Err(GetClosestPeersError::Timeout { key, peers }) => {
943                            debug!(
944                                "Get closest peers query for {} timed out with {} results",
945                                hex::encode(key),
946                                peers.len(),
947                            );
948
949                            for peer in peers {
950                                cancelled = Self::unbounded_send_and_cancel_on_error(
951                                    &mut self.swarm.behaviour_mut().kademlia,
952                                    sender,
953                                    peer.peer_id,
954                                    "GetClosestPeersError::Timeout",
955                                    id,
956                                ) || cancelled;
957                            }
958                        }
959                    }
960                }
961
962                if last || cancelled {
963                    // There will be no more progress
964                    self.query_id_receivers.remove(&id);
965                }
966            }
967            KademliaEvent::OutboundQueryProgressed {
968                step: ProgressStep { last, .. },
969                id,
970                result: QueryResult::GetRecord(result),
971                ..
972            } => {
973                let mut cancelled = false;
974                if let Some(QueryResultSender::Value { sender, .. }) =
975                    self.query_id_receivers.get(&id)
976                {
977                    match result {
978                        Ok(GetRecordOk::FoundRecord(rec)) => {
979                            trace!(
980                                key = hex::encode(&rec.record.key),
981                                "Get record query succeeded",
982                            );
983
984                            cancelled = Self::unbounded_send_and_cancel_on_error(
985                                &mut self.swarm.behaviour_mut().kademlia,
986                                sender,
987                                rec,
988                                "GetRecordOk",
989                                id,
990                            ) || cancelled;
991                        }
992                        Ok(GetRecordOk::FinishedWithNoAdditionalRecord { .. }) => {
993                            trace!("Get record query yielded no results");
994                        }
995                        Err(error) => match error {
996                            GetRecordError::NotFound { key, .. } => {
997                                debug!(
998                                    key = hex::encode(&key),
999                                    "Get record query failed with no results",
1000                                );
1001                            }
1002                            GetRecordError::Timeout { key } => {
1003                                debug!(key = hex::encode(&key), "Get record query timed out");
1004                            }
1005                        },
1006                    }
1007                }
1008
1009                if last || cancelled {
1010                    // There will be no more progress
1011                    self.query_id_receivers.remove(&id);
1012                }
1013            }
1014            KademliaEvent::OutboundQueryProgressed {
1015                step: ProgressStep { last, .. },
1016                id,
1017                result: QueryResult::GetProviders(result),
1018                ..
1019            } => {
1020                let mut cancelled = false;
1021                if let Some(QueryResultSender::Providers { key, sender, .. }) =
1022                    self.query_id_receivers.get(&id)
1023                {
1024                    match result {
1025                        Ok(GetProvidersOk::FoundProviders { key, providers }) => {
1026                            trace!(
1027                                key = hex::encode(&key),
1028                                "Get providers query yielded {} results",
1029                                providers.len(),
1030                            );
1031
1032                            for provider in providers {
1033                                cancelled = Self::unbounded_send_and_cancel_on_error(
1034                                    &mut self.swarm.behaviour_mut().kademlia,
1035                                    sender,
1036                                    provider,
1037                                    "GetProvidersOk",
1038                                    id,
1039                                ) || cancelled;
1040                            }
1041                        }
1042                        Ok(GetProvidersOk::FinishedWithNoAdditionalRecord { closest_peers }) => {
1043                            trace!(
1044                                key = hex::encode(key),
1045                                closest_peers = %closest_peers.len(),
1046                                "Get providers query yielded no results"
1047                            );
1048                        }
1049                        Err(error) => {
1050                            let GetProvidersError::Timeout { key, .. } = error;
1051
1052                            debug!(
1053                                key = hex::encode(&key),
1054                                "Get providers query failed with no results",
1055                            );
1056                        }
1057                    }
1058                }
1059
1060                if last || cancelled {
1061                    // There will be no more progress
1062                    self.query_id_receivers.remove(&id);
1063                }
1064            }
1065            KademliaEvent::OutboundQueryProgressed {
1066                step: ProgressStep { last, .. },
1067                id,
1068                result: QueryResult::PutRecord(result),
1069                ..
1070            } => {
1071                let mut cancelled = false;
1072                if let Some(QueryResultSender::PutValue { sender, .. }) =
1073                    self.query_id_receivers.get(&id)
1074                {
1075                    match result {
1076                        Ok(PutRecordOk { key }) => {
1077                            trace!("Put record query for {} succeeded", hex::encode(&key));
1078
1079                            cancelled = Self::unbounded_send_and_cancel_on_error(
1080                                &mut self.swarm.behaviour_mut().kademlia,
1081                                sender,
1082                                (),
1083                                "PutRecordOk",
1084                                id,
1085                            ) || cancelled;
1086                        }
1087                        Err(error) => {
1088                            debug!(?error, "Put record query failed.");
1089                        }
1090                    }
1091                }
1092
1093                if last || cancelled {
1094                    // There will be no more progress
1095                    self.query_id_receivers.remove(&id);
1096                }
1097            }
1098            KademliaEvent::OutboundQueryProgressed {
1099                step: ProgressStep { last, count },
1100                id,
1101                result: QueryResult::Bootstrap(result),
1102                stats,
1103            } => {
1104                debug!(?stats, %last, %count, ?id, ?result, "Bootstrap OutboundQueryProgressed step.");
1105
1106                let mut cancelled = false;
1107                if let Some(QueryResultSender::Bootstrap { sender }) =
1108                    self.query_id_receivers.get_mut(&id)
1109                {
1110                    match result {
1111                        Ok(BootstrapOk {
1112                            peer,
1113                            num_remaining,
1114                        }) => {
1115                            trace!(%peer, %num_remaining, %last, "Bootstrap query step succeeded");
1116
1117                            cancelled = Self::unbounded_send_and_cancel_on_error(
1118                                &mut self.swarm.behaviour_mut().kademlia,
1119                                sender,
1120                                (),
1121                                "Bootstrap",
1122                                id,
1123                            ) || cancelled;
1124                        }
1125                        Err(error) => {
1126                            debug!(?error, "Bootstrap query failed.");
1127                        }
1128                    }
1129                }
1130
1131                if last || cancelled {
1132                    // There will be no more progress
1133                    self.query_id_receivers.remove(&id);
1134                }
1135            }
1136            _ => {}
1137        }
1138    }
1139
1140    // Returns `true` if query was cancelled
1141    fn unbounded_send_and_cancel_on_error<T>(
1142        kademlia: &mut Kademlia<DummyRecordStore>,
1143        sender: &mpsc::UnboundedSender<T>,
1144        value: T,
1145        channel: &'static str,
1146        id: QueryId,
1147    ) -> bool {
1148        if sender.unbounded_send(value).is_err() {
1149            debug!("{} channel was dropped", channel);
1150
1151            // Cancel query
1152            if let Some(mut query) = kademlia.query_mut(&id) {
1153                query.finish();
1154            }
1155            true
1156        } else {
1157            false
1158        }
1159    }
1160
1161    fn handle_gossipsub_event(&mut self, event: GossipsubEvent) {
1162        if let GossipsubEvent::Message { message, .. } = event
1163            && let Some(senders) = self.topic_subscription_senders.get(&message.topic)
1164        {
1165            let bytes = Bytes::from(message.data);
1166
1167            for sender in senders.values() {
1168                // Doesn't matter if receiver is still listening for messages or not.
1169                let _: Result<(), _> = sender.unbounded_send(bytes.clone());
1170            }
1171        }
1172    }
1173
1174    fn handle_request_response_event(&mut self, event: RequestResponseEvent) {
1175        // No actions on statistics events.
1176        trace!("Request response event: {:?}", event);
1177    }
1178
1179    fn handle_autonat_event(&mut self, event: AutonatEvent) {
1180        trace!(?event, "Autonat event received.");
1181        let autonat = &self.swarm.behaviour().autonat;
1182        debug!(
1183            public_address=?autonat.public_address(),
1184            confidence=%autonat.confidence(),
1185            "Current public address confidence."
1186        );
1187
1188        match event {
1189            AutonatEvent::InboundProbe(_inbound_probe_event) => {
1190                // We do not care about this event
1191            }
1192            AutonatEvent::OutboundProbe(outbound_probe_event) => {
1193                match outbound_probe_event {
1194                    OutboundProbeEvent::Request { peer, .. } => {
1195                        // For outbound probe request add peer to allow list to ensure they can dial
1196                        // us back and not hit global incoming connection limit
1197                        self.swarm
1198                            .behaviour_mut()
1199                            .connection_limits
1200                            // We expect a single successful dial from this peer
1201                            .add_to_incoming_allow_list(
1202                                peer,
1203                                self.peer_ip_addresses
1204                                    .get(&peer)
1205                                    .iter()
1206                                    .flat_map(|ip_addresses| ip_addresses.iter())
1207                                    .copied(),
1208                                1,
1209                            );
1210                    }
1211                    OutboundProbeEvent::Response { peer, .. } => {
1212                        self.swarm
1213                            .behaviour_mut()
1214                            .connection_limits
1215                            .remove_from_incoming_allow_list(&peer, Some(1));
1216                    }
1217                    OutboundProbeEvent::Error { peer, .. } => {
1218                        if let Some(peer) = peer {
1219                            self.swarm
1220                                .behaviour_mut()
1221                                .connection_limits
1222                                .remove_from_incoming_allow_list(&peer, Some(1));
1223                        }
1224                    }
1225                }
1226            }
1227            AutonatEvent::StatusChanged { old, new } => {
1228                debug!(?old, ?new, "Public address status changed.");
1229
1230                // TODO: Remove block once https://github.com/libp2p/rust-libp2p/issues/4863 is resolved
1231                if let (NatStatus::Public(old_address), NatStatus::Private) = (old, new.clone()) {
1232                    self.swarm.remove_external_address(&old_address);
1233                    debug!(
1234                        ?old_address,
1235                        new_status = ?new,
1236                        "Removing old external address...",
1237                    );
1238
1239                    // Trigger potential mode change manually
1240                    self.swarm.behaviour_mut().kademlia.set_mode(None);
1241                }
1242
1243                let connected_peers = self.swarm.connected_peers().copied().collect::<Vec<_>>();
1244                self.swarm.behaviour_mut().identify.push(connected_peers);
1245            }
1246        }
1247    }
1248
1249    fn handle_command(&mut self, command: Command) {
1250        match command {
1251            Command::GetValue {
1252                key,
1253                result_sender,
1254                permit,
1255            } => {
1256                let query_id = self
1257                    .swarm
1258                    .behaviour_mut()
1259                    .kademlia
1260                    .get_record(key.to_bytes().into());
1261
1262                self.query_id_receivers.insert(
1263                    query_id,
1264                    QueryResultSender::Value {
1265                        sender: result_sender,
1266                        _permit: permit,
1267                    },
1268                );
1269            }
1270            Command::PutValue {
1271                key,
1272                value,
1273                result_sender,
1274                permit,
1275            } => {
1276                let local_peer_id = *self.swarm.local_peer_id();
1277
1278                let record = Record {
1279                    key: key.into(),
1280                    value,
1281                    publisher: Some(local_peer_id),
1282                    expires: None, // No time expiration.
1283                };
1284                let query_result = self
1285                    .swarm
1286                    .behaviour_mut()
1287                    .kademlia
1288                    .put_record(record, Quorum::One);
1289
1290                match query_result {
1291                    Ok(query_id) => {
1292                        self.query_id_receivers.insert(
1293                            query_id,
1294                            QueryResultSender::PutValue {
1295                                sender: result_sender,
1296                                _permit: permit,
1297                            },
1298                        );
1299                    }
1300                    Err(err) => {
1301                        warn!(?err, "Failed to put value.");
1302                    }
1303                }
1304            }
1305            Command::Subscribe {
1306                topic,
1307                result_sender,
1308            } => {
1309                assert!(
1310                    self.swarm.behaviour().gossipsub.is_enabled(),
1311                    "Gossipsub protocol is disabled."
1312                );
1313
1314                let topic_hash = topic.hash();
1315                let (sender, receiver) = mpsc::unbounded();
1316
1317                // Unconditionally create subscription ID, code is simpler this way.
1318                let subscription_id = self.next_subscription_id;
1319                self.next_subscription_id += 1;
1320
1321                let created_subscription = CreatedSubscription {
1322                    subscription_id,
1323                    receiver,
1324                };
1325
1326                match self.topic_subscription_senders.entry(topic_hash) {
1327                    Entry::Occupied(mut entry) => {
1328                        // In case subscription already exists, just add one more sender to it.
1329                        if result_sender.send(Ok(created_subscription)).is_ok() {
1330                            entry.get_mut().insert(subscription_id, sender);
1331                        }
1332                    }
1333                    Entry::Vacant(entry) => {
1334                        // Otherwise subscription needs to be created.
1335
1336                        if let Some(gossipsub) = self.swarm.behaviour_mut().gossipsub.as_mut() {
1337                            match gossipsub.subscribe(&topic) {
1338                                Ok(true) => {
1339                                    if result_sender.send(Ok(created_subscription)).is_ok() {
1340                                        entry
1341                                            .insert(IntMap::from_iter([(subscription_id, sender)]));
1342                                    }
1343                                }
1344                                Ok(false) => {
1345                                    panic!(
1346                                        "Logic error, topic subscription wasn't created, this \
1347                                        must never happen"
1348                                    );
1349                                }
1350                                Err(error) => {
1351                                    let _: Result<(), _> = result_sender.send(Err(error));
1352                                }
1353                            }
1354                        }
1355                    }
1356                }
1357            }
1358            Command::Unsubscribe {
1359                topic,
1360                subscription_id,
1361            } => {
1362                assert!(
1363                    self.swarm.behaviour().gossipsub.is_enabled(),
1364                    "Gossipsub protocol is disabled."
1365                );
1366
1367                if let Entry::Occupied(mut entry) =
1368                    self.topic_subscription_senders.entry(topic.hash())
1369                {
1370                    entry.get_mut().remove(&subscription_id);
1371
1372                    // If last sender was removed - unsubscribe.
1373                    if entry.get().is_empty() {
1374                        entry.remove_entry();
1375
1376                        if let Some(gossipsub) = self.swarm.behaviour_mut().gossipsub.as_mut()
1377                            && !gossipsub.unsubscribe(&topic)
1378                        {
1379                            warn!(
1380                                "Can't unsubscribe from topic {topic} because subscription doesn't \
1381                                exist, this is a logic error in the subspace or swarm libraries"
1382                            );
1383                        }
1384                    }
1385                } else {
1386                    error!(
1387                        "Can't unsubscribe from topic {topic} because subscription doesn't exist, \
1388                        this is a logic error in the subspace library"
1389                    );
1390                }
1391            }
1392            Command::Publish {
1393                topic,
1394                message,
1395                result_sender,
1396            } => {
1397                assert!(
1398                    self.swarm.behaviour().gossipsub.is_enabled(),
1399                    "Gossipsub protocol is disabled"
1400                );
1401
1402                if let Some(gossipsub) = self.swarm.behaviour_mut().gossipsub.as_mut() {
1403                    // Doesn't matter if receiver still waits for response.
1404                    let _: Result<(), _> =
1405                        result_sender.send(gossipsub.publish(topic, message).map(|_message_id| ()));
1406                }
1407            }
1408            Command::GetClosestPeers {
1409                key,
1410                result_sender,
1411                permit,
1412            } => {
1413                let query_id = self.swarm.behaviour_mut().kademlia.get_closest_peers(key);
1414
1415                self.query_id_receivers.insert(
1416                    query_id,
1417                    QueryResultSender::ClosestPeers {
1418                        sender: result_sender,
1419                        _permit: permit,
1420                    },
1421                );
1422            }
1423            Command::GetClosestLocalPeers {
1424                key,
1425                source,
1426                result_sender,
1427            } => {
1428                let source = source.unwrap_or_else(|| *self.swarm.local_peer_id());
1429                let result = self
1430                    .swarm
1431                    .behaviour_mut()
1432                    .kademlia
1433                    .find_closest_local_peers(&KBucketKey::from(key), &source)
1434                    .filter(|peer| !peer.multiaddrs.is_empty())
1435                    .map(|peer| (peer.node_id, peer.multiaddrs))
1436                    .collect();
1437
1438                // Doesn't matter if receiver still waits for response.
1439                let _: Result<(), _> = result_sender.send(result);
1440            }
1441            Command::GenericRequest {
1442                peer_id,
1443                addresses,
1444                protocol_name,
1445                request,
1446                result_sender,
1447            } => {
1448                self.swarm.behaviour_mut().request_response.send_request(
1449                    &peer_id,
1450                    protocol_name,
1451                    request,
1452                    result_sender,
1453                    IfDisconnected::TryConnect,
1454                    addresses,
1455                );
1456            }
1457            Command::GetProviders {
1458                key,
1459                result_sender,
1460                permit,
1461            } => {
1462                let query_id = self
1463                    .swarm
1464                    .behaviour_mut()
1465                    .kademlia
1466                    .get_providers(key.clone());
1467
1468                self.query_id_receivers.insert(
1469                    query_id,
1470                    QueryResultSender::Providers {
1471                        key,
1472                        sender: result_sender,
1473                        _permit: permit,
1474                    },
1475                );
1476            }
1477            Command::BanPeer { peer_id } => {
1478                self.ban_peer(peer_id);
1479            }
1480            Command::Dial { address } => {
1481                let _: Result<(), _> = self.swarm.dial(address);
1482            }
1483            Command::ConnectedPeers { result_sender } => {
1484                let connected_peers = self.swarm.connected_peers().copied().collect();
1485
1486                let _: Result<(), _> = result_sender.send(connected_peers);
1487            }
1488            Command::ConnectedServers { result_sender } => {
1489                let connected_servers = self.connected_servers.iter().copied().collect();
1490
1491                let _: Result<(), _> = result_sender.send(connected_servers);
1492            }
1493            Command::Bootstrap { result_sender } => {
1494                let kademlia = &mut self.swarm.behaviour_mut().kademlia;
1495
1496                match kademlia.bootstrap() {
1497                    Ok(query_id) => {
1498                        if let Some(result_sender) = result_sender {
1499                            self.query_id_receivers.insert(
1500                                query_id,
1501                                QueryResultSender::Bootstrap {
1502                                    sender: result_sender,
1503                                },
1504                            );
1505                        }
1506                    }
1507                    Err(err) => {
1508                        debug!(?err, "Bootstrap error.");
1509                    }
1510                }
1511            }
1512        }
1513    }
1514
1515    fn ban_peer(&mut self, peer_id: PeerId) {
1516        // Remove temporary ban if there is one, before creating a permanent one.
1517        self.temporary_bans.lock().remove(&peer_id);
1518
1519        debug!(?peer_id, "Banning peer on network level");
1520
1521        self.swarm.behaviour_mut().block_list.block_peer(peer_id);
1522        self.swarm.behaviour_mut().kademlia.remove_peer(&peer_id);
1523        self.known_peers_registry
1524            .remove_all_known_peer_addresses(peer_id);
1525
1526        // Immediately disconnect the peer to cancel any in-flight requests.
1527        let _: Result<(), ()> = self.swarm.disconnect_peer_id(peer_id);
1528    }
1529
1530    fn register_event_metrics(&mut self, swarm_event: &SwarmEvent<Event>) {
1531        if let Some(ref mut metrics) = self.libp2p_metrics {
1532            match swarm_event {
1533                SwarmEvent::Behaviour(Event::Ping(ping_event)) => {
1534                    metrics.record(ping_event);
1535                }
1536                SwarmEvent::Behaviour(Event::Identify(identify_event)) => {
1537                    metrics.record(identify_event.as_ref());
1538                }
1539                SwarmEvent::Behaviour(Event::Kademlia(kademlia_event)) => {
1540                    metrics.record(kademlia_event);
1541                }
1542                SwarmEvent::Behaviour(Event::Gossipsub(gossipsub_event)) => {
1543                    metrics.record(gossipsub_event);
1544                }
1545                // TODO: implement in the upstream repository
1546                // SwarmEvent::Behaviour(Event::RequestResponse(request_response_event)) => {
1547                //     self.metrics.record(request_response_event);
1548                // }
1549                swarm_event => {
1550                    metrics.record(swarm_event);
1551                }
1552            }
1553        }
1554    }
1555
1556    fn log_kademlia_stats(&mut self) {
1557        let mut peer_counter = 0usize;
1558        let mut peer_with_no_address_counter = 0usize;
1559        for kbucket in self.swarm.behaviour_mut().kademlia.kbuckets() {
1560            for entry in kbucket.iter() {
1561                peer_counter += 1;
1562                if entry.node.value.len() == 0 {
1563                    peer_with_no_address_counter += 1;
1564                }
1565            }
1566        }
1567
1568        debug!(
1569            peers = %peer_counter,
1570            peers_with_no_address = %peer_with_no_address_counter,
1571            "Kademlia stats"
1572        );
1573    }
1574}