Skip to main content

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