1use crate::{
22 config::{
23 FullNetworkConfiguration, IncomingRequest, NodeKeyConfig, NotificationHandshake, Params,
24 SetConfig, TransportConfig,
25 },
26 error::Error,
27 event::{DhtEvent, Event},
28 litep2p::{
29 bitswap::BitswapService,
30 discovery::{Discovery, DiscoveryEvent},
31 ipfs_dht::IpfsDht,
32 peerstore::Peerstore,
33 service::{Litep2pNetworkService, NetworkServiceCommand},
34 shim::{
35 notification::{
36 config::{NotificationProtocolConfig, ProtocolControlHandle},
37 peerset::PeersetCommand,
38 },
39 request_response::{RequestResponseConfig, RequestResponseProtocol},
40 },
41 },
42 peer_store::PeerStoreProvider,
43 service::{
44 metrics::{register_without_sources, MetricSources, Metrics, NotificationMetrics},
45 out_events,
46 traits::{BandwidthSink, NetworkBackend, NetworkService},
47 },
48 webrtc, NetworkStatus, NotificationService, ProtocolName,
49};
50
51use codec::Encode;
52use futures::StreamExt;
53use litep2p::{
54 config::ConfigBuilder,
55 crypto::ed25519::Keypair,
56 error::{DialError, NegotiationError},
57 executor::Executor,
58 protocol::{
59 libp2p::kademlia::{QueryId, Record},
60 request_response::ConfigBuilder as RequestResponseConfigBuilder,
61 },
62 transport::{
63 tcp::config::Config as TcpTransportConfig, webrtc::config::Config as WebRtcTransportConfig,
64 websocket::config::Config as WebSocketTransportConfig, ConnectionLimitsConfig, Endpoint,
65 },
66 types::{
67 multiaddr::{Multiaddr, Protocol},
68 ConnectionId,
69 },
70 Litep2p, Litep2pEvent, ProtocolName as Litep2pProtocolName,
71};
72use prometheus_endpoint::Registry;
73use sc_network_types::kad::{Key as RecordKey, PeerRecord, Record as P2PRecord};
74
75use sc_client_api::BlockBackend;
76use sc_network_common::{role::Roles, ExHashT};
77use sc_network_types::PeerId;
78use sc_utils::mpsc::{tracing_unbounded, TracingUnboundedReceiver};
79use sp_runtime::traits::Block as BlockT;
80
81use std::{
82 cmp,
83 collections::{hash_map::Entry, HashMap, HashSet},
84 fs,
85 future::Future,
86 iter,
87 pin::Pin,
88 sync::{
89 atomic::{AtomicUsize, Ordering},
90 Arc,
91 },
92 time::{Duration, Instant},
93};
94
95mod bitswap;
96mod bitswap_metrics;
97mod discovery;
98mod ipfs_dht;
99mod peerstore;
100mod service;
101mod shim;
102
103struct Litep2pBandwidthSink {
105 sink: litep2p::BandwidthSink,
106}
107
108impl BandwidthSink for Litep2pBandwidthSink {
109 fn total_inbound(&self) -> u64 {
110 self.sink.inbound() as u64
111 }
112
113 fn total_outbound(&self) -> u64 {
114 self.sink.outbound() as u64
115 }
116}
117
118struct Litep2pExecutor {
120 executor: Box<dyn Fn(Pin<Box<dyn Future<Output = ()> + Send>>) + Send + Sync>,
122}
123
124impl Executor for Litep2pExecutor {
125 fn run(&self, future: Pin<Box<dyn Future<Output = ()> + Send>>) {
126 (self.executor)(future)
127 }
128
129 fn run_with_name(&self, _: &'static str, future: Pin<Box<dyn Future<Output = ()> + Send>>) {
130 (self.executor)(future)
131 }
132}
133
134const LOG_TARGET: &str = "sub-libp2p";
136
137struct ConnectionContext {
139 endpoints: HashMap<ConnectionId, Endpoint>,
141
142 num_connections: usize,
144}
145
146#[derive(Debug)]
148enum KadQuery {
149 FindNode(PeerId, Instant),
151 GetValue(RecordKey, Instant),
153 PutValue(RecordKey, Instant),
155 GetProviders(RecordKey, Instant),
157 AddProvider(RecordKey, Instant),
159}
160
161pub struct Litep2pNetworkBackend {
163 litep2p: Litep2p,
165
166 network_service: Arc<dyn NetworkService>,
168
169 cmd_rx: TracingUnboundedReceiver<NetworkServiceCommand>,
171
172 peerset_handles: HashMap<ProtocolName, ProtocolControlHandle>,
174
175 pending_queries: HashMap<QueryId, KadQuery>,
177
178 discovery: Discovery,
180
181 num_connected: Arc<AtomicUsize>,
183
184 peers: HashMap<litep2p::PeerId, ConnectionContext>,
186
187 peerstore_handle: Arc<dyn PeerStoreProvider>,
189
190 block_announce_protocol: ProtocolName,
192
193 event_streams: out_events::OutChannels,
195
196 metrics: Option<Metrics>,
198}
199
200impl Litep2pNetworkBackend {
201 fn parse_addresses(
204 addresses: impl Iterator<Item = Multiaddr>,
205 ) -> HashMap<PeerId, Vec<Multiaddr>> {
206 addresses
207 .into_iter()
208 .filter_map(|address| match address.iter().next() {
209 Some(
210 Protocol::Dns(_) |
211 Protocol::Dns4(_) |
212 Protocol::Dns6(_) |
213 Protocol::Ip6(_) |
214 Protocol::Ip4(_),
215 ) => match address.iter().find(|protocol| std::matches!(protocol, Protocol::P2p(_)))
216 {
217 Some(Protocol::P2p(peer_id)) => Some((peer_id.into(), Some(address))),
218 _ => None,
219 },
220 Some(Protocol::P2p(peer_id)) => Some((peer_id.into(), None)),
221 _ => None,
222 })
223 .fold(HashMap::new(), |mut acc, (peer, maybe_address)| {
224 let entry = acc.entry(peer).or_default();
225 maybe_address.map(|address| entry.push(address));
226
227 acc
228 })
229 }
230
231 fn add_addresses(&mut self, peers: impl Iterator<Item = Multiaddr>) -> HashSet<PeerId> {
233 Self::parse_addresses(peers.into_iter())
234 .into_iter()
235 .filter_map(|(peer, addresses)| {
236 if addresses.is_empty() {
238 return Some(peer);
239 }
240
241 if self.litep2p.add_known_address(peer.into(), addresses.clone().into_iter()) == 0 {
242 log::warn!(
243 target: LOG_TARGET,
244 "couldn't add any addresses for {peer:?} and it won't be added as reserved peer",
245 );
246 return None;
247 }
248
249 self.peerstore_handle.add_known_peer(peer);
250 Some(peer)
251 })
252 .collect()
253 }
254}
255
256impl Litep2pNetworkBackend {
257 fn get_keypair(node_key: &NodeKeyConfig) -> Result<(Keypair, litep2p::PeerId), Error> {
259 let secret: litep2p::crypto::ed25519::SecretKey =
260 node_key.clone().into_keypair()?.secret().into();
261
262 let local_identity = Keypair::from(secret);
263 let local_public = local_identity.public();
264 let local_peer_id = local_public.to_peer_id();
265
266 Ok((local_identity, local_peer_id))
267 }
268
269 fn configure_transport<B: BlockT + 'static, H: ExHashT>(
271 config: &FullNetworkConfiguration<B, H, Self>,
272 keypair: Keypair,
273 ) -> Result<ConfigBuilder, Error> {
274 let _ = match config.network_config.transport {
275 TransportConfig::MemoryOnly => panic!("memory transport not supported"),
276 TransportConfig::Normal { .. } => false,
277 };
278 let config_builder = ConfigBuilder::new();
279
280 let listen_addr_len = config.network_config.listen_addresses.len();
281 let mut tcp_addresses = Vec::with_capacity(listen_addr_len);
282 let mut websocket_addresses = Vec::with_capacity(listen_addr_len);
283 let mut webrtc_addresses = Vec::with_capacity(listen_addr_len);
284
285 for addr in &config.network_config.listen_addresses {
286 use sc_network_types::multiaddr::Protocol;
287
288 let mut iter = addr.iter();
289
290 let ip_version = iter.next();
291 let Some(Protocol::Ip4(_) | Protocol::Ip6(_)) = ip_version else {
292 log::error!(
293 target: LOG_TARGET,
294 "unknown protocol {ip_version:?}, ignoring {addr:?}",
295 );
296 continue;
297 };
298
299 let transport_layer = iter.next();
300 let protocol_type = iter.next();
301
302 match (&transport_layer, &protocol_type) {
303 (Some(Protocol::Tcp(_)), Some(Protocol::P2p(_)) | None) => {
305 tcp_addresses.push(addr.clone());
306 },
307
308 (Some(Protocol::Tcp(_)), Some(Protocol::Ws(_) | Protocol::Wss(_))) => {
310 websocket_addresses.push(addr.clone());
311 },
312 (Some(Protocol::Udp(_)), Some(Protocol::WebRTCDirect)) => {
314 if !config.network_config.experimental_webrtc {
316 log::warn!(
317 target: LOG_TARGET,
318 "WebRTC address provided but --experimental-webrtc flag not enabled, ignoring {addr:?}"
319 );
320 continue;
321 }
322 webrtc_addresses.push(addr.clone());
323 },
324 _ => {
325 log::error!(
326 target: LOG_TARGET,
327 "unknown transport layer {transport_layer:?} and protocol type {protocol_type:?}, ignoring {addr:?}",
328 );
329 },
330 };
331 }
332
333 let mut config_builder = config_builder
334 .with_websocket(WebSocketTransportConfig {
335 listen_addresses: websocket_addresses.into_iter().map(Into::into).collect(),
336 yamux_config: litep2p::yamux::Config::default(),
337 nodelay: true,
338 ..Default::default()
339 })
340 .with_tcp(TcpTransportConfig {
341 listen_addresses: tcp_addresses.into_iter().map(Into::into).collect(),
342 yamux_config: litep2p::yamux::Config::default(),
343 nodelay: true,
344 ..Default::default()
345 });
346
347 if !webrtc_addresses.is_empty() {
348 let certificate =
350 webrtc::derive_certificate(keypair.secret()).map_err(Error::Litep2p)?;
351 config_builder = config_builder.with_webrtc(WebRtcTransportConfig {
352 listen_addresses: webrtc_addresses.into_iter().map(Into::into).collect(),
353 certificate: Some(certificate),
354 ..Default::default()
355 });
356 } else if config.network_config.experimental_webrtc {
357 log::warn!(
358 target: LOG_TARGET,
359 "WebRTC enabled but no listen address specified"
360 );
361 }
362
363 Ok(config_builder.with_keypair(keypair))
364 }
365}
366
367#[async_trait::async_trait]
368impl<B: BlockT + 'static, H: ExHashT> NetworkBackend<B, H> for Litep2pNetworkBackend {
369 type NotificationProtocolConfig = NotificationProtocolConfig;
370 type RequestResponseProtocolConfig = RequestResponseConfig;
371 type NetworkService<Block, Hash> = Arc<Litep2pNetworkService>;
372 type PeerStore = Peerstore;
373 type BitswapConfig = bitswap::BitswapConfig;
374
375 fn new(mut params: Params<B, H, Self>) -> Result<Self, Error>
376 where
377 Self: Sized,
378 {
379 if let Err(err) = rustls::crypto::ring::default_provider().install_default() {
381 log::warn!(
382 target: LOG_TARGET,
383 "failed to install ring CryptoProvider for rustls, another provider might be installed: {err:?}",
384 );
385 }
386
387 let (keypair, local_peer_id) =
388 Self::get_keypair(¶ms.network_config.network_config.node_key)?;
389 let (cmd_tx, cmd_rx) = tracing_unbounded("mpsc_network_worker", 100_000);
390
391 params.network_config.network_config.boot_nodes = params
392 .network_config
393 .network_config
394 .boot_nodes
395 .into_iter()
396 .filter(|boot_node| boot_node.peer_id != local_peer_id.into())
397 .collect();
398 params.network_config.network_config.default_peers_set.reserved_nodes = params
399 .network_config
400 .network_config
401 .default_peers_set
402 .reserved_nodes
403 .into_iter()
404 .filter(|reserved_node| {
405 if reserved_node.peer_id == local_peer_id.into() {
406 log::warn!(
407 target: LOG_TARGET,
408 "Local peer ID used in reserved node, ignoring: {reserved_node}",
409 );
410 false
411 } else {
412 true
413 }
414 })
415 .collect();
416
417 if let Some(path) = ¶ms.network_config.network_config.net_config_path {
418 fs::create_dir_all(path)?;
419 }
420
421 log::info!(target: LOG_TARGET, "Local node identity is: {local_peer_id}");
422 log::info!(target: LOG_TARGET, "Running litep2p network backend");
423
424 params.network_config.sanity_check_addresses()?;
425 params.network_config.sanity_check_bootnodes()?;
426
427 let mut config_builder =
428 Self::configure_transport(¶ms.network_config, keypair.clone())?;
429 let known_addresses = params.network_config.known_addresses();
430 let peer_store_handle = params.network_config.peer_store_handle();
431 let executor = Arc::new(Litep2pExecutor { executor: params.executor });
432
433 let FullNetworkConfiguration {
434 notification_protocols,
435 request_response_protocols,
436 network_config,
437 ..
438 } = params.network_config;
439
440 let block_announce_protocol = params.block_announce_config.protocol_name().clone();
446 let mut notif_protocols = HashMap::from_iter([(
447 params.block_announce_config.protocol_name().clone(),
448 params.block_announce_config.handle,
449 )]);
450
451 config_builder = notification_protocols
453 .into_iter()
454 .fold(config_builder, |config_builder, mut config| {
455 config.config.set_handshake(Roles::from(¶ms.role).encode());
456 notif_protocols.insert(config.protocol_name, config.handle);
457
458 config_builder.with_notification_protocol(config.config)
459 })
460 .with_notification_protocol(params.block_announce_config.config);
461
462 let metrics = match ¶ms.metrics_registry {
464 Some(registry) => Some(register_without_sources(registry)?),
465 None => None,
466 };
467
468 let (mut request_response_receivers, request_response_senders): (
474 HashMap<_, _>,
475 HashMap<_, _>,
476 ) = request_response_protocols
477 .iter()
478 .map(|config| {
479 let (tx, rx) = tracing_unbounded("outbound-requests", 10_000);
480 ((config.protocol_name.clone(), rx), (config.protocol_name.clone(), tx))
481 })
482 .unzip();
483
484 config_builder = request_response_protocols.into_iter().fold(
485 config_builder,
486 |config_builder, config| {
487 let (protocol_config, handle) = RequestResponseConfigBuilder::new(
488 Litep2pProtocolName::from(config.protocol_name.clone()),
489 )
490 .with_max_size(cmp::max(config.max_request_size, config.max_response_size) as usize)
491 .with_fallback_names(config.fallback_names.into_iter().map(From::from).collect())
492 .with_timeout(config.request_timeout)
493 .build();
494
495 let protocol = RequestResponseProtocol::new(
496 config.protocol_name.clone(),
497 handle,
498 Arc::clone(&peer_store_handle),
499 config.inbound_queue,
500 request_response_receivers
501 .remove(&config.protocol_name)
502 .expect("receiver exists as it was just added and there are no duplicate protocols; qed"),
503 request_response_senders.clone(),
504 metrics.clone(),
505 );
506
507 executor.run(Box::pin(async move {
508 protocol.run().await;
509 }));
510
511 config_builder.with_request_response_protocol(protocol_config)
512 },
513 );
514
515 let known_addresses: HashMap<litep2p::PeerId, Vec<Multiaddr>> =
517 known_addresses.into_iter().fold(HashMap::new(), |mut acc, (peer, address)| {
518 use sc_network_types::multiaddr::Protocol;
519
520 let address = match address.iter().last() {
521 Some(Protocol::Ws(_) | Protocol::Wss(_) | Protocol::Tcp(_)) => {
522 address.with(Protocol::P2p(peer.into()))
523 },
524 Some(Protocol::WebRTCDirect | Protocol::Certhash(_)) => {
525 address.with(Protocol::P2p(peer.into()))
526 },
527 Some(Protocol::P2p(_)) => address,
528 _ => return acc,
529 };
530
531 acc.entry(peer.into()).or_default().push(address.into());
532 peer_store_handle.add_known_peer(peer);
533
534 acc
535 });
536
537 let listen_addresses = Arc::new(Default::default());
539 let (discovery, ping_config, identify_config, kademlia_config, maybe_mdns_config) =
540 Discovery::new(
541 local_peer_id,
542 &network_config,
543 params.genesis_hash,
544 params.fork_id.as_deref(),
545 ¶ms.protocol_id,
546 known_addresses.clone(),
547 Arc::clone(&listen_addresses),
548 Arc::clone(&peer_store_handle),
549 );
550
551 let bitswap_cmd_tx = params.ipfs_config.as_ref().map(|c| c.bitswap_config.cmd_tx.clone());
552
553 if let Some(config) = params.ipfs_config {
555 config_builder =
556 config_builder.with_libp2p_bitswap(config.bitswap_config.litep2p_config);
557
558 if !config.bootnodes.is_empty() {
559 let (ipfs_dht, kad_config) = IpfsDht::new(config.bootnodes, config.block_provider);
560 config_builder = config_builder.with_libp2p_kademlia(kad_config);
561 executor.run(Box::pin(ipfs_dht.run()));
562 } else {
563 log::warn!(
564 target: LOG_TARGET,
565 "Not starting IPFS DHT publisher because no IPFS bootnodes are configured. \
566 Only direct Bitswap requests will be handled.",
567 );
568 }
569 }
570
571 config_builder = config_builder
572 .with_known_addresses(known_addresses.clone().into_iter())
573 .with_libp2p_ping(ping_config)
574 .with_libp2p_identify(identify_config)
575 .with_libp2p_kademlia(kademlia_config)
576 .with_connection_limits(ConnectionLimitsConfig::default().max_incoming_connections(
577 Some(crate::MAX_CONNECTIONS_ESTABLISHED_INCOMING as usize),
578 ))
579 .with_keep_alive_timeout(network_config.idle_connection_timeout)
580 .with_system_resolver()
583 .with_executor(executor);
584
585 if let Some(config) = maybe_mdns_config {
586 config_builder = config_builder.with_mdns(config);
587 }
588
589 let litep2p =
590 Litep2p::new(config_builder.build()).map_err(|error| Error::Litep2p(error))?;
591
592 litep2p.listen_addresses().for_each(|address| {
593 log::debug!(target: LOG_TARGET, "listening on: {address}");
594
595 listen_addresses.write().insert(address.clone());
596 });
597
598 let public_addresses = litep2p.public_addresses();
599 for address in network_config.public_addresses.iter() {
600 if let Err(err) = public_addresses.add_address(address.clone().into()) {
601 log::warn!(
602 target: LOG_TARGET,
603 "failed to add public address {address:?}: {err:?}",
604 );
605 }
606 }
607
608 let network_service = Arc::new(Litep2pNetworkService::new(
609 local_peer_id,
610 keypair.clone(),
611 cmd_tx,
612 Arc::clone(&peer_store_handle),
613 notif_protocols.clone(),
614 block_announce_protocol.clone(),
615 request_response_senders,
616 Arc::clone(&listen_addresses),
617 public_addresses,
618 bitswap_cmd_tx,
619 ));
620
621 let num_connected = Arc::new(Default::default());
623 let bandwidth: Arc<dyn BandwidthSink> =
624 Arc::new(Litep2pBandwidthSink { sink: litep2p.bandwidth_sink() });
625
626 if let Some(registry) = ¶ms.metrics_registry {
627 MetricSources::register(registry, bandwidth, Arc::clone(&num_connected))?;
628 }
629
630 Ok(Self {
631 network_service,
632 cmd_rx,
633 metrics,
634 peerset_handles: notif_protocols,
635 num_connected,
636 discovery,
637 pending_queries: HashMap::new(),
638 peerstore_handle: peer_store_handle,
639 block_announce_protocol,
640 event_streams: out_events::OutChannels::new(None)?,
641 peers: HashMap::new(),
642 litep2p,
643 })
644 }
645
646 fn network_service(&self) -> Arc<dyn NetworkService> {
647 Arc::clone(&self.network_service)
648 }
649
650 fn peer_store(
651 bootnodes: Vec<sc_network_types::PeerId>,
652 metrics_registry: Option<Registry>,
653 ) -> Self::PeerStore {
654 Peerstore::new(bootnodes, metrics_registry)
655 }
656
657 fn register_notification_metrics(registry: Option<&Registry>) -> NotificationMetrics {
658 NotificationMetrics::new(registry)
659 }
660
661 fn bitswap_server(
663 client: Arc<dyn BlockBackend<B> + Send + Sync>,
664 metrics_registry: Option<Registry>,
665 ) -> (Pin<Box<dyn Future<Output = ()> + Send>>, Self::BitswapConfig) {
666 BitswapService::new(client, metrics_registry.as_ref())
667 }
668
669 fn notification_config(
671 protocol_name: ProtocolName,
672 fallback_names: Vec<ProtocolName>,
673 max_notification_size: u64,
674 handshake: Option<NotificationHandshake>,
675 set_config: SetConfig,
676 metrics: NotificationMetrics,
677 peerstore_handle: Arc<dyn PeerStoreProvider>,
678 ) -> (Self::NotificationProtocolConfig, Box<dyn NotificationService>) {
679 Self::NotificationProtocolConfig::new(
680 protocol_name,
681 fallback_names,
682 max_notification_size as usize,
683 handshake,
684 set_config,
685 metrics,
686 peerstore_handle,
687 )
688 }
689
690 fn request_response_config(
692 protocol_name: ProtocolName,
693 fallback_names: Vec<ProtocolName>,
694 max_request_size: u64,
695 max_response_size: u64,
696 request_timeout: Duration,
697 inbound_queue: Option<async_channel::Sender<IncomingRequest>>,
698 ) -> Self::RequestResponseProtocolConfig {
699 Self::RequestResponseProtocolConfig::new(
700 protocol_name,
701 fallback_names,
702 max_request_size,
703 max_response_size,
704 request_timeout,
705 inbound_queue,
706 )
707 }
708
709 async fn run(mut self) {
711 log::debug!(target: LOG_TARGET, "starting litep2p network backend");
712
713 loop {
714 let num_connected_peers = self
715 .peerset_handles
716 .get(&self.block_announce_protocol)
717 .map_or(0usize, |handle| handle.connected_peers.load(Ordering::Relaxed));
718 self.num_connected.store(num_connected_peers, Ordering::Relaxed);
719
720 tokio::select! {
721 command = self.cmd_rx.next() => match command {
722 None => return,
723 Some(command) => match command {
724 NetworkServiceCommand::FindClosestPeers { target } => {
725 let query_id = self.discovery.find_node(target.into()).await;
726 self.pending_queries.insert(query_id, KadQuery::FindNode(target, Instant::now()));
727 }
728 NetworkServiceCommand::GetValue{ key } => {
729 let query_id = self.discovery.get_value(key.clone()).await;
730 self.pending_queries.insert(query_id, KadQuery::GetValue(key, Instant::now()));
731 }
732 NetworkServiceCommand::PutValue { key, value } => {
733 let query_id = self.discovery.put_value(key.clone(), value).await;
734 self.pending_queries.insert(query_id, KadQuery::PutValue(key, Instant::now()));
735 }
736 NetworkServiceCommand::PutValueTo { record, peers, update_local_storage} => {
737 let kademlia_key = record.key.clone();
738 let query_id = self.discovery.put_value_to_peers(record.into(), peers, update_local_storage).await;
739 self.pending_queries.insert(query_id, KadQuery::PutValue(kademlia_key, Instant::now()));
740 }
741 NetworkServiceCommand::StoreRecord { key, value, publisher, expires } => {
742 self.discovery.store_record(key, value, publisher.map(Into::into), expires).await;
743 }
744 NetworkServiceCommand::StartProviding { key } => {
745 let query_id = self.discovery.start_providing(key.clone()).await;
746 self.pending_queries.insert(query_id, KadQuery::AddProvider(key, Instant::now()));
747 }
748 NetworkServiceCommand::StopProviding { key } => {
749 self.discovery.stop_providing(key).await;
750 }
751 NetworkServiceCommand::GetProviders { key } => {
752 let query_id = self.discovery.get_providers(key.clone()).await;
753 self.pending_queries.insert(query_id, KadQuery::GetProviders(key, Instant::now()));
754 }
755 NetworkServiceCommand::EventStream { tx } => {
756 self.event_streams.push(tx);
757 }
758 NetworkServiceCommand::Status { tx } => {
759 let _ = tx.send(NetworkStatus {
760 num_connected_peers: self
761 .peerset_handles
762 .get(&self.block_announce_protocol)
763 .map_or(0usize, |handle| handle.connected_peers.load(Ordering::Relaxed)),
764 total_bytes_inbound: self.litep2p.bandwidth_sink().inbound() as u64,
765 total_bytes_outbound: self.litep2p.bandwidth_sink().outbound() as u64,
766 });
767 }
768 NetworkServiceCommand::AddPeersToReservedSet {
769 protocol,
770 peers,
771 } => {
772 let peers = self.add_addresses(peers.into_iter().map(Into::into));
773
774 match self.peerset_handles.get(&protocol) {
775 Some(handle) => {
776 let _ = handle.tx.unbounded_send(PeersetCommand::AddReservedPeers { peers });
777 }
778 None => log::warn!(target: LOG_TARGET, "protocol {protocol} doens't exist"),
779 };
780 }
781 NetworkServiceCommand::AddKnownAddress { peer, address } => {
782 let mut address: Multiaddr = address.into();
783
784 if !address.iter().any(|protocol| std::matches!(protocol, Protocol::P2p(_))) {
785 address.push(Protocol::P2p(litep2p::PeerId::from(peer).into()));
786 }
787
788 if self.litep2p.add_known_address(peer.into(), iter::once(address.clone())) > 0 {
789 self.peerstore_handle.add_known_peer(peer);
793 } else {
794 log::debug!(
795 target: LOG_TARGET,
796 "couldn't add known address ({address}) for {peer:?}, unsupported transport"
797 );
798 }
799 },
800 NetworkServiceCommand::SetReservedPeers { protocol, peers } => {
801 let peers = self.add_addresses(peers.into_iter().map(Into::into));
802
803 match self.peerset_handles.get(&protocol) {
804 Some(handle) => {
805 let _ = handle.tx.unbounded_send(PeersetCommand::SetReservedPeers { peers });
806 }
807 None => log::warn!(target: LOG_TARGET, "protocol {protocol} doens't exist"),
808 }
809
810 },
811 NetworkServiceCommand::DisconnectPeer {
812 protocol,
813 peer,
814 } => {
815 let Some(handle) = self.peerset_handles.get(&protocol) else {
816 log::warn!(target: LOG_TARGET, "protocol {protocol} doens't exist");
817 continue
818 };
819
820 let _ = handle.tx.unbounded_send(PeersetCommand::DisconnectPeer { peer });
821 }
822 NetworkServiceCommand::SetReservedOnly {
823 protocol,
824 reserved_only,
825 } => {
826 let Some(handle) = self.peerset_handles.get(&protocol) else {
827 log::warn!(target: LOG_TARGET, "protocol {protocol} doens't exist");
828 continue
829 };
830
831 let _ = handle.tx.unbounded_send(PeersetCommand::SetReservedOnly { reserved_only });
832 }
833 NetworkServiceCommand::RemoveReservedPeers {
834 protocol,
835 peers,
836 } => {
837 let Some(handle) = self.peerset_handles.get(&protocol) else {
838 log::warn!(target: LOG_TARGET, "protocol {protocol} doens't exist");
839 continue
840 };
841
842 let _ = handle.tx.unbounded_send(PeersetCommand::RemoveReservedPeers { peers });
843 }
844 }
845 },
846 event = self.discovery.next() => match event {
847 None => return,
848 Some(DiscoveryEvent::Discovered { addresses }) => {
849 for (peer, addresses) in Litep2pNetworkBackend::parse_addresses(addresses.into_iter()) {
851 if self.litep2p.add_known_address(peer.into(), addresses.clone().into_iter()) > 0 {
852 self.peerstore_handle.add_known_peer(peer);
853 }
854 }
855 }
856 Some(DiscoveryEvent::RoutingTableUpdate { peers }) => {
857 for peer in peers {
858 self.peerstore_handle.add_known_peer(peer.into());
859 }
860 }
861 Some(DiscoveryEvent::FindNodeSuccess { query_id, target, peers }) => {
862 match self.pending_queries.remove(&query_id) {
863 Some(KadQuery::FindNode(_, started)) => {
864 log::trace!(
865 target: LOG_TARGET,
866 "`FIND_NODE` for {target:?} ({query_id:?}) succeeded",
867 );
868
869 self.event_streams.send(
870 Event::Dht(
871 DhtEvent::ClosestPeersFound(
872 target.into(),
873 peers
874 .into_iter()
875 .map(|(peer, addrs)| (
876 peer.into(),
877 addrs.into_iter().map(Into::into).collect(),
878 ))
879 .collect(),
880 )
881 )
882 );
883
884 if let Some(ref metrics) = self.metrics {
885 metrics
886 .kademlia_query_duration
887 .with_label_values(&["node-find"])
888 .observe(started.elapsed().as_secs_f64());
889 }
890 },
891 query => {
892 log::error!(
893 target: LOG_TARGET,
894 "Missing/invalid pending query for `FIND_NODE`: {query:?}"
895 );
896 debug_assert!(false);
897 }
898 }
899 },
900 Some(DiscoveryEvent::GetRecordPartialResult { query_id, record }) => {
901 if !self.pending_queries.contains_key(&query_id) {
902 log::error!(
903 target: LOG_TARGET,
904 "Missing/invalid pending query for `GET_VALUE` partial result: {query_id:?}"
905 );
906
907 continue
908 }
909
910 let peer_id: sc_network_types::PeerId = record.peer.into();
911 let record = PeerRecord {
912 record: P2PRecord {
913 key: record.record.key.to_vec().into(),
914 value: record.record.value,
915 publisher: record.record.publisher.map(|peer_id| {
916 let peer_id: sc_network_types::PeerId = peer_id.into();
917 peer_id.into()
918 }),
919 expires: record.record.expires,
920 },
921 peer: Some(peer_id.into()),
922 };
923
924 self.event_streams.send(
925 Event::Dht(
926 DhtEvent::ValueFound(
927 record.into()
928 )
929 )
930 );
931 }
932 Some(DiscoveryEvent::GetRecordSuccess { query_id }) => {
933 match self.pending_queries.remove(&query_id) {
934 Some(KadQuery::GetValue(key, started)) => {
935 log::trace!(
936 target: LOG_TARGET,
937 "`GET_VALUE` for {key:?} ({query_id:?}) succeeded",
938 );
939
940 if let Some(ref metrics) = self.metrics {
941 metrics
942 .kademlia_query_duration
943 .with_label_values(&["value-get"])
944 .observe(started.elapsed().as_secs_f64());
945 }
946 },
947 query => {
948 log::error!(
949 target: LOG_TARGET,
950 "Missing/invalid pending query for `GET_VALUE`: {query:?}"
951 );
952 debug_assert!(false);
953 },
954 }
955 }
956 Some(DiscoveryEvent::PutRecordSuccess { query_id }) => {
957 match self.pending_queries.remove(&query_id) {
958 Some(KadQuery::PutValue(key, started)) => {
959 log::trace!(
960 target: LOG_TARGET,
961 "`PUT_VALUE` for {key:?} ({query_id:?}) succeeded",
962 );
963
964 self.event_streams.send(Event::Dht(
965 DhtEvent::ValuePut(key)
966 ));
967
968 if let Some(ref metrics) = self.metrics {
969 metrics
970 .kademlia_query_duration
971 .with_label_values(&["value-put"])
972 .observe(started.elapsed().as_secs_f64());
973 }
974 },
975 query => {
976 log::error!(
977 target: LOG_TARGET,
978 "Missing/invalid pending query for `PUT_VALUE`: {query:?}"
979 );
980 debug_assert!(false);
981 }
982 }
983 }
984 Some(DiscoveryEvent::GetProvidersSuccess { query_id, providers }) => {
985 match self.pending_queries.remove(&query_id) {
986 Some(KadQuery::GetProviders(key, started)) => {
987 log::trace!(
988 target: LOG_TARGET,
989 "`GET_PROVIDERS` for {key:?} ({query_id:?}) succeeded",
990 );
991
992 providers.iter().for_each(|p| {
997 self.litep2p.add_known_address(p.peer, p.addresses.clone().into_iter());
998 });
999
1000 self.event_streams.send(Event::Dht(
1001 DhtEvent::ProvidersFound(
1002 key.clone().into(),
1003 providers.into_iter().map(|p| p.peer.into()).collect()
1004 )
1005 ));
1006
1007 self.event_streams.send(Event::Dht(
1010 DhtEvent::NoMoreProviders(key.into())
1011 ));
1012
1013 if let Some(ref metrics) = self.metrics {
1014 metrics
1015 .kademlia_query_duration
1016 .with_label_values(&["providers-get"])
1017 .observe(started.elapsed().as_secs_f64());
1018 }
1019 },
1020 query => {
1021 log::error!(
1022 target: LOG_TARGET,
1023 "Missing/invalid pending query for `GET_PROVIDERS`: {query:?}"
1024 );
1025 debug_assert!(false);
1026 }
1027 }
1028 }
1029 Some(DiscoveryEvent::AddProviderSuccess { query_id, provided_key }) => {
1030 match self.pending_queries.remove(&query_id) {
1031 Some(KadQuery::AddProvider(key, started)) => {
1032 debug_assert_eq!(key, provided_key.into());
1033
1034 log::trace!(
1035 target: LOG_TARGET,
1036 "`ADD_PROVIDER` for {key:?} ({query_id:?}) succeeded",
1037 );
1038
1039 self.event_streams.send(Event::Dht(
1040 DhtEvent::StartedProviding(key.into())
1041 ));
1042
1043 if let Some(ref metrics) = self.metrics {
1044 metrics
1045 .kademlia_query_duration
1046 .with_label_values(&["provider-add"])
1047 .observe(started.elapsed().as_secs_f64());
1048 }
1049 }
1050 Some(_) => {
1051 log::error!(
1052 target: LOG_TARGET,
1053 "Invalid pending query for `ADD_PROVIDER`: {query_id:?}"
1054 );
1055 debug_assert!(false);
1056 }
1057 None => {
1058 log::trace!(
1059 target: LOG_TARGET,
1060 "`ADD_PROVIDER` for key {provided_key:?} ({query_id:?}) succeeded (republishing)",
1061 );
1062 }
1063 }
1064 }
1065 Some(DiscoveryEvent::QueryFailed { query_id }) => {
1066 match self.pending_queries.remove(&query_id) {
1067 Some(KadQuery::FindNode(peer_id, started)) => {
1068 log::debug!(
1069 target: LOG_TARGET,
1070 "`FIND_NODE` ({query_id:?}) failed for target {peer_id:?}",
1071 );
1072
1073 self.event_streams.send(Event::Dht(
1074 DhtEvent::ClosestPeersNotFound(peer_id.into())
1075 ));
1076
1077 if let Some(ref metrics) = self.metrics {
1078 metrics
1079 .kademlia_query_duration
1080 .with_label_values(&["node-find-failed"])
1081 .observe(started.elapsed().as_secs_f64());
1082 }
1083 },
1084 Some(KadQuery::GetValue(key, started)) => {
1085 log::debug!(
1086 target: LOG_TARGET,
1087 "`GET_VALUE` ({query_id:?}) failed for key {key:?}",
1088 );
1089
1090 self.event_streams.send(Event::Dht(
1091 DhtEvent::ValueNotFound(key)
1092 ));
1093
1094 if let Some(ref metrics) = self.metrics {
1095 metrics
1096 .kademlia_query_duration
1097 .with_label_values(&["value-get-failed"])
1098 .observe(started.elapsed().as_secs_f64());
1099 }
1100 },
1101 Some(KadQuery::PutValue(key, started)) => {
1102 log::debug!(
1103 target: LOG_TARGET,
1104 "`PUT_VALUE` ({query_id:?}) failed for key {key:?}",
1105 );
1106
1107 self.event_streams.send(Event::Dht(
1108 DhtEvent::ValuePutFailed(key)
1109 ));
1110
1111 if let Some(ref metrics) = self.metrics {
1112 metrics
1113 .kademlia_query_duration
1114 .with_label_values(&["value-put-failed"])
1115 .observe(started.elapsed().as_secs_f64());
1116 }
1117 },
1118 Some(KadQuery::GetProviders(key, started)) => {
1119 log::debug!(
1120 target: LOG_TARGET,
1121 "`GET_PROVIDERS` ({query_id:?}) failed for key {key:?}"
1122 );
1123
1124 self.event_streams.send(Event::Dht(
1125 DhtEvent::ProvidersNotFound(key)
1126 ));
1127
1128 if let Some(ref metrics) = self.metrics {
1129 metrics
1130 .kademlia_query_duration
1131 .with_label_values(&["providers-get-failed"])
1132 .observe(started.elapsed().as_secs_f64());
1133 }
1134 },
1135 Some(KadQuery::AddProvider(key, started)) => {
1136 log::debug!(
1137 target: LOG_TARGET,
1138 "`ADD_PROVIDER` ({query_id:?}) failed with key {key:?}",
1139 );
1140
1141 self.event_streams.send(Event::Dht(
1142 DhtEvent::StartProvidingFailed(key)
1143 ));
1144
1145 if let Some(ref metrics) = self.metrics {
1146 metrics
1147 .kademlia_query_duration
1148 .with_label_values(&["provider-add-failed"])
1149 .observe(started.elapsed().as_secs_f64());
1150 }
1151 },
1152 None => {
1153 log::debug!(
1154 target: LOG_TARGET,
1155 "non-existent query (likely republishing a provider) failed ({query_id:?})",
1156 );
1157 }
1158 }
1159 }
1160 Some(DiscoveryEvent::Identified { peer, listen_addresses, supported_protocols, .. }) => {
1161 self.discovery.add_self_reported_address(peer, supported_protocols, listen_addresses).await;
1162 }
1163 Some(DiscoveryEvent::ExternalAddressDiscovered { address }) => {
1164 match self.litep2p.public_addresses().add_address(address.clone().into()) {
1165 Ok(inserted) => if inserted {
1166 log::info!(target: LOG_TARGET, "๐ Discovered new external address for our node: {address}");
1167 },
1168 Err(err) => {
1169 log::warn!(
1170 target: LOG_TARGET,
1171 "๐ Failed to add discovered external address {address:?}: {err:?}",
1172 );
1173 },
1174 }
1175 }
1176 Some(DiscoveryEvent::ExternalAddressExpired{ address }) => {
1177 let local_peer_id = self.litep2p.local_peer_id();
1178
1179 let address = if !std::matches!(address.iter().last(), Some(Protocol::P2p(_))) {
1181 address.with(Protocol::P2p((*local_peer_id).into()))
1182 } else {
1183 address
1184 };
1185
1186 if self.litep2p.public_addresses().remove_address(&address) {
1187 log::info!(target: LOG_TARGET, "๐ Expired external address for our node: {address}");
1188 } else {
1189 log::warn!(
1190 target: LOG_TARGET,
1191 "๐ Failed to remove expired external address {address:?}"
1192 );
1193 }
1194 }
1195 Some(DiscoveryEvent::Ping { peer, rtt }) => {
1196 log::trace!(
1197 target: LOG_TARGET,
1198 "ping time with {peer:?}: {rtt:?}",
1199 );
1200 }
1201 Some(DiscoveryEvent::IncomingRecord { record: Record { key, value, publisher, expires }} ) => {
1202 self.event_streams.send(Event::Dht(
1203 DhtEvent::PutRecordRequest(
1204 key.into(),
1205 value,
1206 publisher.map(Into::into),
1207 expires,
1208 )
1209 ));
1210 },
1211
1212 Some(DiscoveryEvent::RandomKademliaStarted) => {
1213 if let Some(metrics) = self.metrics.as_ref() {
1214 metrics.kademlia_random_queries_total.inc();
1215 }
1216 }
1217 },
1218 event = self.litep2p.next_event() => match event {
1219 Some(Litep2pEvent::ConnectionEstablished { peer, endpoint }) => {
1220 let Some(metrics) = &self.metrics else {
1221 continue;
1222 };
1223
1224 let direction = match endpoint {
1225 Endpoint::Dialer { .. } => "out",
1226 Endpoint::Listener { .. } => {
1227 metrics.incoming_connections_total.inc();
1232
1233 "in"
1234 },
1235 };
1236 metrics.connections_opened_total.with_label_values(&[direction]).inc();
1237
1238 match self.peers.entry(peer) {
1239 Entry::Vacant(entry) => {
1240 entry.insert(ConnectionContext {
1241 endpoints: HashMap::from_iter([(endpoint.connection_id(), endpoint)]),
1242 num_connections: 1usize,
1243 });
1244 metrics.distinct_peers_connections_opened_total.inc();
1245 }
1246 Entry::Occupied(entry) => {
1247 let entry = entry.into_mut();
1248 entry.num_connections += 1;
1249 entry.endpoints.insert(endpoint.connection_id(), endpoint);
1250 }
1251 }
1252 }
1253 Some(Litep2pEvent::ConnectionClosed { peer, connection_id }) => {
1254 let Some(metrics) = &self.metrics else {
1255 continue;
1256 };
1257
1258 let Some(context) = self.peers.get_mut(&peer) else {
1259 log::debug!(target: LOG_TARGET, "unknown peer disconnected: {peer:?} ({connection_id:?})");
1260 continue
1261 };
1262
1263 let direction = match context.endpoints.remove(&connection_id) {
1264 None => {
1265 log::debug!(target: LOG_TARGET, "connection {connection_id:?} doesn't exist for {peer:?} ");
1266 continue
1267 }
1268 Some(endpoint) => {
1269 context.num_connections -= 1;
1270
1271 match endpoint {
1272 Endpoint::Dialer { .. } => "out",
1273 Endpoint::Listener { .. } => "in",
1274 }
1275 }
1276 };
1277
1278 metrics.connections_closed_total.with_label_values(&[direction, "actively-closed"]).inc();
1279
1280 if context.num_connections == 0 {
1281 self.peers.remove(&peer);
1282 metrics.distinct_peers_connections_closed_total.inc();
1283 }
1284 }
1285 Some(Litep2pEvent::DialFailure { address, error }) => {
1286 log::debug!(
1287 target: LOG_TARGET,
1288 "failed to dial peer at {address:?}: {error:?}",
1289 );
1290
1291 if let Some(metrics) = &self.metrics {
1292 let reason = match error {
1293 DialError::Timeout => "timeout",
1294 DialError::AddressError(_) => "invalid-address",
1295 DialError::DnsError(_) => "cannot-resolve-dns",
1296 DialError::NegotiationError(error) => match error {
1297 NegotiationError::Timeout => "timeout",
1298 NegotiationError::PeerIdMissing => "missing-peer-id",
1299 NegotiationError::StateMismatch => "state-mismatch",
1300 NegotiationError::PeerIdMismatch(_,_) => "peer-id-missmatch",
1301 NegotiationError::MultistreamSelectError(_) => "multistream-select-error",
1302 NegotiationError::SnowError(_) => "noise-error",
1303 NegotiationError::ParseError(_) => "parse-error",
1304 NegotiationError::IoError(_) => "io-error",
1305 NegotiationError::WebSocket(_) => "webscoket-error",
1306 NegotiationError::BadSignature => "bad-signature",
1307 }
1308 };
1309
1310 metrics.pending_connections_errors_total.with_label_values(&[&reason]).inc();
1311 }
1312 }
1313 Some(Litep2pEvent::ListDialFailures { errors }) => {
1314 log::debug!(
1315 target: LOG_TARGET,
1316 "failed to dial peer on multiple addresses {errors:?}",
1317 );
1318
1319 if let Some(metrics) = &self.metrics {
1320 metrics.pending_connections_errors_total.with_label_values(&["transport-errors"]).inc();
1321 }
1322 }
1323 None => {
1324 log::error!(
1325 target: LOG_TARGET,
1326 "Litep2p backend terminated"
1327 );
1328 return
1329 }
1330 },
1331 }
1332 }
1333 }
1334}