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		discovery::{Discovery, DiscoveryEvent},
30		ipfs_dht::IpfsDht,
31		peerstore::Peerstore,
32		service::{Litep2pNetworkService, NetworkServiceCommand},
33		shim::{
34			notification::{
35				config::{NotificationProtocolConfig, ProtocolControlHandle},
36				peerset::PeersetCommand,
37			},
38			request_response::{RequestResponseConfig, RequestResponseProtocol},
39		},
40	},
41	peer_store::PeerStoreProvider,
42	service::{
43		metrics::{register_without_sources, MetricSources, Metrics, NotificationMetrics},
44		out_events,
45		traits::{BandwidthSink, NetworkBackend, NetworkService},
46	},
47	webrtc, NetworkStatus, NotificationService, ProtocolName,
48};
49
50use codec::Encode;
51use futures::StreamExt;
52use litep2p::{
53	config::ConfigBuilder,
54	crypto::ed25519::Keypair,
55	error::{DialError, NegotiationError},
56	executor::Executor,
57	protocol::{
58		libp2p::kademlia::{QueryId, Record},
59		request_response::ConfigBuilder as RequestResponseConfigBuilder,
60	},
61	transport::{
62		tcp::config::Config as TcpTransportConfig, webrtc::config::Config as WebRtcTransportConfig,
63		websocket::config::Config as WebSocketTransportConfig, ConnectionLimitsConfig, Endpoint,
64	},
65	types::{
66		multiaddr::{Multiaddr, Protocol},
67		ConnectionId,
68	},
69	Litep2p, Litep2pEvent, ProtocolName as Litep2pProtocolName,
70};
71use prometheus_endpoint::Registry;
72
73use sc_network_common::{role::Roles, ExHashT};
74use sc_network_types::{
75	kad::{Key as RecordKey, PeerRecord, Record as P2PRecord},
76	multiaddr::Protocol as NetworkProtocol,
77	PeerId,
78};
79use sc_utils::mpsc::{tracing_unbounded, TracingUnboundedReceiver};
80use sp_runtime::traits::Block as BlockT;
81
82use std::{
83	cmp,
84	collections::{hash_map::Entry, HashMap, HashSet},
85	fs,
86	future::Future,
87	iter,
88	pin::Pin,
89	sync::{
90		atomic::{AtomicUsize, Ordering},
91		Arc,
92	},
93	time::{Duration, Instant},
94};
95
96mod discovery;
97mod ipfs_dht;
98mod peerstore;
99mod service;
100mod shim;
101
102/// Litep2p bandwidth sink.
103struct Litep2pBandwidthSink {
104	sink: litep2p::BandwidthSink,
105}
106
107impl BandwidthSink for Litep2pBandwidthSink {
108	fn total_inbound(&self) -> u64 {
109		self.sink.inbound() as u64
110	}
111
112	fn total_outbound(&self) -> u64 {
113		self.sink.outbound() as u64
114	}
115}
116
117/// Litep2p task executor.
118struct Litep2pExecutor {
119	/// Executor.
120	executor: Box<dyn Fn(Pin<Box<dyn Future<Output = ()> + Send>>) + Send + Sync>,
121}
122
123impl Executor for Litep2pExecutor {
124	fn run(&self, future: Pin<Box<dyn Future<Output = ()> + Send>>) {
125		(self.executor)(future)
126	}
127
128	fn run_with_name(&self, _: &'static str, future: Pin<Box<dyn Future<Output = ()> + Send>>) {
129		(self.executor)(future)
130	}
131}
132
133/// Logging target for the file.
134const LOG_TARGET: &str = "sub-libp2p";
135
136/// Peer context.
137struct ConnectionContext {
138	/// Peer endpoints.
139	endpoints: HashMap<ConnectionId, Endpoint>,
140
141	/// Number of active connections.
142	num_connections: usize,
143}
144
145/// Kademlia query we are tracking.
146#[derive(Debug)]
147enum KadQuery {
148	/// `FIND_NODE` query for target and when it was initiated.
149	FindNode(PeerId, Instant),
150	/// `GET_VALUE` query for key and when it was initiated.
151	GetValue(RecordKey, Instant),
152	/// `PUT_VALUE` query for key and when it was initiated.
153	PutValue(RecordKey, Instant),
154	/// `GET_PROVIDERS` query for key and when it was initiated.
155	GetProviders(RecordKey, Instant),
156	/// `ADD_PROVIDER` query for key and when it was initiated.
157	AddProvider(RecordKey, Instant),
158}
159
160/// Networking backend for `litep2p`.
161pub struct Litep2pNetworkBackend {
162	/// Main `litep2p` object.
163	litep2p: Litep2p,
164
165	/// `NetworkService` implementation for `Litep2pNetworkBackend`.
166	network_service: Arc<dyn NetworkService>,
167
168	/// RX channel for receiving commands from `Litep2pNetworkService`.
169	cmd_rx: TracingUnboundedReceiver<NetworkServiceCommand>,
170
171	/// `Peerset` handles to notification protocols.
172	peerset_handles: HashMap<ProtocolName, ProtocolControlHandle>,
173
174	/// Pending Kademlia queries.
175	pending_queries: HashMap<QueryId, KadQuery>,
176
177	/// Discovery.
178	discovery: Discovery,
179
180	/// Number of connected peers.
181	num_connected: Arc<AtomicUsize>,
182
183	/// Connected peers.
184	peers: HashMap<litep2p::PeerId, ConnectionContext>,
185
186	/// Peerstore.
187	peerstore_handle: Arc<dyn PeerStoreProvider>,
188
189	/// Block announce protocol name.
190	block_announce_protocol: ProtocolName,
191
192	/// Sender for DHT events.
193	event_streams: out_events::OutChannels,
194
195	/// Prometheus metrics.
196	metrics: Option<Metrics>,
197}
198
199impl Litep2pNetworkBackend {
200	/// From an iterator of multiaddress(es), parse and group all addresses of peers
201	/// so that litep2p can consume the information easily.
202	fn parse_addresses(
203		addresses: impl Iterator<Item = Multiaddr>,
204	) -> HashMap<PeerId, Vec<Multiaddr>> {
205		addresses
206			.into_iter()
207			.filter_map(|address| match address.iter().next() {
208				Some(
209					Protocol::Dns(_) |
210					Protocol::Dns4(_) |
211					Protocol::Dns6(_) |
212					Protocol::Ip6(_) |
213					Protocol::Ip4(_),
214				) => match address.iter().find(|protocol| std::matches!(protocol, Protocol::P2p(_)))
215				{
216					Some(Protocol::P2p(peer_id)) => Some((peer_id.into(), Some(address))),
217					_ => None,
218				},
219				Some(Protocol::P2p(peer_id)) => Some((peer_id.into(), None)),
220				_ => None,
221			})
222			.fold(HashMap::new(), |mut acc, (peer, maybe_address)| {
223				let entry = acc.entry(peer).or_default();
224				maybe_address.map(|address| entry.push(address));
225
226				acc
227			})
228	}
229
230	/// Add new known addresses to `litep2p` and return the parsed peer IDs.
231	fn add_addresses(&mut self, peers: impl Iterator<Item = Multiaddr>) -> HashSet<PeerId> {
232		Self::parse_addresses(peers.into_iter())
233			.into_iter()
234			.filter_map(|(peer, addresses)| {
235				// `peers` contained multiaddress in the form `/p2p/<peer ID>`
236				if addresses.is_empty() {
237					return Some(peer);
238				}
239
240				if self.litep2p.add_known_address(peer.into(), addresses.clone().into_iter()) == 0 {
241					log::warn!(
242						target: LOG_TARGET,
243						"couldn't add any addresses for {peer:?} and it won't be added as reserved peer",
244					);
245					return None;
246				}
247
248				self.peerstore_handle.add_known_peer(peer);
249				Some(peer)
250			})
251			.collect()
252	}
253}
254
255impl Litep2pNetworkBackend {
256	/// Get `litep2p` keypair from `NodeKeyConfig`.
257	fn get_keypair(node_key: &NodeKeyConfig) -> Result<(Keypair, litep2p::PeerId), Error> {
258		let secret: litep2p::crypto::ed25519::SecretKey =
259			node_key.clone().into_keypair()?.secret().into();
260
261		let local_identity = Keypair::from(secret);
262		let local_public = local_identity.public();
263		let local_peer_id = local_public.to_peer_id();
264
265		Ok((local_identity, local_peer_id))
266	}
267
268	/// Configure transport protocols for `Litep2pNetworkBackend`.
269	fn configure_transport<B: BlockT + 'static, H: ExHashT>(
270		config: &FullNetworkConfiguration<B, H, Self>,
271		keypair: Keypair,
272	) -> Result<ConfigBuilder, Error> {
273		let _ = match config.network_config.transport {
274			TransportConfig::MemoryOnly => panic!("memory transport not supported"),
275			TransportConfig::Normal { .. } => false,
276		};
277		let config_builder = ConfigBuilder::new();
278
279		let listen_addr_len = config.network_config.listen_addresses.len();
280		let mut tcp_addresses = Vec::with_capacity(listen_addr_len);
281		let mut websocket_addresses = Vec::with_capacity(listen_addr_len);
282		let mut webrtc_addresses = Vec::with_capacity(listen_addr_len);
283
284		for addr in &config.network_config.listen_addresses {
285			let mut iter = addr.iter();
286
287			let ip_version = iter.next();
288			let Some(NetworkProtocol::Ip4(_) | NetworkProtocol::Ip6(_)) = ip_version else {
289				log::error!(
290					target: LOG_TARGET,
291					"unknown protocol {ip_version:?}, ignoring {addr:?}",
292				);
293				continue;
294			};
295
296			let transport_layer = iter.next();
297			let protocol_type = iter.next();
298
299			match (&transport_layer, &protocol_type) {
300				// Plain TCP address.
301				(Some(NetworkProtocol::Tcp(_)), Some(NetworkProtocol::P2p(_)) | None) => {
302					tcp_addresses.push(addr.clone());
303				},
304				// Websocket address.
305				(
306					Some(NetworkProtocol::Tcp(_)),
307					Some(NetworkProtocol::Ws(_) | NetworkProtocol::Wss(_)),
308				) => {
309					websocket_addresses.push(addr.clone());
310				},
311				// WebRTC-Direct address.
312				(Some(NetworkProtocol::Udp(_)), Some(NetworkProtocol::WebRTCDirect)) => {
313					// An address carrying anything past `webrtc-direct` is rejected.
314					webrtc::validate_listen_address(addr)?;
315					webrtc_addresses.push(addr.clone());
316				},
317				_ => {
318					log::error!(
319						target: LOG_TARGET,
320						"unknown transport layer {transport_layer:?} and protocol type {protocol_type:?}, ignoring {addr:?}",
321					);
322				},
323			};
324		}
325
326		let mut config_builder = config_builder
327			.with_websocket(WebSocketTransportConfig {
328				listen_addresses: websocket_addresses.into_iter().map(Into::into).collect(),
329				yamux_config: litep2p::yamux::Config::default(),
330				nodelay: true,
331				..Default::default()
332			})
333			.with_tcp(TcpTransportConfig {
334				listen_addresses: tcp_addresses.into_iter().map(Into::into).collect(),
335				yamux_config: litep2p::yamux::Config::default(),
336				nodelay: true,
337				..Default::default()
338			});
339
340		if !webrtc_addresses.is_empty() {
341			let certificate =
342				webrtc::derive_certificate(keypair.secret()).map_err(Error::Litep2p)?;
343			log::info!(target: LOG_TARGET, "WebRTC certhash: {}", certificate.certhash_b64());
344			config_builder = config_builder.with_webrtc(WebRtcTransportConfig {
345				listen_addresses: webrtc_addresses.into_iter().map(Into::into).collect(),
346				certificate: Some(certificate),
347				..Default::default()
348			});
349		}
350
351		Ok(config_builder.with_keypair(keypair))
352	}
353}
354
355#[async_trait::async_trait]
356impl<B: BlockT + 'static, H: ExHashT> NetworkBackend<B, H> for Litep2pNetworkBackend {
357	const SUPPORTS_IPFS: bool = true;
358
359	type NotificationProtocolConfig = NotificationProtocolConfig;
360	type RequestResponseProtocolConfig = RequestResponseConfig;
361	type NetworkService<Block, Hash> = Arc<Litep2pNetworkService>;
362	type PeerStore = Peerstore;
363
364	fn new(mut params: Params<B, H, Self>) -> Result<Self, Error>
365	where
366		Self: Sized,
367	{
368		// Install the ring CryptoProvider for rustls before any TLS connections are made.
369		if let Err(err) = rustls::crypto::ring::default_provider().install_default() {
370			log::warn!(
371				target: LOG_TARGET,
372				"failed to install ring CryptoProvider for rustls, another provider might be installed: {err:?}",
373			);
374		}
375
376		let (keypair, local_peer_id) =
377			Self::get_keypair(&params.network_config.network_config.node_key)?;
378		let (cmd_tx, cmd_rx) = tracing_unbounded("mpsc_network_worker", 100_000);
379
380		params.network_config.network_config.boot_nodes = params
381			.network_config
382			.network_config
383			.boot_nodes
384			.into_iter()
385			.filter(|boot_node| boot_node.peer_id != local_peer_id.into())
386			.collect();
387		params.network_config.network_config.default_peers_set.reserved_nodes = params
388			.network_config
389			.network_config
390			.default_peers_set
391			.reserved_nodes
392			.into_iter()
393			.filter(|reserved_node| {
394				if reserved_node.peer_id == local_peer_id.into() {
395					log::warn!(
396						target: LOG_TARGET,
397						"Local peer ID used in reserved node, ignoring: {reserved_node}",
398					);
399					false
400				} else {
401					true
402				}
403			})
404			.collect();
405
406		if let Some(path) = &params.network_config.network_config.net_config_path {
407			fs::create_dir_all(path)?;
408		}
409
410		log::info!(target: LOG_TARGET, "Local node identity is: {local_peer_id}");
411		log::info!(target: LOG_TARGET, "Running litep2p network backend");
412
413		params.network_config.sanity_check_addresses()?;
414		params.network_config.sanity_check_bootnodes()?;
415
416		let mut config_builder =
417			Self::configure_transport(&params.network_config, keypair.clone())?;
418		let known_addresses = params.network_config.known_addresses();
419		let peer_store_handle = params.network_config.peer_store_handle();
420		let executor = Arc::new(Litep2pExecutor { executor: params.executor });
421
422		let FullNetworkConfiguration {
423			notification_protocols,
424			request_response_protocols,
425			network_config,
426			..
427		} = params.network_config;
428
429		// initialize notification protocols
430		//
431		// pass the protocol configuration to `Litep2pConfigBuilder` and save the TX channel
432		// to the protocol's `Peerset` together with the protocol name to allow other subsystems
433		// of Polkadot SDK to control connectivity of the notification protocol
434		let block_announce_protocol = params.block_announce_config.protocol_name().clone();
435		let mut notif_protocols = HashMap::from_iter([(
436			params.block_announce_config.protocol_name().clone(),
437			params.block_announce_config.handle,
438		)]);
439
440		// handshake for all but the syncing protocol is set to node role
441		config_builder = notification_protocols
442			.into_iter()
443			.fold(config_builder, |config_builder, mut config| {
444				config.config.set_handshake(Roles::from(&params.role).encode());
445				notif_protocols.insert(config.protocol_name, config.handle);
446
447				config_builder.with_notification_protocol(config.config)
448			})
449			.with_notification_protocol(params.block_announce_config.config);
450
451		// initialize request-response protocols
452		let metrics = match &params.metrics_registry {
453			Some(registry) => Some(register_without_sources(registry)?),
454			None => None,
455		};
456
457		// create channels that are used to send request before initializing protocols so the
458		// senders can be passed onto all request-response protocols
459		//
460		// all protocols must have each others' senders so they can send the fallback request in
461		// case the main protocol is not supported by the remote peer and user specified a fallback
462		let (mut request_response_receivers, request_response_senders): (
463			HashMap<_, _>,
464			HashMap<_, _>,
465		) = request_response_protocols
466			.iter()
467			.map(|config| {
468				let (tx, rx) = tracing_unbounded("outbound-requests", 10_000);
469				((config.protocol_name.clone(), rx), (config.protocol_name.clone(), tx))
470			})
471			.unzip();
472
473		config_builder = request_response_protocols.into_iter().fold(
474			config_builder,
475			|config_builder, config| {
476				let (protocol_config, handle) = RequestResponseConfigBuilder::new(
477					Litep2pProtocolName::from(config.protocol_name.clone()),
478				)
479				.with_max_size(cmp::max(config.max_request_size, config.max_response_size) as usize)
480				.with_fallback_names(config.fallback_names.into_iter().map(From::from).collect())
481				.with_timeout(config.request_timeout)
482				.build();
483
484				let protocol = RequestResponseProtocol::new(
485					config.protocol_name.clone(),
486					handle,
487					Arc::clone(&peer_store_handle),
488					config.inbound_queue,
489					request_response_receivers
490						.remove(&config.protocol_name)
491						.expect("receiver exists as it was just added and there are no duplicate protocols; qed"),
492					request_response_senders.clone(),
493					metrics.clone(),
494				);
495
496				executor.run(Box::pin(async move {
497					protocol.run().await;
498				}));
499
500				config_builder.with_request_response_protocol(protocol_config)
501			},
502		);
503
504		// collect known addresses
505		let known_addresses: HashMap<litep2p::PeerId, Vec<Multiaddr>> =
506			known_addresses.into_iter().fold(HashMap::new(), |mut acc, (peer, address)| {
507				let address = match address.iter().last() {
508					Some(
509						NetworkProtocol::Ws(_) | NetworkProtocol::Wss(_) | NetworkProtocol::Tcp(_),
510					) => address.with(NetworkProtocol::P2p(peer.into())),
511					Some(NetworkProtocol::WebRTCDirect | NetworkProtocol::Certhash(_)) => {
512						address.with(NetworkProtocol::P2p(peer.into()))
513					},
514					Some(NetworkProtocol::P2p(_)) => address,
515					_ => return acc,
516				};
517
518				acc.entry(peer.into()).or_default().push(address.into());
519				peer_store_handle.add_known_peer(peer);
520
521				acc
522			});
523
524		// enable ipfs ping, identify and kademlia, and potentially mdns if user enabled it
525		let listen_addresses = Arc::new(Default::default());
526		let (discovery, ping_config, identify_config, kademlia_config, maybe_mdns_config) =
527			Discovery::new(
528				local_peer_id,
529				&network_config,
530				params.genesis_hash,
531				params.fork_id.as_deref(),
532				&params.protocol_id,
533				known_addresses.clone(),
534				Arc::clone(&listen_addresses),
535				Arc::clone(&peer_store_handle),
536			);
537
538		if let Some(ipfs) = params.ipfs_config {
539			config_builder = config_builder.with_libp2p_bitswap(ipfs.litep2p_bitswap_config);
540
541			if !ipfs.bootnodes.is_empty() {
542				let (ipfs_dht, kad_config) = IpfsDht::new(ipfs.bootnodes, ipfs.block_provider);
543				config_builder = config_builder.with_libp2p_kademlia(kad_config);
544				executor.run(Box::pin(ipfs_dht.run()));
545			} else {
546				log::warn!(
547					target: LOG_TARGET,
548					"Not starting IPFS DHT publisher because no IPFS bootnodes are configured. \
549					 Only direct Bitswap requests will be handled.",
550				);
551			}
552		}
553
554		config_builder = config_builder
555			.with_known_addresses(known_addresses.clone().into_iter())
556			.with_libp2p_ping(ping_config)
557			.with_libp2p_identify(identify_config)
558			.with_libp2p_kademlia(kademlia_config)
559			.with_connection_limits(ConnectionLimitsConfig::default().max_incoming_connections(
560				Some(crate::MAX_CONNECTIONS_ESTABLISHED_INCOMING as usize),
561			))
562			.with_keep_alive_timeout(network_config.idle_connection_timeout)
563			// Use system DNS resolver to enable intranet domain resolution and administrator
564			// control over DNS lookup.
565			.with_system_resolver()
566			.with_executor(executor);
567
568		if let Some(config) = maybe_mdns_config {
569			config_builder = config_builder.with_mdns(config);
570		}
571
572		let litep2p =
573			Litep2p::new(config_builder.build()).map_err(|error| Error::Litep2p(error))?;
574
575		litep2p.listen_addresses().for_each(|address| {
576			log::debug!(target: LOG_TARGET, "listening on: {address}");
577
578			listen_addresses.write().insert(address.clone());
579		});
580
581		let public_addresses = litep2p.public_addresses();
582		for address in network_config.public_addresses.iter() {
583			if let Err(err) = public_addresses.add_address(address.clone().into()) {
584				log::warn!(
585					target: LOG_TARGET,
586					"failed to add public address {address:?}: {err:?}",
587				);
588			}
589		}
590
591		let network_service = Arc::new(Litep2pNetworkService::new(
592			local_peer_id,
593			keypair.clone(),
594			cmd_tx,
595			Arc::clone(&peer_store_handle),
596			notif_protocols.clone(),
597			block_announce_protocol.clone(),
598			request_response_senders,
599			Arc::clone(&listen_addresses),
600			public_addresses,
601		));
602
603		// register rest of the metrics now that `Litep2p` has been created
604		let num_connected = Arc::new(Default::default());
605		let bandwidth: Arc<dyn BandwidthSink> =
606			Arc::new(Litep2pBandwidthSink { sink: litep2p.bandwidth_sink() });
607
608		if let Some(registry) = &params.metrics_registry {
609			MetricSources::register(registry, bandwidth)?;
610		}
611
612		Ok(Self {
613			network_service,
614			cmd_rx,
615			metrics,
616			peerset_handles: notif_protocols,
617			num_connected,
618			discovery,
619			pending_queries: HashMap::new(),
620			peerstore_handle: peer_store_handle,
621			block_announce_protocol,
622			event_streams: out_events::OutChannels::new(params.metrics_registry.as_ref())?,
623			peers: HashMap::new(),
624			litep2p,
625		})
626	}
627
628	fn network_service(&self) -> Arc<dyn NetworkService> {
629		Arc::clone(&self.network_service)
630	}
631
632	fn peer_store(
633		bootnodes: Vec<sc_network_types::PeerId>,
634		metrics_registry: Option<Registry>,
635	) -> Self::PeerStore {
636		Peerstore::new(bootnodes, metrics_registry)
637	}
638
639	fn register_notification_metrics(registry: Option<&Registry>) -> NotificationMetrics {
640		NotificationMetrics::new(registry)
641	}
642
643	/// Create notification protocol configuration for `protocol`.
644	fn notification_config(
645		protocol_name: ProtocolName,
646		fallback_names: Vec<ProtocolName>,
647		max_notification_size: u64,
648		handshake: Option<NotificationHandshake>,
649		set_config: SetConfig,
650		metrics: NotificationMetrics,
651		peerstore_handle: Arc<dyn PeerStoreProvider>,
652	) -> (Self::NotificationProtocolConfig, Box<dyn NotificationService>) {
653		Self::NotificationProtocolConfig::new(
654			protocol_name,
655			fallback_names,
656			max_notification_size as usize,
657			handshake,
658			set_config,
659			metrics,
660			peerstore_handle,
661		)
662	}
663
664	/// Create request-response protocol configuration.
665	fn request_response_config(
666		protocol_name: ProtocolName,
667		fallback_names: Vec<ProtocolName>,
668		max_request_size: u64,
669		max_response_size: u64,
670		request_timeout: Duration,
671		inbound_queue: Option<async_channel::Sender<IncomingRequest>>,
672	) -> Self::RequestResponseProtocolConfig {
673		Self::RequestResponseProtocolConfig::new(
674			protocol_name,
675			fallback_names,
676			max_request_size,
677			max_response_size,
678			request_timeout,
679			inbound_queue,
680		)
681	}
682
683	/// Start [`Litep2pNetworkBackend`] event loop.
684	async fn run(mut self) {
685		log::debug!(target: LOG_TARGET, "starting litep2p network backend");
686
687		loop {
688			let num_connected_peers = self
689				.peerset_handles
690				.get(&self.block_announce_protocol)
691				.map_or(0usize, |handle| handle.connected_peers.load(Ordering::Relaxed));
692			self.num_connected.store(num_connected_peers, Ordering::Relaxed);
693
694			tokio::select! {
695				command = self.cmd_rx.next() => match command {
696					None => return,
697					Some(command) => match command {
698						NetworkServiceCommand::FindClosestPeers { target } => {
699							let query_id = self.discovery.find_node(target.into()).await;
700							self.pending_queries.insert(query_id, KadQuery::FindNode(target, Instant::now()));
701						}
702						NetworkServiceCommand::GetValue{ key } => {
703							let query_id = self.discovery.get_value(key.clone()).await;
704							self.pending_queries.insert(query_id, KadQuery::GetValue(key, Instant::now()));
705						}
706						NetworkServiceCommand::PutValue { key, value } => {
707							let query_id = self.discovery.put_value(key.clone(), value).await;
708							self.pending_queries.insert(query_id, KadQuery::PutValue(key, Instant::now()));
709						}
710						NetworkServiceCommand::PutValueTo { record, peers, update_local_storage} => {
711							let kademlia_key = record.key.clone();
712							let query_id = self.discovery.put_value_to_peers(record.into(), peers, update_local_storage).await;
713							self.pending_queries.insert(query_id, KadQuery::PutValue(kademlia_key, Instant::now()));
714						}
715						NetworkServiceCommand::StoreRecord { key, value, publisher, expires } => {
716							self.discovery.store_record(key, value, publisher.map(Into::into), expires).await;
717						}
718						NetworkServiceCommand::StartProviding { key } => {
719							let query_id = self.discovery.start_providing(key.clone()).await;
720							self.pending_queries.insert(query_id, KadQuery::AddProvider(key, Instant::now()));
721						}
722						NetworkServiceCommand::StopProviding { key } => {
723							self.discovery.stop_providing(key).await;
724						}
725						NetworkServiceCommand::GetProviders { key } => {
726							let query_id = self.discovery.get_providers(key.clone()).await;
727							self.pending_queries.insert(query_id, KadQuery::GetProviders(key, Instant::now()));
728						}
729						NetworkServiceCommand::EventStream { tx } => {
730							self.event_streams.push(tx);
731						}
732						NetworkServiceCommand::Status { tx } => {
733							let _ = tx.send(NetworkStatus {
734								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								total_bytes_inbound: self.litep2p.bandwidth_sink().inbound() as u64,
739								total_bytes_outbound: self.litep2p.bandwidth_sink().outbound() as u64,
740							});
741						}
742						NetworkServiceCommand::AddPeersToReservedSet {
743							protocol,
744							peers,
745						} => {
746							let peers = self.add_addresses(peers.into_iter().map(Into::into));
747
748							match self.peerset_handles.get(&protocol) {
749								Some(handle) => {
750									let _ = handle.tx.unbounded_send(PeersetCommand::AddReservedPeers { peers });
751								}
752								None => log::warn!(target: LOG_TARGET, "protocol {protocol} doens't exist"),
753							};
754						}
755						NetworkServiceCommand::AddKnownAddress { peer, address } => {
756							let mut address: Multiaddr = address.into();
757
758							if !address.iter().any(|protocol| std::matches!(protocol, Protocol::P2p(_))) {
759								address.push(Protocol::P2p(litep2p::PeerId::from(peer).into()));
760							}
761
762							if self.litep2p.add_known_address(peer.into(), iter::once(address.clone())) > 0 {
763								// libp2p backend generates `DiscoveryOut::Discovered(peer_id)`
764								// event when a new address is added for a peer, which leads to the
765								// peer being added to peerstore. Do the same directly here.
766								self.peerstore_handle.add_known_peer(peer);
767							} else {
768								log::debug!(
769									target: LOG_TARGET,
770									"couldn't add known address ({address}) for {peer:?}, unsupported transport"
771								);
772							}
773						},
774						NetworkServiceCommand::SetReservedPeers { protocol, peers } => {
775							let peers = self.add_addresses(peers.into_iter().map(Into::into));
776
777							match self.peerset_handles.get(&protocol) {
778								Some(handle) => {
779									let _ = handle.tx.unbounded_send(PeersetCommand::SetReservedPeers { peers });
780								}
781								None => log::warn!(target: LOG_TARGET, "protocol {protocol} doens't exist"),
782							}
783
784						},
785						NetworkServiceCommand::DisconnectPeer {
786							protocol,
787							peer,
788						} => {
789							let Some(handle) = self.peerset_handles.get(&protocol) else {
790								log::warn!(target: LOG_TARGET, "protocol {protocol} doens't exist");
791								continue
792							};
793
794							let _ = handle.tx.unbounded_send(PeersetCommand::DisconnectPeer { peer });
795						}
796						NetworkServiceCommand::SetReservedOnly {
797							protocol,
798							reserved_only,
799						} => {
800							let Some(handle) = self.peerset_handles.get(&protocol) else {
801								log::warn!(target: LOG_TARGET, "protocol {protocol} doens't exist");
802								continue
803							};
804
805							let _ = handle.tx.unbounded_send(PeersetCommand::SetReservedOnly { reserved_only });
806						}
807						NetworkServiceCommand::RemoveReservedPeers {
808							protocol,
809							peers,
810						} => {
811							let Some(handle) = self.peerset_handles.get(&protocol) else {
812								log::warn!(target: LOG_TARGET, "protocol {protocol} doens't exist");
813								continue
814							};
815
816							let _ = handle.tx.unbounded_send(PeersetCommand::RemoveReservedPeers { peers });
817						}
818					}
819				},
820				event = self.discovery.next() => match event {
821					None => return,
822					Some(DiscoveryEvent::Discovered { addresses }) => {
823						// if at least one address was added for the peer, report the peer to `Peerstore`
824						for (peer, addresses) in Litep2pNetworkBackend::parse_addresses(addresses.into_iter()) {
825							if self.litep2p.add_known_address(peer.into(), addresses.clone().into_iter()) > 0 {
826								self.peerstore_handle.add_known_peer(peer);
827							}
828						}
829					}
830					Some(DiscoveryEvent::RoutingTableUpdate { peers }) => {
831						let peers = peers.into_iter().map(Into::into).collect::<Vec<_>>();
832
833						for peer in &peers {
834							self.peerstore_handle.add_known_peer(*peer);
835						}
836
837						if !peers.is_empty() {
838							self.event_streams.send(Event::PeerRoutingTableUpdate(peers));
839						}
840					}
841					Some(DiscoveryEvent::FindNodeSuccess { query_id, target, peers }) => {
842						match self.pending_queries.remove(&query_id) {
843							Some(KadQuery::FindNode(_, started)) => {
844								log::trace!(
845									target: LOG_TARGET,
846									"`FIND_NODE` for {target:?} ({query_id:?}) succeeded",
847								);
848
849								self.event_streams.send(
850									Event::Dht(
851										DhtEvent::ClosestPeersFound(
852											target.into(),
853											peers
854												.into_iter()
855												.map(|(peer, addrs)| (
856													peer.into(),
857													addrs.into_iter().map(Into::into).collect(),
858												))
859												.collect(),
860										)
861									)
862								);
863
864								if let Some(ref metrics) = self.metrics {
865									metrics
866										.kademlia_query_duration
867										.with_label_values(&["node-find"])
868										.observe(started.elapsed().as_secs_f64());
869								}
870							},
871							query => {
872								log::error!(
873									target: LOG_TARGET,
874									"Missing/invalid pending query for `FIND_NODE`: {query:?}"
875								);
876								debug_assert!(false);
877							}
878						}
879					},
880					Some(DiscoveryEvent::GetRecordPartialResult { query_id, record }) => {
881						if !self.pending_queries.contains_key(&query_id) {
882							log::error!(
883								target: LOG_TARGET,
884								"Missing/invalid pending query for `GET_VALUE` partial result: {query_id:?}"
885							);
886
887							continue
888						}
889
890						let peer_id: sc_network_types::PeerId = record.peer.into();
891						let record = PeerRecord {
892							record: P2PRecord {
893								key: record.record.key.to_vec().into(),
894								value: record.record.value,
895								publisher: record.record.publisher.map(|peer_id| {
896									let peer_id: sc_network_types::PeerId = peer_id.into();
897									peer_id.into()
898								}),
899								expires: record.record.expires,
900							},
901							peer: Some(peer_id.into()),
902						};
903
904						self.event_streams.send(
905							Event::Dht(
906								DhtEvent::ValueFound(
907									record.into()
908								)
909							)
910						);
911					}
912					Some(DiscoveryEvent::GetRecordSuccess { query_id }) => {
913						match self.pending_queries.remove(&query_id) {
914							Some(KadQuery::GetValue(key, started)) => {
915								log::trace!(
916									target: LOG_TARGET,
917									"`GET_VALUE` for {key:?} ({query_id:?}) succeeded",
918								);
919
920								if let Some(ref metrics) = self.metrics {
921									metrics
922										.kademlia_query_duration
923										.with_label_values(&["value-get"])
924										.observe(started.elapsed().as_secs_f64());
925								}
926							},
927							query => {
928								log::error!(
929									target: LOG_TARGET,
930									"Missing/invalid pending query for `GET_VALUE`: {query:?}"
931								);
932								debug_assert!(false);
933							},
934						}
935					}
936					Some(DiscoveryEvent::PutRecordSuccess { query_id }) => {
937						match self.pending_queries.remove(&query_id) {
938							Some(KadQuery::PutValue(key, started)) => {
939								log::trace!(
940									target: LOG_TARGET,
941									"`PUT_VALUE` for {key:?} ({query_id:?}) succeeded",
942								);
943
944								self.event_streams.send(Event::Dht(
945									DhtEvent::ValuePut(key)
946								));
947
948								if let Some(ref metrics) = self.metrics {
949									metrics
950										.kademlia_query_duration
951										.with_label_values(&["value-put"])
952										.observe(started.elapsed().as_secs_f64());
953								}
954							},
955							query => {
956								log::error!(
957									target: LOG_TARGET,
958									"Missing/invalid pending query for `PUT_VALUE`: {query:?}"
959								);
960								debug_assert!(false);
961							}
962						}
963					}
964					Some(DiscoveryEvent::GetProvidersSuccess { query_id, providers }) => {
965						match self.pending_queries.remove(&query_id) {
966							Some(KadQuery::GetProviders(key, started)) => {
967								log::trace!(
968									target: LOG_TARGET,
969									"`GET_PROVIDERS` for {key:?} ({query_id:?}) succeeded",
970								);
971
972								// We likely requested providers to connect to them,
973								// so let's add their addresses to litep2p's transport manager.
974								// Consider also looking the addresses of providers up with `FIND_NODE`
975								// query, as it can yield more up to date addresses.
976								providers.iter().for_each(|p| {
977									self.litep2p.add_known_address(p.peer, p.addresses.clone().into_iter());
978								});
979
980								self.event_streams.send(Event::Dht(
981									DhtEvent::ProvidersFound(
982										key.clone().into(),
983										providers.into_iter().map(|p| p.peer.into()).collect()
984									)
985								));
986
987								// litep2p returns all providers in a single event, so we let
988								// subscribers know no more providers will be yielded.
989								self.event_streams.send(Event::Dht(
990									DhtEvent::NoMoreProviders(key.into())
991								));
992
993								if let Some(ref metrics) = self.metrics {
994									metrics
995										.kademlia_query_duration
996										.with_label_values(&["providers-get"])
997										.observe(started.elapsed().as_secs_f64());
998								}
999							},
1000							query => {
1001								log::error!(
1002									target: LOG_TARGET,
1003									"Missing/invalid pending query for `GET_PROVIDERS`: {query:?}"
1004								);
1005								debug_assert!(false);
1006							}
1007						}
1008					}
1009					Some(DiscoveryEvent::AddProviderSuccess { query_id, provided_key }) => {
1010						match self.pending_queries.remove(&query_id) {
1011							Some(KadQuery::AddProvider(key, started)) => {
1012								debug_assert_eq!(key, provided_key.into());
1013
1014								log::trace!(
1015									target: LOG_TARGET,
1016									"`ADD_PROVIDER` for {key:?} ({query_id:?}) succeeded",
1017								);
1018
1019								self.event_streams.send(Event::Dht(
1020									DhtEvent::StartedProviding(key.into())
1021								));
1022
1023								if let Some(ref metrics) = self.metrics {
1024									metrics
1025										.kademlia_query_duration
1026										.with_label_values(&["provider-add"])
1027										.observe(started.elapsed().as_secs_f64());
1028								}
1029							}
1030							Some(_) => {
1031								log::error!(
1032									target: LOG_TARGET,
1033									"Invalid pending query for `ADD_PROVIDER`: {query_id:?}"
1034								);
1035								debug_assert!(false);
1036							}
1037							None => {
1038								log::trace!(
1039									target: LOG_TARGET,
1040									"`ADD_PROVIDER` for key {provided_key:?} ({query_id:?}) succeeded (republishing)",
1041								);
1042							}
1043						}
1044					}
1045					Some(DiscoveryEvent::QueryFailed { query_id }) => {
1046						match self.pending_queries.remove(&query_id) {
1047							Some(KadQuery::FindNode(peer_id, started)) => {
1048								log::debug!(
1049									target: LOG_TARGET,
1050									"`FIND_NODE` ({query_id:?}) failed for target {peer_id:?}",
1051								);
1052
1053								self.event_streams.send(Event::Dht(
1054									DhtEvent::ClosestPeersNotFound(peer_id.into())
1055								));
1056
1057								if let Some(ref metrics) = self.metrics {
1058									metrics
1059										.kademlia_query_duration
1060										.with_label_values(&["node-find-failed"])
1061										.observe(started.elapsed().as_secs_f64());
1062								}
1063							},
1064							Some(KadQuery::GetValue(key, started)) => {
1065								log::debug!(
1066									target: LOG_TARGET,
1067									"`GET_VALUE` ({query_id:?}) failed for key {key:?}",
1068								);
1069
1070								self.event_streams.send(Event::Dht(
1071									DhtEvent::ValueNotFound(key)
1072								));
1073
1074								if let Some(ref metrics) = self.metrics {
1075									metrics
1076										.kademlia_query_duration
1077										.with_label_values(&["value-get-failed"])
1078										.observe(started.elapsed().as_secs_f64());
1079								}
1080							},
1081							Some(KadQuery::PutValue(key, started)) => {
1082								log::debug!(
1083									target: LOG_TARGET,
1084									"`PUT_VALUE` ({query_id:?}) failed for key {key:?}",
1085								);
1086
1087								self.event_streams.send(Event::Dht(
1088									DhtEvent::ValuePutFailed(key)
1089								));
1090
1091								if let Some(ref metrics) = self.metrics {
1092									metrics
1093										.kademlia_query_duration
1094										.with_label_values(&["value-put-failed"])
1095										.observe(started.elapsed().as_secs_f64());
1096								}
1097							},
1098							Some(KadQuery::GetProviders(key, started)) => {
1099								log::debug!(
1100									target: LOG_TARGET,
1101									"`GET_PROVIDERS` ({query_id:?}) failed for key {key:?}"
1102								);
1103
1104								self.event_streams.send(Event::Dht(
1105									DhtEvent::ProvidersNotFound(key)
1106								));
1107
1108								if let Some(ref metrics) = self.metrics {
1109									metrics
1110										.kademlia_query_duration
1111										.with_label_values(&["providers-get-failed"])
1112										.observe(started.elapsed().as_secs_f64());
1113								}
1114							},
1115							Some(KadQuery::AddProvider(key, started)) => {
1116								log::debug!(
1117									target: LOG_TARGET,
1118									"`ADD_PROVIDER` ({query_id:?}) failed with key {key:?}",
1119								);
1120
1121								self.event_streams.send(Event::Dht(
1122									DhtEvent::StartProvidingFailed(key)
1123								));
1124
1125								if let Some(ref metrics) = self.metrics {
1126									metrics
1127										.kademlia_query_duration
1128										.with_label_values(&["provider-add-failed"])
1129										.observe(started.elapsed().as_secs_f64());
1130								}
1131							},
1132							None => {
1133								log::debug!(
1134									target: LOG_TARGET,
1135									"non-existent query (likely republishing a provider) failed ({query_id:?})",
1136								);
1137							}
1138						}
1139					}
1140					Some(DiscoveryEvent::Identified { peer, listen_addresses, supported_protocols, .. }) => {
1141						self.event_streams.send(Event::PeerIdentified {
1142							peer: peer.into(),
1143							supported_protocols: supported_protocols.iter().cloned().map(Into::into).collect(),
1144						});
1145						self.discovery.add_self_reported_address(peer, supported_protocols, listen_addresses).await;
1146					}
1147					Some(DiscoveryEvent::ExternalAddressDiscovered { address }) => {
1148						match self.litep2p.public_addresses().add_address(address.clone().into()) {
1149							Ok(inserted) => if inserted {
1150								log::info!(target: LOG_TARGET, "๐Ÿ” Discovered new external address for our node: {address}");
1151							},
1152							Err(err) => {
1153								log::warn!(
1154									target: LOG_TARGET,
1155									"๐Ÿ” Failed to add discovered external address {address:?}: {err:?}",
1156								);
1157							},
1158						}
1159					}
1160					Some(DiscoveryEvent::ExternalAddressExpired{ address }) => {
1161						let local_peer_id = self.litep2p.local_peer_id();
1162
1163						// Litep2p requires the peer ID to be present in the address.
1164						let address = if !std::matches!(address.iter().last(), Some(Protocol::P2p(_))) {
1165							address.with(Protocol::P2p((*local_peer_id).into()))
1166						} else {
1167							address
1168						};
1169
1170						if self.litep2p.public_addresses().remove_address(&address) {
1171							log::info!(target: LOG_TARGET, "๐Ÿ” Expired external address for our node: {address}");
1172						} else {
1173							log::warn!(
1174								target: LOG_TARGET,
1175								"๐Ÿ” Failed to remove expired external address {address:?}"
1176							);
1177						}
1178					}
1179					Some(DiscoveryEvent::Ping { peer, rtt }) => {
1180						log::trace!(
1181							target: LOG_TARGET,
1182							"ping time with {peer:?}: {rtt:?}",
1183						);
1184					}
1185					Some(DiscoveryEvent::IncomingRecord { record: Record { key, value, publisher, expires }} ) => {
1186						self.event_streams.send(Event::Dht(
1187							DhtEvent::PutRecordRequest(
1188								key.into(),
1189								value,
1190								publisher.map(Into::into),
1191								expires,
1192							)
1193						));
1194					},
1195
1196					Some(DiscoveryEvent::RandomKademliaStarted) => {
1197						if let Some(metrics) = self.metrics.as_ref() {
1198							metrics.kademlia_random_queries_total.inc();
1199						}
1200					}
1201				},
1202				event = self.litep2p.next_event() => match event {
1203					Some(Litep2pEvent::ConnectionEstablished { peer, endpoint }) => {
1204						let Some(metrics) = &self.metrics else {
1205							continue;
1206						};
1207
1208						let direction = match endpoint {
1209							Endpoint::Dialer { .. } => "out",
1210							Endpoint::Listener { .. } => {
1211								// Increment incoming connections counter.
1212								//
1213								// Note: For litep2p these are represented by established negotiated connections,
1214								// while for libp2p (legacy) these represent not-yet-negotiated connections.
1215								metrics.incoming_connections_total.inc();
1216
1217								"in"
1218							},
1219						};
1220						metrics.connections_opened_total.with_label_values(&[direction]).inc();
1221
1222						match self.peers.entry(peer) {
1223							Entry::Vacant(entry) => {
1224								entry.insert(ConnectionContext {
1225									endpoints: HashMap::from_iter([(endpoint.connection_id(), endpoint)]),
1226									num_connections: 1usize,
1227								});
1228								metrics.distinct_peers_connections_opened_total.inc();
1229							}
1230							Entry::Occupied(entry) => {
1231								let entry = entry.into_mut();
1232								entry.num_connections += 1;
1233								entry.endpoints.insert(endpoint.connection_id(), endpoint);
1234							}
1235						}
1236					}
1237					Some(Litep2pEvent::ConnectionClosed { peer, connection_id }) => {
1238						let Some(metrics) = &self.metrics else {
1239							continue;
1240						};
1241
1242						let Some(context) = self.peers.get_mut(&peer) else {
1243							log::debug!(target: LOG_TARGET, "unknown peer disconnected: {peer:?} ({connection_id:?})");
1244							continue
1245						};
1246
1247						let direction = match context.endpoints.remove(&connection_id) {
1248							None => {
1249								log::debug!(target: LOG_TARGET, "connection {connection_id:?} doesn't exist for {peer:?} ");
1250								continue
1251							}
1252							Some(endpoint) => {
1253								context.num_connections -= 1;
1254
1255								match endpoint {
1256									Endpoint::Dialer { .. } => "out",
1257									Endpoint::Listener { .. } => "in",
1258								}
1259							}
1260						};
1261
1262						metrics.connections_closed_total.with_label_values(&[direction, "actively-closed"]).inc();
1263
1264						if context.num_connections == 0 {
1265							self.peers.remove(&peer);
1266							metrics.distinct_peers_connections_closed_total.inc();
1267						}
1268					}
1269					Some(Litep2pEvent::DialFailure { address, error }) => {
1270						log::debug!(
1271							target: LOG_TARGET,
1272							"failed to dial peer at {address:?}: {error:?}",
1273						);
1274
1275						if let Some(metrics) = &self.metrics {
1276							let reason = match error {
1277								DialError::Timeout => "timeout",
1278								DialError::AddressError(_) => "invalid-address",
1279								DialError::DnsError(_) => "cannot-resolve-dns",
1280								DialError::NegotiationError(error) => match error {
1281									NegotiationError::Timeout => "timeout",
1282									NegotiationError::PeerIdMissing => "missing-peer-id",
1283									NegotiationError::StateMismatch => "state-mismatch",
1284									NegotiationError::PeerIdMismatch(_,_) => "peer-id-missmatch",
1285									NegotiationError::MultistreamSelectError(_) => "multistream-select-error",
1286									NegotiationError::SnowError(_) => "noise-error",
1287									NegotiationError::ParseError(_) => "parse-error",
1288									NegotiationError::IoError(_) => "io-error",
1289									NegotiationError::WebSocket(_) => "webscoket-error",
1290									NegotiationError::BadSignature => "bad-signature",
1291								}
1292							};
1293
1294							metrics.pending_connections_errors_total.with_label_values(&[&reason]).inc();
1295						}
1296					}
1297					Some(Litep2pEvent::ListDialFailures { errors }) => {
1298						log::debug!(
1299							target: LOG_TARGET,
1300							"failed to dial peer on multiple addresses {errors:?}",
1301						);
1302
1303						if let Some(metrics) = &self.metrics {
1304							metrics.pending_connections_errors_total.with_label_values(&["transport-errors"]).inc();
1305						}
1306					}
1307					None => {
1308						log::error!(
1309								target: LOG_TARGET,
1310								"Litep2p backend terminated"
1311						);
1312						return
1313					}
1314				},
1315			}
1316		}
1317	}
1318}
1319
1320#[cfg(test)]
1321mod tests {
1322	use super::*;
1323	use crate::{
1324		config::{ed25519, NetworkConfiguration, ProtocolId, Role, Secret},
1325		service::traits::NetworkStateInfo,
1326	};
1327	use sc_network_types::{
1328		multiaddr::Multiaddr as NetworkMultiaddr, multihash::Multihash as NetworkMultihash,
1329	};
1330	use sp_core::H256;
1331	use substrate_test_runtime_client::runtime::Block;
1332
1333	/// `--listen-addr` of a node behind a NAT. Port `0` so concurrent tests don't collide.
1334	const WEBRTC_LISTEN_ADDRESS: &str = "/ip4/127.0.0.1/udp/0/webrtc-direct";
1335
1336	/// `--public-addr` of that node: the routable coordinate peers are told to dial.
1337	/// This is the shape an operator supplies when the node sits behind a proxy.
1338	const WEBRTC_PUBLIC_ADDRESS: &str = "/ip4/203.0.113.9/udp/31234/webrtc-direct";
1339
1340	/// The `/certhash` component of `address`.
1341	fn certhash(address: &NetworkMultiaddr) -> Option<NetworkMultihash> {
1342		address.iter().find_map(|protocol| match protocol {
1343			NetworkProtocol::Certhash(hash) => Some(hash),
1344			_ => None,
1345		})
1346	}
1347
1348	/// Bring up the backend from `network_config` through its real entry point.
1349	fn start_backend(
1350		network_config: &NetworkConfiguration,
1351	) -> Result<Litep2pNetworkBackend, Error> {
1352		let config = FullNetworkConfiguration::<Block, H256, Litep2pNetworkBackend>::new(
1353			network_config,
1354			None,
1355		);
1356
1357		let (block_announce_config, _notification_service) =
1358			<Litep2pNetworkBackend as NetworkBackend<Block, H256>>::notification_config(
1359				"/block-announces/1".into(),
1360				vec![],
1361				1024,
1362				None,
1363				SetConfig::default(),
1364				NotificationMetrics::new(None),
1365				config.peer_store_handle(),
1366			);
1367
1368		<Litep2pNetworkBackend as NetworkBackend<Block, H256>>::new(Params {
1369			role: Role::Full,
1370			executor: Box::new(|future| {
1371				tokio::spawn(future);
1372			}),
1373			network_config: config,
1374			protocol_id: ProtocolId::from("test"),
1375			genesis_hash: H256::zero(),
1376			fork_id: None,
1377			metrics_registry: None,
1378			block_announce_config,
1379			ipfs_config: None,
1380			notification_metrics: NotificationMetrics::new(None),
1381		})
1382	}
1383
1384	/// Both the address the node binds and the one it tells peers to dial must carry the node's
1385	/// `/certhash`: a `webrtc-direct` dialer has no other way to verify the DTLS handshake.
1386	#[tokio::test]
1387	async fn webrtc_addresses_advertised_with_certhash() {
1388		let mut network_config = NetworkConfiguration::new_local();
1389		network_config.listen_addresses = vec![WEBRTC_LISTEN_ADDRESS.parse().unwrap()];
1390		network_config.public_addresses = vec![WEBRTC_PUBLIC_ADDRESS.parse().unwrap()];
1391		// Fixed, so the certificate the node will present can be derived here as well.
1392		network_config.node_key = NodeKeyConfig::Ed25519(Secret::Input(
1393			ed25519::SecretKey::try_from_bytes([7u8; 32]).unwrap(),
1394		));
1395		network_config.validate_and_complete_addresses().unwrap();
1396
1397		let (keypair, _peer_id) =
1398			Litep2pNetworkBackend::get_keypair(&network_config.node_key).unwrap();
1399		let node_certhash: NetworkMultihash =
1400			webrtc::derive_certificate(keypair.secret()).unwrap().certhash().into();
1401
1402		// Held for the duration of the test: dropping it closes the node's sockets.
1403		let backend = start_backend(&network_config).unwrap();
1404		let network_service =
1405			<Litep2pNetworkBackend as NetworkBackend<Block, H256>>::network_service(&backend);
1406
1407		let advertised = network_service.listen_addresses();
1408		assert_eq!(advertised.len(), 1);
1409		assert_eq!(certhash(&advertised[0]), Some(node_certhash));
1410
1411		let advertised = network_service.external_addresses();
1412		assert_eq!(advertised.len(), 1);
1413		assert_eq!(certhash(&advertised[0]), Some(node_certhash));
1414	}
1415
1416	/// The default node key (`Secret::New`) generates a fresh key per `into_keypair()` call:
1417	/// completing the addresses must pin the resolved key, or the backend would serve a
1418	/// certificate matching neither the advertised `/certhash` nor each other's.
1419	#[tokio::test]
1420	async fn webrtc_certhash_consistent_with_default_node_key() {
1421		let mut network_config = NetworkConfiguration::new_local();
1422		network_config.listen_addresses = vec![WEBRTC_LISTEN_ADDRESS.parse().unwrap()];
1423		network_config.public_addresses = vec![WEBRTC_PUBLIC_ADDRESS.parse().unwrap()];
1424		network_config.validate_and_complete_addresses().unwrap();
1425
1426		// Completing the addresses pinned the key, so this is the key the backend serves.
1427		let (keypair, _peer_id) =
1428			Litep2pNetworkBackend::get_keypair(&network_config.node_key).unwrap();
1429		let node_certhash: NetworkMultihash =
1430			webrtc::derive_certificate(keypair.secret()).unwrap().certhash().into();
1431
1432		// Held for the duration of the test: dropping it closes the node's sockets.
1433		let backend = start_backend(&network_config).unwrap();
1434		let network_service =
1435			<Litep2pNetworkBackend as NetworkBackend<Block, H256>>::network_service(&backend);
1436
1437		let advertised = network_service.listen_addresses();
1438		assert_eq!(advertised.len(), 1);
1439		assert_eq!(certhash(&advertised[0]), Some(node_certhash));
1440
1441		let advertised = network_service.external_addresses();
1442		assert_eq!(advertised.len(), 1);
1443		assert_eq!(certhash(&advertised[0]), Some(node_certhash));
1444	}
1445
1446	#[tokio::test]
1447	async fn webrtc_public_address_completed_at_config_creation_accepted() {
1448		let mut network_config = NetworkConfiguration::new_local();
1449		network_config.listen_addresses = vec![WEBRTC_LISTEN_ADDRESS.parse().unwrap()];
1450		network_config.public_addresses = vec![WEBRTC_PUBLIC_ADDRESS.parse().unwrap()];
1451		network_config.node_key = NodeKeyConfig::Ed25519(Secret::Input(
1452			ed25519::SecretKey::try_from_bytes([7u8; 32]).unwrap(),
1453		));
1454		network_config.validate_and_complete_addresses().unwrap();
1455
1456		// Held for the duration of the test: dropping it closes the node's sockets.
1457		let backend = start_backend(&network_config).unwrap();
1458		let network_service =
1459			<Litep2pNetworkBackend as NetworkBackend<Block, H256>>::network_service(&backend);
1460		let local_peer_id = network_service.local_peer_id();
1461
1462		assert_eq!(
1463			network_service.external_addresses(),
1464			vec![network_config.public_addresses[0]
1465				.clone()
1466				.with(NetworkProtocol::P2p(local_peer_id.into()))],
1467		);
1468	}
1469
1470	/// A configured `/certhash` is checked against the node's own and removed by
1471	/// `NetworkConfiguration::validate_and_complete_addresses`, which this test deliberately skips:
1472	/// the backend is the guard for a configuration built without it.
1473	#[tokio::test]
1474	async fn webrtc_listen_address_with_certhash_refuses_to_start() {
1475		let listen_address: NetworkMultiaddr = WEBRTC_LISTEN_ADDRESS.parse().unwrap();
1476
1477		let mut network_config = NetworkConfiguration::new_local();
1478		network_config.listen_addresses =
1479			vec![listen_address.with(NetworkProtocol::Certhash(a_certhash()))];
1480
1481		assert!(matches!(start_backend(&network_config), Err(Error::InvalidWebRtcAddress { .. }),));
1482	}
1483
1484	#[tokio::test]
1485	async fn non_webrtc_public_address_untouched() {
1486		let public_address: NetworkMultiaddr = "/ip4/203.0.113.9/tcp/31234".parse().unwrap();
1487
1488		let mut network_config = NetworkConfiguration::new_local();
1489		network_config.listen_addresses = vec!["/ip4/127.0.0.1/tcp/0".parse().unwrap()];
1490		network_config.public_addresses = vec![public_address.clone()];
1491
1492		// Held for the duration of the test: dropping it closes the node's sockets.
1493		let backend = start_backend(&network_config).unwrap();
1494		let network_service =
1495			<Litep2pNetworkBackend as NetworkBackend<Block, H256>>::network_service(&backend);
1496		let local_peer_id = network_service.local_peer_id();
1497
1498		// litep2p appends the local `/p2p` to every public address; nothing else is added.
1499		assert_eq!(
1500			network_service.external_addresses(),
1501			vec![public_address.with(NetworkProtocol::P2p(local_peer_id.into()))],
1502		);
1503	}
1504
1505	/// A certhash standing in for the node's own.
1506	fn a_certhash() -> NetworkMultihash {
1507		sc_network_types::multihash::Code::Sha2_256.digest(b"certificate")
1508	}
1509}