referrerpolicy=no-referrer-when-downgrade

sc_network/
service.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//! Main entry point of the sc-network crate.
20//!
21//! There are two main structs in this module: [`NetworkWorker`] and [`NetworkService`].
22//! The [`NetworkWorker`] *is* the network. Network is driven by [`NetworkWorker::run`] future that
23//! terminates only when all instances of the control handles [`NetworkService`] were dropped.
24//! The [`NetworkService`] is merely a shared version of the [`NetworkWorker`]. You can obtain an
25//! `Arc<NetworkService>` by calling [`NetworkWorker::service`].
26//!
27//! The methods of the [`NetworkService`] are implemented by sending a message over a channel,
28//! which is then processed by [`NetworkWorker::next_action`].
29
30use crate::{
31	behaviour::{self, Behaviour, BehaviourOut},
32	config::{
33		parse_addr, FullNetworkConfiguration, IncomingRequest, MultiaddrWithPeerId,
34		NonDefaultSetConfig, NotificationHandshake, Params, SetConfig, TransportConfig,
35	},
36	discovery::DiscoveryConfig,
37	error::Error,
38	event::{DhtEvent, Event},
39	network_state::{
40		NetworkState, NotConnectedPeer as NetworkStateNotConnectedPeer, Peer as NetworkStatePeer,
41	},
42	peer_store::{PeerStore, PeerStoreProvider},
43	protocol::{self, Protocol, Ready},
44	protocol_controller::{self, ProtoSetConfig, ProtocolController, SetId},
45	request_responses::{IfDisconnected, ProtocolConfig as RequestResponseConfig, RequestFailure},
46	service::{
47		signature::{Signature, SigningError},
48		traits::{
49			BandwidthSink, NetworkBackend, NetworkDHTProvider, NetworkEventStream, NetworkPeers,
50			NetworkRequest, NetworkService as NetworkServiceT, NetworkSigner, NetworkStateInfo,
51			NetworkStatus, NetworkStatusProvider, NotificationSender as NotificationSenderT,
52			NotificationSenderError, NotificationSenderReady as NotificationSenderReadyT,
53		},
54	},
55	transport,
56	types::ProtocolName,
57	NotificationService, ReputationChange,
58};
59
60use codec::DecodeAll;
61use futures::{channel::oneshot, prelude::*};
62use libp2p::{
63	connection_limits::{ConnectionLimits, Exceeded},
64	core::{upgrade, ConnectedPoint, Endpoint},
65	identify::Info as IdentifyInfo,
66	identity::ed25519,
67	multiaddr::{self, Multiaddr},
68	swarm::{
69		Config as SwarmConfig, ConnectionError, ConnectionId, DialError, Executor, ListenError,
70		NetworkBehaviour, Swarm, SwarmEvent,
71	},
72	PeerId,
73};
74use log::{debug, error, info, trace, warn};
75use metrics::{Histogram, MetricSources, Metrics};
76use parking_lot::Mutex;
77use prometheus_endpoint::Registry;
78use sc_network_types::kad::{Key as KademliaKey, Record};
79
80use sc_network_common::{
81	role::{ObservedRole, Roles},
82	ExHashT,
83};
84use sc_utils::mpsc::{tracing_unbounded, TracingUnboundedReceiver, TracingUnboundedSender};
85use sp_runtime::traits::Block as BlockT;
86
87pub use behaviour::{InboundFailure, OutboundFailure, ResponseFailure};
88pub use libp2p::identity::{DecodingError, Keypair, PublicKey};
89pub use metrics::NotificationMetrics;
90pub use protocol::NotificationsSink;
91use std::{
92	collections::{HashMap, HashSet},
93	fs, iter,
94	marker::PhantomData,
95	num::NonZeroUsize,
96	pin::Pin,
97	str,
98	sync::{
99		atomic::{AtomicUsize, Ordering},
100		Arc,
101	},
102	time::{Duration, Instant},
103};
104
105pub(crate) mod metrics;
106pub(crate) mod out_events;
107
108pub mod signature;
109pub mod traits;
110
111/// Logging target for the file.
112const LOG_TARGET: &str = "sub-libp2p";
113
114struct Libp2pBandwidthSink {
115	#[allow(deprecated)]
116	sink: Arc<transport::BandwidthSinks>,
117}
118
119impl BandwidthSink for Libp2pBandwidthSink {
120	fn total_inbound(&self) -> u64 {
121		self.sink.total_inbound()
122	}
123
124	fn total_outbound(&self) -> u64 {
125		self.sink.total_outbound()
126	}
127}
128
129/// Substrate network service. Handles network IO and manages connectivity.
130pub struct NetworkService<B: BlockT + 'static, H: ExHashT> {
131	/// Number of peers we're connected to.
132	num_connected: Arc<AtomicUsize>,
133	/// The local external addresses.
134	external_addresses: Arc<Mutex<HashSet<Multiaddr>>>,
135	/// Listen addresses. Do **NOT** include a trailing `/p2p/` with our `PeerId`.
136	listen_addresses: Arc<Mutex<HashSet<Multiaddr>>>,
137	/// Local copy of the `PeerId` of the local node.
138	local_peer_id: PeerId,
139	/// The `KeyPair` that defines the `PeerId` of the local node.
140	local_identity: Keypair,
141	/// Bandwidth logging system. Can be queried to know the average bandwidth consumed.
142	bandwidth: Arc<dyn BandwidthSink>,
143	/// Channel that sends messages to the actual worker.
144	to_worker: TracingUnboundedSender<ServiceToWorkerMsg>,
145	/// Protocol name -> `SetId` mapping for notification protocols. The map never changes after
146	/// initialization.
147	notification_protocol_ids: HashMap<ProtocolName, SetId>,
148	/// Handles to manage peer connections on notification protocols. The vector never changes
149	/// after initialization.
150	protocol_handles: Vec<protocol_controller::ProtocolHandle>,
151	/// Shortcut to sync protocol handle (`protocol_handles[0]`).
152	sync_protocol_handle: protocol_controller::ProtocolHandle,
153	/// Handle to `PeerStore`.
154	peer_store_handle: Arc<dyn PeerStoreProvider>,
155	/// Marker to pin the `H` generic. Serves no purpose except to not break backwards
156	/// compatibility.
157	_marker: PhantomData<H>,
158	/// Marker for block type
159	_block: PhantomData<B>,
160}
161
162#[async_trait::async_trait]
163impl<B, H> NetworkBackend<B, H> for NetworkWorker<B, H>
164where
165	B: BlockT + 'static,
166	H: ExHashT,
167{
168	type NotificationProtocolConfig = NonDefaultSetConfig;
169	type RequestResponseProtocolConfig = RequestResponseConfig;
170	type NetworkService<Block, Hash> = Arc<NetworkService<B, H>>;
171	type PeerStore = PeerStore;
172
173	fn new(params: Params<B, H, Self>) -> Result<Self, Error>
174	where
175		Self: Sized,
176	{
177		NetworkWorker::new(params)
178	}
179
180	/// Get handle to `NetworkService` of the `NetworkBackend`.
181	fn network_service(&self) -> Arc<dyn NetworkServiceT> {
182		self.service.clone()
183	}
184
185	/// Create `PeerStore`.
186	fn peer_store(
187		bootnodes: Vec<sc_network_types::PeerId>,
188		metrics_registry: Option<Registry>,
189	) -> Self::PeerStore {
190		PeerStore::new(bootnodes.into_iter().map(From::from).collect(), metrics_registry)
191	}
192
193	fn register_notification_metrics(registry: Option<&Registry>) -> NotificationMetrics {
194		NotificationMetrics::new(registry)
195	}
196
197	/// Create notification protocol configuration.
198	fn notification_config(
199		protocol_name: ProtocolName,
200		fallback_names: Vec<ProtocolName>,
201		max_notification_size: u64,
202		handshake: Option<NotificationHandshake>,
203		set_config: SetConfig,
204		_metrics: NotificationMetrics,
205		_peerstore_handle: Arc<dyn PeerStoreProvider>,
206	) -> (Self::NotificationProtocolConfig, Box<dyn NotificationService>) {
207		NonDefaultSetConfig::new(
208			protocol_name,
209			fallback_names,
210			max_notification_size,
211			handshake,
212			set_config,
213		)
214	}
215
216	/// Create request-response protocol configuration.
217	fn request_response_config(
218		protocol_name: ProtocolName,
219		fallback_names: Vec<ProtocolName>,
220		max_request_size: u64,
221		max_response_size: u64,
222		request_timeout: Duration,
223		inbound_queue: Option<async_channel::Sender<IncomingRequest>>,
224	) -> Self::RequestResponseProtocolConfig {
225		Self::RequestResponseProtocolConfig {
226			name: protocol_name,
227			fallback_names,
228			max_request_size,
229			max_response_size,
230			request_timeout,
231			inbound_queue,
232		}
233	}
234
235	/// Start [`NetworkBackend`] event loop.
236	async fn run(mut self) {
237		self.run().await
238	}
239}
240
241impl<B, H> NetworkWorker<B, H>
242where
243	B: BlockT + 'static,
244	H: ExHashT,
245{
246	/// Creates the network service.
247	///
248	/// Returns a `NetworkWorker` that implements `Future` and must be regularly polled in order
249	/// for the network processing to advance. From it, you can extract a `NetworkService` using
250	/// `worker.service()`. The `NetworkService` can be shared through the codebase.
251	pub fn new(params: Params<B, H, Self>) -> Result<Self, Error> {
252		let peer_store_handle = params.network_config.peer_store_handle();
253		let FullNetworkConfiguration {
254			notification_protocols,
255			request_response_protocols,
256			mut network_config,
257			..
258		} = params.network_config;
259
260		// Private and public keys configuration.
261		let local_identity = network_config.node_key.clone().into_keypair()?;
262		let local_public = local_identity.public();
263		let local_peer_id = local_public.to_peer_id();
264
265		// Convert to libp2p types.
266		let local_identity: ed25519::Keypair = local_identity.into();
267		let local_public: ed25519::PublicKey = local_public.into();
268		let local_peer_id: PeerId = local_peer_id.into();
269
270		network_config.boot_nodes = network_config
271			.boot_nodes
272			.into_iter()
273			.filter(|boot_node| boot_node.peer_id != local_peer_id.into())
274			.collect();
275		network_config.default_peers_set.reserved_nodes = network_config
276			.default_peers_set
277			.reserved_nodes
278			.into_iter()
279			.filter(|reserved_node| {
280				if reserved_node.peer_id == local_peer_id.into() {
281					warn!(
282						target: LOG_TARGET,
283						"Local peer ID used in reserved node, ignoring: {}",
284						reserved_node,
285					);
286					false
287				} else {
288					true
289				}
290			})
291			.collect();
292
293		// Ensure the listen addresses are consistent with the transport.
294		ensure_addresses_consistent_with_transport(
295			network_config.listen_addresses.iter(),
296			&network_config.transport,
297		)?;
298		ensure_addresses_consistent_with_transport(
299			network_config.boot_nodes.iter().map(|x| &x.multiaddr),
300			&network_config.transport,
301		)?;
302		ensure_addresses_consistent_with_transport(
303			network_config.default_peers_set.reserved_nodes.iter().map(|x| &x.multiaddr),
304			&network_config.transport,
305		)?;
306		for notification_protocol in &notification_protocols {
307			ensure_addresses_consistent_with_transport(
308				notification_protocol.set_config().reserved_nodes.iter().map(|x| &x.multiaddr),
309				&network_config.transport,
310			)?;
311		}
312		ensure_addresses_consistent_with_transport(
313			network_config.public_addresses.iter(),
314			&network_config.transport,
315		)?;
316
317		let (to_worker, from_service) = tracing_unbounded("mpsc_network_worker", 100_000);
318
319		if let Some(path) = &network_config.net_config_path {
320			fs::create_dir_all(path)?;
321		}
322
323		info!(
324			target: LOG_TARGET,
325			"๐Ÿท  Local node identity is: {}",
326			local_peer_id.to_base58(),
327		);
328		info!(target: LOG_TARGET, "Running libp2p network backend");
329
330		let (transport, bandwidth) = {
331			let config_mem = match network_config.transport {
332				TransportConfig::MemoryOnly => true,
333				TransportConfig::Normal { .. } => false,
334			};
335
336			transport::build_transport(local_identity.clone().into(), config_mem)
337		};
338
339		let (to_notifications, from_protocol_controllers) =
340			tracing_unbounded("mpsc_protocol_controllers_to_notifications", 10_000);
341
342		// We must prepend a hardcoded default peer set to notification protocols.
343		let all_peer_sets_iter = iter::once(&network_config.default_peers_set)
344			.chain(notification_protocols.iter().map(|protocol| protocol.set_config()));
345
346		let (protocol_handles, protocol_controllers): (Vec<_>, Vec<_>) = all_peer_sets_iter
347			.enumerate()
348			.map(|(set_id, set_config)| {
349				let proto_set_config = ProtoSetConfig {
350					in_peers: set_config.in_peers,
351					out_peers: set_config.out_peers,
352					reserved_nodes: set_config
353						.reserved_nodes
354						.iter()
355						.map(|node| node.peer_id.into())
356						.collect(),
357					reserved_only: set_config.non_reserved_mode.is_reserved_only(),
358				};
359
360				ProtocolController::new(
361					SetId::from(set_id),
362					proto_set_config,
363					to_notifications.clone(),
364					Arc::clone(&peer_store_handle),
365				)
366			})
367			.unzip();
368
369		// Shortcut to default (sync) peer set protocol handle.
370		let sync_protocol_handle = protocol_handles[0].clone();
371
372		// Spawn `ProtocolController` runners.
373		protocol_controllers
374			.into_iter()
375			.for_each(|controller| (params.executor)(controller.run().boxed()));
376
377		// Protocol name to protocol id mapping. The first protocol is always block announce (sync)
378		// protocol, aka default (hardcoded) peer set.
379		let notification_protocol_ids: HashMap<ProtocolName, SetId> =
380			iter::once(&params.block_announce_config)
381				.chain(notification_protocols.iter())
382				.enumerate()
383				.map(|(index, protocol)| (protocol.protocol_name().clone(), SetId::from(index)))
384				.collect();
385
386		let known_addresses = {
387			// Collect all reserved nodes and bootnodes addresses.
388			let mut addresses: Vec<_> = network_config
389				.default_peers_set
390				.reserved_nodes
391				.iter()
392				.map(|reserved| (reserved.peer_id, reserved.multiaddr.clone()))
393				.chain(notification_protocols.iter().flat_map(|protocol| {
394					protocol
395						.set_config()
396						.reserved_nodes
397						.iter()
398						.map(|reserved| (reserved.peer_id, reserved.multiaddr.clone()))
399				}))
400				.chain(
401					network_config
402						.boot_nodes
403						.iter()
404						.map(|bootnode| (bootnode.peer_id, bootnode.multiaddr.clone())),
405				)
406				.collect();
407
408			// Remove possible duplicates.
409			addresses.sort();
410			addresses.dedup();
411
412			addresses
413		};
414
415		// Check for duplicate bootnodes.
416		network_config.boot_nodes.iter().try_for_each(|bootnode| {
417			if let Some(other) = network_config
418				.boot_nodes
419				.iter()
420				.filter(|o| o.multiaddr == bootnode.multiaddr)
421				.find(|o| o.peer_id != bootnode.peer_id)
422			{
423				Err(Error::DuplicateBootnode {
424					address: bootnode.multiaddr.clone().into(),
425					first_id: bootnode.peer_id.into(),
426					second_id: other.peer_id.into(),
427				})
428			} else {
429				Ok(())
430			}
431		})?;
432
433		// List of bootnode multiaddresses.
434		let mut boot_node_ids = HashMap::<PeerId, Vec<Multiaddr>>::new();
435
436		for bootnode in network_config.boot_nodes.iter() {
437			boot_node_ids
438				.entry(bootnode.peer_id.into())
439				.or_default()
440				.push(bootnode.multiaddr.clone().into());
441		}
442
443		let boot_node_ids = Arc::new(boot_node_ids);
444
445		let num_connected = Arc::new(AtomicUsize::new(0));
446		let external_addresses = Arc::new(Mutex::new(HashSet::new()));
447
448		let (protocol, notif_protocol_handles) = Protocol::new(
449			From::from(&params.role),
450			params.notification_metrics,
451			notification_protocols,
452			params.block_announce_config,
453			Arc::clone(&peer_store_handle),
454			protocol_handles.clone(),
455			from_protocol_controllers,
456		)?;
457
458		// Build the swarm.
459		let (mut swarm, bandwidth): (Swarm<Behaviour<B>>, _) = {
460			let user_agent =
461				format!("{} ({})", network_config.client_version, network_config.node_name);
462
463			let discovery_config = {
464				let mut config = DiscoveryConfig::new(local_peer_id);
465				config.with_permanent_addresses(
466					known_addresses
467						.iter()
468						.map(|(peer, address)| (peer.into(), address.clone().into()))
469						.collect::<Vec<_>>(),
470				);
471				config.discovery_limit(u64::from(network_config.default_peers_set.out_peers) + 15);
472				config.with_kademlia(
473					params.genesis_hash,
474					params.fork_id.as_deref(),
475					&params.protocol_id,
476				);
477				config.with_dht_random_walk(network_config.enable_dht_random_walk);
478				config.allow_non_globals_in_dht(network_config.allow_non_globals_in_dht);
479				config.use_kademlia_disjoint_query_paths(
480					network_config.kademlia_disjoint_query_paths,
481				);
482				config.with_kademlia_replication_factor(network_config.kademlia_replication_factor);
483
484				match network_config.transport {
485					TransportConfig::MemoryOnly => {
486						config.with_mdns(false);
487						config.allow_private_ip(false);
488					},
489					TransportConfig::Normal {
490						enable_mdns,
491						allow_private_ip: allow_private_ipv4,
492						..
493					} => {
494						config.with_mdns(enable_mdns);
495						config.allow_private_ip(allow_private_ipv4);
496					},
497				}
498
499				config
500			};
501
502			let behaviour = {
503				let result = Behaviour::new(
504					protocol,
505					user_agent,
506					local_public.into(),
507					discovery_config,
508					request_response_protocols,
509					Arc::clone(&peer_store_handle),
510					external_addresses.clone(),
511					network_config.public_addresses.iter().cloned().map(Into::into).collect(),
512					ConnectionLimits::default()
513						.with_max_established_per_peer(Some(crate::MAX_CONNECTIONS_PER_PEER as u32))
514						.with_max_established_incoming(Some(
515							crate::MAX_CONNECTIONS_ESTABLISHED_INCOMING,
516						)),
517				);
518
519				match result {
520					Ok(b) => b,
521					Err(crate::request_responses::RegisterError::DuplicateProtocol(proto)) => {
522						return Err(Error::DuplicateRequestResponseProtocol { protocol: proto })
523					},
524				}
525			};
526
527			let swarm = {
528				struct SpawnImpl<F>(F);
529				impl<F: Fn(Pin<Box<dyn Future<Output = ()> + Send>>)> Executor for SpawnImpl<F> {
530					fn exec(&self, f: Pin<Box<dyn Future<Output = ()> + Send>>) {
531						(self.0)(f)
532					}
533				}
534
535				let config = SwarmConfig::with_executor(SpawnImpl(params.executor))
536					.with_substream_upgrade_protocol_override(upgrade::Version::V1)
537					.with_notify_handler_buffer_size(NonZeroUsize::new(32).expect("32 != 0; qed"))
538					// NOTE: 24 is somewhat arbitrary and should be tuned in the future if
539					// necessary. See <https://github.com/paritytech/substrate/pull/6080>
540					.with_per_connection_event_buffer_size(24)
541					.with_max_negotiating_inbound_streams(2048)
542					.with_idle_connection_timeout(network_config.idle_connection_timeout);
543
544				Swarm::new(transport, behaviour, local_peer_id, config)
545			};
546
547			(swarm, Arc::new(Libp2pBandwidthSink { sink: bandwidth }))
548		};
549
550		// Initialize the metrics.
551		let metrics = match &params.metrics_registry {
552			Some(registry) => Some(metrics::register(
553				registry,
554				MetricSources {
555					bandwidth: bandwidth.clone(),
556					connected_peers: num_connected.clone(),
557				},
558			)?),
559			None => None,
560		};
561
562		// Listen on multiaddresses.
563		for addr in &network_config.listen_addresses {
564			if let Err(err) = Swarm::<Behaviour<B>>::listen_on(&mut swarm, addr.clone().into()) {
565				warn!(target: LOG_TARGET, "Can't listen on {} because: {:?}", addr, err)
566			}
567		}
568
569		// Add external addresses.
570		for addr in &network_config.public_addresses {
571			Swarm::<Behaviour<B>>::add_external_address(&mut swarm, addr.clone().into());
572		}
573
574		let listen_addresses_set = Arc::new(Mutex::new(HashSet::new()));
575
576		let service = Arc::new(NetworkService {
577			bandwidth,
578			external_addresses,
579			listen_addresses: listen_addresses_set.clone(),
580			num_connected: num_connected.clone(),
581			local_peer_id,
582			local_identity: local_identity.into(),
583			to_worker,
584			notification_protocol_ids,
585			protocol_handles,
586			sync_protocol_handle,
587			peer_store_handle: Arc::clone(&peer_store_handle),
588			_marker: PhantomData,
589			_block: Default::default(),
590		});
591
592		Ok(NetworkWorker {
593			listen_addresses: listen_addresses_set,
594			num_connected,
595			network_service: swarm,
596			service,
597			from_service,
598			event_streams: out_events::OutChannels::new(params.metrics_registry.as_ref())?,
599			metrics,
600			boot_node_ids,
601			reported_invalid_boot_nodes: Default::default(),
602			peer_store_handle: Arc::clone(&peer_store_handle),
603			notif_protocol_handles,
604			_marker: Default::default(),
605			_block: Default::default(),
606		})
607	}
608
609	/// High-level network status information.
610	pub fn status(&self) -> NetworkStatus {
611		NetworkStatus {
612			num_connected_peers: self.num_connected_peers(),
613			total_bytes_inbound: self.total_bytes_inbound(),
614			total_bytes_outbound: self.total_bytes_outbound(),
615		}
616	}
617
618	/// Returns the total number of bytes received so far.
619	pub fn total_bytes_inbound(&self) -> u64 {
620		self.service.bandwidth.total_inbound()
621	}
622
623	/// Returns the total number of bytes sent so far.
624	pub fn total_bytes_outbound(&self) -> u64 {
625		self.service.bandwidth.total_outbound()
626	}
627
628	/// Returns the number of peers we're connected to.
629	pub fn num_connected_peers(&self) -> usize {
630		self.network_service.behaviour().user_protocol().num_sync_peers()
631	}
632
633	/// Adds an address for a node.
634	pub fn add_known_address(&mut self, peer_id: PeerId, addr: Multiaddr) {
635		self.network_service.behaviour_mut().add_known_address(peer_id, addr);
636	}
637
638	/// Return a `NetworkService` that can be shared through the code base and can be used to
639	/// manipulate the worker.
640	pub fn service(&self) -> &Arc<NetworkService<B, H>> {
641		&self.service
642	}
643
644	/// Returns the local `PeerId`.
645	pub fn local_peer_id(&self) -> &PeerId {
646		Swarm::<Behaviour<B>>::local_peer_id(&self.network_service)
647	}
648
649	/// Returns the list of addresses we are listening on.
650	///
651	/// Does **NOT** include a trailing `/p2p/` with our `PeerId`.
652	pub fn listen_addresses(&self) -> impl Iterator<Item = &Multiaddr> {
653		Swarm::<Behaviour<B>>::listeners(&self.network_service)
654	}
655
656	/// Get network state.
657	///
658	/// **Note**: Use this only for debugging. This API is unstable. There are warnings literally
659	/// everywhere about this. Please don't use this function to retrieve actual information.
660	pub fn network_state(&mut self) -> NetworkState {
661		let swarm = &mut self.network_service;
662		let open = swarm.behaviour_mut().user_protocol().open_peers().cloned().collect::<Vec<_>>();
663		let connected_peers = {
664			let swarm = &mut *swarm;
665			open.iter()
666				.filter_map(move |peer_id| {
667					let known_addresses = if let Ok(addrs) =
668						NetworkBehaviour::handle_pending_outbound_connection(
669							swarm.behaviour_mut(),
670							ConnectionId::new_unchecked(0), // dummy value
671							Some(*peer_id),
672							&vec![],
673							Endpoint::Listener,
674						) {
675						addrs.into_iter().collect()
676					} else {
677						error!(target: LOG_TARGET, "Was not able to get known addresses for {:?}", peer_id);
678						return None;
679					};
680
681					let endpoint = if let Some(e) =
682						swarm.behaviour_mut().node(peer_id).and_then(|i| i.endpoint())
683					{
684						e.clone().into()
685					} else {
686						error!(target: LOG_TARGET, "Found state inconsistency between custom protocol \
687						and debug information about {:?}", peer_id);
688						return None;
689					};
690
691					Some((
692						peer_id.to_base58(),
693						NetworkStatePeer {
694							endpoint,
695							version_string: swarm
696								.behaviour_mut()
697								.node(peer_id)
698								.and_then(|i| i.client_version().map(|s| s.to_owned())),
699							latest_ping_time: swarm
700								.behaviour_mut()
701								.node(peer_id)
702								.and_then(|i| i.latest_ping()),
703							known_addresses,
704						},
705					))
706				})
707				.collect()
708		};
709
710		let not_connected_peers = {
711			let swarm = &mut *swarm;
712			swarm
713				.behaviour_mut()
714				.known_peers()
715				.into_iter()
716				.filter(|p| open.iter().all(|n| n != p))
717				.map(move |peer_id| {
718					let known_addresses = if let Ok(addrs) =
719						NetworkBehaviour::handle_pending_outbound_connection(
720							swarm.behaviour_mut(),
721							ConnectionId::new_unchecked(0), // dummy value
722							Some(peer_id),
723							&vec![],
724							Endpoint::Listener,
725						) {
726						addrs.into_iter().collect()
727					} else {
728						error!(target: LOG_TARGET, "Was not able to get known addresses for {:?}", peer_id);
729						Default::default()
730					};
731
732					(
733						peer_id.to_base58(),
734						NetworkStateNotConnectedPeer {
735							version_string: swarm
736								.behaviour_mut()
737								.node(&peer_id)
738								.and_then(|i| i.client_version().map(|s| s.to_owned())),
739							latest_ping_time: swarm
740								.behaviour_mut()
741								.node(&peer_id)
742								.and_then(|i| i.latest_ping()),
743							known_addresses,
744						},
745					)
746				})
747				.collect()
748		};
749
750		let peer_id = Swarm::<Behaviour<B>>::local_peer_id(swarm).to_base58();
751		let listened_addresses = swarm.listeners().cloned().collect();
752		let external_addresses = swarm.external_addresses().cloned().collect();
753
754		NetworkState {
755			peer_id,
756			listened_addresses,
757			external_addresses,
758			connected_peers,
759			not_connected_peers,
760			// TODO: Check what info we can include here.
761			//       Issue reference: https://github.com/paritytech/substrate/issues/14160.
762			peerset: serde_json::json!(
763				"Unimplemented. See https://github.com/paritytech/substrate/issues/14160."
764			),
765		}
766	}
767
768	/// Removes a `PeerId` from the list of reserved peers.
769	pub fn remove_reserved_peer(&self, peer: PeerId) {
770		self.service.remove_reserved_peer(peer.into());
771	}
772
773	/// Adds a `PeerId` and its `Multiaddr` as reserved.
774	pub fn add_reserved_peer(&self, peer: MultiaddrWithPeerId) -> Result<(), String> {
775		self.service.add_reserved_peer(peer)
776	}
777}
778
779impl<B: BlockT + 'static, H: ExHashT> NetworkService<B, H> {
780	/// Get network state.
781	///
782	/// **Note**: Use this only for debugging. This API is unstable. There are warnings literally
783	/// everywhere about this. Please don't use this function to retrieve actual information.
784	///
785	/// Returns an error if the `NetworkWorker` is no longer running.
786	pub async fn network_state(&self) -> Result<NetworkState, ()> {
787		let (tx, rx) = oneshot::channel();
788
789		let _ = self
790			.to_worker
791			.unbounded_send(ServiceToWorkerMsg::NetworkState { pending_response: tx });
792
793		match rx.await {
794			Ok(v) => v.map_err(|_| ()),
795			// The channel can only be closed if the network worker no longer exists.
796			Err(_) => Err(()),
797		}
798	}
799
800	/// Utility function to extract `PeerId` from each `Multiaddr` for peer set updates.
801	///
802	/// Returns an `Err` if one of the given addresses is invalid or contains an
803	/// invalid peer ID (which includes the local peer ID).
804	fn split_multiaddr_and_peer_id(
805		&self,
806		peers: HashSet<Multiaddr>,
807	) -> Result<Vec<(PeerId, Multiaddr)>, String> {
808		peers
809			.into_iter()
810			.map(|mut addr| {
811				let peer = match addr.pop() {
812					Some(multiaddr::Protocol::P2p(peer_id)) => peer_id,
813					_ => return Err("Missing PeerId from address".to_string()),
814				};
815
816				// Make sure the local peer ID is never added to the PSM
817				// or added as a "known address", even if given.
818				if peer == self.local_peer_id {
819					Err("Local peer ID in peer set.".to_string())
820				} else {
821					Ok((peer, addr))
822				}
823			})
824			.collect::<Result<Vec<(PeerId, Multiaddr)>, String>>()
825	}
826}
827
828impl<B, H> NetworkStateInfo for NetworkService<B, H>
829where
830	B: sp_runtime::traits::Block,
831	H: ExHashT,
832{
833	/// Returns the local external addresses.
834	fn external_addresses(&self) -> Vec<sc_network_types::multiaddr::Multiaddr> {
835		self.external_addresses.lock().iter().cloned().map(Into::into).collect()
836	}
837
838	/// Returns the listener addresses (without trailing `/p2p/` with our `PeerId`).
839	fn listen_addresses(&self) -> Vec<sc_network_types::multiaddr::Multiaddr> {
840		self.listen_addresses.lock().iter().cloned().map(Into::into).collect()
841	}
842
843	/// Returns the local Peer ID.
844	fn local_peer_id(&self) -> sc_network_types::PeerId {
845		self.local_peer_id.into()
846	}
847}
848
849impl<B, H> NetworkSigner for NetworkService<B, H>
850where
851	B: sp_runtime::traits::Block,
852	H: ExHashT,
853{
854	fn sign_with_local_identity(&self, msg: Vec<u8>) -> Result<Signature, SigningError> {
855		let public_key = self.local_identity.public();
856		let bytes = self.local_identity.sign(msg.as_ref())?;
857
858		Ok(Signature {
859			public_key: crate::service::signature::PublicKey::Libp2p(public_key),
860			bytes,
861		})
862	}
863
864	fn verify(
865		&self,
866		peer_id: sc_network_types::PeerId,
867		public_key: &Vec<u8>,
868		signature: &Vec<u8>,
869		message: &Vec<u8>,
870	) -> Result<bool, String> {
871		let public_key =
872			PublicKey::try_decode_protobuf(&public_key).map_err(|error| error.to_string())?;
873		let peer_id: PeerId = peer_id.into();
874		let remote: libp2p::PeerId = public_key.to_peer_id();
875
876		Ok(peer_id == remote && public_key.verify(message, signature))
877	}
878}
879
880impl<B, H> NetworkDHTProvider for NetworkService<B, H>
881where
882	B: BlockT + 'static,
883	H: ExHashT,
884{
885	/// Start finding closest peerst to the target peer ID in the DHT.
886	///
887	/// This will generate either a `ClosestPeersFound` or a `ClosestPeersNotFound` event and pass
888	/// it as an item on the [`NetworkWorker`] stream.
889	fn find_closest_peers(&self, target: sc_network_types::PeerId) {
890		let _ = self
891			.to_worker
892			.unbounded_send(ServiceToWorkerMsg::FindClosestPeers(target.into()));
893	}
894
895	/// Start getting a value from the DHT.
896	///
897	/// This will generate either a `ValueFound` or a `ValueNotFound` event and pass it as an
898	/// item on the [`NetworkWorker`] stream.
899	fn get_value(&self, key: &KademliaKey) {
900		let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::GetValue(key.clone()));
901	}
902
903	/// Start putting a value in the DHT.
904	///
905	/// This will generate either a `ValuePut` or a `ValuePutFailed` event and pass it as an
906	/// item on the [`NetworkWorker`] stream.
907	fn put_value(&self, key: KademliaKey, value: Vec<u8>) {
908		let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::PutValue(key, value));
909	}
910
911	fn put_record_to(
912		&self,
913		record: Record,
914		peers: HashSet<sc_network_types::PeerId>,
915		update_local_storage: bool,
916	) {
917		let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::PutRecordTo {
918			record,
919			peers,
920			update_local_storage,
921		});
922	}
923
924	fn store_record(
925		&self,
926		key: KademliaKey,
927		value: Vec<u8>,
928		publisher: Option<sc_network_types::PeerId>,
929		expires: Option<Instant>,
930	) {
931		let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::StoreRecord(
932			key,
933			value,
934			publisher.map(Into::into),
935			expires,
936		));
937	}
938
939	fn start_providing(&self, key: KademliaKey) {
940		let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::StartProviding(key));
941	}
942
943	fn stop_providing(&self, key: KademliaKey) {
944		let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::StopProviding(key));
945	}
946
947	fn get_providers(&self, key: KademliaKey) {
948		let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::GetProviders(key));
949	}
950}
951
952#[async_trait::async_trait]
953impl<B, H> NetworkStatusProvider for NetworkService<B, H>
954where
955	B: BlockT + 'static,
956	H: ExHashT,
957{
958	async fn status(&self) -> Result<NetworkStatus, ()> {
959		let (tx, rx) = oneshot::channel();
960
961		let _ = self
962			.to_worker
963			.unbounded_send(ServiceToWorkerMsg::NetworkStatus { pending_response: tx });
964
965		match rx.await {
966			Ok(v) => v.map_err(|_| ()),
967			// The channel can only be closed if the network worker no longer exists.
968			Err(_) => Err(()),
969		}
970	}
971
972	async fn network_state(&self) -> Result<NetworkState, ()> {
973		let (tx, rx) = oneshot::channel();
974
975		let _ = self
976			.to_worker
977			.unbounded_send(ServiceToWorkerMsg::NetworkState { pending_response: tx });
978
979		match rx.await {
980			Ok(v) => v.map_err(|_| ()),
981			// The channel can only be closed if the network worker no longer exists.
982			Err(_) => Err(()),
983		}
984	}
985}
986
987#[async_trait::async_trait]
988impl<B, H> NetworkPeers for NetworkService<B, H>
989where
990	B: BlockT + 'static,
991	H: ExHashT,
992{
993	fn set_authorized_peers(&self, peers: HashSet<sc_network_types::PeerId>) {
994		self.sync_protocol_handle
995			.set_reserved_peers(peers.iter().map(|peer| (*peer).into()).collect());
996	}
997
998	fn set_authorized_only(&self, reserved_only: bool) {
999		self.sync_protocol_handle.set_reserved_only(reserved_only);
1000	}
1001
1002	fn add_known_address(
1003		&self,
1004		peer_id: sc_network_types::PeerId,
1005		addr: sc_network_types::multiaddr::Multiaddr,
1006	) {
1007		let _ = self
1008			.to_worker
1009			.unbounded_send(ServiceToWorkerMsg::AddKnownAddress(peer_id.into(), addr.into()));
1010	}
1011
1012	fn report_peer(&self, peer_id: sc_network_types::PeerId, cost_benefit: ReputationChange) {
1013		self.peer_store_handle.report_peer(peer_id, cost_benefit);
1014	}
1015
1016	fn peer_reputation(&self, peer_id: &sc_network_types::PeerId) -> i32 {
1017		self.peer_store_handle.peer_reputation(peer_id)
1018	}
1019
1020	fn disconnect_peer(&self, peer_id: sc_network_types::PeerId, protocol: ProtocolName) {
1021		let _ = self
1022			.to_worker
1023			.unbounded_send(ServiceToWorkerMsg::DisconnectPeer(peer_id.into(), protocol));
1024	}
1025
1026	fn accept_unreserved_peers(&self) {
1027		self.sync_protocol_handle.set_reserved_only(false);
1028	}
1029
1030	fn deny_unreserved_peers(&self) {
1031		self.sync_protocol_handle.set_reserved_only(true);
1032	}
1033
1034	fn add_reserved_peer(&self, peer: MultiaddrWithPeerId) -> Result<(), String> {
1035		// Make sure the local peer ID is never added as a reserved peer.
1036		if peer.peer_id == self.local_peer_id.into() {
1037			return Err("Local peer ID cannot be added as a reserved peer.".to_string());
1038		}
1039
1040		let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::AddKnownAddress(
1041			peer.peer_id.into(),
1042			peer.multiaddr.into(),
1043		));
1044		self.sync_protocol_handle.add_reserved_peer(peer.peer_id.into());
1045
1046		Ok(())
1047	}
1048
1049	fn remove_reserved_peer(&self, peer_id: sc_network_types::PeerId) {
1050		self.sync_protocol_handle.remove_reserved_peer(peer_id.into());
1051	}
1052
1053	fn set_reserved_peers(
1054		&self,
1055		protocol: ProtocolName,
1056		peers: HashSet<sc_network_types::multiaddr::Multiaddr>,
1057	) -> Result<(), String> {
1058		let Some(set_id) = self.notification_protocol_ids.get(&protocol) else {
1059			return Err(format!("Cannot set reserved peers for unknown protocol: {}", protocol));
1060		};
1061
1062		let peers: HashSet<Multiaddr> = peers.into_iter().map(Into::into).collect();
1063		let peers_addrs = self.split_multiaddr_and_peer_id(peers)?;
1064
1065		let mut peers: HashSet<PeerId> = HashSet::with_capacity(peers_addrs.len());
1066
1067		for (peer_id, addr) in peers_addrs.into_iter() {
1068			// Make sure the local peer ID is never added to the PSM.
1069			if peer_id == self.local_peer_id {
1070				return Err("Local peer ID cannot be added as a reserved peer.".to_string());
1071			}
1072
1073			peers.insert(peer_id.into());
1074
1075			if !addr.is_empty() {
1076				let _ = self
1077					.to_worker
1078					.unbounded_send(ServiceToWorkerMsg::AddKnownAddress(peer_id, addr));
1079			}
1080		}
1081
1082		self.protocol_handles[usize::from(*set_id)].set_reserved_peers(peers);
1083
1084		Ok(())
1085	}
1086
1087	fn add_peers_to_reserved_set(
1088		&self,
1089		protocol: ProtocolName,
1090		peers: HashSet<sc_network_types::multiaddr::Multiaddr>,
1091	) -> Result<(), String> {
1092		let Some(set_id) = self.notification_protocol_ids.get(&protocol) else {
1093			return Err(format!(
1094				"Cannot add peers to reserved set of unknown protocol: {}",
1095				protocol
1096			));
1097		};
1098
1099		let peers: HashSet<Multiaddr> = peers.into_iter().map(Into::into).collect();
1100		let peers = self.split_multiaddr_and_peer_id(peers)?;
1101
1102		for (peer_id, addr) in peers.into_iter() {
1103			// Make sure the local peer ID is never added to the PSM.
1104			if peer_id == self.local_peer_id {
1105				return Err("Local peer ID cannot be added as a reserved peer.".to_string());
1106			}
1107
1108			if !addr.is_empty() {
1109				let _ = self
1110					.to_worker
1111					.unbounded_send(ServiceToWorkerMsg::AddKnownAddress(peer_id, addr));
1112			}
1113
1114			self.protocol_handles[usize::from(*set_id)].add_reserved_peer(peer_id);
1115		}
1116
1117		Ok(())
1118	}
1119
1120	fn remove_peers_from_reserved_set(
1121		&self,
1122		protocol: ProtocolName,
1123		peers: Vec<sc_network_types::PeerId>,
1124	) -> Result<(), String> {
1125		let Some(set_id) = self.notification_protocol_ids.get(&protocol) else {
1126			return Err(format!(
1127				"Cannot remove peers from reserved set of unknown protocol: {}",
1128				protocol
1129			));
1130		};
1131
1132		for peer_id in peers.into_iter() {
1133			self.protocol_handles[usize::from(*set_id)].remove_reserved_peer(peer_id.into());
1134		}
1135
1136		Ok(())
1137	}
1138
1139	fn sync_num_connected(&self) -> usize {
1140		self.num_connected.load(Ordering::Relaxed)
1141	}
1142
1143	fn peer_role(
1144		&self,
1145		peer_id: sc_network_types::PeerId,
1146		handshake: Vec<u8>,
1147	) -> Option<ObservedRole> {
1148		match Roles::decode_all(&mut &handshake[..]) {
1149			Ok(role) => Some(role.into()),
1150			Err(_) => {
1151				log::debug!(target: LOG_TARGET, "handshake doesn't contain peer role: {handshake:?}");
1152				self.peer_store_handle.peer_role(&(peer_id.into()))
1153			},
1154		}
1155	}
1156
1157	/// Get the list of reserved peers.
1158	///
1159	/// Returns an error if the `NetworkWorker` is no longer running.
1160	async fn reserved_peers(&self) -> Result<Vec<sc_network_types::PeerId>, ()> {
1161		let (tx, rx) = oneshot::channel();
1162
1163		self.sync_protocol_handle.reserved_peers(tx);
1164
1165		// The channel can only be closed if `ProtocolController` no longer exists.
1166		rx.await
1167			.map(|peers| peers.into_iter().map(From::from).collect())
1168			.map_err(|_| ())
1169	}
1170}
1171
1172impl<B, H> NetworkEventStream for NetworkService<B, H>
1173where
1174	B: BlockT + 'static,
1175	H: ExHashT,
1176{
1177	fn event_stream(&self, name: &'static str) -> Pin<Box<dyn Stream<Item = Event> + Send>> {
1178		let (tx, rx) = out_events::channel(name, 100_000);
1179		let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::EventStream(tx));
1180		Box::pin(rx)
1181	}
1182}
1183
1184#[async_trait::async_trait]
1185impl<B, H> NetworkRequest for NetworkService<B, H>
1186where
1187	B: BlockT + 'static,
1188	H: ExHashT,
1189{
1190	async fn request(
1191		&self,
1192		target: sc_network_types::PeerId,
1193		protocol: ProtocolName,
1194		request: Vec<u8>,
1195		fallback_request: Option<(Vec<u8>, ProtocolName)>,
1196		connect: IfDisconnected,
1197	) -> Result<(Vec<u8>, ProtocolName), RequestFailure> {
1198		let (tx, rx) = oneshot::channel();
1199
1200		self.start_request(target.into(), protocol, request, fallback_request, tx, connect);
1201
1202		match rx.await {
1203			Ok(v) => v,
1204			// The channel can only be closed if the network worker no longer exists. If the
1205			// network worker no longer exists, then all connections to `target` are necessarily
1206			// closed, and we legitimately report this situation as a "ConnectionClosed".
1207			Err(_) => Err(RequestFailure::Network(OutboundFailure::ConnectionClosed)),
1208		}
1209	}
1210
1211	fn start_request(
1212		&self,
1213		target: sc_network_types::PeerId,
1214		protocol: ProtocolName,
1215		request: Vec<u8>,
1216		fallback_request: Option<(Vec<u8>, ProtocolName)>,
1217		tx: oneshot::Sender<Result<(Vec<u8>, ProtocolName), RequestFailure>>,
1218		connect: IfDisconnected,
1219	) {
1220		let _ = self.to_worker.unbounded_send(ServiceToWorkerMsg::Request {
1221			target: target.into(),
1222			protocol: protocol.into(),
1223			request,
1224			fallback_request,
1225			pending_response: tx,
1226			connect,
1227		});
1228	}
1229}
1230
1231/// A `NotificationSender` allows for sending notifications to a peer with a chosen protocol.
1232#[must_use]
1233pub struct NotificationSender {
1234	sink: NotificationsSink,
1235
1236	/// Name of the protocol on the wire.
1237	protocol_name: ProtocolName,
1238
1239	/// Field extracted from the [`Metrics`] struct and necessary to report the
1240	/// notifications-related metrics.
1241	notification_size_metric: Option<Histogram>,
1242}
1243
1244#[async_trait::async_trait]
1245impl NotificationSenderT for NotificationSender {
1246	async fn ready(
1247		&self,
1248	) -> Result<Box<dyn NotificationSenderReadyT + '_>, NotificationSenderError> {
1249		Ok(Box::new(NotificationSenderReady {
1250			ready: match self.sink.reserve_notification().await {
1251				Ok(r) => Some(r),
1252				Err(()) => return Err(NotificationSenderError::Closed),
1253			},
1254			peer_id: self.sink.peer_id(),
1255			protocol_name: &self.protocol_name,
1256			notification_size_metric: self.notification_size_metric.clone(),
1257		}))
1258	}
1259}
1260
1261/// Reserved slot in the notifications buffer, ready to accept data.
1262#[must_use]
1263pub struct NotificationSenderReady<'a> {
1264	ready: Option<Ready<'a>>,
1265
1266	/// Target of the notification.
1267	peer_id: &'a PeerId,
1268
1269	/// Name of the protocol on the wire.
1270	protocol_name: &'a ProtocolName,
1271
1272	/// Field extracted from the [`Metrics`] struct and necessary to report the
1273	/// notifications-related metrics.
1274	notification_size_metric: Option<Histogram>,
1275}
1276
1277impl<'a> NotificationSenderReadyT for NotificationSenderReady<'a> {
1278	fn send(&mut self, notification: Vec<u8>) -> Result<(), NotificationSenderError> {
1279		if let Some(notification_size_metric) = &self.notification_size_metric {
1280			notification_size_metric.observe(notification.len() as f64);
1281		}
1282
1283		trace!(
1284			target: LOG_TARGET,
1285			"External API => Notification({:?}, {}, {} bytes)",
1286			self.peer_id, self.protocol_name, notification.len(),
1287		);
1288		trace!(target: LOG_TARGET, "Handler({:?}) <= Async notification", self.peer_id);
1289
1290		self.ready
1291			.take()
1292			.ok_or(NotificationSenderError::Closed)?
1293			.send(notification)
1294			.map_err(|()| NotificationSenderError::Closed)
1295	}
1296}
1297
1298/// Messages sent from the `NetworkService` to the `NetworkWorker`.
1299///
1300/// Each entry corresponds to a method of `NetworkService`.
1301enum ServiceToWorkerMsg {
1302	FindClosestPeers(PeerId),
1303	GetValue(KademliaKey),
1304	PutValue(KademliaKey, Vec<u8>),
1305	PutRecordTo {
1306		record: Record,
1307		peers: HashSet<sc_network_types::PeerId>,
1308		update_local_storage: bool,
1309	},
1310	StoreRecord(KademliaKey, Vec<u8>, Option<PeerId>, Option<Instant>),
1311	StartProviding(KademliaKey),
1312	StopProviding(KademliaKey),
1313	GetProviders(KademliaKey),
1314	AddKnownAddress(PeerId, Multiaddr),
1315	EventStream(out_events::Sender),
1316	Request {
1317		target: PeerId,
1318		protocol: ProtocolName,
1319		request: Vec<u8>,
1320		fallback_request: Option<(Vec<u8>, ProtocolName)>,
1321		pending_response: oneshot::Sender<Result<(Vec<u8>, ProtocolName), RequestFailure>>,
1322		connect: IfDisconnected,
1323	},
1324	NetworkStatus {
1325		pending_response: oneshot::Sender<Result<NetworkStatus, RequestFailure>>,
1326	},
1327	NetworkState {
1328		pending_response: oneshot::Sender<Result<NetworkState, RequestFailure>>,
1329	},
1330	DisconnectPeer(PeerId, ProtocolName),
1331}
1332
1333/// Main network worker. Must be polled in order for the network to advance.
1334///
1335/// You are encouraged to poll this in a separate background thread or task.
1336#[must_use = "The NetworkWorker must be polled in order for the network to advance"]
1337pub struct NetworkWorker<B, H>
1338where
1339	B: BlockT + 'static,
1340	H: ExHashT,
1341{
1342	/// Updated by the `NetworkWorker` and loaded by the `NetworkService`.
1343	listen_addresses: Arc<Mutex<HashSet<Multiaddr>>>,
1344	/// Updated by the `NetworkWorker` and loaded by the `NetworkService`.
1345	num_connected: Arc<AtomicUsize>,
1346	/// The network service that can be extracted and shared through the codebase.
1347	service: Arc<NetworkService<B, H>>,
1348	/// The *actual* network.
1349	network_service: Swarm<Behaviour<B>>,
1350	/// Messages from the [`NetworkService`] that must be processed.
1351	from_service: TracingUnboundedReceiver<ServiceToWorkerMsg>,
1352	/// Senders for events that happen on the network.
1353	event_streams: out_events::OutChannels,
1354	/// Prometheus network metrics.
1355	metrics: Option<Metrics>,
1356	/// The `PeerId`'s of all boot nodes mapped to the registered addresses.
1357	boot_node_ids: Arc<HashMap<PeerId, Vec<Multiaddr>>>,
1358	/// Boot nodes that we already have reported as invalid.
1359	reported_invalid_boot_nodes: HashSet<PeerId>,
1360	/// Peer reputation store handle.
1361	peer_store_handle: Arc<dyn PeerStoreProvider>,
1362	/// Notification protocol handles.
1363	notif_protocol_handles: Vec<protocol::ProtocolHandle>,
1364	/// Marker to pin the `H` generic. Serves no purpose except to not break backwards
1365	/// compatibility.
1366	_marker: PhantomData<H>,
1367	/// Marker for block type
1368	_block: PhantomData<B>,
1369}
1370
1371impl<B, H> NetworkWorker<B, H>
1372where
1373	B: BlockT + 'static,
1374	H: ExHashT,
1375{
1376	/// Run the network.
1377	pub async fn run(mut self) {
1378		while self.next_action().await {}
1379	}
1380
1381	/// Perform one action on the network.
1382	///
1383	/// Returns `false` when the worker should be shutdown.
1384	/// Use in tests only.
1385	pub async fn next_action(&mut self) -> bool {
1386		futures::select! {
1387			// Next message from the service.
1388			msg = self.from_service.next() => {
1389				if let Some(msg) = msg {
1390					self.handle_worker_message(msg);
1391				} else {
1392					return false
1393				}
1394			},
1395			// Next event from `Swarm` (the stream guaranteed to never terminate).
1396			event = self.network_service.select_next_some() => {
1397				self.handle_swarm_event(event);
1398			},
1399		};
1400
1401		// Update the `num_connected` count shared with the `NetworkService`.
1402		let num_connected_peers = self.network_service.behaviour().user_protocol().num_sync_peers();
1403		self.num_connected.store(num_connected_peers, Ordering::Relaxed);
1404
1405		if let Some(metrics) = self.metrics.as_ref() {
1406			if let Some(buckets) = self.network_service.behaviour_mut().num_entries_per_kbucket() {
1407				for (lower_ilog2_bucket_bound, num_entries) in buckets {
1408					metrics
1409						.kbuckets_num_nodes
1410						.with_label_values(&[&lower_ilog2_bucket_bound.to_string()])
1411						.set(num_entries as u64);
1412				}
1413			}
1414			if let Some(num_entries) = self.network_service.behaviour_mut().num_kademlia_records() {
1415				metrics.kademlia_records_count.set(num_entries as u64);
1416			}
1417			if let Some(num_entries) =
1418				self.network_service.behaviour_mut().kademlia_records_total_size()
1419			{
1420				metrics.kademlia_records_sizes_total.set(num_entries as u64);
1421			}
1422
1423			metrics.pending_connections.set(
1424				Swarm::network_info(&self.network_service).connection_counters().num_pending()
1425					as u64,
1426			);
1427		}
1428
1429		true
1430	}
1431
1432	/// Process the next message coming from the `NetworkService`.
1433	fn handle_worker_message(&mut self, msg: ServiceToWorkerMsg) {
1434		match msg {
1435			ServiceToWorkerMsg::FindClosestPeers(target) => {
1436				self.network_service.behaviour_mut().find_closest_peers(target)
1437			},
1438			ServiceToWorkerMsg::GetValue(key) => {
1439				self.network_service.behaviour_mut().get_value(key.into())
1440			},
1441			ServiceToWorkerMsg::PutValue(key, value) => {
1442				self.network_service.behaviour_mut().put_value(key.into(), value)
1443			},
1444			ServiceToWorkerMsg::PutRecordTo { record, peers, update_local_storage } => self
1445				.network_service
1446				.behaviour_mut()
1447				.put_record_to(record.into(), peers, update_local_storage),
1448			ServiceToWorkerMsg::StoreRecord(key, value, publisher, expires) => self
1449				.network_service
1450				.behaviour_mut()
1451				.store_record(key.into(), value, publisher, expires),
1452			ServiceToWorkerMsg::StartProviding(key) => {
1453				self.network_service.behaviour_mut().start_providing(key.into())
1454			},
1455			ServiceToWorkerMsg::StopProviding(key) => {
1456				self.network_service.behaviour_mut().stop_providing(&key.into())
1457			},
1458			ServiceToWorkerMsg::GetProviders(key) => {
1459				self.network_service.behaviour_mut().get_providers(key.into())
1460			},
1461			ServiceToWorkerMsg::AddKnownAddress(peer_id, addr) => {
1462				self.network_service.behaviour_mut().add_known_address(peer_id, addr)
1463			},
1464			ServiceToWorkerMsg::EventStream(sender) => self.event_streams.push(sender),
1465			ServiceToWorkerMsg::Request {
1466				target,
1467				protocol,
1468				request,
1469				fallback_request,
1470				pending_response,
1471				connect,
1472			} => {
1473				self.network_service.behaviour_mut().send_request(
1474					&target,
1475					protocol,
1476					request,
1477					fallback_request,
1478					pending_response,
1479					connect,
1480				);
1481			},
1482			ServiceToWorkerMsg::NetworkStatus { pending_response } => {
1483				let _ = pending_response.send(Ok(self.status()));
1484			},
1485			ServiceToWorkerMsg::NetworkState { pending_response } => {
1486				let _ = pending_response.send(Ok(self.network_state()));
1487			},
1488			ServiceToWorkerMsg::DisconnectPeer(who, protocol_name) => self
1489				.network_service
1490				.behaviour_mut()
1491				.user_protocol_mut()
1492				.disconnect_peer(&who, protocol_name),
1493		}
1494	}
1495
1496	/// Process the next event coming from `Swarm`.
1497	fn handle_swarm_event(&mut self, event: SwarmEvent<BehaviourOut>) {
1498		match event {
1499			SwarmEvent::Behaviour(BehaviourOut::InboundRequest { protocol, result, .. }) => {
1500				if let Some(metrics) = self.metrics.as_ref() {
1501					match result {
1502						Ok(serve_time) => {
1503							metrics
1504								.requests_in_success_total
1505								.with_label_values(&[&protocol])
1506								.observe(serve_time.as_secs_f64());
1507						},
1508						Err(err) => {
1509							let reason = match err {
1510								ResponseFailure::Network(InboundFailure::Timeout) => {
1511									Some("timeout")
1512								},
1513								ResponseFailure::Network(InboundFailure::UnsupportedProtocols) =>
1514								// `UnsupportedProtocols` is reported for every single
1515								// inbound request whenever a request with an unsupported
1516								// protocol is received. This is not reported in order to
1517								// avoid confusions.
1518								{
1519									None
1520								},
1521								ResponseFailure::Network(InboundFailure::ResponseOmission) => {
1522									Some("busy-omitted")
1523								},
1524								ResponseFailure::Network(InboundFailure::ConnectionClosed) => {
1525									Some("connection-closed")
1526								},
1527								ResponseFailure::Network(InboundFailure::Io(_)) => Some("io"),
1528							};
1529
1530							if let Some(reason) = reason {
1531								metrics
1532									.requests_in_failure_total
1533									.with_label_values(&[&protocol, reason])
1534									.inc();
1535							}
1536						},
1537					}
1538				}
1539			},
1540			SwarmEvent::Behaviour(BehaviourOut::RequestFinished {
1541				protocol,
1542				duration,
1543				result,
1544				..
1545			}) => {
1546				if let Some(metrics) = self.metrics.as_ref() {
1547					match result {
1548						Ok(_) => {
1549							metrics
1550								.requests_out_success_total
1551								.with_label_values(&[&protocol])
1552								.observe(duration.as_secs_f64());
1553						},
1554						Err(err) => {
1555							let reason = match err {
1556								RequestFailure::NotConnected => "not-connected",
1557								RequestFailure::UnknownProtocol => "unknown-protocol",
1558								RequestFailure::InvalidRequest => "invalid-request",
1559								RequestFailure::Refused => "refused",
1560								RequestFailure::Obsolete => "obsolete",
1561								RequestFailure::Network(OutboundFailure::DialFailure) => {
1562									"dial-failure"
1563								},
1564								RequestFailure::Network(OutboundFailure::Timeout) => "timeout",
1565								RequestFailure::Network(OutboundFailure::ConnectionClosed) => {
1566									"connection-closed"
1567								},
1568								RequestFailure::Network(OutboundFailure::UnsupportedProtocols) => {
1569									"unsupported"
1570								},
1571								RequestFailure::Network(OutboundFailure::Io(_)) => "io",
1572							};
1573
1574							metrics
1575								.requests_out_failure_total
1576								.with_label_values(&[&protocol, reason])
1577								.inc();
1578						},
1579					}
1580				}
1581			},
1582			SwarmEvent::Behaviour(BehaviourOut::ReputationChanges { peer, changes }) => {
1583				for change in changes {
1584					self.peer_store_handle.report_peer(peer.into(), change);
1585				}
1586			},
1587			SwarmEvent::Behaviour(BehaviourOut::PeerIdentify {
1588				peer_id,
1589				info:
1590					IdentifyInfo {
1591						protocol_version, agent_version, mut listen_addrs, protocols, ..
1592					},
1593			}) => {
1594				if listen_addrs.len() > 30 {
1595					debug!(
1596						target: LOG_TARGET,
1597						"Node {:?} has reported more than 30 addresses; it is identified by {:?} and {:?}",
1598						peer_id, protocol_version, agent_version
1599					);
1600					listen_addrs.truncate(30);
1601				}
1602				for addr in listen_addrs {
1603					self.network_service.behaviour_mut().add_self_reported_address_to_dht(
1604						&peer_id,
1605						&protocols,
1606						addr.clone(),
1607					);
1608				}
1609				self.peer_store_handle.add_known_peer(peer_id.into());
1610			},
1611			SwarmEvent::Behaviour(BehaviourOut::Discovered(peer_id)) => {
1612				self.peer_store_handle.add_known_peer(peer_id.into());
1613			},
1614			SwarmEvent::Behaviour(BehaviourOut::RandomKademliaStarted) => {
1615				if let Some(metrics) = self.metrics.as_ref() {
1616					metrics.kademlia_random_queries_total.inc();
1617				}
1618			},
1619			SwarmEvent::Behaviour(BehaviourOut::NotificationStreamOpened {
1620				remote,
1621				set_id,
1622				direction,
1623				negotiated_fallback,
1624				notifications_sink,
1625				received_handshake,
1626			}) => {
1627				let _ = self.notif_protocol_handles[usize::from(set_id)].report_substream_opened(
1628					remote,
1629					direction,
1630					received_handshake,
1631					negotiated_fallback,
1632					notifications_sink,
1633				);
1634			},
1635			SwarmEvent::Behaviour(BehaviourOut::NotificationStreamReplaced {
1636				remote,
1637				set_id,
1638				notifications_sink,
1639			}) => {
1640				let _ = self.notif_protocol_handles[usize::from(set_id)]
1641					.report_notification_sink_replaced(remote, notifications_sink);
1642
1643				// TODO: Notifications might have been lost as a result of the previous
1644				// connection being dropped, and as a result it would be preferable to notify
1645				// the users of this fact by simulating the substream being closed then
1646				// reopened.
1647				// The code below doesn't compile because `role` is unknown. Propagating the
1648				// handshake of the secondary connections is quite an invasive change and
1649				// would conflict with https://github.com/paritytech/substrate/issues/6403.
1650				// Considering that dropping notifications is generally regarded as
1651				// acceptable, this bug is at the moment intentionally left there and is
1652				// intended to be fixed at the same time as
1653				// https://github.com/paritytech/substrate/issues/6403.
1654				// self.event_streams.send(Event::NotificationStreamClosed {
1655				// remote,
1656				// protocol,
1657				// });
1658				// self.event_streams.send(Event::NotificationStreamOpened {
1659				// remote,
1660				// protocol,
1661				// role,
1662				// });
1663			},
1664			SwarmEvent::Behaviour(BehaviourOut::NotificationStreamClosed { remote, set_id }) => {
1665				let _ = self.notif_protocol_handles[usize::from(set_id)]
1666					.report_substream_closed(remote);
1667			},
1668			SwarmEvent::Behaviour(BehaviourOut::NotificationsReceived {
1669				remote,
1670				set_id,
1671				notification,
1672			}) => {
1673				let _ = self.notif_protocol_handles[usize::from(set_id)]
1674					.report_notification_received(remote, notification);
1675			},
1676			SwarmEvent::Behaviour(BehaviourOut::Dht(event, duration)) => {
1677				match (self.metrics.as_ref(), duration) {
1678					(Some(metrics), Some(duration)) => {
1679						let query_type = match event {
1680							DhtEvent::ClosestPeersFound(_, _) => "peers-found",
1681							DhtEvent::ClosestPeersNotFound(_) => "peers-not-found",
1682							DhtEvent::ValueFound(_) => "value-found",
1683							DhtEvent::ValueNotFound(_) => "value-not-found",
1684							DhtEvent::ValuePut(_) => "value-put",
1685							DhtEvent::ValuePutFailed(_) => "value-put-failed",
1686							DhtEvent::PutRecordRequest(_, _, _, _) => "put-record-request",
1687							DhtEvent::StartedProviding(_) => "started-providing",
1688							DhtEvent::StartProvidingFailed(_) => "start-providing-failed",
1689							DhtEvent::ProvidersFound(_, _) => "providers-found",
1690							DhtEvent::NoMoreProviders(_) => "no-more-providers",
1691							DhtEvent::ProvidersNotFound(_) => "providers-not-found",
1692						};
1693						metrics
1694							.kademlia_query_duration
1695							.with_label_values(&[query_type])
1696							.observe(duration.as_secs_f64());
1697					},
1698					_ => {},
1699				}
1700
1701				self.event_streams.send(Event::Dht(event));
1702			},
1703			SwarmEvent::Behaviour(BehaviourOut::None) => {
1704				// Ignored event from lower layers.
1705			},
1706			SwarmEvent::ConnectionEstablished {
1707				peer_id,
1708				endpoint,
1709				num_established,
1710				concurrent_dial_errors,
1711				..
1712			} => {
1713				if let Some(errors) = concurrent_dial_errors {
1714					debug!(target: LOG_TARGET, "Libp2p => Connected({:?}) with errors: {:?}", peer_id, errors);
1715				} else {
1716					debug!(target: LOG_TARGET, "Libp2p => Connected({:?})", peer_id);
1717				}
1718
1719				if let Some(metrics) = self.metrics.as_ref() {
1720					let direction = match endpoint {
1721						ConnectedPoint::Dialer { .. } => "out",
1722						ConnectedPoint::Listener { .. } => "in",
1723					};
1724					metrics.connections_opened_total.with_label_values(&[direction]).inc();
1725
1726					if num_established.get() == 1 {
1727						metrics.distinct_peers_connections_opened_total.inc();
1728					}
1729				}
1730			},
1731			SwarmEvent::ConnectionClosed {
1732				connection_id,
1733				peer_id,
1734				cause,
1735				endpoint,
1736				num_established,
1737			} => {
1738				debug!(target: LOG_TARGET, "Libp2p => Disconnected({peer_id:?} via {connection_id:?}, {cause:?})");
1739				if let Some(metrics) = self.metrics.as_ref() {
1740					let direction = match endpoint {
1741						ConnectedPoint::Dialer { .. } => "out",
1742						ConnectedPoint::Listener { .. } => "in",
1743					};
1744					let reason = match cause {
1745						Some(ConnectionError::IO(_)) => "transport-error",
1746						Some(ConnectionError::KeepAliveTimeout) => "keep-alive-timeout",
1747						None => "actively-closed",
1748					};
1749					metrics.connections_closed_total.with_label_values(&[direction, reason]).inc();
1750
1751					// `num_established` represents the number of *remaining* connections.
1752					if num_established == 0 {
1753						metrics.distinct_peers_connections_closed_total.inc();
1754					}
1755				}
1756			},
1757			SwarmEvent::NewListenAddr { address, .. } => {
1758				trace!(target: LOG_TARGET, "Libp2p => NewListenAddr({})", address);
1759				if let Some(metrics) = self.metrics.as_ref() {
1760					metrics.listeners_local_addresses.inc();
1761				}
1762				self.listen_addresses.lock().insert(address.clone());
1763			},
1764			SwarmEvent::ExpiredListenAddr { address, .. } => {
1765				info!(target: LOG_TARGET, "๐Ÿ“ช No longer listening on {}", address);
1766				if let Some(metrics) = self.metrics.as_ref() {
1767					metrics.listeners_local_addresses.dec();
1768				}
1769				self.listen_addresses.lock().remove(&address);
1770			},
1771			SwarmEvent::OutgoingConnectionError { connection_id, peer_id, error } => {
1772				if let Some(peer_id) = peer_id {
1773					trace!(
1774						target: LOG_TARGET,
1775						"Libp2p => Failed to reach {peer_id:?} via {connection_id:?}: {error}",
1776					);
1777
1778					let not_reported = !self.reported_invalid_boot_nodes.contains(&peer_id);
1779
1780					if let Some(addresses) =
1781						not_reported.then(|| self.boot_node_ids.get(&peer_id)).flatten()
1782					{
1783						if let DialError::WrongPeerId { obtained, endpoint } = &error {
1784							if let ConnectedPoint::Dialer {
1785								address,
1786								role_override: _,
1787								port_use: _,
1788							} = endpoint
1789							{
1790								let address_without_peer_id = parse_addr(address.clone().into())
1791									.map_or_else(|_| address.clone(), |r| r.1.into());
1792
1793								// Only report for address of boot node that was added at startup of
1794								// the node and not for any address that the node learned of the
1795								// boot node.
1796								if addresses.iter().any(|a| address_without_peer_id == *a) {
1797									warn!(
1798										"๐Ÿ’” The bootnode you want to connect to at `{address}` provided a \
1799										 different peer ID `{obtained}` than the one you expect `{peer_id}`.",
1800									);
1801
1802									self.reported_invalid_boot_nodes.insert(peer_id);
1803								}
1804							}
1805						}
1806					}
1807				}
1808
1809				if let Some(metrics) = self.metrics.as_ref() {
1810					let reason = match error {
1811						DialError::Denied { cause } => {
1812							if cause.downcast::<Exceeded>().is_ok() {
1813								Some("limit-reached")
1814							} else {
1815								None
1816							}
1817						},
1818						DialError::LocalPeerId { .. } => Some("local-peer-id"),
1819						DialError::WrongPeerId { .. } => Some("invalid-peer-id"),
1820						DialError::Transport(_) => Some("transport-error"),
1821						DialError::NoAddresses |
1822						DialError::DialPeerConditionFalse(_) |
1823						DialError::Aborted => None, // ignore them
1824					};
1825					if let Some(reason) = reason {
1826						metrics.pending_connections_errors_total.with_label_values(&[reason]).inc();
1827					}
1828				}
1829			},
1830			SwarmEvent::Dialing { connection_id, peer_id } => {
1831				trace!(target: LOG_TARGET, "Libp2p => Dialing({peer_id:?}) via {connection_id:?}")
1832			},
1833			SwarmEvent::IncomingConnection { connection_id, local_addr, send_back_addr } => {
1834				trace!(target: LOG_TARGET, "Libp2p => IncomingConnection({local_addr},{send_back_addr} via {connection_id:?}))");
1835				if let Some(metrics) = self.metrics.as_ref() {
1836					metrics.incoming_connections_total.inc();
1837				}
1838			},
1839			SwarmEvent::IncomingConnectionError {
1840				connection_id,
1841				local_addr,
1842				send_back_addr,
1843				error,
1844			} => {
1845				debug!(
1846					target: LOG_TARGET,
1847					"Libp2p => IncomingConnectionError({local_addr},{send_back_addr} via {connection_id:?}): {error}"
1848				);
1849				if let Some(metrics) = self.metrics.as_ref() {
1850					let reason = match error {
1851						ListenError::Denied { cause } => {
1852							if cause.downcast::<Exceeded>().is_ok() {
1853								Some("limit-reached")
1854							} else {
1855								None
1856							}
1857						},
1858						ListenError::WrongPeerId { .. } | ListenError::LocalPeerId { .. } => {
1859							Some("invalid-peer-id")
1860						},
1861						ListenError::Transport(_) => Some("transport-error"),
1862						ListenError::Aborted => None, // ignore it
1863					};
1864
1865					if let Some(reason) = reason {
1866						metrics
1867							.incoming_connections_errors_total
1868							.with_label_values(&[reason])
1869							.inc();
1870					}
1871				}
1872			},
1873			SwarmEvent::ListenerClosed { reason, addresses, .. } => {
1874				if let Some(metrics) = self.metrics.as_ref() {
1875					metrics.listeners_local_addresses.sub(addresses.len() as u64);
1876				}
1877				let mut listen_addresses = self.listen_addresses.lock();
1878				for addr in &addresses {
1879					listen_addresses.remove(addr);
1880				}
1881				drop(listen_addresses);
1882
1883				let addrs =
1884					addresses.into_iter().map(|a| a.to_string()).collect::<Vec<_>>().join(", ");
1885				match reason {
1886					Ok(()) => error!(
1887						target: LOG_TARGET,
1888						"๐Ÿ“ช Libp2p listener ({}) closed gracefully",
1889						addrs
1890					),
1891					Err(e) => error!(
1892						target: LOG_TARGET,
1893						"๐Ÿ“ช Libp2p listener ({}) closed: {}",
1894						addrs, e
1895					),
1896				}
1897			},
1898			SwarmEvent::ListenerError { error, .. } => {
1899				debug!(target: LOG_TARGET, "Libp2p => ListenerError: {}", error);
1900				if let Some(metrics) = self.metrics.as_ref() {
1901					metrics.listeners_errors_total.inc();
1902				}
1903			},
1904			SwarmEvent::NewExternalAddrCandidate { address } => {
1905				trace!(target: LOG_TARGET, "Libp2p => NewExternalAddrCandidate: {address:?}");
1906			},
1907			SwarmEvent::ExternalAddrConfirmed { address } => {
1908				trace!(target: LOG_TARGET, "Libp2p => ExternalAddrConfirmed: {address:?}");
1909			},
1910			SwarmEvent::ExternalAddrExpired { address } => {
1911				trace!(target: LOG_TARGET, "Libp2p => ExternalAddrExpired: {address:?}");
1912			},
1913			SwarmEvent::NewExternalAddrOfPeer { peer_id, address } => {
1914				trace!(target: LOG_TARGET, "Libp2p => NewExternalAddrOfPeer({peer_id:?}): {address:?}")
1915			},
1916			event => {
1917				warn!(target: LOG_TARGET, "New unknown SwarmEvent libp2p event: {event:?}");
1918			},
1919		}
1920	}
1921}
1922
1923impl<B, H> Unpin for NetworkWorker<B, H>
1924where
1925	B: BlockT + 'static,
1926	H: ExHashT,
1927{
1928}
1929
1930pub(crate) fn ensure_addresses_consistent_with_transport<'a>(
1931	addresses: impl Iterator<Item = &'a sc_network_types::multiaddr::Multiaddr>,
1932	transport: &TransportConfig,
1933) -> Result<(), Error> {
1934	use sc_network_types::multiaddr::Protocol;
1935
1936	if matches!(transport, TransportConfig::MemoryOnly) {
1937		let addresses: Vec<_> = addresses
1938			.filter(|x| x.iter().any(|y| !matches!(y, Protocol::Memory(_))))
1939			.cloned()
1940			.collect();
1941
1942		if !addresses.is_empty() {
1943			return Err(Error::AddressesForAnotherTransport {
1944				transport: transport.clone(),
1945				addresses,
1946			});
1947		}
1948	} else {
1949		let addresses: Vec<_> = addresses
1950			.filter(|x| x.iter().any(|y| matches!(y, Protocol::Memory(_))))
1951			.cloned()
1952			.collect();
1953
1954		if !addresses.is_empty() {
1955			return Err(Error::AddressesForAnotherTransport {
1956				transport: transport.clone(),
1957				addresses,
1958			});
1959		}
1960	}
1961
1962	Ok(())
1963}