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 _permit: OwnedSemaphorePermit,
58 },
59 ClosestPeers {
60 sender: mpsc::UnboundedSender<PeerId>,
61 _permit: Option<OwnedSemaphorePermit>,
63 },
64 Providers {
65 key: RecordKey,
66 sender: mpsc::UnboundedSender<PeerId>,
67 _permit: Option<OwnedSemaphorePermit>,
69 },
70 PutValue {
71 sender: mpsc::UnboundedSender<()>,
72 _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#[must_use = "Node does not function properly unless its runner is driven forward"]
90pub struct NodeRunner {
91 allow_non_global_addresses_in_dht: bool,
93 is_listening: bool,
95 command_receiver: mpsc::Receiver<Command>,
96 swarm: Swarm<Behavior>,
97 shared_weak: Weak<Shared>,
98 next_random_query_interval: Duration,
100 query_id_receivers: HashMap<QueryId, QueryResultSender>,
101 next_subscription_id: usize,
104 topic_subscription_senders: HashMap<TopicHash, IntMap<usize, mpsc::UnboundedSender<Bytes>>>,
107 random_query_timeout: Pin<Box<Fuse<Sleep>>>,
108 periodical_tasks_interval: Pin<Box<Fuse<Sleep>>>,
110 known_peers_registry: Box<dyn KnownPeersRegistry>,
112 connected_servers: HashSet<PeerId>,
113 reserved_peers: HashMap<PeerId, Multiaddr>,
115 temporary_bans: Arc<Mutex<TemporaryBans>>,
117 libp2p_metrics: Option<Metrics>,
119 metrics: Option<SubspaceMetrics>,
121 peer_ip_addresses: HashMap<PeerId, HashSet<IpAddr>>,
123 protocol_version: String,
125 bootstrap_addresses: Vec<Multiaddr>,
127 bootstrap_command_state: Arc<AsyncMutex<BootstrapCommandState>>,
129 removed_addresses_rx: mpsc::UnboundedReceiver<PeerAddressRemovedEvent>,
131 _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
143pub(crate) struct NodeRunnerConfig {
145 pub(crate) allow_non_global_addresses_in_dht: bool,
146 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 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 random_query_timeout: Box::pin(tokio::time::sleep(Duration::from_secs(0)).fuse()),
204 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 pub async fn run(&mut self) {
223 if self.is_listening {
224 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 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 yield_now().await;
283 }
284 }
285
286 async fn bootstrap(&mut self) {
288 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 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 fn handle_periodical_tasks(&mut self) {
388 let network_info = self.swarm.network_info();
390 let connections = network_info.connection_counters();
391
392 debug!(?connections, "Current connections and limits.");
393
394 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 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 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 if let ConnectedPoint::Dialer { address, .. } = &endpoint {
504 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 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 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 if true {
669 return false;
670 }
671
672 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 if self.temporary_bans.lock().is_banned(peer_id) {
687 return false;
688 }
689 }
690 }
691 }
692 true
694 }
695 DialError::LocalPeerId { .. } => {
696 debug!("Local peer dial attempt detected.");
698
699 false
700 }
701 DialError::NoAddresses => {
702 true
704 }
705 DialError::DialPeerConditionFalse(_) => {
706 false
708 }
709 DialError::Aborted => {
710 false
712 }
713 DialError::WrongPeerId { .. } => {
714 false
716 }
717 DialError::Denied { .. } => {
718 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 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 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 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 && 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 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 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 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 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 self.query_id_receivers.remove(&id);
1134 }
1135 }
1136 _ => {}
1137 }
1138 }
1139
1140 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 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 let _: Result<(), _> = sender.unbounded_send(bytes.clone());
1170 }
1171 }
1172 }
1173
1174 fn handle_request_response_event(&mut self, event: RequestResponseEvent) {
1175 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 }
1192 AutonatEvent::OutboundProbe(outbound_probe_event) => {
1193 match outbound_probe_event {
1194 OutboundProbeEvent::Request { peer, .. } => {
1195 self.swarm
1198 .behaviour_mut()
1199 .connection_limits
1200 .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 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 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, };
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 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 if result_sender.send(Ok(created_subscription)).is_ok() {
1330 entry.get_mut().insert(subscription_id, sender);
1331 }
1332 }
1333 Entry::Vacant(entry) => {
1334 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 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 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 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 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 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 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}