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