1use crate::{
31 behaviour::{self, Behaviour, BehaviourOut},
32 config::{
33 parse_addr, FullNetworkConfiguration, IncomingRequest, MultiaddrWithPeerId,
34 NonDefaultSetConfig, NotificationHandshake, Params, SetConfig, TransportConfig,
35 },
36 discovery::DiscoveryConfig,
37 error::Error,
38 event::{DhtEvent, Event},
39 network_state::{
40 NetworkState, NotConnectedPeer as NetworkStateNotConnectedPeer, Peer as NetworkStatePeer,
41 },
42 peer_store::{PeerStore, PeerStoreProvider},
43 protocol::{self, Protocol, Ready},
44 protocol_controller::{self, ProtoSetConfig, ProtocolController, SetId},
45 request_responses::{IfDisconnected, ProtocolConfig as RequestResponseConfig, RequestFailure},
46 service::{
47 signature::{Signature, SigningError},
48 traits::{
49 BandwidthSink, NetworkBackend, NetworkDHTProvider, NetworkEventStream, NetworkPeers,
50 NetworkRequest, NetworkService as NetworkServiceT, NetworkSigner, NetworkStateInfo,
51 NetworkStatus, NetworkStatusProvider, NotificationSender as NotificationSenderT,
52 NotificationSenderError, NotificationSenderReady as NotificationSenderReadyT,
53 },
54 },
55 transport,
56 types::ProtocolName,
57 NotificationService, ReputationChange,
58};
59
60use codec::DecodeAll;
61use futures::{channel::oneshot, prelude::*};
62use libp2p::{
63 connection_limits::{ConnectionLimits, Exceeded},
64 core::{upgrade, ConnectedPoint, Endpoint},
65 identify::Info as IdentifyInfo,
66 identity::ed25519,
67 multiaddr::{self, Multiaddr},
68 swarm::{
69 Config as SwarmConfig, ConnectionError, ConnectionId, DialError, Executor, ListenError,
70 NetworkBehaviour, Swarm, SwarmEvent,
71 },
72 PeerId,
73};
74use log::{debug, error, info, trace, warn};
75use metrics::{Histogram, MetricSources, Metrics};
76use parking_lot::Mutex;
77use prometheus_endpoint::Registry;
78use sc_network_types::kad::{Key as KademliaKey, Record};
79
80use sc_network_common::{
81 role::{ObservedRole, Roles},
82 ExHashT,
83};
84use sc_utils::mpsc::{tracing_unbounded, TracingUnboundedReceiver, TracingUnboundedSender};
85use sp_runtime::traits::Block as BlockT;
86
87pub use behaviour::{InboundFailure, OutboundFailure, ResponseFailure};
88pub use libp2p::identity::{DecodingError, Keypair, PublicKey};
89pub use metrics::NotificationMetrics;
90pub use protocol::NotificationsSink;
91use std::{
92 collections::{HashMap, HashSet},
93 fs, iter,
94 marker::PhantomData,
95 num::NonZeroUsize,
96 pin::Pin,
97 str,
98 sync::{
99 atomic::{AtomicUsize, Ordering},
100 Arc,
101 },
102 time::{Duration, Instant},
103};
104
105pub(crate) mod metrics;
106pub(crate) mod out_events;
107
108pub mod signature;
109pub mod traits;
110
111const LOG_TARGET: &str = "sub-libp2p";
113
114struct Libp2pBandwidthSink {
115 #[allow(deprecated)]
116 sink: Arc<transport::BandwidthSinks>,
117}
118
119impl BandwidthSink for Libp2pBandwidthSink {
120 fn total_inbound(&self) -> u64 {
121 self.sink.total_inbound()
122 }
123
124 fn total_outbound(&self) -> u64 {
125 self.sink.total_outbound()
126 }
127}
128
129pub struct NetworkService<B: BlockT + 'static, H: ExHashT> {
131 num_connected: Arc<AtomicUsize>,
133 external_addresses: Arc<Mutex<HashSet<Multiaddr>>>,
135 listen_addresses: Arc<Mutex<HashSet<Multiaddr>>>,
137 local_peer_id: PeerId,
139 local_identity: Keypair,
141 bandwidth: Arc<dyn BandwidthSink>,
143 to_worker: TracingUnboundedSender<ServiceToWorkerMsg>,
145 notification_protocol_ids: HashMap<ProtocolName, SetId>,
148 protocol_handles: Vec<protocol_controller::ProtocolHandle>,
151 sync_protocol_handle: protocol_controller::ProtocolHandle,
153 peer_store_handle: Arc<dyn PeerStoreProvider>,
155 _marker: PhantomData<H>,
158 _block: PhantomData<B>,
160}
161
162#[async_trait::async_trait]
163impl<B, H> NetworkBackend<B, H> for NetworkWorker<B, H>
164where
165 B: BlockT + 'static,
166 H: ExHashT,
167{
168 type NotificationProtocolConfig = NonDefaultSetConfig;
169 type RequestResponseProtocolConfig = RequestResponseConfig;
170 type NetworkService<Block, Hash> = Arc<NetworkService<B, H>>;
171 type PeerStore = PeerStore;
172
173 fn new(params: Params<B, H, Self>) -> Result<Self, Error>
174 where
175 Self: Sized,
176 {
177 NetworkWorker::new(params)
178 }
179
180 fn network_service(&self) -> Arc<dyn NetworkServiceT> {
182 self.service.clone()
183 }
184
185 fn peer_store(
187 bootnodes: Vec<sc_network_types::PeerId>,
188 metrics_registry: Option<Registry>,
189 ) -> Self::PeerStore {
190 PeerStore::new(bootnodes.into_iter().map(From::from).collect(), metrics_registry)
191 }
192
193 fn register_notification_metrics(registry: Option<&Registry>) -> NotificationMetrics {
194 NotificationMetrics::new(registry)
195 }
196
197 fn notification_config(
199 protocol_name: ProtocolName,
200 fallback_names: Vec<ProtocolName>,
201 max_notification_size: u64,
202 handshake: Option<NotificationHandshake>,
203 set_config: SetConfig,
204 _metrics: NotificationMetrics,
205 _peerstore_handle: Arc<dyn PeerStoreProvider>,
206 ) -> (Self::NotificationProtocolConfig, Box<dyn NotificationService>) {
207 NonDefaultSetConfig::new(
208 protocol_name,
209 fallback_names,
210 max_notification_size,
211 handshake,
212 set_config,
213 )
214 }
215
216 fn request_response_config(
218 protocol_name: ProtocolName,
219 fallback_names: Vec<ProtocolName>,
220 max_request_size: u64,
221 max_response_size: u64,
222 request_timeout: Duration,
223 inbound_queue: Option<async_channel::Sender<IncomingRequest>>,
224 ) -> Self::RequestResponseProtocolConfig {
225 Self::RequestResponseProtocolConfig {
226 name: protocol_name,
227 fallback_names,
228 max_request_size,
229 max_response_size,
230 request_timeout,
231 inbound_queue,
232 }
233 }
234
235 async fn run(mut self) {
237 self.run().await
238 }
239}
240
241impl<B, H> NetworkWorker<B, H>
242where
243 B: BlockT + 'static,
244 H: ExHashT,
245{
246 pub fn new(params: Params<B, H, Self>) -> Result<Self, Error> {
252 let peer_store_handle = params.network_config.peer_store_handle();
253 let FullNetworkConfiguration {
254 notification_protocols,
255 request_response_protocols,
256 mut network_config,
257 ..
258 } = params.network_config;
259
260 let local_identity = network_config.node_key.clone().into_keypair()?;
262 let local_public = local_identity.public();
263 let local_peer_id = local_public.to_peer_id();
264
265 let local_identity: ed25519::Keypair = local_identity.into();
267 let local_public: ed25519::PublicKey = local_public.into();
268 let local_peer_id: PeerId = local_peer_id.into();
269
270 network_config.boot_nodes = network_config
271 .boot_nodes
272 .into_iter()
273 .filter(|boot_node| boot_node.peer_id != local_peer_id.into())
274 .collect();
275 network_config.default_peers_set.reserved_nodes = network_config
276 .default_peers_set
277 .reserved_nodes
278 .into_iter()
279 .filter(|reserved_node| {
280 if reserved_node.peer_id == local_peer_id.into() {
281 warn!(
282 target: LOG_TARGET,
283 "Local peer ID used in reserved node, ignoring: {}",
284 reserved_node,
285 );
286 false
287 } else {
288 true
289 }
290 })
291 .collect();
292
293 ensure_addresses_consistent_with_transport(
295 network_config.listen_addresses.iter(),
296 &network_config.transport,
297 )?;
298 ensure_addresses_consistent_with_transport(
299 network_config.boot_nodes.iter().map(|x| &x.multiaddr),
300 &network_config.transport,
301 )?;
302 ensure_addresses_consistent_with_transport(
303 network_config.default_peers_set.reserved_nodes.iter().map(|x| &x.multiaddr),
304 &network_config.transport,
305 )?;
306 for notification_protocol in ¬ification_protocols {
307 ensure_addresses_consistent_with_transport(
308 notification_protocol.set_config().reserved_nodes.iter().map(|x| &x.multiaddr),
309 &network_config.transport,
310 )?;
311 }
312 ensure_addresses_consistent_with_transport(
313 network_config.public_addresses.iter(),
314 &network_config.transport,
315 )?;
316
317 let (to_worker, from_service) = tracing_unbounded("mpsc_network_worker", 100_000);
318
319 if let Some(path) = &network_config.net_config_path {
320 fs::create_dir_all(path)?;
321 }
322
323 info!(
324 target: LOG_TARGET,
325 "๐ท Local node identity is: {}",
326 local_peer_id.to_base58(),
327 );
328 info!(target: LOG_TARGET, "Running libp2p network backend");
329
330 let (transport, bandwidth) = {
331 let config_mem = match network_config.transport {
332 TransportConfig::MemoryOnly => true,
333 TransportConfig::Normal { .. } => false,
334 };
335
336 transport::build_transport(local_identity.clone().into(), config_mem)
337 };
338
339 let (to_notifications, from_protocol_controllers) =
340 tracing_unbounded("mpsc_protocol_controllers_to_notifications", 10_000);
341
342 let all_peer_sets_iter = iter::once(&network_config.default_peers_set)
344 .chain(notification_protocols.iter().map(|protocol| protocol.set_config()));
345
346 let (protocol_handles, protocol_controllers): (Vec<_>, Vec<_>) = all_peer_sets_iter
347 .enumerate()
348 .map(|(set_id, set_config)| {
349 let proto_set_config = ProtoSetConfig {
350 in_peers: set_config.in_peers,
351 out_peers: set_config.out_peers,
352 reserved_nodes: set_config
353 .reserved_nodes
354 .iter()
355 .map(|node| node.peer_id.into())
356 .collect(),
357 reserved_only: set_config.non_reserved_mode.is_reserved_only(),
358 };
359
360 ProtocolController::new(
361 SetId::from(set_id),
362 proto_set_config,
363 to_notifications.clone(),
364 Arc::clone(&peer_store_handle),
365 )
366 })
367 .unzip();
368
369 let sync_protocol_handle = protocol_handles[0].clone();
371
372 protocol_controllers
374 .into_iter()
375 .for_each(|controller| (params.executor)(controller.run().boxed()));
376
377 let notification_protocol_ids: HashMap<ProtocolName, SetId> =
380 iter::once(¶ms.block_announce_config)
381 .chain(notification_protocols.iter())
382 .enumerate()
383 .map(|(index, protocol)| (protocol.protocol_name().clone(), SetId::from(index)))
384 .collect();
385
386 let known_addresses = {
387 let mut addresses: Vec<_> = network_config
389 .default_peers_set
390 .reserved_nodes
391 .iter()
392 .map(|reserved| (reserved.peer_id, reserved.multiaddr.clone()))
393 .chain(notification_protocols.iter().flat_map(|protocol| {
394 protocol
395 .set_config()
396 .reserved_nodes
397 .iter()
398 .map(|reserved| (reserved.peer_id, reserved.multiaddr.clone()))
399 }))
400 .chain(
401 network_config
402 .boot_nodes
403 .iter()
404 .map(|bootnode| (bootnode.peer_id, bootnode.multiaddr.clone())),
405 )
406 .collect();
407
408 addresses.sort();
410 addresses.dedup();
411
412 addresses
413 };
414
415 network_config.boot_nodes.iter().try_for_each(|bootnode| {
417 if let Some(other) = network_config
418 .boot_nodes
419 .iter()
420 .filter(|o| o.multiaddr == bootnode.multiaddr)
421 .find(|o| o.peer_id != bootnode.peer_id)
422 {
423 Err(Error::DuplicateBootnode {
424 address: bootnode.multiaddr.clone().into(),
425 first_id: bootnode.peer_id.into(),
426 second_id: other.peer_id.into(),
427 })
428 } else {
429 Ok(())
430 }
431 })?;
432
433 let mut boot_node_ids = HashMap::<PeerId, Vec<Multiaddr>>::new();
435
436 for bootnode in network_config.boot_nodes.iter() {
437 boot_node_ids
438 .entry(bootnode.peer_id.into())
439 .or_default()
440 .push(bootnode.multiaddr.clone().into());
441 }
442
443 let boot_node_ids = Arc::new(boot_node_ids);
444
445 let num_connected = Arc::new(AtomicUsize::new(0));
446 let external_addresses = Arc::new(Mutex::new(HashSet::new()));
447
448 let (protocol, notif_protocol_handles) = Protocol::new(
449 From::from(¶ms.role),
450 params.notification_metrics,
451 notification_protocols,
452 params.block_announce_config,
453 Arc::clone(&peer_store_handle),
454 protocol_handles.clone(),
455 from_protocol_controllers,
456 )?;
457
458 let (mut swarm, bandwidth): (Swarm<Behaviour<B>>, _) = {
460 let user_agent =
461 format!("{} ({})", network_config.client_version, network_config.node_name);
462
463 let discovery_config = {
464 let mut config = DiscoveryConfig::new(local_peer_id);
465 config.with_permanent_addresses(
466 known_addresses
467 .iter()
468 .map(|(peer, address)| (peer.into(), address.clone().into()))
469 .collect::<Vec<_>>(),
470 );
471 config.discovery_limit(u64::from(network_config.default_peers_set.out_peers) + 15);
472 config.with_kademlia(
473 params.genesis_hash,
474 params.fork_id.as_deref(),
475 ¶ms.protocol_id,
476 );
477 config.with_dht_random_walk(network_config.enable_dht_random_walk);
478 config.allow_non_globals_in_dht(network_config.allow_non_globals_in_dht);
479 config.use_kademlia_disjoint_query_paths(
480 network_config.kademlia_disjoint_query_paths,
481 );
482 config.with_kademlia_replication_factor(network_config.kademlia_replication_factor);
483
484 match network_config.transport {
485 TransportConfig::MemoryOnly => {
486 config.with_mdns(false);
487 config.allow_private_ip(false);
488 },
489 TransportConfig::Normal {
490 enable_mdns,
491 allow_private_ip: allow_private_ipv4,
492 ..
493 } => {
494 config.with_mdns(enable_mdns);
495 config.allow_private_ip(allow_private_ipv4);
496 },
497 }
498
499 config
500 };
501
502 let behaviour = {
503 let result = Behaviour::new(
504 protocol,
505 user_agent,
506 local_public.into(),
507 discovery_config,
508 request_response_protocols,
509 Arc::clone(&peer_store_handle),
510 external_addresses.clone(),
511 network_config.public_addresses.iter().cloned().map(Into::into).collect(),
512 ConnectionLimits::default()
513 .with_max_established_per_peer(Some(crate::MAX_CONNECTIONS_PER_PEER as u32))
514 .with_max_established_incoming(Some(
515 crate::MAX_CONNECTIONS_ESTABLISHED_INCOMING,
516 )),
517 );
518
519 match result {
520 Ok(b) => b,
521 Err(crate::request_responses::RegisterError::DuplicateProtocol(proto)) => {
522 return Err(Error::DuplicateRequestResponseProtocol { protocol: proto })
523 },
524 }
525 };
526
527 let swarm = {
528 struct SpawnImpl<F>(F);
529 impl<F: Fn(Pin<Box<dyn Future<Output = ()> + Send>>)> Executor for SpawnImpl<F> {
530 fn exec(&self, f: Pin<Box<dyn Future<Output = ()> + Send>>) {
531 (self.0)(f)
532 }
533 }
534
535 let config = SwarmConfig::with_executor(SpawnImpl(params.executor))
536 .with_substream_upgrade_protocol_override(upgrade::Version::V1)
537 .with_notify_handler_buffer_size(NonZeroUsize::new(32).expect("32 != 0; qed"))
538 .with_per_connection_event_buffer_size(24)
541 .with_max_negotiating_inbound_streams(2048)
542 .with_idle_connection_timeout(network_config.idle_connection_timeout);
543
544 Swarm::new(transport, behaviour, local_peer_id, config)
545 };
546
547 (swarm, Arc::new(Libp2pBandwidthSink { sink: bandwidth }))
548 };
549
550 let metrics = match ¶ms.metrics_registry {
552 Some(registry) => Some(metrics::register(
553 registry,
554 MetricSources {
555 bandwidth: bandwidth.clone(),
556 connected_peers: num_connected.clone(),
557 },
558 )?),
559 None => None,
560 };
561
562 for addr in &network_config.listen_addresses {
564 if let Err(err) = Swarm::<Behaviour<B>>::listen_on(&mut swarm, addr.clone().into()) {
565 warn!(target: LOG_TARGET, "Can't listen on {} because: {:?}", addr, err)
566 }
567 }
568
569 for addr in &network_config.public_addresses {
571 Swarm::<Behaviour<B>>::add_external_address(&mut swarm, addr.clone().into());
572 }
573
574 let listen_addresses_set = Arc::new(Mutex::new(HashSet::new()));
575
576 let service = Arc::new(NetworkService {
577 bandwidth,
578 external_addresses,
579 listen_addresses: listen_addresses_set.clone(),
580 num_connected: num_connected.clone(),
581 local_peer_id,
582 local_identity: local_identity.into(),
583 to_worker,
584 notification_protocol_ids,
585 protocol_handles,
586 sync_protocol_handle,
587 peer_store_handle: Arc::clone(&peer_store_handle),
588 _marker: PhantomData,
589 _block: Default::default(),
590 });
591
592 Ok(NetworkWorker {
593 listen_addresses: listen_addresses_set,
594 num_connected,
595 network_service: swarm,
596 service,
597 from_service,
598 event_streams: out_events::OutChannels::new(params.metrics_registry.as_ref())?,
599 metrics,
600 boot_node_ids,
601 reported_invalid_boot_nodes: Default::default(),
602 peer_store_handle: Arc::clone(&peer_store_handle),
603 notif_protocol_handles,
604 _marker: Default::default(),
605 _block: Default::default(),
606 })
607 }
608
609 pub fn status(&self) -> NetworkStatus {
611 NetworkStatus {
612 num_connected_peers: self.num_connected_peers(),
613 total_bytes_inbound: self.total_bytes_inbound(),
614 total_bytes_outbound: self.total_bytes_outbound(),
615 }
616 }
617
618 pub fn total_bytes_inbound(&self) -> u64 {
620 self.service.bandwidth.total_inbound()
621 }
622
623 pub fn total_bytes_outbound(&self) -> u64 {
625 self.service.bandwidth.total_outbound()
626 }
627
628 pub fn num_connected_peers(&self) -> usize {
630 self.network_service.behaviour().user_protocol().num_sync_peers()
631 }
632
633 pub fn add_known_address(&mut self, peer_id: PeerId, addr: Multiaddr) {
635 self.network_service.behaviour_mut().add_known_address(peer_id, addr);
636 }
637
638 pub fn service(&self) -> &Arc<NetworkService<B, H>> {
641 &self.service
642 }
643
644 pub fn local_peer_id(&self) -> &PeerId {
646 Swarm::<Behaviour<B>>::local_peer_id(&self.network_service)
647 }
648
649 pub fn listen_addresses(&self) -> impl Iterator<Item = &Multiaddr> {
653 Swarm::<Behaviour<B>>::listeners(&self.network_service)
654 }
655
656 pub fn network_state(&mut self) -> NetworkState {
661 let swarm = &mut self.network_service;
662 let open = swarm.behaviour_mut().user_protocol().open_peers().cloned().collect::<Vec<_>>();
663 let connected_peers = {
664 let swarm = &mut *swarm;
665 open.iter()
666 .filter_map(move |peer_id| {
667 let known_addresses = if let Ok(addrs) =
668 NetworkBehaviour::handle_pending_outbound_connection(
669 swarm.behaviour_mut(),
670 ConnectionId::new_unchecked(0), Some(*peer_id),
672 &vec![],
673 Endpoint::Listener,
674 ) {
675 addrs.into_iter().collect()
676 } else {
677 error!(target: LOG_TARGET, "Was not able to get known addresses for {:?}", peer_id);
678 return None;
679 };
680
681 let endpoint = if let Some(e) =
682 swarm.behaviour_mut().node(peer_id).and_then(|i| i.endpoint())
683 {
684 e.clone().into()
685 } else {
686 error!(target: LOG_TARGET, "Found state inconsistency between custom protocol \
687 and debug information about {:?}", peer_id);
688 return None;
689 };
690
691 Some((
692 peer_id.to_base58(),
693 NetworkStatePeer {
694 endpoint,
695 version_string: swarm
696 .behaviour_mut()
697 .node(peer_id)
698 .and_then(|i| i.client_version().map(|s| s.to_owned())),
699 latest_ping_time: swarm
700 .behaviour_mut()
701 .node(peer_id)
702 .and_then(|i| i.latest_ping()),
703 known_addresses,
704 },
705 ))
706 })
707 .collect()
708 };
709
710 let not_connected_peers = {
711 let swarm = &mut *swarm;
712 swarm
713 .behaviour_mut()
714 .known_peers()
715 .into_iter()
716 .filter(|p| open.iter().all(|n| n != p))
717 .map(move |peer_id| {
718 let known_addresses = if let Ok(addrs) =
719 NetworkBehaviour::handle_pending_outbound_connection(
720 swarm.behaviour_mut(),
721 ConnectionId::new_unchecked(0), Some(peer_id),
723 &vec![],
724 Endpoint::Listener,
725 ) {
726 addrs.into_iter().collect()
727 } else {
728 error!(target: LOG_TARGET, "Was not able to get known addresses for {:?}", peer_id);
729 Default::default()
730 };
731
732 (
733 peer_id.to_base58(),
734 NetworkStateNotConnectedPeer {
735 version_string: swarm
736 .behaviour_mut()
737 .node(&peer_id)
738 .and_then(|i| i.client_version().map(|s| s.to_owned())),
739 latest_ping_time: swarm
740 .behaviour_mut()
741 .node(&peer_id)
742 .and_then(|i| i.latest_ping()),
743 known_addresses,
744 },
745 )
746 })
747 .collect()
748 };
749
750 let peer_id = Swarm::<Behaviour<B>>::local_peer_id(swarm).to_base58();
751 let listened_addresses = swarm.listeners().cloned().collect();
752 let external_addresses = swarm.external_addresses().cloned().collect();
753
754 NetworkState {
755 peer_id,
756 listened_addresses,
757 external_addresses,
758 connected_peers,
759 not_connected_peers,
760 peerset: serde_json::json!(
763 "Unimplemented. See https://github.com/paritytech/substrate/issues/14160."
764 ),
765 }
766 }
767
768 pub fn remove_reserved_peer(&self, peer: PeerId) {
770 self.service.remove_reserved_peer(peer.into());
771 }
772
773 pub fn add_reserved_peer(&self, peer: MultiaddrWithPeerId) -> Result<(), String> {
775 self.service.add_reserved_peer(peer)
776 }
777}
778
779impl<B: BlockT + 'static, H: ExHashT> NetworkService<B, H> {
780 pub async fn network_state(&self) -> Result<NetworkState, ()> {
787 let (tx, rx) = oneshot::channel();
788
789 let _ = self
790 .to_worker
791 .unbounded_send(ServiceToWorkerMsg::NetworkState { pending_response: tx });
792
793 match rx.await {
794 Ok(v) => v.map_err(|_| ()),
795 Err(_) => Err(()),
797 }
798 }
799
800 fn split_multiaddr_and_peer_id(
805 &self,
806 peers: HashSet<Multiaddr>,
807 ) -> Result<Vec<(PeerId, Multiaddr)>, String> {
808 peers
809 .into_iter()
810 .map(|mut addr| {
811 let peer = match addr.pop() {
812 Some(multiaddr::Protocol::P2p(peer_id)) => peer_id,
813 _ => return Err("Missing PeerId from address".to_string()),
814 };
815
816 if peer == self.local_peer_id {
819 Err("Local peer ID in peer set.".to_string())
820 } else {
821 Ok((peer, addr))
822 }
823 })
824 .collect::<Result<Vec<(PeerId, Multiaddr)>, String>>()
825 }
826}
827
828impl<B, H> NetworkStateInfo for NetworkService<B, H>
829where
830 B: sp_runtime::traits::Block,
831 H: ExHashT,
832{
833 fn external_addresses(&self) -> Vec<sc_network_types::multiaddr::Multiaddr> {
835 self.external_addresses.lock().iter().cloned().map(Into::into).collect()
836 }
837
838 fn listen_addresses(&self) -> Vec<sc_network_types::multiaddr::Multiaddr> {
840 self.listen_addresses.lock().iter().cloned().map(Into::into).collect()
841 }
842
843 fn local_peer_id(&self) -> sc_network_types::PeerId {
845 self.local_peer_id.into()
846 }
847}
848
849impl<B, H> NetworkSigner for NetworkService<B, H>
850where
851 B: sp_runtime::traits::Block,
852 H: ExHashT,
853{
854 fn sign_with_local_identity(&self, msg: Vec<u8>) -> Result<Signature, SigningError> {
855 let public_key = self.local_identity.public();
856 let bytes = self.local_identity.sign(msg.as_ref())?;
857
858 Ok(Signature {
859 public_key: crate::service::signature::PublicKey::Libp2p(public_key),
860 bytes,
861 })
862 }
863
864 fn verify(
865 &self,
866 peer_id: sc_network_types::PeerId,
867 public_key: &Vec<u8>,
868 signature: &Vec<u8>,
869 message: &Vec<u8>,
870 ) -> Result<bool, String> {
871 let public_key =
872 PublicKey::try_decode_protobuf(&public_key).map_err(|error| error.to_string())?;
873 let peer_id: PeerId = peer_id.into();
874 let remote: libp2p::PeerId = public_key.to_peer_id();
875
876 Ok(peer_id == remote && public_key.verify(message, signature))
877 }
878}
879
880impl<B, H> NetworkDHTProvider for NetworkService<B, H>
881where
882 B: BlockT + 'static,
883 H: ExHashT,
884{
885 fn find_closest_peers(&self, target: sc_network_types::PeerId) {
890 let _ = self
891 .to_worker
892 .unbounded_send(ServiceToWorkerMsg::FindClosestPeers(target.into()));
893 }
894
895 fn get_value(&self, key: &KademliaKey) {
900 let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::GetValue(key.clone()));
901 }
902
903 fn put_value(&self, key: KademliaKey, value: Vec<u8>) {
908 let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::PutValue(key, value));
909 }
910
911 fn put_record_to(
912 &self,
913 record: Record,
914 peers: HashSet<sc_network_types::PeerId>,
915 update_local_storage: bool,
916 ) {
917 let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::PutRecordTo {
918 record,
919 peers,
920 update_local_storage,
921 });
922 }
923
924 fn store_record(
925 &self,
926 key: KademliaKey,
927 value: Vec<u8>,
928 publisher: Option<sc_network_types::PeerId>,
929 expires: Option<Instant>,
930 ) {
931 let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::StoreRecord(
932 key,
933 value,
934 publisher.map(Into::into),
935 expires,
936 ));
937 }
938
939 fn start_providing(&self, key: KademliaKey) {
940 let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::StartProviding(key));
941 }
942
943 fn stop_providing(&self, key: KademliaKey) {
944 let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::StopProviding(key));
945 }
946
947 fn get_providers(&self, key: KademliaKey) {
948 let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::GetProviders(key));
949 }
950}
951
952#[async_trait::async_trait]
953impl<B, H> NetworkStatusProvider for NetworkService<B, H>
954where
955 B: BlockT + 'static,
956 H: ExHashT,
957{
958 async fn status(&self) -> Result<NetworkStatus, ()> {
959 let (tx, rx) = oneshot::channel();
960
961 let _ = self
962 .to_worker
963 .unbounded_send(ServiceToWorkerMsg::NetworkStatus { pending_response: tx });
964
965 match rx.await {
966 Ok(v) => v.map_err(|_| ()),
967 Err(_) => Err(()),
969 }
970 }
971
972 async fn network_state(&self) -> Result<NetworkState, ()> {
973 let (tx, rx) = oneshot::channel();
974
975 let _ = self
976 .to_worker
977 .unbounded_send(ServiceToWorkerMsg::NetworkState { pending_response: tx });
978
979 match rx.await {
980 Ok(v) => v.map_err(|_| ()),
981 Err(_) => Err(()),
983 }
984 }
985}
986
987#[async_trait::async_trait]
988impl<B, H> NetworkPeers for NetworkService<B, H>
989where
990 B: BlockT + 'static,
991 H: ExHashT,
992{
993 fn set_authorized_peers(&self, peers: HashSet<sc_network_types::PeerId>) {
994 self.sync_protocol_handle
995 .set_reserved_peers(peers.iter().map(|peer| (*peer).into()).collect());
996 }
997
998 fn set_authorized_only(&self, reserved_only: bool) {
999 self.sync_protocol_handle.set_reserved_only(reserved_only);
1000 }
1001
1002 fn add_known_address(
1003 &self,
1004 peer_id: sc_network_types::PeerId,
1005 addr: sc_network_types::multiaddr::Multiaddr,
1006 ) {
1007 let _ = self
1008 .to_worker
1009 .unbounded_send(ServiceToWorkerMsg::AddKnownAddress(peer_id.into(), addr.into()));
1010 }
1011
1012 fn report_peer(&self, peer_id: sc_network_types::PeerId, cost_benefit: ReputationChange) {
1013 self.peer_store_handle.report_peer(peer_id, cost_benefit);
1014 }
1015
1016 fn peer_reputation(&self, peer_id: &sc_network_types::PeerId) -> i32 {
1017 self.peer_store_handle.peer_reputation(peer_id)
1018 }
1019
1020 fn disconnect_peer(&self, peer_id: sc_network_types::PeerId, protocol: ProtocolName) {
1021 let _ = self
1022 .to_worker
1023 .unbounded_send(ServiceToWorkerMsg::DisconnectPeer(peer_id.into(), protocol));
1024 }
1025
1026 fn accept_unreserved_peers(&self) {
1027 self.sync_protocol_handle.set_reserved_only(false);
1028 }
1029
1030 fn deny_unreserved_peers(&self) {
1031 self.sync_protocol_handle.set_reserved_only(true);
1032 }
1033
1034 fn add_reserved_peer(&self, peer: MultiaddrWithPeerId) -> Result<(), String> {
1035 if peer.peer_id == self.local_peer_id.into() {
1037 return Err("Local peer ID cannot be added as a reserved peer.".to_string());
1038 }
1039
1040 let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::AddKnownAddress(
1041 peer.peer_id.into(),
1042 peer.multiaddr.into(),
1043 ));
1044 self.sync_protocol_handle.add_reserved_peer(peer.peer_id.into());
1045
1046 Ok(())
1047 }
1048
1049 fn remove_reserved_peer(&self, peer_id: sc_network_types::PeerId) {
1050 self.sync_protocol_handle.remove_reserved_peer(peer_id.into());
1051 }
1052
1053 fn set_reserved_peers(
1054 &self,
1055 protocol: ProtocolName,
1056 peers: HashSet<sc_network_types::multiaddr::Multiaddr>,
1057 ) -> Result<(), String> {
1058 let Some(set_id) = self.notification_protocol_ids.get(&protocol) else {
1059 return Err(format!("Cannot set reserved peers for unknown protocol: {}", protocol));
1060 };
1061
1062 let peers: HashSet<Multiaddr> = peers.into_iter().map(Into::into).collect();
1063 let peers_addrs = self.split_multiaddr_and_peer_id(peers)?;
1064
1065 let mut peers: HashSet<PeerId> = HashSet::with_capacity(peers_addrs.len());
1066
1067 for (peer_id, addr) in peers_addrs.into_iter() {
1068 if peer_id == self.local_peer_id {
1070 return Err("Local peer ID cannot be added as a reserved peer.".to_string());
1071 }
1072
1073 peers.insert(peer_id.into());
1074
1075 if !addr.is_empty() {
1076 let _ = self
1077 .to_worker
1078 .unbounded_send(ServiceToWorkerMsg::AddKnownAddress(peer_id, addr));
1079 }
1080 }
1081
1082 self.protocol_handles[usize::from(*set_id)].set_reserved_peers(peers);
1083
1084 Ok(())
1085 }
1086
1087 fn add_peers_to_reserved_set(
1088 &self,
1089 protocol: ProtocolName,
1090 peers: HashSet<sc_network_types::multiaddr::Multiaddr>,
1091 ) -> Result<(), String> {
1092 let Some(set_id) = self.notification_protocol_ids.get(&protocol) else {
1093 return Err(format!(
1094 "Cannot add peers to reserved set of unknown protocol: {}",
1095 protocol
1096 ));
1097 };
1098
1099 let peers: HashSet<Multiaddr> = peers.into_iter().map(Into::into).collect();
1100 let peers = self.split_multiaddr_and_peer_id(peers)?;
1101
1102 for (peer_id, addr) in peers.into_iter() {
1103 if peer_id == self.local_peer_id {
1105 return Err("Local peer ID cannot be added as a reserved peer.".to_string());
1106 }
1107
1108 if !addr.is_empty() {
1109 let _ = self
1110 .to_worker
1111 .unbounded_send(ServiceToWorkerMsg::AddKnownAddress(peer_id, addr));
1112 }
1113
1114 self.protocol_handles[usize::from(*set_id)].add_reserved_peer(peer_id);
1115 }
1116
1117 Ok(())
1118 }
1119
1120 fn remove_peers_from_reserved_set(
1121 &self,
1122 protocol: ProtocolName,
1123 peers: Vec<sc_network_types::PeerId>,
1124 ) -> Result<(), String> {
1125 let Some(set_id) = self.notification_protocol_ids.get(&protocol) else {
1126 return Err(format!(
1127 "Cannot remove peers from reserved set of unknown protocol: {}",
1128 protocol
1129 ));
1130 };
1131
1132 for peer_id in peers.into_iter() {
1133 self.protocol_handles[usize::from(*set_id)].remove_reserved_peer(peer_id.into());
1134 }
1135
1136 Ok(())
1137 }
1138
1139 fn sync_num_connected(&self) -> usize {
1140 self.num_connected.load(Ordering::Relaxed)
1141 }
1142
1143 fn peer_role(
1144 &self,
1145 peer_id: sc_network_types::PeerId,
1146 handshake: Vec<u8>,
1147 ) -> Option<ObservedRole> {
1148 match Roles::decode_all(&mut &handshake[..]) {
1149 Ok(role) => Some(role.into()),
1150 Err(_) => {
1151 log::debug!(target: LOG_TARGET, "handshake doesn't contain peer role: {handshake:?}");
1152 self.peer_store_handle.peer_role(&(peer_id.into()))
1153 },
1154 }
1155 }
1156
1157 async fn reserved_peers(&self) -> Result<Vec<sc_network_types::PeerId>, ()> {
1161 let (tx, rx) = oneshot::channel();
1162
1163 self.sync_protocol_handle.reserved_peers(tx);
1164
1165 rx.await
1167 .map(|peers| peers.into_iter().map(From::from).collect())
1168 .map_err(|_| ())
1169 }
1170}
1171
1172impl<B, H> NetworkEventStream for NetworkService<B, H>
1173where
1174 B: BlockT + 'static,
1175 H: ExHashT,
1176{
1177 fn event_stream(&self, name: &'static str) -> Pin<Box<dyn Stream<Item = Event> + Send>> {
1178 let (tx, rx) = out_events::channel(name, 100_000);
1179 let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::EventStream(tx));
1180 Box::pin(rx)
1181 }
1182}
1183
1184#[async_trait::async_trait]
1185impl<B, H> NetworkRequest for NetworkService<B, H>
1186where
1187 B: BlockT + 'static,
1188 H: ExHashT,
1189{
1190 async fn request(
1191 &self,
1192 target: sc_network_types::PeerId,
1193 protocol: ProtocolName,
1194 request: Vec<u8>,
1195 fallback_request: Option<(Vec<u8>, ProtocolName)>,
1196 connect: IfDisconnected,
1197 ) -> Result<(Vec<u8>, ProtocolName), RequestFailure> {
1198 let (tx, rx) = oneshot::channel();
1199
1200 self.start_request(target.into(), protocol, request, fallback_request, tx, connect);
1201
1202 match rx.await {
1203 Ok(v) => v,
1204 Err(_) => Err(RequestFailure::Network(OutboundFailure::ConnectionClosed)),
1208 }
1209 }
1210
1211 fn start_request(
1212 &self,
1213 target: sc_network_types::PeerId,
1214 protocol: ProtocolName,
1215 request: Vec<u8>,
1216 fallback_request: Option<(Vec<u8>, ProtocolName)>,
1217 tx: oneshot::Sender<Result<(Vec<u8>, ProtocolName), RequestFailure>>,
1218 connect: IfDisconnected,
1219 ) {
1220 let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::Request {
1221 target: target.into(),
1222 protocol: protocol.into(),
1223 request,
1224 fallback_request,
1225 pending_response: tx,
1226 connect,
1227 });
1228 }
1229}
1230
1231#[must_use]
1233pub struct NotificationSender {
1234 sink: NotificationsSink,
1235
1236 protocol_name: ProtocolName,
1238
1239 notification_size_metric: Option<Histogram>,
1242}
1243
1244#[async_trait::async_trait]
1245impl NotificationSenderT for NotificationSender {
1246 async fn ready(
1247 &self,
1248 ) -> Result<Box<dyn NotificationSenderReadyT + '_>, NotificationSenderError> {
1249 Ok(Box::new(NotificationSenderReady {
1250 ready: match self.sink.reserve_notification().await {
1251 Ok(r) => Some(r),
1252 Err(()) => return Err(NotificationSenderError::Closed),
1253 },
1254 peer_id: self.sink.peer_id(),
1255 protocol_name: &self.protocol_name,
1256 notification_size_metric: self.notification_size_metric.clone(),
1257 }))
1258 }
1259}
1260
1261#[must_use]
1263pub struct NotificationSenderReady<'a> {
1264 ready: Option<Ready<'a>>,
1265
1266 peer_id: &'a PeerId,
1268
1269 protocol_name: &'a ProtocolName,
1271
1272 notification_size_metric: Option<Histogram>,
1275}
1276
1277impl<'a> NotificationSenderReadyT for NotificationSenderReady<'a> {
1278 fn send(&mut self, notification: Vec<u8>) -> Result<(), NotificationSenderError> {
1279 if let Some(notification_size_metric) = &self.notification_size_metric {
1280 notification_size_metric.observe(notification.len() as f64);
1281 }
1282
1283 trace!(
1284 target: LOG_TARGET,
1285 "External API => Notification({:?}, {}, {} bytes)",
1286 self.peer_id, self.protocol_name, notification.len(),
1287 );
1288 trace!(target: LOG_TARGET, "Handler({:?}) <= Async notification", self.peer_id);
1289
1290 self.ready
1291 .take()
1292 .ok_or(NotificationSenderError::Closed)?
1293 .send(notification)
1294 .map_err(|()| NotificationSenderError::Closed)
1295 }
1296}
1297
1298enum ServiceToWorkerMsg {
1302 FindClosestPeers(PeerId),
1303 GetValue(KademliaKey),
1304 PutValue(KademliaKey, Vec<u8>),
1305 PutRecordTo {
1306 record: Record,
1307 peers: HashSet<sc_network_types::PeerId>,
1308 update_local_storage: bool,
1309 },
1310 StoreRecord(KademliaKey, Vec<u8>, Option<PeerId>, Option<Instant>),
1311 StartProviding(KademliaKey),
1312 StopProviding(KademliaKey),
1313 GetProviders(KademliaKey),
1314 AddKnownAddress(PeerId, Multiaddr),
1315 EventStream(out_events::Sender),
1316 Request {
1317 target: PeerId,
1318 protocol: ProtocolName,
1319 request: Vec<u8>,
1320 fallback_request: Option<(Vec<u8>, ProtocolName)>,
1321 pending_response: oneshot::Sender<Result<(Vec<u8>, ProtocolName), RequestFailure>>,
1322 connect: IfDisconnected,
1323 },
1324 NetworkStatus {
1325 pending_response: oneshot::Sender<Result<NetworkStatus, RequestFailure>>,
1326 },
1327 NetworkState {
1328 pending_response: oneshot::Sender<Result<NetworkState, RequestFailure>>,
1329 },
1330 DisconnectPeer(PeerId, ProtocolName),
1331}
1332
1333#[must_use = "The NetworkWorker must be polled in order for the network to advance"]
1337pub struct NetworkWorker<B, H>
1338where
1339 B: BlockT + 'static,
1340 H: ExHashT,
1341{
1342 listen_addresses: Arc<Mutex<HashSet<Multiaddr>>>,
1344 num_connected: Arc<AtomicUsize>,
1346 service: Arc<NetworkService<B, H>>,
1348 network_service: Swarm<Behaviour<B>>,
1350 from_service: TracingUnboundedReceiver<ServiceToWorkerMsg>,
1352 event_streams: out_events::OutChannels,
1354 metrics: Option<Metrics>,
1356 boot_node_ids: Arc<HashMap<PeerId, Vec<Multiaddr>>>,
1358 reported_invalid_boot_nodes: HashSet<PeerId>,
1360 peer_store_handle: Arc<dyn PeerStoreProvider>,
1362 notif_protocol_handles: Vec<protocol::ProtocolHandle>,
1364 _marker: PhantomData<H>,
1367 _block: PhantomData<B>,
1369}
1370
1371impl<B, H> NetworkWorker<B, H>
1372where
1373 B: BlockT + 'static,
1374 H: ExHashT,
1375{
1376 pub async fn run(mut self) {
1378 while self.next_action().await {}
1379 }
1380
1381 pub async fn next_action(&mut self) -> bool {
1386 futures::select! {
1387 msg = self.from_service.next() => {
1389 if let Some(msg) = msg {
1390 self.handle_worker_message(msg);
1391 } else {
1392 return false
1393 }
1394 },
1395 event = self.network_service.select_next_some() => {
1397 self.handle_swarm_event(event);
1398 },
1399 };
1400
1401 let num_connected_peers = self.network_service.behaviour().user_protocol().num_sync_peers();
1403 self.num_connected.store(num_connected_peers, Ordering::Relaxed);
1404
1405 if let Some(metrics) = self.metrics.as_ref() {
1406 if let Some(buckets) = self.network_service.behaviour_mut().num_entries_per_kbucket() {
1407 for (lower_ilog2_bucket_bound, num_entries) in buckets {
1408 metrics
1409 .kbuckets_num_nodes
1410 .with_label_values(&[&lower_ilog2_bucket_bound.to_string()])
1411 .set(num_entries as u64);
1412 }
1413 }
1414 if let Some(num_entries) = self.network_service.behaviour_mut().num_kademlia_records() {
1415 metrics.kademlia_records_count.set(num_entries as u64);
1416 }
1417 if let Some(num_entries) =
1418 self.network_service.behaviour_mut().kademlia_records_total_size()
1419 {
1420 metrics.kademlia_records_sizes_total.set(num_entries as u64);
1421 }
1422
1423 metrics.pending_connections.set(
1424 Swarm::network_info(&self.network_service).connection_counters().num_pending()
1425 as u64,
1426 );
1427 }
1428
1429 true
1430 }
1431
1432 fn handle_worker_message(&mut self, msg: ServiceToWorkerMsg) {
1434 match msg {
1435 ServiceToWorkerMsg::FindClosestPeers(target) => {
1436 self.network_service.behaviour_mut().find_closest_peers(target)
1437 },
1438 ServiceToWorkerMsg::GetValue(key) => {
1439 self.network_service.behaviour_mut().get_value(key.into())
1440 },
1441 ServiceToWorkerMsg::PutValue(key, value) => {
1442 self.network_service.behaviour_mut().put_value(key.into(), value)
1443 },
1444 ServiceToWorkerMsg::PutRecordTo { record, peers, update_local_storage } => self
1445 .network_service
1446 .behaviour_mut()
1447 .put_record_to(record.into(), peers, update_local_storage),
1448 ServiceToWorkerMsg::StoreRecord(key, value, publisher, expires) => self
1449 .network_service
1450 .behaviour_mut()
1451 .store_record(key.into(), value, publisher, expires),
1452 ServiceToWorkerMsg::StartProviding(key) => {
1453 self.network_service.behaviour_mut().start_providing(key.into())
1454 },
1455 ServiceToWorkerMsg::StopProviding(key) => {
1456 self.network_service.behaviour_mut().stop_providing(&key.into())
1457 },
1458 ServiceToWorkerMsg::GetProviders(key) => {
1459 self.network_service.behaviour_mut().get_providers(key.into())
1460 },
1461 ServiceToWorkerMsg::AddKnownAddress(peer_id, addr) => {
1462 self.network_service.behaviour_mut().add_known_address(peer_id, addr)
1463 },
1464 ServiceToWorkerMsg::EventStream(sender) => self.event_streams.push(sender),
1465 ServiceToWorkerMsg::Request {
1466 target,
1467 protocol,
1468 request,
1469 fallback_request,
1470 pending_response,
1471 connect,
1472 } => {
1473 self.network_service.behaviour_mut().send_request(
1474 &target,
1475 protocol,
1476 request,
1477 fallback_request,
1478 pending_response,
1479 connect,
1480 );
1481 },
1482 ServiceToWorkerMsg::NetworkStatus { pending_response } => {
1483 let _ = pending_response.send(Ok(self.status()));
1484 },
1485 ServiceToWorkerMsg::NetworkState { pending_response } => {
1486 let _ = pending_response.send(Ok(self.network_state()));
1487 },
1488 ServiceToWorkerMsg::DisconnectPeer(who, protocol_name) => self
1489 .network_service
1490 .behaviour_mut()
1491 .user_protocol_mut()
1492 .disconnect_peer(&who, protocol_name),
1493 }
1494 }
1495
1496 fn handle_swarm_event(&mut self, event: SwarmEvent<BehaviourOut>) {
1498 match event {
1499 SwarmEvent::Behaviour(BehaviourOut::InboundRequest { protocol, result, .. }) => {
1500 if let Some(metrics) = self.metrics.as_ref() {
1501 match result {
1502 Ok(serve_time) => {
1503 metrics
1504 .requests_in_success_total
1505 .with_label_values(&[&protocol])
1506 .observe(serve_time.as_secs_f64());
1507 },
1508 Err(err) => {
1509 let reason = match err {
1510 ResponseFailure::Network(InboundFailure::Timeout) => {
1511 Some("timeout")
1512 },
1513 ResponseFailure::Network(InboundFailure::UnsupportedProtocols) =>
1514 {
1519 None
1520 },
1521 ResponseFailure::Network(InboundFailure::ResponseOmission) => {
1522 Some("busy-omitted")
1523 },
1524 ResponseFailure::Network(InboundFailure::ConnectionClosed) => {
1525 Some("connection-closed")
1526 },
1527 ResponseFailure::Network(InboundFailure::Io(_)) => Some("io"),
1528 };
1529
1530 if let Some(reason) = reason {
1531 metrics
1532 .requests_in_failure_total
1533 .with_label_values(&[&protocol, reason])
1534 .inc();
1535 }
1536 },
1537 }
1538 }
1539 },
1540 SwarmEvent::Behaviour(BehaviourOut::RequestFinished {
1541 protocol,
1542 duration,
1543 result,
1544 ..
1545 }) => {
1546 if let Some(metrics) = self.metrics.as_ref() {
1547 match result {
1548 Ok(_) => {
1549 metrics
1550 .requests_out_success_total
1551 .with_label_values(&[&protocol])
1552 .observe(duration.as_secs_f64());
1553 },
1554 Err(err) => {
1555 let reason = match err {
1556 RequestFailure::NotConnected => "not-connected",
1557 RequestFailure::UnknownProtocol => "unknown-protocol",
1558 RequestFailure::InvalidRequest => "invalid-request",
1559 RequestFailure::Refused => "refused",
1560 RequestFailure::Obsolete => "obsolete",
1561 RequestFailure::Network(OutboundFailure::DialFailure) => {
1562 "dial-failure"
1563 },
1564 RequestFailure::Network(OutboundFailure::Timeout) => "timeout",
1565 RequestFailure::Network(OutboundFailure::ConnectionClosed) => {
1566 "connection-closed"
1567 },
1568 RequestFailure::Network(OutboundFailure::UnsupportedProtocols) => {
1569 "unsupported"
1570 },
1571 RequestFailure::Network(OutboundFailure::Io(_)) => "io",
1572 };
1573
1574 metrics
1575 .requests_out_failure_total
1576 .with_label_values(&[&protocol, reason])
1577 .inc();
1578 },
1579 }
1580 }
1581 },
1582 SwarmEvent::Behaviour(BehaviourOut::ReputationChanges { peer, changes }) => {
1583 for change in changes {
1584 self.peer_store_handle.report_peer(peer.into(), change);
1585 }
1586 },
1587 SwarmEvent::Behaviour(BehaviourOut::PeerIdentify {
1588 peer_id,
1589 info:
1590 IdentifyInfo {
1591 protocol_version, agent_version, mut listen_addrs, protocols, ..
1592 },
1593 }) => {
1594 if listen_addrs.len() > 30 {
1595 debug!(
1596 target: LOG_TARGET,
1597 "Node {:?} has reported more than 30 addresses; it is identified by {:?} and {:?}",
1598 peer_id, protocol_version, agent_version
1599 );
1600 listen_addrs.truncate(30);
1601 }
1602 for addr in listen_addrs {
1603 self.network_service.behaviour_mut().add_self_reported_address_to_dht(
1604 &peer_id,
1605 &protocols,
1606 addr.clone(),
1607 );
1608 }
1609 self.peer_store_handle.add_known_peer(peer_id.into());
1610 },
1611 SwarmEvent::Behaviour(BehaviourOut::Discovered(peer_id)) => {
1612 self.peer_store_handle.add_known_peer(peer_id.into());
1613 },
1614 SwarmEvent::Behaviour(BehaviourOut::RandomKademliaStarted) => {
1615 if let Some(metrics) = self.metrics.as_ref() {
1616 metrics.kademlia_random_queries_total.inc();
1617 }
1618 },
1619 SwarmEvent::Behaviour(BehaviourOut::NotificationStreamOpened {
1620 remote,
1621 set_id,
1622 direction,
1623 negotiated_fallback,
1624 notifications_sink,
1625 received_handshake,
1626 }) => {
1627 let _ = self.notif_protocol_handles[usize::from(set_id)].report_substream_opened(
1628 remote,
1629 direction,
1630 received_handshake,
1631 negotiated_fallback,
1632 notifications_sink,
1633 );
1634 },
1635 SwarmEvent::Behaviour(BehaviourOut::NotificationStreamReplaced {
1636 remote,
1637 set_id,
1638 notifications_sink,
1639 }) => {
1640 let _ = self.notif_protocol_handles[usize::from(set_id)]
1641 .report_notification_sink_replaced(remote, notifications_sink);
1642
1643 },
1664 SwarmEvent::Behaviour(BehaviourOut::NotificationStreamClosed { remote, set_id }) => {
1665 let _ = self.notif_protocol_handles[usize::from(set_id)]
1666 .report_substream_closed(remote);
1667 },
1668 SwarmEvent::Behaviour(BehaviourOut::NotificationsReceived {
1669 remote,
1670 set_id,
1671 notification,
1672 }) => {
1673 let _ = self.notif_protocol_handles[usize::from(set_id)]
1674 .report_notification_received(remote, notification);
1675 },
1676 SwarmEvent::Behaviour(BehaviourOut::Dht(event, duration)) => {
1677 match (self.metrics.as_ref(), duration) {
1678 (Some(metrics), Some(duration)) => {
1679 let query_type = match event {
1680 DhtEvent::ClosestPeersFound(_, _) => "peers-found",
1681 DhtEvent::ClosestPeersNotFound(_) => "peers-not-found",
1682 DhtEvent::ValueFound(_) => "value-found",
1683 DhtEvent::ValueNotFound(_) => "value-not-found",
1684 DhtEvent::ValuePut(_) => "value-put",
1685 DhtEvent::ValuePutFailed(_) => "value-put-failed",
1686 DhtEvent::PutRecordRequest(_, _, _, _) => "put-record-request",
1687 DhtEvent::StartedProviding(_) => "started-providing",
1688 DhtEvent::StartProvidingFailed(_) => "start-providing-failed",
1689 DhtEvent::ProvidersFound(_, _) => "providers-found",
1690 DhtEvent::NoMoreProviders(_) => "no-more-providers",
1691 DhtEvent::ProvidersNotFound(_) => "providers-not-found",
1692 };
1693 metrics
1694 .kademlia_query_duration
1695 .with_label_values(&[query_type])
1696 .observe(duration.as_secs_f64());
1697 },
1698 _ => {},
1699 }
1700
1701 self.event_streams.send(Event::Dht(event));
1702 },
1703 SwarmEvent::Behaviour(BehaviourOut::None) => {
1704 },
1706 SwarmEvent::ConnectionEstablished {
1707 peer_id,
1708 endpoint,
1709 num_established,
1710 concurrent_dial_errors,
1711 ..
1712 } => {
1713 if let Some(errors) = concurrent_dial_errors {
1714 debug!(target: LOG_TARGET, "Libp2p => Connected({:?}) with errors: {:?}", peer_id, errors);
1715 } else {
1716 debug!(target: LOG_TARGET, "Libp2p => Connected({:?})", peer_id);
1717 }
1718
1719 if let Some(metrics) = self.metrics.as_ref() {
1720 let direction = match endpoint {
1721 ConnectedPoint::Dialer { .. } => "out",
1722 ConnectedPoint::Listener { .. } => "in",
1723 };
1724 metrics.connections_opened_total.with_label_values(&[direction]).inc();
1725
1726 if num_established.get() == 1 {
1727 metrics.distinct_peers_connections_opened_total.inc();
1728 }
1729 }
1730 },
1731 SwarmEvent::ConnectionClosed {
1732 connection_id,
1733 peer_id,
1734 cause,
1735 endpoint,
1736 num_established,
1737 } => {
1738 debug!(target: LOG_TARGET, "Libp2p => Disconnected({peer_id:?} via {connection_id:?}, {cause:?})");
1739 if let Some(metrics) = self.metrics.as_ref() {
1740 let direction = match endpoint {
1741 ConnectedPoint::Dialer { .. } => "out",
1742 ConnectedPoint::Listener { .. } => "in",
1743 };
1744 let reason = match cause {
1745 Some(ConnectionError::IO(_)) => "transport-error",
1746 Some(ConnectionError::KeepAliveTimeout) => "keep-alive-timeout",
1747 None => "actively-closed",
1748 };
1749 metrics.connections_closed_total.with_label_values(&[direction, reason]).inc();
1750
1751 if num_established == 0 {
1753 metrics.distinct_peers_connections_closed_total.inc();
1754 }
1755 }
1756 },
1757 SwarmEvent::NewListenAddr { address, .. } => {
1758 trace!(target: LOG_TARGET, "Libp2p => NewListenAddr({})", address);
1759 if let Some(metrics) = self.metrics.as_ref() {
1760 metrics.listeners_local_addresses.inc();
1761 }
1762 self.listen_addresses.lock().insert(address.clone());
1763 },
1764 SwarmEvent::ExpiredListenAddr { address, .. } => {
1765 info!(target: LOG_TARGET, "๐ช No longer listening on {}", address);
1766 if let Some(metrics) = self.metrics.as_ref() {
1767 metrics.listeners_local_addresses.dec();
1768 }
1769 self.listen_addresses.lock().remove(&address);
1770 },
1771 SwarmEvent::OutgoingConnectionError { connection_id, peer_id, error } => {
1772 if let Some(peer_id) = peer_id {
1773 trace!(
1774 target: LOG_TARGET,
1775 "Libp2p => Failed to reach {peer_id:?} via {connection_id:?}: {error}",
1776 );
1777
1778 let not_reported = !self.reported_invalid_boot_nodes.contains(&peer_id);
1779
1780 if let Some(addresses) =
1781 not_reported.then(|| self.boot_node_ids.get(&peer_id)).flatten()
1782 {
1783 if let DialError::WrongPeerId { obtained, endpoint } = &error {
1784 if let ConnectedPoint::Dialer {
1785 address,
1786 role_override: _,
1787 port_use: _,
1788 } = endpoint
1789 {
1790 let address_without_peer_id = parse_addr(address.clone().into())
1791 .map_or_else(|_| address.clone(), |r| r.1.into());
1792
1793 if addresses.iter().any(|a| address_without_peer_id == *a) {
1797 warn!(
1798 "๐ The bootnode you want to connect to at `{address}` provided a \
1799 different peer ID `{obtained}` than the one you expect `{peer_id}`.",
1800 );
1801
1802 self.reported_invalid_boot_nodes.insert(peer_id);
1803 }
1804 }
1805 }
1806 }
1807 }
1808
1809 if let Some(metrics) = self.metrics.as_ref() {
1810 let reason = match error {
1811 DialError::Denied { cause } => {
1812 if cause.downcast::<Exceeded>().is_ok() {
1813 Some("limit-reached")
1814 } else {
1815 None
1816 }
1817 },
1818 DialError::LocalPeerId { .. } => Some("local-peer-id"),
1819 DialError::WrongPeerId { .. } => Some("invalid-peer-id"),
1820 DialError::Transport(_) => Some("transport-error"),
1821 DialError::NoAddresses |
1822 DialError::DialPeerConditionFalse(_) |
1823 DialError::Aborted => None, };
1825 if let Some(reason) = reason {
1826 metrics.pending_connections_errors_total.with_label_values(&[reason]).inc();
1827 }
1828 }
1829 },
1830 SwarmEvent::Dialing { connection_id, peer_id } => {
1831 trace!(target: LOG_TARGET, "Libp2p => Dialing({peer_id:?}) via {connection_id:?}")
1832 },
1833 SwarmEvent::IncomingConnection { connection_id, local_addr, send_back_addr } => {
1834 trace!(target: LOG_TARGET, "Libp2p => IncomingConnection({local_addr},{send_back_addr} via {connection_id:?}))");
1835 if let Some(metrics) = self.metrics.as_ref() {
1836 metrics.incoming_connections_total.inc();
1837 }
1838 },
1839 SwarmEvent::IncomingConnectionError {
1840 connection_id,
1841 local_addr,
1842 send_back_addr,
1843 error,
1844 } => {
1845 debug!(
1846 target: LOG_TARGET,
1847 "Libp2p => IncomingConnectionError({local_addr},{send_back_addr} via {connection_id:?}): {error}"
1848 );
1849 if let Some(metrics) = self.metrics.as_ref() {
1850 let reason = match error {
1851 ListenError::Denied { cause } => {
1852 if cause.downcast::<Exceeded>().is_ok() {
1853 Some("limit-reached")
1854 } else {
1855 None
1856 }
1857 },
1858 ListenError::WrongPeerId { .. } | ListenError::LocalPeerId { .. } => {
1859 Some("invalid-peer-id")
1860 },
1861 ListenError::Transport(_) => Some("transport-error"),
1862 ListenError::Aborted => None, };
1864
1865 if let Some(reason) = reason {
1866 metrics
1867 .incoming_connections_errors_total
1868 .with_label_values(&[reason])
1869 .inc();
1870 }
1871 }
1872 },
1873 SwarmEvent::ListenerClosed { reason, addresses, .. } => {
1874 if let Some(metrics) = self.metrics.as_ref() {
1875 metrics.listeners_local_addresses.sub(addresses.len() as u64);
1876 }
1877 let mut listen_addresses = self.listen_addresses.lock();
1878 for addr in &addresses {
1879 listen_addresses.remove(addr);
1880 }
1881 drop(listen_addresses);
1882
1883 let addrs =
1884 addresses.into_iter().map(|a| a.to_string()).collect::<Vec<_>>().join(", ");
1885 match reason {
1886 Ok(()) => error!(
1887 target: LOG_TARGET,
1888 "๐ช Libp2p listener ({}) closed gracefully",
1889 addrs
1890 ),
1891 Err(e) => error!(
1892 target: LOG_TARGET,
1893 "๐ช Libp2p listener ({}) closed: {}",
1894 addrs, e
1895 ),
1896 }
1897 },
1898 SwarmEvent::ListenerError { error, .. } => {
1899 debug!(target: LOG_TARGET, "Libp2p => ListenerError: {}", error);
1900 if let Some(metrics) = self.metrics.as_ref() {
1901 metrics.listeners_errors_total.inc();
1902 }
1903 },
1904 SwarmEvent::NewExternalAddrCandidate { address } => {
1905 trace!(target: LOG_TARGET, "Libp2p => NewExternalAddrCandidate: {address:?}");
1906 },
1907 SwarmEvent::ExternalAddrConfirmed { address } => {
1908 trace!(target: LOG_TARGET, "Libp2p => ExternalAddrConfirmed: {address:?}");
1909 },
1910 SwarmEvent::ExternalAddrExpired { address } => {
1911 trace!(target: LOG_TARGET, "Libp2p => ExternalAddrExpired: {address:?}");
1912 },
1913 SwarmEvent::NewExternalAddrOfPeer { peer_id, address } => {
1914 trace!(target: LOG_TARGET, "Libp2p => NewExternalAddrOfPeer({peer_id:?}): {address:?}")
1915 },
1916 event => {
1917 warn!(target: LOG_TARGET, "New unknown SwarmEvent libp2p event: {event:?}");
1918 },
1919 }
1920 }
1921}
1922
1923impl<B, H> Unpin for NetworkWorker<B, H>
1924where
1925 B: BlockT + 'static,
1926 H: ExHashT,
1927{
1928}
1929
1930pub(crate) fn ensure_addresses_consistent_with_transport<'a>(
1931 addresses: impl Iterator<Item = &'a sc_network_types::multiaddr::Multiaddr>,
1932 transport: &TransportConfig,
1933) -> Result<(), Error> {
1934 use sc_network_types::multiaddr::Protocol;
1935
1936 if matches!(transport, TransportConfig::MemoryOnly) {
1937 let addresses: Vec<_> = addresses
1938 .filter(|x| x.iter().any(|y| !matches!(y, Protocol::Memory(_))))
1939 .cloned()
1940 .collect();
1941
1942 if !addresses.is_empty() {
1943 return Err(Error::AddressesForAnotherTransport {
1944 transport: transport.clone(),
1945 addresses,
1946 });
1947 }
1948 } else {
1949 let addresses: Vec<_> = addresses
1950 .filter(|x| x.iter().any(|y| matches!(y, Protocol::Memory(_))))
1951 .cloned()
1952 .collect();
1953
1954 if !addresses.is_empty() {
1955 return Err(Error::AddressesForAnotherTransport {
1956 transport: transport.clone(),
1957 addresses,
1958 });
1959 }
1960 }
1961
1962 Ok(())
1963}