referrerpolicy=no-referrer-when-downgrade

sc_network/litep2p/
mod.rs

1// This file is part of Substrate.
2
3// Copyright (C) Parity Technologies (UK) Ltd.
4// SPDX-License-Identifier: GPL-3.0-or-later WITH Classpath-exception-2.0
5
6// This program is free software: you can redistribute it and/or modify
7// it under the terms of the GNU General Public License as published by
8// the Free Software Foundation, either version 3 of the License, or
9// (at your option) any later version.
10
11// This program is distributed in the hope that it will be useful,
12// but WITHOUT ANY WARRANTY; without even the implied warranty of
13// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
14// GNU General Public License for more details.
15
16// You should have received a copy of the GNU General Public License
17// along with this program. If not, see <https://www.gnu.org/licenses/>.
18
19//! `NetworkBackend` implementation for `litep2p`.
20
21use 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
103/// Litep2p bandwidth sink.
104struct 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
118/// Litep2p task executor.
119struct Litep2pExecutor {
120	/// Executor.
121	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
134/// Logging target for the file.
135const LOG_TARGET: &str = "sub-libp2p";
136
137/// Peer context.
138struct ConnectionContext {
139	/// Peer endpoints.
140	endpoints: HashMap<ConnectionId, Endpoint>,
141
142	/// Number of active connections.
143	num_connections: usize,
144}
145
146/// Kademlia query we are tracking.
147#[derive(Debug)]
148enum KadQuery {
149	/// `FIND_NODE` query for target and when it was initiated.
150	FindNode(PeerId, Instant),
151	/// `GET_VALUE` query for key and when it was initiated.
152	GetValue(RecordKey, Instant),
153	/// `PUT_VALUE` query for key and when it was initiated.
154	PutValue(RecordKey, Instant),
155	/// `GET_PROVIDERS` query for key and when it was initiated.
156	GetProviders(RecordKey, Instant),
157	/// `ADD_PROVIDER` query for key and when it was initiated.
158	AddProvider(RecordKey, Instant),
159}
160
161/// Networking backend for `litep2p`.
162pub struct Litep2pNetworkBackend {
163	/// Main `litep2p` object.
164	litep2p: Litep2p,
165
166	/// `NetworkService` implementation for `Litep2pNetworkBackend`.
167	network_service: Arc<dyn NetworkService>,
168
169	/// RX channel for receiving commands from `Litep2pNetworkService`.
170	cmd_rx: TracingUnboundedReceiver<NetworkServiceCommand>,
171
172	/// `Peerset` handles to notification protocols.
173	peerset_handles: HashMap<ProtocolName, ProtocolControlHandle>,
174
175	/// Pending Kademlia queries.
176	pending_queries: HashMap<QueryId, KadQuery>,
177
178	/// Discovery.
179	discovery: Discovery,
180
181	/// Number of connected peers.
182	num_connected: Arc<AtomicUsize>,
183
184	/// Connected peers.
185	peers: HashMap<litep2p::PeerId, ConnectionContext>,
186
187	/// Peerstore.
188	peerstore_handle: Arc<dyn PeerStoreProvider>,
189
190	/// Block announce protocol name.
191	block_announce_protocol: ProtocolName,
192
193	/// Sender for DHT events.
194	event_streams: out_events::OutChannels,
195
196	/// Prometheus metrics.
197	metrics: Option<Metrics>,
198}
199
200impl Litep2pNetworkBackend {
201	/// From an iterator of multiaddress(es), parse and group all addresses of peers
202	/// so that litep2p can consume the information easily.
203	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	/// Add new known addresses to `litep2p` and return the parsed peer IDs.
232	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				// `peers` contained multiaddress in the form `/p2p/<peer ID>`
237				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	/// Get `litep2p` keypair from `NodeKeyConfig`.
258	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	/// Configure transport protocols for `Litep2pNetworkBackend`.
270	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				// Plain TCP address.
304				(Some(Protocol::Tcp(_)), Some(Protocol::P2p(_)) | None) => {
305					tcp_addresses.push(addr.clone());
306				},
307
308				// Websocket address.
309				(Some(Protocol::Tcp(_)), Some(Protocol::Ws(_) | Protocol::Wss(_))) => {
310					websocket_addresses.push(addr.clone());
311				},
312				// WebRTCDirecet address.
313				(Some(Protocol::Udp(_)), Some(Protocol::WebRTCDirect)) => {
314					// Ignore WebRTC addresses unless the experimental feature is enabled.
315					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			// WebRTC cert/key are unambiguously defined by the node key.
349			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		// Install the ring CryptoProvider for rustls before any TLS connections are made.
380		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(&params.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) = &params.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(&params.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		// initialize notification protocols
441		//
442		// pass the protocol configuration to `Litep2pConfigBuilder` and save the TX channel
443		// to the protocol's `Peerset` together with the protocol name to allow other subsystems
444		// of Polkadot SDK to control connectivity of the notification protocol
445		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		// handshake for all but the syncing protocol is set to node role
452		config_builder = notification_protocols
453			.into_iter()
454			.fold(config_builder, |config_builder, mut config| {
455				config.config.set_handshake(Roles::from(&params.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		// initialize request-response protocols
463		let metrics = match &params.metrics_registry {
464			Some(registry) => Some(register_without_sources(registry)?),
465			None => None,
466		};
467
468		// create channels that are used to send request before initializing protocols so the
469		// senders can be passed onto all request-response protocols
470		//
471		// all protocols must have each others' senders so they can send the fallback request in
472		// case the main protocol is not supported by the remote peer and user specified a fallback
473		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		// collect known addresses
516		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		// enable ipfs ping, identify and kademlia, and potentially mdns if user enabled it
538		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				&params.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		// enable Bitswap & IPFS DHT
554		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			// Use system DNS resolver to enable intranet domain resolution and administrator
581			// control over DNS lookup.
582			.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		// register rest of the metrics now that `Litep2p` has been created
622		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) = &params.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	/// Create Bitswap server.
662	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	/// Create notification protocol configuration for `protocol`.
670	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	/// Create request-response protocol configuration.
691	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	/// Start [`Litep2pNetworkBackend`] event loop.
710	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								// libp2p backend generates `DiscoveryOut::Discovered(peer_id)`
790								// event when a new address is added for a peer, which leads to the
791								// peer being added to peerstore. Do the same directly here.
792								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						// if at least one address was added for the peer, report the peer to `Peerstore`
850						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								// We likely requested providers to connect to them,
993								// so let's add their addresses to litep2p's transport manager.
994								// Consider also looking the addresses of providers up with `FIND_NODE`
995								// query, as it can yield more up to date addresses.
996								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								// litep2p returns all providers in a single event, so we let
1008								// subscribers know no more providers will be yielded.
1009								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						// Litep2p requires the peer ID to be present in the address.
1180						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								// Increment incoming connections counter.
1228								//
1229								// Note: For litep2p these are represented by established negotiated connections,
1230								// while for libp2p (legacy) these represent not-yet-negotiated connections.
1231								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}