smoldot_light/
network_service.rs

1// Smoldot
2// Copyright (C) 2019-2022  Parity Technologies (UK) Ltd.
3// SPDX-License-Identifier: GPL-3.0-or-later WITH Classpath-exception-2.0
4
5// This program is free software: you can redistribute it and/or modify
6// it under the terms of the GNU General Public License as published by
7// the Free Software Foundation, either version 3 of the License, or
8// (at your option) any later version.
9
10// This program is distributed in the hope that it will be useful,
11// but WITHOUT ANY WARRANTY; without even the implied warranty of
12// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
13// GNU General Public License for more details.
14
15// You should have received a copy of the GNU General Public License
16// along with this program.  If not, see <http://www.gnu.org/licenses/>.
17
18//! Background network service.
19//!
20//! The [`NetworkService`] manages background tasks dedicated to connecting to other nodes.
21//! Importantly, its design is oriented towards the particular use case of the light client.
22//!
23//! The [`NetworkService`] spawns one background task (using [`PlatformRef::spawn_task`]) for
24//! each active connection.
25//!
26//! The objective of the [`NetworkService`] in general is to try stay connected as much as
27//! possible to the nodes of the peer-to-peer network of the chain, and maintain open substreams
28//! with them in order to send out requests (e.g. block requests) and notifications (e.g. block
29//! announces).
30//!
31//! Connectivity to the network is performed in the background as an implementation detail of
32//! the service. The public API only allows emitting requests and notifications towards the
33//! already-connected nodes.
34//!
35//! After a [`NetworkService`] is created, one can add chains using [`NetworkService::add_chain`].
36//! If all references to a [`NetworkServiceChain`] are destroyed, the chain is automatically
37//! purged.
38//!
39//! An important part of the API is the list of channel receivers of [`Event`] returned by
40//! [`NetworkServiceChain::subscribe`]. These channels inform the foreground about updates to the
41//! network connectivity.
42
43use crate::{
44    log, metrics,
45    platform::{self, PlatformRef, address_parse},
46};
47
48use alloc::{
49    borrow::ToOwned as _,
50    boxed::Box,
51    collections::{BTreeMap, VecDeque},
52    format,
53    string::{String, ToString as _},
54    sync::Arc,
55    vec::{self, Vec},
56};
57use core::{cmp, mem, num::NonZero, pin::Pin, time::Duration};
58use futures_channel::oneshot;
59use futures_lite::FutureExt as _;
60use futures_util::{StreamExt as _, future, stream};
61use hashbrown::{HashMap, HashSet};
62use rand::seq::IteratorRandom as _;
63use rand_chacha::rand_core::SeedableRng as _;
64use smoldot::{
65    header,
66    informant::{BytesDisplay, HashDisplay},
67    libp2p::{
68        connection,
69        multiaddr::{self, Multiaddr},
70        peer_id,
71    },
72    network::{basic_peering_strategy, bitswap_peering_strategy, codec, service},
73};
74
75pub use codec::{AffinityFilter, CallProofRequestConfig, Role};
76use service::SendTopicAffinityError;
77pub use service::{
78    ChainId, EncodedMerkleProof, PeerId, QueueNotificationError, SendBitswapMessageError,
79};
80
81mod tasks;
82
83/// Configuration for a [`NetworkService`].
84pub struct Config<TPlat> {
85    /// Access to the platform's capabilities.
86    pub platform: TPlat,
87
88    /// Value sent back for the agent version when receiving an identification request.
89    pub identify_agent_version: String,
90
91    /// Capacity to allocate for the list of chains.
92    pub chains_capacity: usize,
93
94    /// Maximum number of connections that the service can open simultaneously. After this value
95    /// has been reached, a new connection can be opened after each
96    /// [`Config::connections_open_pool_restore_delay`].
97    pub connections_open_pool_size: u32,
98
99    /// Delay after which the service can open a new connection.
100    /// The delay is cumulative. If no connection has been opened for example for twice this
101    /// duration, then two connections can be opened at the same time, up to a maximum of
102    /// [`Config::connections_open_pool_size`].
103    pub connections_open_pool_restore_delay: Duration,
104}
105
106/// See [`NetworkService::add_chain`].
107///
108/// Note that this configuration is intentionally missing a field containing the bootstrap
109/// nodes of the chain. Bootstrap nodes are supposed to be added afterwards by calling
110/// [`NetworkServiceChain::discover`].
111pub struct ConfigChain {
112    /// Name of the chain, for logging purposes.
113    pub log_name: String,
114
115    /// Number of "out slots" of this chain. We establish simultaneously gossip substreams up to
116    /// this number of peers.
117    pub num_out_slots: usize,
118
119    /// Hash of the genesis block of the chain. Sent to other nodes in order to determine whether
120    /// the chains match.
121    ///
122    /// > **Note**: Be aware that this *must* be the *genesis* block, not any block known to be
123    /// >           in the chain.
124    pub genesis_block_hash: [u8; 32],
125
126    /// Number and hash of the current best block. Can later be updated with
127    /// [`NetworkServiceChain::set_local_best_block`].
128    pub best_block: (u64, [u8; 32]),
129
130    /// Optional identifier to insert into the networking protocol names. Used to differentiate
131    /// between chains with the same genesis hash.
132    pub fork_id: Option<String>,
133
134    /// Number of bytes of the block number in the networking protocol.
135    pub block_number_bytes: usize,
136
137    /// Must be `Some` if and only if the chain uses the GrandPa networking protocol. Contains the
138    /// number of the finalized block at the time of the initialization.
139    pub grandpa_protocol_finalized_block_height: Option<u64>,
140
141    /// If `true`, enables the statement store protocol.
142    pub enable_statement_protocol: bool,
143
144    /// Metrics of the chain.
145    pub metrics: Arc<metrics::ChainMetrics>,
146}
147
148pub struct NetworkService<TPlat: PlatformRef> {
149    /// Channel connected to the background service.
150    messages_tx: async_channel::Sender<ToBackground<TPlat>>,
151
152    /// See [`Config::platform`].
153    platform: TPlat,
154
155    /// Process-wide network metrics, shared with the background service.
156    metrics: Arc<metrics::NetworkMetrics>,
157}
158
159impl<TPlat: PlatformRef> NetworkService<TPlat> {
160    /// Initializes the network service with the given configuration.
161    pub fn new(config: Config<TPlat>) -> Arc<Self> {
162        let (main_messages_tx, main_messages_rx) = async_channel::bounded(4);
163
164        let metrics = Arc::new(metrics::NetworkMetrics::default());
165
166        let network = service::ChainNetwork::new(service::Config {
167            chains_capacity: config.chains_capacity,
168            connections_capacity: 32,
169            // Shortened from 8s: parallel dials hold slots until this fires.
170            handshake_timeout: Duration::from_secs(4),
171            randomness_seed: {
172                let mut seed = [0; 32];
173                config.platform.fill_random_bytes(&mut seed);
174                seed
175            },
176        });
177
178        // Spawn main task that processes the network service.
179        let (tasks_messages_tx, tasks_messages_rx) = async_channel::bounded(32);
180        let task = Box::pin(background_task(BackgroundTask {
181            randomness: rand_chacha::ChaCha20Rng::from_seed({
182                let mut seed = [0; 32];
183                config.platform.fill_random_bytes(&mut seed);
184                seed
185            }),
186            identify_agent_version: config.identify_agent_version,
187            tasks_messages_tx,
188            tasks_messages_rx: Box::pin(tasks_messages_rx),
189            peering_strategy: basic_peering_strategy::BasicPeeringStrategy::new(
190                basic_peering_strategy::Config {
191                    randomness_seed: {
192                        let mut seed = [0; 32];
193                        config.platform.fill_random_bytes(&mut seed);
194                        seed
195                    },
196                    peers_capacity: 50, // TODO: ?
197                    chains_capacity: config.chains_capacity,
198                },
199            ),
200            bitswap_peering_strategy: bitswap_peering_strategy::BitswapPeeringStrategy::new(
201                bitswap_peering_strategy::Config {
202                    randomness_seed: {
203                        let mut seed = [0; 32];
204                        config.platform.fill_random_bytes(&mut seed);
205                        seed
206                    },
207                    peers_capacity: 50, // TODO: hardcoded to the same value as `peering_strategy`.
208                },
209            ),
210            network,
211            connections_open_pool_size: config.connections_open_pool_size,
212            connections_open_pool_restore_delay: config.connections_open_pool_restore_delay,
213            num_recent_connection_opening: 0,
214            next_recent_connection_restore: None,
215            platform: config.platform.clone(),
216            metrics: metrics.clone(),
217            open_gossip_links: BTreeMap::new(),
218            chains_ever_gossip_connected: HashSet::with_capacity_and_hasher(4, Default::default()),
219            v2_statement_peers: HashMap::with_capacity_and_hasher(4, Default::default()),
220            current_affinity_filter: HashMap::with_capacity_and_hasher(4, Default::default()),
221            events_pending_send: VecDeque::with_capacity(4),
222            event_senders: either::Left(Vec::new()),
223            pending_new_subscriptions: Vec::new(),
224            bitswap_event_pending_send: None,
225            bitswap_connected_peers: 0,
226            bitswap_event_senders: either::Left(Vec::new()),
227            pending_new_bitswap_subscriptions: Vec::new(),
228            statement_event_pending_send: None,
229            statement_event_senders: either::Left(Vec::new()),
230            pending_new_statement_subscriptions: Vec::new(),
231            important_nodes: HashMap::with_capacity_and_hasher(16, Default::default()),
232            main_messages_rx: Box::pin(main_messages_rx),
233            messages_rx: stream::SelectAll::new(),
234            blocks_requests: HashMap::with_capacity_and_hasher(8, Default::default()),
235            grandpa_warp_sync_requests: HashMap::with_capacity_and_hasher(8, Default::default()),
236            storage_proof_requests: HashMap::with_capacity_and_hasher(8, Default::default()),
237            call_proof_requests: HashMap::with_capacity_and_hasher(8, Default::default()),
238            child_storage_proof_requests: HashMap::with_capacity_and_hasher(8, Default::default()),
239            chains_by_next_discovery: BTreeMap::new(),
240        }));
241
242        config.platform.spawn_task("network-service".into(), {
243            let platform = config.platform.clone();
244            async move {
245                task.await;
246                log!(&platform, Debug, "network", "shutdown");
247            }
248        });
249
250        Arc::new(NetworkService {
251            messages_tx: main_messages_tx,
252            platform: config.platform,
253            metrics,
254        })
255    }
256
257    /// Returns the process-wide network metrics.
258    pub fn metrics(&self) -> Arc<metrics::NetworkMetrics> {
259        self.metrics.clone()
260    }
261
262    /// Adds a chain to the list of chains that the network service connects to.
263    ///
264    /// Returns an object representing the chain and that allows interacting with it. If all
265    /// references to [`NetworkServiceChain`] are destroyed, the network service automatically
266    /// purges that chain.
267    pub fn add_chain(&self, config: ConfigChain) -> Arc<NetworkServiceChain<TPlat>> {
268        let (messages_tx, messages_rx) = async_channel::bounded(32);
269        let chain_metrics = config.metrics.clone();
270
271        // TODO: this code is hacky because we don't want to make `add_chain` async at the moment, because it's not convenient for lib.rs
272        self.platform.spawn_task("add-chain-message-send".into(), {
273            let config = service::ChainConfig {
274                grandpa_protocol_config: config.grandpa_protocol_finalized_block_height.map(
275                    |commit_finalized_height| service::GrandpaState {
276                        commit_finalized_height,
277                        round_number: 1,
278                        set_id: 0,
279                    },
280                ),
281                enable_statement_protocol: config.enable_statement_protocol,
282                fork_id: config.fork_id.clone(),
283                block_number_bytes: config.block_number_bytes,
284                best_hash: config.best_block.1,
285                best_number: config.best_block.0,
286                genesis_hash: config.genesis_block_hash,
287                role: Role::Light,
288                allow_inbound_block_requests: false,
289                user_data: Chain {
290                    log_name: config.log_name,
291                    block_number_bytes: config.block_number_bytes,
292                    num_out_slots: config.num_out_slots,
293                    num_references: NonZero::<usize>::new(1).unwrap(),
294                    next_discovery_period: Duration::from_secs(2),
295                    next_discovery_when: self.platform.now(),
296                    metrics: config.metrics,
297                },
298            };
299
300            let messages_tx = self.messages_tx.clone();
301            async move {
302                let _ = messages_tx
303                    .send(ToBackground::AddChain {
304                        messages_rx,
305                        config,
306                    })
307                    .await;
308            }
309        });
310
311        Arc::new(NetworkServiceChain {
312            _keep_alive_messages_tx: self.messages_tx.clone(),
313            messages_tx,
314            metrics: chain_metrics,
315            platform: self.platform.clone(),
316        })
317    }
318}
319
320pub struct NetworkServiceChain<TPlat: PlatformRef> {
321    /// Copy of [`NetworkService::messages_tx`]. Used in order to maintain the network service
322    /// background task alive.
323    _keep_alive_messages_tx: async_channel::Sender<ToBackground<TPlat>>,
324
325    /// Channel to send messages to the background task.
326    messages_tx: async_channel::Sender<ToBackgroundChain>,
327
328    /// Metrics of the chain.
329    metrics: Arc<metrics::ChainMetrics>,
330
331    /// See [`Config::platform`].
332    platform: TPlat,
333}
334
335/// Severity of a ban. See [`NetworkServiceChain::ban_and_disconnect`].
336#[derive(Debug, Copy, Clone, PartialEq, Eq)]
337pub enum BanSeverity {
338    Low,
339    High,
340}
341
342/// Reason for banning a peer. Printed in the logs and used as the `reason` label
343/// of the `networkPeerBansTotal` metric.
344// Strum derives are explained on [`crate::metrics::MetricLabel`]; `Display`
345// additionally prints the same string in log messages.
346#[derive(
347    Debug, Copy, Clone, PartialEq, Eq, strum::Display, strum::EnumIter, strum::IntoStaticStr,
348)]
349#[strum(serialize_all = "kebab-case")]
350pub enum BanReason {
351    BadBlockAnnounce,
352    BadChildTrieRoot,
353    BadGrandpaCommit,
354    BadJustification,
355    BadMerkleProof,
356    BadWarpSyncFragment,
357    InvalidCallProof,
358    BlocksRequestFailed,
359    CallProofRequestFailed,
360    ChildStorageRequestFailed,
361    StorageRequestFailed,
362    WarpSyncRequestFailed,
363}
364
365impl<TPlat: PlatformRef> NetworkServiceChain<TPlat> {
366    /// Subscribes to the networking events that happen on the given chain.
367    ///
368    /// Calling this function returns a `Receiver` that receives events about the chain.
369    /// The new channel will immediately receive events about all the existing connections, so
370    /// that it is able to maintain a coherent view of the network.
371    ///
372    /// Note that this function is `async`, but it should return very quickly.
373    ///
374    /// The `Receiver` **must** be polled continuously. When the channel is full, the networking
375    /// connections will be back-pressured until the channel isn't full anymore.
376    ///
377    /// The `Receiver` never yields `None` unless the [`NetworkService`] crashes or is destroyed.
378    /// If `None` is yielded and the [`NetworkService`] is still alive, you should call
379    /// [`NetworkServiceChain::subscribe`] again to obtain a new `Receiver`.
380    ///
381    // TODO: consider not killing the background until the channel is destroyed, as that would be a more sensical behaviour
382    pub async fn subscribe(&self) -> async_channel::Receiver<Event> {
383        let (tx, rx) = async_channel::bounded(128);
384
385        self.messages_tx
386            .send(ToBackgroundChain::Subscribe { sender: tx })
387            .await
388            .unwrap();
389
390        rx
391    }
392
393    /// Subscribes to the Bitswap events that happen on the network. Bitswap events subscription is
394    /// separate from other network service events, because Bitswap events are big and are not
395    /// needed by the most of subscribers.
396    ///
397    /// Note that this function is `async`, but it should return very quickly.
398    ///
399    /// The `Receiver` **must** be polled continuously. When the channel is full, the networking
400    /// connections will be back-pressured until the channel isn't full anymore.
401    ///
402    /// The `Receiver` never yields `None` unless the [`NetworkService`] crashes or is destroyed.
403    /// If `None` is yielded and the [`NetworkService`] is still alive, you should call
404    /// [`NetworkServiceChain::subscribe_bitswap`] again to obtain a new `Receiver`.
405    ///
406    // TODO: the last section of the doc seem to contradict itself.
407    pub async fn subscribe_bitswap(&self) -> async_channel::Receiver<BitswapEvent> {
408        let (tx, rx) = async_channel::bounded(128);
409
410        self.messages_tx
411            .send(ToBackgroundChain::SubscribeBitswap { sender: tx })
412            .await
413            .unwrap();
414
415        rx
416    }
417
418    /// Subscribes to the statement notifications that happen on the network.
419    ///
420    /// Note that this function is `async`, but it should return very quickly.
421    ///
422    /// The `Receiver` **must** be polled continuously. When the channel is full, the networking
423    /// connections will be back-pressured until the channel isn't full anymore.
424    ///
425    /// The `Receiver` never yields `None` unless the [`NetworkService`] crashes or is destroyed.
426    /// If `None` is yielded and the [`NetworkService`] is still alive, you should call
427    /// [`NetworkServiceChain::subscribe_statements`] again to obtain a new `Receiver`.
428    pub async fn subscribe_statements(&self) -> async_channel::Receiver<StatementEvent> {
429        let (tx, rx) = async_channel::bounded(128);
430
431        self.messages_tx
432            .send(ToBackgroundChain::SubscribeStatements { sender: tx })
433            .await
434            .unwrap();
435
436        rx
437    }
438
439    /// Starts asynchronously disconnecting the given peer. A [`Event::Disconnected`] will later be
440    /// generated. Prevents a new gossip link with the same peer from being reopened for a
441    /// little while.
442    ///
443    /// `reason` is printed in the logs and counted in the metrics.
444    ///
445    /// Due to race conditions, it is possible to reconnect to the peer soon after, in case the
446    /// reconnection was already happening as the call to this function is still being processed.
447    /// If that happens another [`Event::Disconnected`] will be delivered afterwards. In other
448    /// words, this function guarantees that we will be disconnected in the future rather than
449    /// guarantees that we will disconnect.
450    pub async fn ban_and_disconnect(
451        &self,
452        peer_id: PeerId,
453        severity: BanSeverity,
454        reason: BanReason,
455    ) {
456        self.metrics.peer_bans.inc(reason);
457
458        let _ = self
459            .messages_tx
460            .send(ToBackgroundChain::DisconnectAndBan {
461                peer_id,
462                severity,
463                reason,
464            })
465            .await;
466    }
467
468    /// Sends a blocks request to the given peer.
469    // TODO: more docs
470    pub async fn blocks_request(
471        self: Arc<Self>,
472        target: PeerId,
473        config: codec::BlocksRequestConfig,
474        timeout: Duration,
475    ) -> Result<Vec<codec::BlockData>, BlocksRequestError> {
476        let when_started = self.platform.now();
477        let (tx, rx) = oneshot::channel();
478
479        self.messages_tx
480            .send(ToBackgroundChain::StartBlocksRequest {
481                target: target.clone(),
482                config,
483                timeout,
484                result: tx,
485            })
486            .await
487            .unwrap();
488
489        let result = rx.await.unwrap();
490        self.metrics
491            .blocks_requests
492            .observe(result.is_ok(), self.platform.now() - when_started);
493        result
494    }
495
496    /// Sends a grandpa warp sync request to the given peer.
497    // TODO: more docs
498    pub async fn grandpa_warp_sync_request(
499        self: Arc<Self>,
500        target: PeerId,
501        begin_hash: [u8; 32],
502        timeout: Duration,
503    ) -> Result<service::EncodedGrandpaWarpSyncResponse, WarpSyncRequestError> {
504        let when_started = self.platform.now();
505        let (tx, rx) = oneshot::channel();
506
507        self.messages_tx
508            .send(ToBackgroundChain::StartWarpSyncRequest {
509                target: target.clone(),
510                begin_hash,
511                timeout,
512                result: tx,
513            })
514            .await
515            .unwrap();
516
517        let result = rx.await.unwrap();
518        self.metrics
519            .warp_sync_requests
520            .observe(result.is_ok(), self.platform.now() - when_started);
521        result
522    }
523
524    pub async fn set_local_best_block(&self, best_hash: [u8; 32], best_number: u64) {
525        self.messages_tx
526            .send(ToBackgroundChain::SetLocalBestBlock {
527                best_hash,
528                best_number,
529            })
530            .await
531            .unwrap();
532    }
533
534    pub async fn set_local_grandpa_state(&self, grandpa_state: service::GrandpaState) {
535        self.messages_tx
536            .send(ToBackgroundChain::SetLocalGrandpaState { grandpa_state })
537            .await
538            .unwrap();
539    }
540
541    /// Sends a storage proof request to the given peer.
542    // TODO: more docs
543    pub async fn storage_proof_request(
544        self: Arc<Self>,
545        target: PeerId, // TODO: takes by value because of futures longevity issue
546        config: codec::StorageProofRequestConfig<impl Iterator<Item = impl AsRef<[u8]> + Clone>>,
547        timeout: Duration,
548    ) -> Result<service::EncodedMerkleProof, StorageProofRequestError> {
549        let when_started = self.platform.now();
550        let (tx, rx) = oneshot::channel();
551
552        self.messages_tx
553            .send(ToBackgroundChain::StartStorageProofRequest {
554                target: target.clone(),
555                config: codec::StorageProofRequestConfig {
556                    block_hash: config.block_hash,
557                    keys: config
558                        .keys
559                        .map(|key| key.as_ref().to_vec()) // TODO: to_vec() overhead
560                        .collect::<Vec<_>>()
561                        .into_iter(),
562                },
563                timeout,
564                result: tx,
565            })
566            .await
567            .unwrap();
568
569        let result = rx.await.unwrap();
570        self.metrics
571            .storage_proof_requests
572            .observe(result.is_ok(), self.platform.now() - when_started);
573        result
574    }
575
576    /// Sends a call proof request to the given peer.
577    ///
578    /// See also [`NetworkServiceChain::call_proof_request`].
579    // TODO: more docs
580    pub async fn call_proof_request(
581        self: Arc<Self>,
582        target: PeerId, // TODO: takes by value because of futures longevity issue
583        config: codec::CallProofRequestConfig<'_, impl Iterator<Item = impl AsRef<[u8]>>>,
584        timeout: Duration,
585    ) -> Result<EncodedMerkleProof, CallProofRequestError> {
586        let when_started = self.platform.now();
587        let (tx, rx) = oneshot::channel();
588
589        self.messages_tx
590            .send(ToBackgroundChain::StartCallProofRequest {
591                target: target.clone(),
592                config: codec::CallProofRequestConfig {
593                    block_hash: config.block_hash,
594                    method: config.method.into_owned().into(),
595                    parameter_vectored: config
596                        .parameter_vectored
597                        .map(|v| v.as_ref().to_vec()) // TODO: to_vec() overhead
598                        .collect::<Vec<_>>()
599                        .into_iter(),
600                },
601                timeout,
602                result: tx,
603            })
604            .await
605            .unwrap();
606
607        let result = rx.await.unwrap();
608        self.metrics
609            .call_proof_requests
610            .observe(result.is_ok(), self.platform.now() - when_started);
611        result
612    }
613
614    /// Sends a child storage proof request to the given peer.
615    pub async fn child_storage_proof_request(
616        self: Arc<Self>,
617        target: PeerId,
618        config: codec::ChildStorageProofRequestConfig<
619            impl AsRef<[u8]> + Clone,
620            impl Iterator<Item = impl AsRef<[u8]> + Clone>,
621        >,
622        timeout: Duration,
623    ) -> Result<service::EncodedMerkleProof, ChildStorageProofRequestError> {
624        let when_started = self.platform.now();
625        let (tx, rx) = oneshot::channel();
626
627        self.messages_tx
628            .send(ToBackgroundChain::StartChildStorageProofRequest {
629                target: target.clone(),
630                config: ChildStorageProofRequestConfigOwned {
631                    block_hash: config.block_hash,
632                    child_trie: config.child_trie.as_ref().to_vec(),
633                    keys: config
634                        .keys
635                        .map(|key| key.as_ref().to_vec())
636                        .collect::<Vec<_>>(),
637                },
638                timeout,
639                result: tx,
640            })
641            .await
642            .unwrap();
643
644        let result = rx.await.unwrap();
645        self.metrics
646            .child_storage_proof_requests
647            .observe(result.is_ok(), self.platform.now() - when_started);
648        result
649    }
650
651    /// Announces transaction to the peers we are connected to.
652    ///
653    /// Returns a list of peers that we have sent the transaction to. Can return an empty `Vec`
654    /// if we didn't send the transaction to any peer.
655    ///
656    /// Note that the remote doesn't confirm that it has received the transaction. Because
657    /// networking is inherently unreliable, successfully sending a transaction to a peer doesn't
658    /// necessarily mean that the remote has received it. In practice, however, the likelihood of
659    /// a transaction not being received are extremely low. This can be considered as known flaw.
660    pub async fn announce_transaction(self: Arc<Self>, transaction: &[u8]) -> Vec<PeerId> {
661        let (tx, rx) = oneshot::channel();
662
663        self.messages_tx
664            .send(ToBackgroundChain::AnnounceTransaction {
665                transaction: transaction.to_vec(), // TODO: ovheread
666                result: tx,
667            })
668            .await
669            .unwrap();
670
671        rx.await.unwrap()
672    }
673
674    /// See [`service::ChainNetwork::gossip_send_block_announce`].
675    pub async fn send_block_announce(
676        self: Arc<Self>,
677        target: &PeerId,
678        scale_encoded_header: &[u8],
679        is_best: bool,
680    ) -> Result<(), QueueNotificationError> {
681        let (tx, rx) = oneshot::channel();
682
683        self.messages_tx
684            .send(ToBackgroundChain::SendBlockAnnounce {
685                target: target.clone(),                              // TODO: overhead
686                scale_encoded_header: scale_encoded_header.to_vec(), // TODO: overhead
687                is_best,
688                result: tx,
689            })
690            .await
691            .unwrap();
692
693        rx.await.unwrap()
694    }
695
696    /// Send Bitswap message to the given peer.
697    pub async fn send_bitswap_message(
698        &self,
699        target: PeerId,
700        message: Vec<u8>,
701    ) -> Result<(), SendBitswapMessageError> {
702        let (tx, rx) = oneshot::channel();
703
704        self.messages_tx
705            .send(ToBackgroundChain::SendBitswapMessage {
706                target,
707                message,
708                result: tx,
709            })
710            .await
711            .unwrap();
712
713        rx.await.unwrap()
714    }
715
716    /// Broadcast Bitswap message to all [`service::ChainNetwork::established_bitswap_desired`]
717    /// peers.
718    ///
719    /// Returns the peers message was broadcast to or an error if the message cannot be sent
720    /// to at least one peer.
721    // TODO: better use a dedicated error type instead of reusing a lower-level
722    // `SendBitswapMessageErorr`.
723    pub async fn broadcast_bitswap_message(
724        &self,
725        message: Vec<u8>,
726    ) -> Result<Vec<PeerId>, SendBitswapMessageError> {
727        let (tx, rx) = oneshot::channel();
728
729        self.messages_tx
730            .send(ToBackgroundChain::BroadcastBitswapMessage {
731                message,
732                result: tx,
733            })
734            .await
735            .unwrap();
736
737        rx.await.unwrap()
738    }
739
740    /// Broadcast a statement notification to all gossip-connected peers.
741    pub async fn broadcast_statement(
742        self: Arc<Self>,
743        statement: Vec<u8>,
744    ) -> BroadcastStatementResult {
745        let (tx, rx) = oneshot::channel();
746
747        self.messages_tx
748            .send(ToBackgroundChain::BroadcastStatement {
749                statement,
750                result: tx,
751            })
752            .await
753            .unwrap();
754
755        rx.await.unwrap()
756    }
757
758    pub async fn update_topic_affinity(&self, filter: AffinityFilter) {
759        self.messages_tx
760            .send(ToBackgroundChain::UpdateTopicAffinity { filter })
761            .await
762            .unwrap();
763    }
764
765    /// Marks the given peers as belonging to the given chain, and adds some addresses to these
766    /// peers to the address book.
767    ///
768    /// The `important_nodes` parameter indicates whether these nodes are considered note-worthy
769    /// and should have additional logging.
770    pub async fn discover(
771        &self,
772        list: impl IntoIterator<Item = (PeerId, impl IntoIterator<Item = Multiaddr>)>,
773        important_nodes: bool,
774    ) {
775        self.messages_tx
776            .send(ToBackgroundChain::Discover {
777                // TODO: overhead
778                list: list
779                    .into_iter()
780                    .map(|(peer_id, addrs)| {
781                        (peer_id, addrs.into_iter().collect::<Vec<_>>().into_iter())
782                    })
783                    .collect::<Vec<_>>()
784                    .into_iter(),
785                important_nodes,
786            })
787            .await
788            .unwrap();
789    }
790
791    /// Returns a list of nodes (their [`PeerId`] and multiaddresses) that we know are part of
792    /// the network.
793    ///
794    /// Nodes that are discovered might disappear over time. In other words, there is no guarantee
795    /// that a node that has been added through [`NetworkServiceChain::discover`] will later be
796    /// returned by [`NetworkServiceChain::discovered_nodes`].
797    pub async fn discovered_nodes(
798        &self,
799    ) -> impl Iterator<Item = (PeerId, impl Iterator<Item = Multiaddr>)> {
800        let (tx, rx) = oneshot::channel();
801
802        self.messages_tx
803            .send(ToBackgroundChain::DiscoveredNodes { result: tx })
804            .await
805            .unwrap();
806
807        rx.await
808            .unwrap()
809            .into_iter()
810            .map(|(peer_id, addrs)| (peer_id, addrs.into_iter()))
811    }
812
813    /// Returns an iterator to the list of [`PeerId`]s that we have an established connection
814    /// with.
815    pub async fn peers_list(&self) -> impl Iterator<Item = PeerId> {
816        let (tx, rx) = oneshot::channel();
817        self.messages_tx
818            .send(ToBackgroundChain::PeersList { result: tx })
819            .await
820            .unwrap();
821        rx.await.unwrap().into_iter()
822    }
823}
824
825#[derive(Debug, Clone)]
826pub struct BroadcastStatementResult {
827    pub sent: usize,
828    pub total: usize,
829}
830
831/// Event that can happen on the network service.
832#[derive(Debug, Clone)]
833pub enum Event {
834    Connected {
835        peer_id: PeerId,
836        role: Role,
837        best_block_number: u64,
838        best_block_hash: [u8; 32],
839    },
840    Disconnected {
841        peer_id: PeerId,
842    },
843    BlockAnnounce {
844        peer_id: PeerId,
845        announce: service::EncodedBlockAnnounce,
846    },
847    GrandpaNeighborPacket {
848        peer_id: PeerId,
849        finalized_block_height: u64,
850    },
851    /// Received a GrandPa commit message from the network.
852    GrandpaCommitMessage {
853        peer_id: PeerId,
854        message: service::EncodedGrandpaCommitMessage,
855    },
856}
857
858/// Bitswap event that can be generated by the network service. Because Bitswap messages are big
859/// (up to 2 MiB) and can be delivered at high rate, we use a dedicated subscriber to not copy them
860/// to all network service subscribers.
861#[derive(Debug, Clone)]
862pub enum BitswapEvent {
863    BitswapMessage {
864        peer_id: PeerId,
865        message: service::EncodedBitswapMessage,
866    },
867}
868
869/// Statement event that can be generated by the network service.
870#[derive(Debug, Clone)]
871pub enum StatementEvent {
872    /// Received a statement notification from the network.
873    StatementsNotification {
874        peer_id: PeerId,
875        statements: Vec<([u8; 32], codec::Statement)>,
876    },
877}
878
879/// Error returned by [`NetworkServiceChain::blocks_request`].
880#[derive(Debug, derive_more::Display, derive_more::Error)]
881pub enum BlocksRequestError {
882    /// No established connection with the target.
883    NoConnection,
884    /// Error during the request.
885    #[display("{_0}")]
886    Request(service::BlocksRequestError),
887}
888
889/// Error returned by [`NetworkServiceChain::grandpa_warp_sync_request`].
890#[derive(Debug, derive_more::Display, derive_more::Error)]
891pub enum WarpSyncRequestError {
892    /// No established connection with the target.
893    NoConnection,
894    /// Error during the request.
895    #[display("{_0}")]
896    Request(service::GrandpaWarpSyncRequestError),
897}
898
899/// Error returned by [`NetworkServiceChain::storage_proof_request`].
900#[derive(Debug, derive_more::Display, derive_more::Error, Clone)]
901pub enum StorageProofRequestError {
902    /// No established connection with the target.
903    NoConnection,
904    /// Storage proof request is too large and can't be sent.
905    RequestTooLarge,
906    /// Error during the request.
907    #[display("{_0}")]
908    Request(service::StorageProofRequestError),
909}
910
911/// Error returned by [`NetworkServiceChain::call_proof_request`].
912#[derive(Debug, derive_more::Display, derive_more::Error, Clone)]
913pub enum CallProofRequestError {
914    /// No established connection with the target.
915    NoConnection,
916    /// Call proof request is too large and can't be sent.
917    RequestTooLarge,
918    /// Error during the request.
919    #[display("{_0}")]
920    Request(service::CallProofRequestError),
921}
922
923impl CallProofRequestError {
924    /// Returns `true` if this is caused by networking issues, as opposed to a consensus-related
925    /// issue.
926    pub fn is_network_problem(&self) -> bool {
927        match self {
928            CallProofRequestError::Request(err) => err.is_network_problem(),
929            CallProofRequestError::RequestTooLarge => false,
930            CallProofRequestError::NoConnection => true,
931        }
932    }
933}
934
935/// Error returned by [`NetworkServiceChain::child_storage_proof_request`].
936#[derive(Debug, derive_more::Display, derive_more::Error, Clone)]
937pub enum ChildStorageProofRequestError {
938    /// No established connection with the target.
939    NoConnection,
940    /// Child storage proof request is too large and can't be sent.
941    RequestTooLarge,
942    /// Error during the request.
943    #[display("{_0}")]
944    Request(service::StorageProofRequestError),
945}
946
947impl ChildStorageProofRequestError {
948    /// Returns `true` if this is caused by networking issues, as opposed to a consensus-related
949    /// issue.
950    pub fn is_network_problem(&self) -> bool {
951        match self {
952            ChildStorageProofRequestError::Request(err) => err.is_network_problem(),
953            ChildStorageProofRequestError::RequestTooLarge => false,
954            ChildStorageProofRequestError::NoConnection => true,
955        }
956    }
957}
958
959/// Why an address obtained through discovery was discarded. The `reason` label of
960/// the `networkDiscoveryAddressesDroppedTotal` metric.
961// Strum derives are explained on [`crate::metrics::MetricLabel`]; `Display`
962// additionally prints the same string in log messages.
963#[derive(
964    Debug, Clone, Copy, PartialEq, Eq, strum::Display, strum::EnumIter, strum::IntoStaticStr,
965)]
966#[strum(serialize_all = "kebab-case")]
967pub(crate) enum DiscoveredAddressDropReason {
968    /// Address contains a `/p2p/` component that doesn't match the peer it was announced for.
969    PeerIdMismatch,
970    /// Platform doesn't support connecting to this type of address.
971    NotSupported,
972    /// Address failed to parse.
973    Invalid,
974}
975
976/// Owned version of [`codec::ChildStorageProofRequestConfig`] for sending across channel.
977struct ChildStorageProofRequestConfigOwned {
978    block_hash: [u8; 32],
979    child_trie: Vec<u8>,
980    keys: Vec<Vec<u8>>,
981}
982
983enum ToBackground<TPlat: PlatformRef> {
984    AddChain {
985        messages_rx: async_channel::Receiver<ToBackgroundChain>,
986        config: service::ChainConfig<Chain<TPlat>>,
987    },
988}
989
990enum ToBackgroundChain {
991    RemoveChain,
992    Subscribe {
993        sender: async_channel::Sender<Event>,
994    },
995    SubscribeBitswap {
996        sender: async_channel::Sender<BitswapEvent>,
997    },
998    SubscribeStatements {
999        sender: async_channel::Sender<StatementEvent>,
1000    },
1001    DisconnectAndBan {
1002        peer_id: PeerId,
1003        severity: BanSeverity,
1004        reason: BanReason,
1005    },
1006    // TODO: serialize the request before sending over channel
1007    StartBlocksRequest {
1008        target: PeerId, // TODO: takes by value because of future longevity issue
1009        config: codec::BlocksRequestConfig,
1010        timeout: Duration,
1011        result: oneshot::Sender<Result<Vec<codec::BlockData>, BlocksRequestError>>,
1012    },
1013    // TODO: serialize the request before sending over channel
1014    StartWarpSyncRequest {
1015        target: PeerId,
1016        begin_hash: [u8; 32],
1017        timeout: Duration,
1018        result:
1019            oneshot::Sender<Result<service::EncodedGrandpaWarpSyncResponse, WarpSyncRequestError>>,
1020    },
1021    // TODO: serialize the request before sending over channel
1022    StartStorageProofRequest {
1023        target: PeerId,
1024        config: codec::StorageProofRequestConfig<vec::IntoIter<Vec<u8>>>,
1025        timeout: Duration,
1026        result: oneshot::Sender<Result<service::EncodedMerkleProof, StorageProofRequestError>>,
1027    },
1028    // TODO: serialize the request before sending over channel
1029    StartCallProofRequest {
1030        target: PeerId, // TODO: takes by value because of futures longevity issue
1031        config: codec::CallProofRequestConfig<'static, vec::IntoIter<Vec<u8>>>,
1032        timeout: Duration,
1033        result: oneshot::Sender<Result<service::EncodedMerkleProof, CallProofRequestError>>,
1034    },
1035    // TODO: serialize the request before sending over channel
1036    StartChildStorageProofRequest {
1037        target: PeerId,
1038        config: ChildStorageProofRequestConfigOwned,
1039        timeout: Duration,
1040        result: oneshot::Sender<Result<service::EncodedMerkleProof, ChildStorageProofRequestError>>,
1041    },
1042    SetLocalBestBlock {
1043        best_hash: [u8; 32],
1044        best_number: u64,
1045    },
1046    SetLocalGrandpaState {
1047        grandpa_state: service::GrandpaState,
1048    },
1049    AnnounceTransaction {
1050        transaction: Vec<u8>,
1051        result: oneshot::Sender<Vec<PeerId>>,
1052    },
1053    SendBlockAnnounce {
1054        target: PeerId,
1055        scale_encoded_header: Vec<u8>,
1056        is_best: bool,
1057        result: oneshot::Sender<Result<(), QueueNotificationError>>,
1058    },
1059    SendBitswapMessage {
1060        target: PeerId,
1061        message: Vec<u8>,
1062        result: oneshot::Sender<Result<(), SendBitswapMessageError>>,
1063    },
1064    BroadcastBitswapMessage {
1065        message: Vec<u8>,
1066        result: oneshot::Sender<Result<Vec<PeerId>, SendBitswapMessageError>>,
1067    },
1068    BroadcastStatement {
1069        statement: Vec<u8>,
1070        result: oneshot::Sender<BroadcastStatementResult>,
1071    },
1072    UpdateTopicAffinity {
1073        filter: AffinityFilter,
1074    },
1075    Discover {
1076        list: vec::IntoIter<(PeerId, vec::IntoIter<Multiaddr>)>,
1077        important_nodes: bool,
1078    },
1079    DiscoveredNodes {
1080        result: oneshot::Sender<Vec<(PeerId, Vec<Multiaddr>)>>,
1081    },
1082    PeersList {
1083        result: oneshot::Sender<Vec<PeerId>>,
1084    },
1085}
1086
1087struct BackgroundTask<TPlat: PlatformRef> {
1088    /// See [`Config::platform`].
1089    platform: TPlat,
1090
1091    /// Process-wide network metrics. Also accessible through [`NetworkService::metrics`].
1092    metrics: Arc<metrics::NetworkMetrics>,
1093
1094    /// Random number generator.
1095    randomness: rand_chacha::ChaCha20Rng,
1096
1097    /// Value provided through [`Config::identify_agent_version`].
1098    identify_agent_version: String,
1099
1100    /// Channel to send messages to the background task.
1101    tasks_messages_tx:
1102        async_channel::Sender<(service::ConnectionId, service::ConnectionToCoordinator)>,
1103
1104    /// Channel to receive messages destined to the background task.
1105    tasks_messages_rx: Pin<
1106        Box<async_channel::Receiver<(service::ConnectionId, service::ConnectionToCoordinator)>>,
1107    >,
1108
1109    /// Data structure holding the entire state of the networking.
1110    network: service::ChainNetwork<
1111        Chain<TPlat>,
1112        async_channel::Sender<service::CoordinatorToConnection>,
1113        TPlat::Instant,
1114    >,
1115
1116    /// All known peers and their addresses.
1117    peering_strategy: basic_peering_strategy::BasicPeeringStrategy<ChainId, TPlat::Instant>,
1118
1119    /// Bitswap slot assignment strategy.
1120    bitswap_peering_strategy: bitswap_peering_strategy::BitswapPeeringStrategy<TPlat::Instant>,
1121
1122    /// See [`Config::connections_open_pool_size`].
1123    connections_open_pool_size: u32,
1124
1125    /// See [`Config::connections_open_pool_restore_delay`].
1126    connections_open_pool_restore_delay: Duration,
1127
1128    /// Every time a connection is opened, the value in this field is increased by one. After
1129    /// [`BackgroundTask::next_recent_connection_restore`] has yielded, the value is reduced by
1130    /// one.
1131    num_recent_connection_opening: u32,
1132
1133    /// Delay after which [`BackgroundTask::num_recent_connection_opening`] is increased by one.
1134    next_recent_connection_restore: Option<Pin<Box<TPlat::Delay>>>,
1135
1136    /// List of all open gossip links.
1137    // TODO: using this data structure unfortunately means that PeerIds are cloned a lot, maybe some user data in ChainNetwork is better? not sure
1138    open_gossip_links: BTreeMap<(ChainId, PeerId), OpenGossipLinkState>,
1139
1140    /// Chains for which a gossip link has been opened at least once. Used to prefer bootnodes for
1141    /// out slots only until the chain first connects.
1142    chains_ever_gossip_connected: HashSet<ChainId, fnv::FnvBuildHasher>,
1143
1144    /// Connected peers using statement protocol V2, per chain.
1145    v2_statement_peers: HashMap<ChainId, HashSet<PeerId, fnv::FnvBuildHasher>, fnv::FnvBuildHasher>,
1146
1147    /// Current topic affinity filter per chain, sent to V2 peers on connect.
1148    current_affinity_filter: HashMap<ChainId, AffinityFilter, fnv::FnvBuildHasher>,
1149
1150    /// Important nodes per chain (in practice the bootnodes; see [`NetworkServiceChain::discover`]).
1151    /// They get extra logging, and slot preference until the chain first connects.
1152    // TODO: should also detect whenever we fail to open a block announces substream with any of these peers
1153    important_nodes: HashMap<ChainId, HashSet<PeerId, fnv::FnvBuildHasher>, fnv::FnvBuildHasher>,
1154
1155    /// Events about to be sent on the senders of [`BackgroundTask::event_senders`].
1156    ///
1157    /// Network events are only pulled when this queue is empty, keeping it small. A queue is
1158    /// nonetheless necessary, as processing a [`ToBackgroundChain::DisconnectAndBan`] message can
1159    /// generate an event while another event is already waiting to be dispatched.
1160    events_pending_send: VecDeque<(ChainId, Event)>,
1161
1162    /// Bitswap event about to be sent on the senders of [`BackgroundTask::bitswap_event_senders`].
1163    bitswap_event_pending_send: Option<BitswapEvent>,
1164
1165    /// Running count of peers with an open Bitswap substream. Maintained from
1166    /// `service::Event::BitswapConnected` / `BitswapDisconnected`. Used for diagnostic logging
1167    /// only; the authoritative per-peer state lives in
1168    /// [`BackgroundTask::bitswap_peering_strategy`].
1169    bitswap_connected_peers: usize,
1170
1171    /// Sending events through the public API.
1172    ///
1173    /// Contains either senders, or a `Future` that is currently sending an event and will yield
1174    /// the senders back once it is finished.
1175    // TODO: sort by ChainId instead of using a Vec?
1176    event_senders: either::Either<
1177        Vec<(ChainId, async_channel::Sender<Event>)>,
1178        Pin<Box<dyn Future<Output = Vec<(ChainId, async_channel::Sender<Event>)>> + Send>>,
1179    >,
1180
1181    /// Whenever [`NetworkServiceChain::subscribe`] is called, the new sender is added to this list.
1182    /// Once [`BackgroundTask::event_senders`] is ready, we properly initialize these senders.
1183    pending_new_subscriptions: Vec<(ChainId, async_channel::Sender<Event>)>,
1184
1185    /// Sending Bitswap events through the public API. We use separate channels for Bitswap events,
1186    /// because Bitswap messages are big and only few of event subscribers are interested in them.
1187    ///
1188    /// Contains either senders, or a `Future` that is currently sending an event and will yield
1189    /// the senders back once it is finished.
1190    ///
1191    /// Note that compared to `event_senders`, `bitswap_event_senders` are not associated with
1192    /// chains, because Bitswap messages coming from the network do not have the information about
1193    /// what chain they are coming from.
1194    bitswap_event_senders: either::Either<
1195        Vec<async_channel::Sender<BitswapEvent>>,
1196        Pin<Box<dyn Future<Output = Vec<async_channel::Sender<BitswapEvent>>> + Send>>,
1197    >,
1198
1199    /// Whenever [`NetworkServiceChain::subscribe_bitswap`] is called, the new sender is added to
1200    /// this list. Once [`BackgroundTask::bitswap_event_senders`] is ready, we properly initialize
1201    /// these senders.
1202    pending_new_bitswap_subscriptions: Vec<async_channel::Sender<BitswapEvent>>,
1203
1204    /// Statement event about to be sent on the senders of
1205    /// [`BackgroundTask::statement_event_senders`].
1206    statement_event_pending_send: Option<(ChainId, StatementEvent)>,
1207
1208    /// Sending statement events through the public API.
1209    ///
1210    /// Contains either senders, or a `Future` that is currently sending an event and will yield
1211    /// the senders back once it is finished.
1212    statement_event_senders: either::Either<
1213        Vec<(ChainId, async_channel::Sender<StatementEvent>)>,
1214        Pin<Box<dyn Future<Output = Vec<(ChainId, async_channel::Sender<StatementEvent>)>> + Send>>,
1215    >,
1216
1217    /// Whenever [`NetworkServiceChain::subscribe_statements`] is called, the new sender is added to
1218    /// this list. Once [`BackgroundTask::statement_event_senders`] is ready, we properly initialize
1219    /// these senders.
1220    pending_new_statement_subscriptions: Vec<(ChainId, async_channel::Sender<StatementEvent>)>,
1221
1222    main_messages_rx: Pin<Box<async_channel::Receiver<ToBackground<TPlat>>>>,
1223
1224    messages_rx:
1225        stream::SelectAll<Pin<Box<dyn stream::Stream<Item = (ChainId, ToBackgroundChain)> + Send>>>,
1226
1227    blocks_requests: HashMap<
1228        service::SubstreamId,
1229        oneshot::Sender<Result<Vec<codec::BlockData>, BlocksRequestError>>,
1230        fnv::FnvBuildHasher,
1231    >,
1232
1233    grandpa_warp_sync_requests: HashMap<
1234        service::SubstreamId,
1235        oneshot::Sender<Result<service::EncodedGrandpaWarpSyncResponse, WarpSyncRequestError>>,
1236        fnv::FnvBuildHasher,
1237    >,
1238
1239    storage_proof_requests: HashMap<
1240        service::SubstreamId,
1241        oneshot::Sender<Result<service::EncodedMerkleProof, StorageProofRequestError>>,
1242        fnv::FnvBuildHasher,
1243    >,
1244
1245    call_proof_requests: HashMap<
1246        service::SubstreamId,
1247        oneshot::Sender<Result<service::EncodedMerkleProof, CallProofRequestError>>,
1248        fnv::FnvBuildHasher,
1249    >,
1250
1251    child_storage_proof_requests: HashMap<
1252        service::SubstreamId,
1253        oneshot::Sender<Result<service::EncodedMerkleProof, ChildStorageProofRequestError>>,
1254        fnv::FnvBuildHasher,
1255    >,
1256
1257    /// All chains, indexed by the value of [`Chain::next_discovery_when`].
1258    chains_by_next_discovery: BTreeMap<(TPlat::Instant, ChainId), Pin<Box<TPlat::Delay>>>,
1259}
1260
1261struct Chain<TPlat: PlatformRef> {
1262    log_name: String,
1263
1264    // TODO: this field is a hack due to the fact that `add_chain` can't be `async`; should eventually be fixed after a lib.rs refactor
1265    num_references: NonZero<usize>,
1266
1267    /// See [`ConfigChain::block_number_bytes`].
1268    // TODO: redundant with ChainNetwork? since we might not need to know this in the future i'm reluctant to add a getter to ChainNetwork
1269    block_number_bytes: usize,
1270
1271    /// See [`ConfigChain::num_out_slots`].
1272    num_out_slots: usize,
1273
1274    /// See [`ConfigChain::metrics`].
1275    metrics: Arc<metrics::ChainMetrics>,
1276
1277    /// When the next discovery should be started for this chain.
1278    next_discovery_when: TPlat::Instant,
1279
1280    /// After [`Chain::next_discovery_when`] is reached, the following discovery happens after
1281    /// the given duration.
1282    next_discovery_period: Duration,
1283}
1284
1285#[derive(Clone)]
1286struct OpenGossipLinkState {
1287    role: Role,
1288    best_block_number: u64,
1289    best_block_hash: [u8; 32],
1290    /// `None` if unknown.
1291    finalized_block_height: Option<u64>,
1292}
1293
1294async fn background_task<TPlat: PlatformRef>(mut task: BackgroundTask<TPlat>) {
1295    loop {
1296        // Yield at every loop in order to provide better tasks granularity.
1297        futures_lite::future::yield_now().await;
1298
1299        enum WakeUpReason<TPlat: PlatformRef> {
1300            ForegroundClosed,
1301            Message(ToBackground<TPlat>),
1302            MessageForChain(ChainId, ToBackgroundChain),
1303            NetworkEvent(service::Event<async_channel::Sender<service::CoordinatorToConnection>>),
1304            CanAssignSlot(PeerId, ChainId),
1305            CanAssignBitswapSlot(PeerId),
1306            NextRecentConnectionRestore,
1307            CanStartConnect(PeerId),
1308            CanOpenGossip(PeerId, ChainId),
1309            CanOpenBitswap(PeerId),
1310            MessageFromConnection {
1311                connection_id: service::ConnectionId,
1312                message: service::ConnectionToCoordinator,
1313            },
1314            MessageToConnection {
1315                connection_id: service::ConnectionId,
1316                message: service::CoordinatorToConnection,
1317            },
1318            EventSendersReady,
1319            BitswapEventSendersReady,
1320            StatementEventSendersReady,
1321            StartDiscovery(ChainId),
1322        }
1323
1324        let wake_up_reason = {
1325            let message_received = async {
1326                task.main_messages_rx
1327                    .next()
1328                    .await
1329                    .map_or(WakeUpReason::ForegroundClosed, WakeUpReason::Message)
1330            };
1331            let message_for_chain_received = async {
1332                // Note that when the last entry of `messages_rx` yields `None`, `messages_rx`
1333                // itself will yield `None`. For this reason, we can't use
1334                // `task.messages_rx.is_empty()` to determine whether `messages_rx` will
1335                // yield `None`.
1336                let Some((chain_id, message)) = task.messages_rx.next().await else {
1337                    future::pending().await
1338                };
1339                WakeUpReason::MessageForChain(chain_id, message)
1340            };
1341            let message_from_task_received = async {
1342                let (connection_id, message) = task.tasks_messages_rx.next().await.unwrap();
1343                WakeUpReason::MessageFromConnection {
1344                    connection_id,
1345                    message,
1346                }
1347            };
1348            let service_event = async {
1349                if let Some(event) = (task.events_pending_send.is_empty()
1350                    && task.bitswap_event_pending_send.is_none()
1351                    && task.statement_event_pending_send.is_none()
1352                    && task.pending_new_subscriptions.is_empty()
1353                    && task.pending_new_bitswap_subscriptions.is_empty()
1354                    && task.pending_new_statement_subscriptions.is_empty())
1355                .then(|| task.network.next_event())
1356                .flatten()
1357                {
1358                    WakeUpReason::NetworkEvent(event)
1359                } else if let Some(start_connect) = {
1360                    let x = (task.num_recent_connection_opening < task.connections_open_pool_size)
1361                        .then(|| {
1362                            task.network
1363                                .unconnected_desired()
1364                                .choose(&mut task.randomness)
1365                                .cloned()
1366                        })
1367                        .flatten();
1368                    x
1369                } {
1370                    WakeUpReason::CanStartConnect(start_connect)
1371                } else if let Some((peer_id, chain_id)) = {
1372                    let x = task
1373                        .network
1374                        .connected_unopened_gossip_desired()
1375                        .choose(&mut task.randomness)
1376                        .map(|(peer_id, chain_id, _)| (peer_id.clone(), chain_id));
1377                    x
1378                } {
1379                    WakeUpReason::CanOpenGossip(peer_id, chain_id)
1380                } else if let Some(peer_id) = {
1381                    let x = task
1382                        .network
1383                        .connected_unopened_bitswap_desired()
1384                        .choose(&mut task.randomness)
1385                        .cloned();
1386                    x
1387                } {
1388                    WakeUpReason::CanOpenBitswap(peer_id)
1389                } else if let Some((connection_id, message)) =
1390                    task.network.pull_message_to_connection()
1391                {
1392                    WakeUpReason::MessageToConnection {
1393                        connection_id,
1394                        message,
1395                    }
1396                } else {
1397                    'search: loop {
1398                        let mut earlier_unban = None;
1399
1400                        for chain_id in task.network.chains().collect::<Vec<_>>() {
1401                            if task.network.gossip_desired_num(
1402                                chain_id,
1403                                service::GossipKind::ConsensusTransactions,
1404                            ) >= task.network[chain_id].num_out_slots
1405                            {
1406                                continue;
1407                            }
1408
1409                            let now = task.platform.now();
1410
1411                            // Until the chain first connects, prefer slots for important nodes
1412                            // (the bootnodes); otherwise use the general pool.
1413                            if !task.chains_ever_gossip_connected.contains(&chain_id) {
1414                                if let basic_peering_strategy::AssignablePeer::Assignable(peer_id) =
1415                                    task.peering_strategy.pick_assignable_peer_filtered(
1416                                        &chain_id,
1417                                        &now,
1418                                        |peer_id| {
1419                                            task.important_nodes
1420                                                .get(&chain_id)
1421                                                .map_or(false, |nodes| nodes.contains(peer_id))
1422                                        },
1423                                    )
1424                                {
1425                                    break 'search WakeUpReason::CanAssignSlot(
1426                                        peer_id.clone(),
1427                                        chain_id,
1428                                    );
1429                                }
1430                            }
1431
1432                            match task.peering_strategy.pick_assignable_peer(&chain_id, &now) {
1433                                basic_peering_strategy::AssignablePeer::Assignable(peer_id) => {
1434                                    break 'search WakeUpReason::CanAssignSlot(
1435                                        peer_id.clone(),
1436                                        chain_id,
1437                                    );
1438                                }
1439                                basic_peering_strategy::AssignablePeer::AllPeersBanned {
1440                                    next_unban,
1441                                } => {
1442                                    if earlier_unban.as_ref().map_or(true, |b| b > next_unban) {
1443                                        earlier_unban = Some(next_unban.clone());
1444                                    }
1445                                }
1446                                basic_peering_strategy::AssignablePeer::NoPeer => continue,
1447                            }
1448                        }
1449
1450                        match task
1451                            .bitswap_peering_strategy
1452                            .pick_assignable_peer(&task.platform.now())
1453                        {
1454                            bitswap_peering_strategy::AssignablePeer::Assignable(peer_id) => {
1455                                break 'search WakeUpReason::CanAssignBitswapSlot(peer_id.clone());
1456                            }
1457                            bitswap_peering_strategy::AssignablePeer::AllPeersBanned {
1458                                next_unban,
1459                            } => {
1460                                if earlier_unban.as_ref().map_or(true, |b| b > next_unban) {
1461                                    earlier_unban = Some(next_unban.clone());
1462                                }
1463                            }
1464                            bitswap_peering_strategy::AssignablePeer::NoPeer => {}
1465                        }
1466
1467                        if let Some(earlier_unban) = earlier_unban {
1468                            task.platform.sleep_until(earlier_unban).await;
1469                        } else {
1470                            future::pending::<()>().await;
1471                        }
1472                    }
1473                }
1474            };
1475            let next_recent_connection_restore = async {
1476                if task.num_recent_connection_opening != 0
1477                    && task.next_recent_connection_restore.is_none()
1478                {
1479                    task.next_recent_connection_restore = Some(Box::pin(
1480                        task.platform
1481                            .sleep(task.connections_open_pool_restore_delay),
1482                    ));
1483                }
1484                if let Some(delay) = task.next_recent_connection_restore.as_mut() {
1485                    delay.await;
1486                    task.next_recent_connection_restore = None;
1487                    WakeUpReason::NextRecentConnectionRestore
1488                } else {
1489                    future::pending().await
1490                }
1491            };
1492            let finished_sending_event = async {
1493                if let either::Right(event_sending_future) = &mut task.event_senders {
1494                    let event_senders = event_sending_future.await;
1495                    task.event_senders = either::Left(event_senders);
1496                    WakeUpReason::EventSendersReady
1497                } else if !task.events_pending_send.is_empty()
1498                    || !task.pending_new_subscriptions.is_empty()
1499                {
1500                    WakeUpReason::EventSendersReady
1501                } else {
1502                    future::pending().await
1503                }
1504            };
1505            let finished_sending_bitswap_event = async {
1506                if let either::Right(bitswap_event_sending_future) = &mut task.bitswap_event_senders
1507                {
1508                    let bitswap_event_senders = bitswap_event_sending_future.await;
1509                    task.bitswap_event_senders = either::Left(bitswap_event_senders);
1510                    WakeUpReason::BitswapEventSendersReady
1511                } else if task.bitswap_event_pending_send.is_some()
1512                    || !task.pending_new_bitswap_subscriptions.is_empty()
1513                {
1514                    WakeUpReason::BitswapEventSendersReady
1515                } else {
1516                    future::pending().await
1517                }
1518            };
1519            let finished_sending_statement_event = async {
1520                if let either::Right(statement_event_sending_future) =
1521                    &mut task.statement_event_senders
1522                {
1523                    let statement_event_senders = statement_event_sending_future.await;
1524                    task.statement_event_senders = either::Left(statement_event_senders);
1525                    WakeUpReason::StatementEventSendersReady
1526                } else if task.statement_event_pending_send.is_some()
1527                    || !task.pending_new_statement_subscriptions.is_empty()
1528                {
1529                    WakeUpReason::StatementEventSendersReady
1530                } else {
1531                    future::pending().await
1532                }
1533            };
1534            let start_discovery = async {
1535                let Some(mut next_discovery) = task.chains_by_next_discovery.first_entry() else {
1536                    future::pending().await
1537                };
1538                next_discovery.get_mut().await;
1539                let ((_, chain_id), _) = next_discovery.remove_entry();
1540                WakeUpReason::StartDiscovery(chain_id)
1541            };
1542
1543            message_for_chain_received
1544                .or(message_received)
1545                .or(message_from_task_received)
1546                .or(service_event)
1547                .or(next_recent_connection_restore)
1548                .or(finished_sending_event)
1549                .or(finished_sending_bitswap_event)
1550                .or(finished_sending_statement_event)
1551                .or(start_discovery)
1552                .await
1553        };
1554
1555        match wake_up_reason {
1556            WakeUpReason::ForegroundClosed => {
1557                // End the task.
1558                return;
1559            }
1560            WakeUpReason::Message(ToBackground::AddChain {
1561                messages_rx,
1562                config,
1563            }) => {
1564                // TODO: this is not a completely clean way of handling duplicate chains, because the existing chain might have a different best block and role and all ; also, multiple sync services will call set_best_block and set_finalized_block
1565                let chain_id = match task.network.add_chain(config) {
1566                    Ok(id) => id,
1567                    Err(service::AddChainError::Duplicate { existing_identical }) => {
1568                        task.network[existing_identical].num_references = task.network
1569                            [existing_identical]
1570                            .num_references
1571                            .checked_add(1)
1572                            .unwrap();
1573                        existing_identical
1574                    }
1575                };
1576
1577                task.chains_by_next_discovery.insert(
1578                    (task.network[chain_id].next_discovery_when.clone(), chain_id),
1579                    Box::pin(
1580                        task.platform
1581                            .sleep_until(task.network[chain_id].next_discovery_when.clone()),
1582                    ),
1583                );
1584
1585                task.messages_rx
1586                    .push(Box::pin(
1587                        messages_rx
1588                            .map(move |msg| (chain_id, msg))
1589                            .chain(stream::once(future::ready((
1590                                chain_id,
1591                                ToBackgroundChain::RemoveChain,
1592                            )))),
1593                    ) as Pin<Box<_>>);
1594
1595                log!(
1596                    &task.platform,
1597                    Debug,
1598                    "network",
1599                    "chain-added",
1600                    id = task.network[chain_id].log_name
1601                );
1602            }
1603            WakeUpReason::EventSendersReady => {
1604                // Dispatch the pending event, if any, to the various senders.
1605
1606                // We made sure that the senders were ready before generating an event.
1607                let either::Left(event_senders) = &mut task.event_senders else {
1608                    unreachable!()
1609                };
1610
1611                if let Some((event_to_dispatch_chain_id, event_to_dispatch)) =
1612                    task.events_pending_send.pop_front()
1613                {
1614                    let mut event_senders = mem::take(event_senders);
1615                    task.event_senders = either::Right(Box::pin(async move {
1616                        // Elements in `event_senders` are removed one by one and inserted
1617                        // back if the channel is still open.
1618                        for index in (0..event_senders.len()).rev() {
1619                            let (event_sender_chain_id, event_sender) =
1620                                event_senders.swap_remove(index);
1621                            if event_sender_chain_id == event_to_dispatch_chain_id {
1622                                if event_sender.send(event_to_dispatch.clone()).await.is_err() {
1623                                    continue;
1624                                }
1625                            }
1626                            event_senders.push((event_sender_chain_id, event_sender));
1627                        }
1628                        event_senders
1629                    }));
1630                } else if !task.pending_new_subscriptions.is_empty() {
1631                    let pending_new_subscriptions = mem::take(&mut task.pending_new_subscriptions);
1632                    let mut event_senders = mem::take(event_senders);
1633                    // TODO: cloning :-/
1634                    let open_gossip_links = task.open_gossip_links.clone();
1635                    task.event_senders = either::Right(Box::pin(async move {
1636                        for (chain_id, new_subscription) in pending_new_subscriptions {
1637                            for ((link_chain_id, peer_id), state) in &open_gossip_links {
1638                                // TODO: optimize? this is O(n) by chain
1639                                if *link_chain_id != chain_id {
1640                                    continue;
1641                                }
1642
1643                                let _ = new_subscription
1644                                    .send(Event::Connected {
1645                                        peer_id: peer_id.clone(),
1646                                        role: state.role,
1647                                        best_block_number: state.best_block_number,
1648                                        best_block_hash: state.best_block_hash,
1649                                    })
1650                                    .await;
1651
1652                                if let Some(finalized_block_height) = state.finalized_block_height {
1653                                    let _ = new_subscription
1654                                        .send(Event::GrandpaNeighborPacket {
1655                                            peer_id: peer_id.clone(),
1656                                            finalized_block_height,
1657                                        })
1658                                        .await;
1659                                }
1660                            }
1661
1662                            event_senders.push((chain_id, new_subscription));
1663                        }
1664
1665                        event_senders
1666                    }));
1667                }
1668            }
1669            WakeUpReason::BitswapEventSendersReady => {
1670                // We made sure that the senders were ready before generating an event.
1671                let either::Left(bitswap_event_senders) = &mut task.bitswap_event_senders else {
1672                    unreachable!()
1673                };
1674
1675                if let Some(event_to_dispatch) = task.bitswap_event_pending_send.take() {
1676                    let mut bitswap_event_senders = mem::take(bitswap_event_senders);
1677                    task.bitswap_event_senders = either::Right(Box::pin(async move {
1678                        // Elements in `bitswap_event_senders` are removed one by one and
1679                        // inserted back if the channel is still open.
1680                        for index in (0..bitswap_event_senders.len()).rev() {
1681                            let event_sender = bitswap_event_senders.swap_remove(index);
1682                            if event_sender.send(event_to_dispatch.clone()).await.is_err() {
1683                                continue;
1684                            }
1685                            bitswap_event_senders.push(event_sender);
1686                        }
1687                        bitswap_event_senders
1688                    }));
1689                } else if !task.pending_new_bitswap_subscriptions.is_empty() {
1690                    bitswap_event_senders.append(&mut task.pending_new_bitswap_subscriptions);
1691                }
1692            }
1693            WakeUpReason::StatementEventSendersReady => {
1694                // We made sure that the senders were ready before generating an event.
1695                let either::Left(statement_event_senders) = &mut task.statement_event_senders
1696                else {
1697                    unreachable!()
1698                };
1699
1700                if let Some((event_to_dispatch_chain_id, event_to_dispatch)) =
1701                    task.statement_event_pending_send.take()
1702                {
1703                    let mut statement_event_senders = mem::take(statement_event_senders);
1704                    task.statement_event_senders = either::Right(Box::pin(async move {
1705                        // Elements in `statement_event_senders` are removed one by one and
1706                        // inserted back if the channel is still open.
1707                        for index in (0..statement_event_senders.len()).rev() {
1708                            let (event_sender_chain_id, event_sender) =
1709                                statement_event_senders.swap_remove(index);
1710                            if event_sender_chain_id == event_to_dispatch_chain_id {
1711                                if event_sender.send(event_to_dispatch.clone()).await.is_err() {
1712                                    continue;
1713                                }
1714                            }
1715                            statement_event_senders.push((event_sender_chain_id, event_sender));
1716                        }
1717                        statement_event_senders
1718                    }));
1719                } else if !task.pending_new_statement_subscriptions.is_empty() {
1720                    statement_event_senders.append(&mut task.pending_new_statement_subscriptions);
1721                }
1722            }
1723            WakeUpReason::MessageFromConnection {
1724                connection_id,
1725                message,
1726            } => {
1727                task.network
1728                    .inject_connection_message(connection_id, message);
1729            }
1730            WakeUpReason::MessageForChain(chain_id, ToBackgroundChain::RemoveChain) => {
1731                if let Some(new_ref) =
1732                    NonZero::<usize>::new(task.network[chain_id].num_references.get() - 1)
1733                {
1734                    task.network[chain_id].num_references = new_ref;
1735                    continue;
1736                }
1737
1738                for peer_id in task
1739                    .network
1740                    .gossip_connected_peers(chain_id, service::GossipKind::ConsensusTransactions)
1741                    .cloned()
1742                    .collect::<Vec<_>>()
1743                {
1744                    task.network
1745                        .gossip_close(
1746                            chain_id,
1747                            &peer_id,
1748                            service::GossipKind::ConsensusTransactions,
1749                        )
1750                        .unwrap();
1751
1752                    let _was_in = task.open_gossip_links.remove(&(chain_id, peer_id));
1753                    debug_assert!(_was_in.is_some());
1754                }
1755
1756                let _was_in = task
1757                    .chains_by_next_discovery
1758                    .remove(&(task.network[chain_id].next_discovery_when.clone(), chain_id));
1759                debug_assert!(_was_in.is_some());
1760
1761                log!(
1762                    &task.platform,
1763                    Debug,
1764                    "network",
1765                    "chain-removed",
1766                    id = task.network[chain_id].log_name
1767                );
1768                task.v2_statement_peers.remove(&chain_id);
1769                task.current_affinity_filter.remove(&chain_id);
1770                task.important_nodes.remove(&chain_id);
1771                task.chains_ever_gossip_connected.remove(&chain_id);
1772                task.network.remove_chain(chain_id).unwrap();
1773                task.peering_strategy.remove_chain_peers(&chain_id);
1774            }
1775            WakeUpReason::MessageForChain(chain_id, ToBackgroundChain::Subscribe { sender }) => {
1776                task.pending_new_subscriptions.push((chain_id, sender));
1777            }
1778            WakeUpReason::MessageForChain(
1779                _chain_id,
1780                ToBackgroundChain::SubscribeBitswap { sender },
1781            ) => {
1782                task.pending_new_bitswap_subscriptions.push(sender);
1783            }
1784            WakeUpReason::MessageForChain(
1785                chain_id,
1786                ToBackgroundChain::SubscribeStatements { sender },
1787            ) => {
1788                task.pending_new_statement_subscriptions
1789                    .push((chain_id, sender));
1790            }
1791            WakeUpReason::MessageForChain(
1792                chain_id,
1793                ToBackgroundChain::DisconnectAndBan {
1794                    peer_id,
1795                    severity,
1796                    reason,
1797                },
1798            ) => {
1799                let ban_duration = Duration::from_secs(match severity {
1800                    BanSeverity::Low => 10,
1801                    BanSeverity::High => 40,
1802                });
1803
1804                let had_slot = matches!(
1805                    task.peering_strategy.unassign_slot_and_ban(
1806                        &chain_id,
1807                        &peer_id,
1808                        task.platform.now() + ban_duration,
1809                    ),
1810                    basic_peering_strategy::UnassignSlotAndBan::Banned { had_slot: true }
1811                );
1812
1813                if had_slot {
1814                    log!(
1815                        &task.platform,
1816                        Debug,
1817                        "network",
1818                        "slot-unassigned",
1819                        chain = &task.network[chain_id].log_name,
1820                        peer_id,
1821                        ?ban_duration,
1822                        reason = "user-ban",
1823                        user_reason = reason
1824                    );
1825                    task.network.gossip_remove_desired(
1826                        chain_id,
1827                        &peer_id,
1828                        service::GossipKind::ConsensusTransactions,
1829                    );
1830                }
1831
1832                if task.network.gossip_is_connected(
1833                    chain_id,
1834                    &peer_id,
1835                    service::GossipKind::ConsensusTransactions,
1836                ) {
1837                    let _closed_result = task.network.gossip_close(
1838                        chain_id,
1839                        &peer_id,
1840                        service::GossipKind::ConsensusTransactions,
1841                    );
1842                    debug_assert!(_closed_result.is_ok());
1843
1844                    log!(
1845                        &task.platform,
1846                        Debug,
1847                        "network",
1848                        "gossip-closed",
1849                        chain = &task.network[chain_id].log_name,
1850                        peer_id,
1851                    );
1852
1853                    let _was_in = task.open_gossip_links.remove(&(chain_id, peer_id.clone()));
1854                    debug_assert!(_was_in.is_some());
1855                    task.network[chain_id].metrics.gossip_peers_connected.dec();
1856
1857                    if let Some(peers) = task.v2_statement_peers.get_mut(&chain_id) {
1858                        peers.remove(&peer_id);
1859                    }
1860
1861                    // Unlike the network-event handlers below, this message handler can run
1862                    // while another event is already queued, hence the push to a queue.
1863                    task.events_pending_send
1864                        .push_back((chain_id, Event::Disconnected { peer_id }));
1865                }
1866            }
1867            WakeUpReason::MessageForChain(
1868                chain_id,
1869                ToBackgroundChain::StartBlocksRequest {
1870                    target,
1871                    config,
1872                    timeout,
1873                    result,
1874                },
1875            ) => {
1876                match &config.start {
1877                    codec::BlocksRequestConfigStart::Hash(hash) => {
1878                        log!(
1879                            &task.platform,
1880                            Debug,
1881                            "network",
1882                            "blocks-request-started",
1883                            chain = task.network[chain_id].log_name, target,
1884                            start = HashDisplay(hash),
1885                            num = config.desired_count.get(),
1886                            descending = ?matches!(config.direction, codec::BlocksRequestDirection::Descending),
1887                            header = ?config.fields.header, body = ?config.fields.body,
1888                            justifications = ?config.fields.justifications
1889                        );
1890                    }
1891                    codec::BlocksRequestConfigStart::Number(number) => {
1892                        log!(
1893                            &task.platform,
1894                            Debug,
1895                            "network",
1896                            "blocks-request-started",
1897                            chain = task.network[chain_id].log_name, target, start = number,
1898                            num = config.desired_count.get(),
1899                            descending = ?matches!(config.direction, codec::BlocksRequestDirection::Descending),
1900                            header = ?config.fields.header, body = ?config.fields.body, justifications = ?config.fields.justifications
1901                        );
1902                    }
1903                }
1904
1905                match task
1906                    .network
1907                    .start_blocks_request(&target, chain_id, config.clone(), timeout)
1908                {
1909                    Ok(substream_id) => {
1910                        task.blocks_requests.insert(substream_id, result);
1911                    }
1912                    Err(service::StartRequestError::NoConnection) => {
1913                        log!(
1914                            &task.platform,
1915                            Debug,
1916                            "network",
1917                            "blocks-request-error",
1918                            chain = task.network[chain_id].log_name,
1919                            target,
1920                            error = "NoConnection"
1921                        );
1922                        let _ = result.send(Err(BlocksRequestError::NoConnection));
1923                    }
1924                }
1925            }
1926            WakeUpReason::MessageForChain(
1927                chain_id,
1928                ToBackgroundChain::StartWarpSyncRequest {
1929                    target,
1930                    begin_hash,
1931                    timeout,
1932                    result,
1933                },
1934            ) => {
1935                log!(
1936                    &task.platform,
1937                    Debug,
1938                    "network",
1939                    "warp-sync-request-started",
1940                    chain = task.network[chain_id].log_name,
1941                    target,
1942                    start = HashDisplay(&begin_hash)
1943                );
1944
1945                match task
1946                    .network
1947                    .start_grandpa_warp_sync_request(&target, chain_id, begin_hash, timeout)
1948                {
1949                    Ok(substream_id) => {
1950                        task.grandpa_warp_sync_requests.insert(substream_id, result);
1951                    }
1952                    Err(service::StartRequestError::NoConnection) => {
1953                        log!(
1954                            &task.platform,
1955                            Debug,
1956                            "network",
1957                            "warp-sync-request-error",
1958                            chain = task.network[chain_id].log_name,
1959                            target,
1960                            error = "NoConnection"
1961                        );
1962                        let _ = result.send(Err(WarpSyncRequestError::NoConnection));
1963                    }
1964                }
1965            }
1966            WakeUpReason::MessageForChain(
1967                chain_id,
1968                ToBackgroundChain::StartStorageProofRequest {
1969                    target,
1970                    config,
1971                    timeout,
1972                    result,
1973                },
1974            ) => {
1975                log!(
1976                    &task.platform,
1977                    Debug,
1978                    "network",
1979                    "storage-proof-request-started",
1980                    chain = task.network[chain_id].log_name,
1981                    target,
1982                    block_hash = HashDisplay(&config.block_hash)
1983                );
1984
1985                match task.network.start_storage_proof_request(
1986                    &target,
1987                    chain_id,
1988                    config.clone(),
1989                    timeout,
1990                ) {
1991                    Ok(substream_id) => {
1992                        task.storage_proof_requests.insert(substream_id, result);
1993                    }
1994                    Err(service::StartRequestMaybeTooLargeError::NoConnection) => {
1995                        log!(
1996                            &task.platform,
1997                            Debug,
1998                            "network",
1999                            "storage-proof-request-error",
2000                            chain = task.network[chain_id].log_name,
2001                            target,
2002                            error = "NoConnection"
2003                        );
2004                        let _ = result.send(Err(StorageProofRequestError::NoConnection));
2005                    }
2006                    Err(service::StartRequestMaybeTooLargeError::RequestTooLarge) => {
2007                        log!(
2008                            &task.platform,
2009                            Debug,
2010                            "network",
2011                            "storage-proof-request-error",
2012                            chain = task.network[chain_id].log_name,
2013                            target,
2014                            error = "RequestTooLarge"
2015                        );
2016                        let _ = result.send(Err(StorageProofRequestError::RequestTooLarge));
2017                    }
2018                };
2019            }
2020            WakeUpReason::MessageForChain(
2021                chain_id,
2022                ToBackgroundChain::StartCallProofRequest {
2023                    target,
2024                    config,
2025                    timeout,
2026                    result,
2027                },
2028            ) => {
2029                log!(
2030                    &task.platform,
2031                    Debug,
2032                    "network",
2033                    "call-proof-request-started",
2034                    chain = task.network[chain_id].log_name,
2035                    target,
2036                    block_hash = HashDisplay(&config.block_hash),
2037                    function = config.method
2038                );
2039                // TODO: log parameter
2040
2041                match task.network.start_call_proof_request(
2042                    &target,
2043                    chain_id,
2044                    config.clone(),
2045                    timeout,
2046                ) {
2047                    Ok(substream_id) => {
2048                        task.call_proof_requests.insert(substream_id, result);
2049                    }
2050                    Err(service::StartRequestMaybeTooLargeError::NoConnection) => {
2051                        log!(
2052                            &task.platform,
2053                            Debug,
2054                            "network",
2055                            "call-proof-request-error",
2056                            chain = task.network[chain_id].log_name,
2057                            target,
2058                            error = "NoConnection"
2059                        );
2060                        let _ = result.send(Err(CallProofRequestError::NoConnection));
2061                    }
2062                    Err(service::StartRequestMaybeTooLargeError::RequestTooLarge) => {
2063                        log!(
2064                            &task.platform,
2065                            Debug,
2066                            "network",
2067                            "call-proof-request-error",
2068                            chain = task.network[chain_id].log_name,
2069                            target,
2070                            error = "RequestTooLarge"
2071                        );
2072                        let _ = result.send(Err(CallProofRequestError::RequestTooLarge));
2073                    }
2074                };
2075            }
2076            WakeUpReason::MessageForChain(
2077                chain_id,
2078                ToBackgroundChain::StartChildStorageProofRequest {
2079                    target,
2080                    config,
2081                    timeout,
2082                    result,
2083                },
2084            ) => {
2085                log!(
2086                    &task.platform,
2087                    Debug,
2088                    "network",
2089                    "child-storage-proof-request-started",
2090                    chain = task.network[chain_id].log_name,
2091                    target,
2092                    block_hash = HashDisplay(&config.block_hash)
2093                );
2094
2095                match task.network.start_child_storage_proof_request(
2096                    &target,
2097                    chain_id,
2098                    codec::ChildStorageProofRequestConfig {
2099                        block_hash: config.block_hash,
2100                        child_trie: &config.child_trie,
2101                        keys: config.keys.iter().map(|k| k.as_slice()),
2102                    },
2103                    timeout,
2104                ) {
2105                    Ok(substream_id) => {
2106                        task.child_storage_proof_requests
2107                            .insert(substream_id, result);
2108                    }
2109                    Err(service::StartRequestMaybeTooLargeError::NoConnection) => {
2110                        log!(
2111                            &task.platform,
2112                            Debug,
2113                            "network",
2114                            "child-storage-proof-request-error",
2115                            chain = task.network[chain_id].log_name,
2116                            target,
2117                            error = "NoConnection"
2118                        );
2119                        let _ = result.send(Err(ChildStorageProofRequestError::NoConnection));
2120                    }
2121                    Err(service::StartRequestMaybeTooLargeError::RequestTooLarge) => {
2122                        log!(
2123                            &task.platform,
2124                            Debug,
2125                            "network",
2126                            "child-storage-proof-request-error",
2127                            chain = task.network[chain_id].log_name,
2128                            target,
2129                            error = "RequestTooLarge"
2130                        );
2131                        let _ = result.send(Err(ChildStorageProofRequestError::RequestTooLarge));
2132                    }
2133                };
2134            }
2135            WakeUpReason::MessageForChain(
2136                chain_id,
2137                ToBackgroundChain::SetLocalBestBlock {
2138                    best_hash,
2139                    best_number,
2140                },
2141            ) => {
2142                task.network
2143                    .set_chain_local_best_block(chain_id, best_hash, best_number);
2144            }
2145            WakeUpReason::MessageForChain(
2146                chain_id,
2147                ToBackgroundChain::SetLocalGrandpaState { grandpa_state },
2148            ) => {
2149                log!(
2150                    &task.platform,
2151                    Debug,
2152                    "network",
2153                    "local-grandpa-state-announced",
2154                    chain = task.network[chain_id].log_name,
2155                    set_id = grandpa_state.set_id,
2156                    commit_finalized_height = grandpa_state.commit_finalized_height,
2157                );
2158
2159                // TODO: log the list of peers we sent the packet to
2160
2161                task.network
2162                    .gossip_broadcast_grandpa_state_and_update(chain_id, grandpa_state);
2163            }
2164            WakeUpReason::MessageForChain(
2165                chain_id,
2166                ToBackgroundChain::AnnounceTransaction {
2167                    transaction,
2168                    result,
2169                },
2170            ) => {
2171                // TODO: keep track of which peer knows about which transaction, and don't send it again
2172
2173                let peers_to_send = task
2174                    .network
2175                    .gossip_connected_peers(chain_id, service::GossipKind::ConsensusTransactions)
2176                    .cloned()
2177                    .collect::<Vec<_>>();
2178
2179                let mut peers_sent = Vec::with_capacity(peers_to_send.len());
2180                let mut peers_queue_full = Vec::with_capacity(peers_to_send.len());
2181                for peer in &peers_to_send {
2182                    match task
2183                        .network
2184                        .gossip_send_transaction(peer, chain_id, &transaction)
2185                    {
2186                        Ok(()) => peers_sent.push(peer.to_base58()),
2187                        Err(QueueNotificationError::QueueFull) => {
2188                            peers_queue_full.push(peer.to_base58())
2189                        }
2190                        Err(QueueNotificationError::NoConnection) => unreachable!(),
2191                    }
2192                }
2193
2194                log!(
2195                    &task.platform,
2196                    Debug,
2197                    "network",
2198                    "transaction-announced",
2199                    chain = task.network[chain_id].log_name,
2200                    transaction =
2201                        hex::encode(blake2_rfc::blake2b::blake2b(32, &[], &transaction).as_bytes()),
2202                    size = transaction.len(),
2203                    peers_sent = peers_sent.join(", "),
2204                    peers_queue_full = peers_queue_full.join(", "),
2205                );
2206
2207                let _ = result.send(peers_to_send);
2208            }
2209            WakeUpReason::MessageForChain(
2210                chain_id,
2211                ToBackgroundChain::SendBlockAnnounce {
2212                    target,
2213                    scale_encoded_header,
2214                    is_best,
2215                    result,
2216                },
2217            ) => {
2218                // TODO: log who the announce was sent to
2219                let _ = result.send(task.network.gossip_send_block_announce(
2220                    &target,
2221                    chain_id,
2222                    &scale_encoded_header,
2223                    is_best,
2224                ));
2225            }
2226            WakeUpReason::MessageForChain(
2227                _chain_id,
2228                ToBackgroundChain::SendBitswapMessage {
2229                    target,
2230                    message,
2231                    result,
2232                },
2233            ) => {
2234                let _ = result.send(task.network.bitswap_send_message(&target, message));
2235            }
2236            WakeUpReason::MessageForChain(
2237                _chain_id,
2238                ToBackgroundChain::BroadcastBitswapMessage { message, result },
2239            ) => {
2240                let peers = task
2241                    .network
2242                    .established_bitswap_desired()
2243                    .cloned()
2244                    .collect::<Vec<_>>();
2245                let results = peers
2246                    .iter()
2247                    .map(|peer| {
2248                        (
2249                            peer,
2250                            task.network.bitswap_send_message(peer, message.clone()),
2251                        )
2252                    })
2253                    .collect::<Vec<_>>(); // we must collect first to send all messages
2254
2255                let succeeded_peers = results
2256                    .iter()
2257                    .filter_map(|(peer, r)| r.is_ok().then(|| (*peer).clone()))
2258                    .collect::<Vec<_>>();
2259
2260                // TODO: introspecting a third-party error type below doesn't seem good.
2261                let r = if !succeeded_peers.is_empty() {
2262                    Ok(succeeded_peers)
2263                } else if results
2264                    .iter()
2265                    .any(|(_peer, r)| matches!(r, Err(SendBitswapMessageError::QueueFull)))
2266                {
2267                    // `QueueFull` has higher priority than `NoConnection` for possible
2268                    // back-pressure in higher level code.
2269                    Err(SendBitswapMessageError::QueueFull)
2270                } else {
2271                    // This is only emitted if all peers fail with `NoConnection` or there is no
2272                    // peers at all.
2273                    Err(SendBitswapMessageError::NoConnection)
2274                };
2275
2276                let _ = result.send(r);
2277            }
2278            WakeUpReason::MessageForChain(
2279                chain_id,
2280                ToBackgroundChain::BroadcastStatement { statement, result },
2281            ) => {
2282                let peers_to_send = task
2283                    .network
2284                    .gossip_connected_peers(chain_id, service::GossipKind::ConsensusTransactions)
2285                    .cloned()
2286                    .collect::<Vec<_>>();
2287
2288                let total = peers_to_send.len();
2289                let mut sent = 0;
2290                for peer in &peers_to_send {
2291                    if task
2292                        .network
2293                        .gossip_send_statement(peer, chain_id, statement.clone())
2294                        .is_ok()
2295                    {
2296                        sent += 1;
2297                    }
2298                }
2299
2300                log!(
2301                    &task.platform,
2302                    Debug,
2303                    "network",
2304                    "statement-broadcast",
2305                    chain = task.network[chain_id].log_name,
2306                    sent,
2307                    total,
2308                );
2309
2310                let _ = result.send(BroadcastStatementResult { sent, total });
2311            }
2312            WakeUpReason::MessageForChain(
2313                chain_id,
2314                ToBackgroundChain::UpdateTopicAffinity { filter },
2315            ) => {
2316                task.current_affinity_filter
2317                    .insert(chain_id, filter.clone());
2318                if let Some(peers) = task.v2_statement_peers.get_mut(&chain_id) {
2319                    let mut to_remove = Vec::new();
2320                    for peer_id in peers.iter() {
2321                        if let Err(
2322                            SendTopicAffinityError::NoConnection
2323                            | SendTopicAffinityError::ProtocolV1,
2324                        ) = task.network.send_topic_affinity(peer_id, chain_id, &filter)
2325                        {
2326                            to_remove.push(peer_id.clone());
2327                        }
2328                    }
2329                    for peer_id in &to_remove {
2330                        peers.remove(peer_id);
2331                    }
2332                }
2333            }
2334            WakeUpReason::MessageForChain(
2335                chain_id,
2336                ToBackgroundChain::Discover {
2337                    list,
2338                    important_nodes,
2339                },
2340            ) => {
2341                for (peer_id, addrs) in list {
2342                    if important_nodes {
2343                        task.important_nodes
2344                            .entry(chain_id)
2345                            .or_default()
2346                            .insert(peer_id.clone());
2347                    }
2348
2349                    // Note that we must call this function before `insert_address`, as documented
2350                    // in `basic_peering_strategy`.
2351                    task.peering_strategy
2352                        .insert_chain_peer(chain_id, peer_id.clone(), 30); // TODO: constant
2353
2354                    for addr in addrs {
2355                        let _ =
2356                            task.peering_strategy
2357                                .insert_address(&peer_id, addr.into_bytes(), 10);
2358                        // TODO: constant
2359                    }
2360                }
2361            }
2362            WakeUpReason::MessageForChain(
2363                chain_id,
2364                ToBackgroundChain::DiscoveredNodes { result },
2365            ) => {
2366                // TODO: consider returning Vec<u8>s for the addresses?
2367                let _ = result.send(
2368                    task.peering_strategy
2369                        .chain_peers_unordered(&chain_id)
2370                        .map(|peer_id| {
2371                            let addrs = task
2372                                .peering_strategy
2373                                .peer_addresses(peer_id)
2374                                .map(|a| Multiaddr::from_bytes(a.to_owned()).unwrap())
2375                                .collect::<Vec<_>>();
2376                            (peer_id.clone(), addrs)
2377                        })
2378                        .collect::<Vec<_>>(),
2379                );
2380            }
2381            WakeUpReason::MessageForChain(chain_id, ToBackgroundChain::PeersList { result }) => {
2382                let _ = result.send(
2383                    task.network
2384                        .gossip_connected_peers(
2385                            chain_id,
2386                            service::GossipKind::ConsensusTransactions,
2387                        )
2388                        .cloned()
2389                        .collect(),
2390                );
2391            }
2392            WakeUpReason::StartDiscovery(chain_id) => {
2393                // Re-insert the chain in `chains_by_next_discovery`.
2394                let chain = &mut task.network[chain_id];
2395                chain.next_discovery_when = task.platform.now() + chain.next_discovery_period;
2396                chain.next_discovery_period =
2397                    cmp::min(chain.next_discovery_period * 2, Duration::from_secs(120));
2398                task.chains_by_next_discovery.insert(
2399                    (chain.next_discovery_when.clone(), chain_id),
2400                    Box::pin(
2401                        task.platform
2402                            .sleep(task.network[chain_id].next_discovery_period),
2403                    ),
2404                );
2405
2406                // Iterative-style discovery: instead of a single FindNode against one peer,
2407                // dispatch up to ALPHA=3 FindNode requests in parallel to distinct peers,
2408                // each with a distinct random target. This gives substantially better DHT
2409                // coverage per discovery round (more peers asked, more diverse keyspace
2410                // walked) without requiring a per-query state machine.
2411                //
2412                // Order of preference for the target peer pool:
2413                //  1. Peers we know speak this chain's Kad protocol (from Identify).
2414                //  2. Peers with an open block-announces gossip substream (best-effort:
2415                //     they're connected and likely speak Kad even if we haven't gotten
2416                //     Identify yet).
2417                const PARALLEL_FIND_NODE_PER_ROUND: usize = 3;
2418
2419                let mut candidates: Vec<PeerId> = task
2420                    .network
2421                    .kademlia_capable_peers(chain_id)
2422                    .cloned()
2423                    .collect();
2424                for p in task
2425                    .network
2426                    .gossip_connected_peers(chain_id, service::GossipKind::ConsensusTransactions)
2427                    .cloned()
2428                {
2429                    if !candidates.contains(&p) {
2430                        candidates.push(p);
2431                    }
2432                }
2433
2434                let started = dispatch_find_node_requests(
2435                    &mut task.network,
2436                    &mut task.randomness,
2437                    chain_id,
2438                    &candidates,
2439                    PARALLEL_FIND_NODE_PER_ROUND,
2440                );
2441
2442                let chain_log_name = &task.network[chain_id].log_name;
2443                for (request_target, requested_peer_id) in &started {
2444                    log!(
2445                        &task.platform,
2446                        Debug,
2447                        "network",
2448                        "discovery-find-node-started",
2449                        chain = chain_log_name,
2450                        request_target,
2451                        requested_peer_id
2452                    );
2453                }
2454                if started.is_empty() {
2455                    log!(
2456                        &task.platform,
2457                        Debug,
2458                        "network",
2459                        "discovery-skipped-no-peer",
2460                        chain = chain_log_name
2461                    );
2462                }
2463            }
2464            WakeUpReason::NetworkEvent(service::Event::HandshakeFinished {
2465                peer_id,
2466                expected_peer_id,
2467                id,
2468            }) => {
2469                let remote_addr =
2470                    Multiaddr::from_bytes(task.network.connection_remote_addr(id)).unwrap(); // TODO: review this unwrap
2471                if let Some(expected_peer_id) = expected_peer_id.as_ref().filter(|p| **p != peer_id)
2472                {
2473                    log!(
2474                        &task.platform,
2475                        Debug,
2476                        "network",
2477                        "handshake-finished-peer-id-mismatch",
2478                        remote_addr,
2479                        expected_peer_id,
2480                        actual_peer_id = peer_id
2481                    );
2482
2483                    let _was_in = task
2484                        .peering_strategy
2485                        .decrease_address_connections_and_remove_if_zero(
2486                            expected_peer_id,
2487                            remote_addr.as_ref(),
2488                        );
2489                    debug_assert!(_was_in.is_ok());
2490                    let _ = task.peering_strategy.increase_address_connections(
2491                        &peer_id,
2492                        remote_addr.into_bytes().to_vec(),
2493                        10,
2494                    );
2495                } else {
2496                    log!(
2497                        &task.platform,
2498                        Debug,
2499                        "network",
2500                        "handshake-finished",
2501                        remote_addr,
2502                        peer_id
2503                    );
2504                }
2505
2506                task.metrics.connections_handshakes_finished.inc();
2507                task.bitswap_peering_strategy
2508                    .increase_peer_connections(&peer_id);
2509            }
2510            WakeUpReason::NetworkEvent(service::Event::PreHandshakeDisconnected {
2511                expected_peer_id: Some(_),
2512                ..
2513            })
2514            | WakeUpReason::NetworkEvent(service::Event::Disconnected { .. }) => {
2515                let (address, peer_id, handshake_finished) = match wake_up_reason {
2516                    WakeUpReason::NetworkEvent(service::Event::PreHandshakeDisconnected {
2517                        address,
2518                        expected_peer_id: Some(peer_id),
2519                        ..
2520                    }) => (address, peer_id, false),
2521                    WakeUpReason::NetworkEvent(service::Event::Disconnected {
2522                        address,
2523                        peer_id,
2524                        ..
2525                    }) => (address, peer_id, true),
2526                    _ => unreachable!(),
2527                };
2528
2529                task.peering_strategy
2530                    .decrease_address_connections(&peer_id, &address)
2531                    .unwrap();
2532                let address = Multiaddr::from_bytes(address).unwrap();
2533                log!(
2534                    &task.platform,
2535                    Debug,
2536                    "network",
2537                    "connection-shutdown",
2538                    peer_id,
2539                    address,
2540                    ?handshake_finished
2541                );
2542                task.metrics.connections_shutdowns.inc();
2543
2544                // Ban the peer in order to avoid trying over and over again the same address(es).
2545                // Even if the handshake was finished, it is possible that the peer simply shuts
2546                // down connections immediately after it has been opened, hence the ban.
2547                // Due to race conditions and peerid mismatches, it is possible that there is
2548                // another existing connection or connection attempt with that same peer. However,
2549                // it is not possible to be sure that we will reach 0 connections or connection
2550                // attempts, and thus we ban the peer every time.
2551                // Pre-handshake failures get a shorter ban: many parallel dials time out
2552                // before any handshake completes, and a long slot-hold there dominates
2553                // peer-discovery latency on restarts.
2554                let ban_duration = if handshake_finished {
2555                    Duration::from_secs(5)
2556                } else {
2557                    Duration::from_secs(2)
2558                };
2559                task.network.gossip_remove_desired_all(
2560                    &peer_id,
2561                    service::GossipKind::ConsensusTransactions,
2562                );
2563                for (&chain_id, what_happened) in task
2564                    .peering_strategy
2565                    .unassign_slots_and_ban(&peer_id, task.platform.now() + ban_duration)
2566                {
2567                    if matches!(
2568                        what_happened,
2569                        basic_peering_strategy::UnassignSlotsAndBan::Banned { had_slot: true }
2570                    ) {
2571                        log!(
2572                            &task.platform,
2573                            Debug,
2574                            "network",
2575                            "slot-unassigned",
2576                            chain = &task.network[chain_id].log_name,
2577                            peer_id,
2578                            ?ban_duration,
2579                            // TODO: `reason` might be wrong, `handshake_finished` is not checked.
2580                            reason = "pre-handshake-disconnect"
2581                        );
2582                    }
2583                }
2584
2585                if handshake_finished {
2586                    task.network.bitswap_remove_desired(&peer_id);
2587                    let what_happened = task
2588                        .bitswap_peering_strategy
2589                        .unassign_slot_and_ban(&peer_id, task.platform.now() + ban_duration);
2590                    if matches!(
2591                        what_happened,
2592                        bitswap_peering_strategy::UnassignSlotAndBan::Banned { had_slot: true },
2593                    ) {
2594                        log!(
2595                            &task.platform,
2596                            Debug,
2597                            "network",
2598                            "bitswap-slot-unassigned",
2599                            peer_id,
2600                            ?ban_duration,
2601                            reason = "disconnect",
2602                        );
2603                    }
2604                    let _ = task
2605                        .bitswap_peering_strategy
2606                        .decrease_peer_connections(&peer_id);
2607                }
2608            }
2609            WakeUpReason::NetworkEvent(service::Event::PreHandshakeDisconnected {
2610                expected_peer_id: None,
2611                ..
2612            }) => {
2613                // This path can't be reached as we always set an expected peer id when creating
2614                // a connection.
2615                debug_assert!(false);
2616            }
2617            WakeUpReason::NetworkEvent(service::Event::PingOutSuccess {
2618                id,
2619                peer_id,
2620                ping_time,
2621            }) => {
2622                let remote_addr =
2623                    Multiaddr::from_bytes(task.network.connection_remote_addr(id)).unwrap(); // TODO: review this unwrap
2624                log!(
2625                    &task.platform,
2626                    Debug,
2627                    "network",
2628                    "pong",
2629                    peer_id,
2630                    remote_addr,
2631                    ?ping_time
2632                );
2633            }
2634            WakeUpReason::NetworkEvent(service::Event::BlockAnnounce {
2635                chain_id,
2636                peer_id,
2637                announce,
2638            }) => {
2639                log!(
2640                    &task.platform,
2641                    Debug,
2642                    "network",
2643                    "block-announce-received",
2644                    chain = &task.network[chain_id].log_name,
2645                    peer_id,
2646                    block_hash = HashDisplay(&header::hash_from_scale_encoded_header(
2647                        announce.decode().scale_encoded_header
2648                    )),
2649                    is_best = announce.decode().is_best
2650                );
2651
2652                let decoded_announce = announce.decode();
2653                if decoded_announce.is_best {
2654                    let link = task
2655                        .open_gossip_links
2656                        .get_mut(&(chain_id, peer_id.clone()))
2657                        .unwrap();
2658                    if let Ok(decoded) = header::decode(
2659                        decoded_announce.scale_encoded_header,
2660                        task.network[chain_id].block_number_bytes,
2661                    ) {
2662                        link.best_block_hash = header::hash_from_scale_encoded_header(
2663                            decoded_announce.scale_encoded_header,
2664                        );
2665                        link.best_block_number = decoded.number;
2666                    }
2667                }
2668
2669                debug_assert!(task.events_pending_send.is_empty());
2670                task.events_pending_send
2671                    .push_back((chain_id, Event::BlockAnnounce { peer_id, announce }));
2672            }
2673            WakeUpReason::NetworkEvent(service::Event::GossipConnected {
2674                peer_id,
2675                chain_id,
2676                role,
2677                best_number,
2678                best_hash,
2679                kind: service::GossipKind::ConsensusTransactions,
2680            }) => {
2681                log!(
2682                    &task.platform,
2683                    Debug,
2684                    "network",
2685                    "gossip-open-success",
2686                    chain = &task.network[chain_id].log_name,
2687                    peer_id,
2688                    best_number,
2689                    best_hash = HashDisplay(&best_hash)
2690                );
2691
2692                let _prev_value = task.open_gossip_links.insert(
2693                    (chain_id, peer_id.clone()),
2694                    OpenGossipLinkState {
2695                        best_block_number: best_number,
2696                        best_block_hash: best_hash,
2697                        role,
2698                        finalized_block_height: None,
2699                    },
2700                );
2701                debug_assert!(_prev_value.is_none());
2702                task.network[chain_id].metrics.gossip_peers_connected.inc();
2703
2704                task.chains_ever_gossip_connected.insert(chain_id);
2705
2706                debug_assert!(task.events_pending_send.is_empty());
2707                task.events_pending_send.push_back((
2708                    chain_id,
2709                    Event::Connected {
2710                        peer_id,
2711                        role,
2712                        best_block_number: best_number,
2713                        best_block_hash: best_hash,
2714                    },
2715                ));
2716            }
2717            WakeUpReason::NetworkEvent(service::Event::GossipOpenFailed {
2718                peer_id,
2719                chain_id,
2720                error,
2721                kind: service::GossipKind::ConsensusTransactions,
2722            }) => {
2723                log!(
2724                    &task.platform,
2725                    Debug,
2726                    "network",
2727                    "gossip-open-error",
2728                    chain = &task.network[chain_id].log_name,
2729                    peer_id,
2730                    ?error,
2731                );
2732                // Must exceed polkadot-sdk's 5s notification-reject ban; otherwise we retry
2733                // into a still-active remote ban. 0.5s margin covers network delay and
2734                // clock skew between the two sides' ban timers.
2735                let ban_duration = Duration::from_millis(5500);
2736
2737                // Note that peer doesn't necessarily have an out slot, as this event might happen
2738                // as a result of an inbound gossip connection.
2739                let had_slot = if let service::GossipConnectError::GenesisMismatch { .. } = error {
2740                    matches!(
2741                        task.peering_strategy
2742                            .unassign_slot_and_remove_chain_peer(&chain_id, &peer_id),
2743                        basic_peering_strategy::UnassignSlotAndRemoveChainPeer::HadSlot
2744                    )
2745                } else {
2746                    matches!(
2747                        task.peering_strategy.unassign_slot_and_ban(
2748                            &chain_id,
2749                            &peer_id,
2750                            task.platform.now() + ban_duration,
2751                        ),
2752                        basic_peering_strategy::UnassignSlotAndBan::Banned { had_slot: true }
2753                    )
2754                };
2755
2756                if had_slot {
2757                    log!(
2758                        &task.platform,
2759                        Debug,
2760                        "network",
2761                        "slot-unassigned",
2762                        chain = &task.network[chain_id].log_name,
2763                        peer_id,
2764                        ?ban_duration,
2765                        reason = "gossip-open-failed"
2766                    );
2767                    task.network.gossip_remove_desired(
2768                        chain_id,
2769                        &peer_id,
2770                        service::GossipKind::ConsensusTransactions,
2771                    );
2772                }
2773            }
2774            WakeUpReason::NetworkEvent(service::Event::GossipDisconnected {
2775                peer_id,
2776                chain_id,
2777                kind: service::GossipKind::ConsensusTransactions,
2778            }) => {
2779                log!(
2780                    &task.platform,
2781                    Debug,
2782                    "network",
2783                    "gossip-closed",
2784                    chain = &task.network[chain_id].log_name,
2785                    peer_id,
2786                );
2787                let ban_duration = Duration::from_secs(10);
2788
2789                let _was_in = task.open_gossip_links.remove(&(chain_id, peer_id.clone()));
2790                debug_assert!(_was_in.is_some());
2791                task.network[chain_id].metrics.gossip_peers_connected.dec();
2792
2793                // Note that peer doesn't necessarily have an out slot, as this event might happen
2794                // as a result of an inbound gossip connection.
2795                if matches!(
2796                    task.peering_strategy.unassign_slot_and_ban(
2797                        &chain_id,
2798                        &peer_id,
2799                        task.platform.now() + ban_duration,
2800                    ),
2801                    basic_peering_strategy::UnassignSlotAndBan::Banned { had_slot: true }
2802                ) {
2803                    log!(
2804                        &task.platform,
2805                        Debug,
2806                        "network",
2807                        "slot-unassigned",
2808                        chain = &task.network[chain_id].log_name,
2809                        peer_id,
2810                        ?ban_duration,
2811                        reason = "gossip-closed"
2812                    );
2813                    task.network.gossip_remove_desired(
2814                        chain_id,
2815                        &peer_id,
2816                        service::GossipKind::ConsensusTransactions,
2817                    );
2818                }
2819
2820                if let Some(peers) = task.v2_statement_peers.get_mut(&chain_id) {
2821                    peers.remove(&peer_id);
2822                }
2823
2824                debug_assert!(task.events_pending_send.is_empty());
2825                task.events_pending_send
2826                    .push_back((chain_id, Event::Disconnected { peer_id }));
2827            }
2828            WakeUpReason::NetworkEvent(service::Event::BitswapConnected { peer_id }) => {
2829                task.bitswap_connected_peers = task.bitswap_connected_peers.saturating_add(1);
2830                log!(
2831                    &task.platform,
2832                    Debug,
2833                    "network",
2834                    "bitswap-open-success",
2835                    peer_id,
2836                    total = task.bitswap_connected_peers
2837                );
2838            }
2839            WakeUpReason::NetworkEvent(service::Event::BitswapOpenFailed { peer_id, error }) => {
2840                log!(
2841                    &task.platform,
2842                    Debug,
2843                    "network",
2844                    "bitswap-open-error",
2845                    peer_id,
2846                    ?error
2847                );
2848                let ban_duration = if error.is_protocol_not_available() {
2849                    Duration::from_secs(600)
2850                } else {
2851                    Duration::from_secs(15)
2852                };
2853                if matches!(
2854                    task.bitswap_peering_strategy
2855                        .unassign_slot_and_ban(&peer_id, task.platform.now() + ban_duration,),
2856                    bitswap_peering_strategy::UnassignSlotAndBan::Banned { had_slot: true }
2857                ) {
2858                    log!(
2859                        &task.platform,
2860                        Debug,
2861                        "network",
2862                        "bitswap-slot-unassigned",
2863                        peer_id,
2864                        ?ban_duration,
2865                        reason = "bitswap-open-failed"
2866                    );
2867                    task.network.bitswap_remove_desired(&peer_id);
2868                }
2869            }
2870            WakeUpReason::NetworkEvent(service::Event::BitswapMessage { peer_id, message }) => {
2871                log!(
2872                    &task.platform,
2873                    Debug,
2874                    "network",
2875                    "bitswap-message-received",
2876                    peer_id
2877                );
2878                debug_assert!(task.bitswap_event_pending_send.is_none());
2879                task.bitswap_event_pending_send =
2880                    Some(BitswapEvent::BitswapMessage { peer_id, message });
2881            }
2882            WakeUpReason::NetworkEvent(service::Event::BitswapDisconnected { peer_id }) => {
2883                debug_assert!(task.bitswap_connected_peers > 0);
2884                task.bitswap_connected_peers = task.bitswap_connected_peers.saturating_sub(1);
2885                log!(
2886                    &task.platform,
2887                    Debug,
2888                    "network",
2889                    "bitswap-closed",
2890                    peer_id,
2891                    total = task.bitswap_connected_peers
2892                );
2893                let ban_duration = Duration::from_secs(10);
2894                if matches!(
2895                    task.bitswap_peering_strategy
2896                        .unassign_slot_and_ban(&peer_id, task.platform.now() + ban_duration,),
2897                    bitswap_peering_strategy::UnassignSlotAndBan::Banned { had_slot: true }
2898                ) {
2899                    log!(
2900                        &task.platform,
2901                        Debug,
2902                        "network",
2903                        "bitswap-slot-unassigned",
2904                        peer_id,
2905                        ?ban_duration,
2906                        reason = "bitswap-closed"
2907                    );
2908                    task.network.bitswap_remove_desired(&peer_id);
2909                }
2910            }
2911            WakeUpReason::NetworkEvent(service::Event::RequestResult {
2912                substream_id,
2913                peer_id,
2914                chain_id,
2915                response: service::RequestResult::Blocks(response),
2916            }) => {
2917                match &response {
2918                    Ok(blocks) => {
2919                        log!(
2920                            &task.platform,
2921                            Debug,
2922                            "network",
2923                            "blocks-request-success",
2924                            chain = task.network[chain_id].log_name,
2925                            target = peer_id,
2926                            num_blocks = blocks.len(),
2927                            block_data_total_size =
2928                                BytesDisplay(blocks.iter().fold(0, |sum, block| {
2929                                    let block_size = block.header.as_ref().map_or(0, |h| h.len())
2930                                        + block
2931                                            .body
2932                                            .as_ref()
2933                                            .map_or(0, |b| b.iter().fold(0, |s, e| s + e.len()))
2934                                        + block
2935                                            .justifications
2936                                            .as_ref()
2937                                            .into_iter()
2938                                            .flat_map(|l| l.iter())
2939                                            .fold(0, |s, j| s + j.justification.len());
2940                                    sum + u64::try_from(block_size).unwrap()
2941                                }))
2942                        );
2943                    }
2944                    Err(error) => {
2945                        log!(
2946                            &task.platform,
2947                            Debug,
2948                            "network",
2949                            "blocks-request-error",
2950                            chain = task.network[chain_id].log_name,
2951                            target = peer_id,
2952                            ?error
2953                        );
2954                    }
2955                }
2956
2957                match &response {
2958                    Ok(_) => {}
2959                    Err(service::BlocksRequestError::Request(err)) if !err.is_protocol_error() => {}
2960                    Err(err) => {
2961                        log!(
2962                            &task.platform,
2963                            Debug,
2964                            "network",
2965                            format!(
2966                                "Error in block request with {}. This might indicate an \
2967                                incompatibility. Error: {}",
2968                                peer_id, err
2969                            )
2970                        );
2971                    }
2972                }
2973
2974                let _ = task
2975                    .blocks_requests
2976                    .remove(&substream_id)
2977                    .unwrap()
2978                    .send(response.map_err(BlocksRequestError::Request));
2979            }
2980            WakeUpReason::NetworkEvent(service::Event::RequestResult {
2981                substream_id,
2982                peer_id,
2983                chain_id,
2984                response: service::RequestResult::GrandpaWarpSync(response),
2985            }) => {
2986                match &response {
2987                    Ok(response) => {
2988                        // TODO: print total bytes size
2989                        let decoded = response.decode();
2990                        log!(
2991                            &task.platform,
2992                            Debug,
2993                            "network",
2994                            "warp-sync-request-success",
2995                            chain = task.network[chain_id].log_name,
2996                            target = peer_id,
2997                            num_fragments = decoded.fragments.len(),
2998                            is_finished = ?decoded.is_finished,
2999                        );
3000                    }
3001                    Err(error) => {
3002                        log!(
3003                            &task.platform,
3004                            Debug,
3005                            "network",
3006                            "warp-sync-request-error",
3007                            chain = task.network[chain_id].log_name,
3008                            target = peer_id,
3009                            ?error,
3010                        );
3011                    }
3012                }
3013
3014                let _ = task
3015                    .grandpa_warp_sync_requests
3016                    .remove(&substream_id)
3017                    .unwrap()
3018                    .send(response.map_err(WarpSyncRequestError::Request));
3019            }
3020            WakeUpReason::NetworkEvent(service::Event::RequestResult {
3021                substream_id,
3022                peer_id,
3023                chain_id,
3024                response: service::RequestResult::StorageProof(response),
3025            }) => {
3026                match &response {
3027                    Ok(items) => {
3028                        let decoded = items.decode();
3029                        log!(
3030                            &task.platform,
3031                            Debug,
3032                            "network",
3033                            "storage-proof-request-success",
3034                            chain = task.network[chain_id].log_name,
3035                            target = peer_id,
3036                            total_size = BytesDisplay(u64::try_from(decoded.len()).unwrap()),
3037                        );
3038                    }
3039                    Err(error) => {
3040                        log!(
3041                            &task.platform,
3042                            Debug,
3043                            "network",
3044                            "storage-proof-request-error",
3045                            chain = task.network[chain_id].log_name,
3046                            target = peer_id,
3047                            ?error
3048                        );
3049                    }
3050                }
3051
3052                // Both regular storage proof and child storage proof use the same protocol,
3053                // so check both HashMaps for the request.
3054                if let Some(sender) = task.storage_proof_requests.remove(&substream_id) {
3055                    let _ = sender.send(response.map_err(StorageProofRequestError::Request));
3056                } else if let Some(sender) = task.child_storage_proof_requests.remove(&substream_id)
3057                {
3058                    let _ = sender.send(response.map_err(ChildStorageProofRequestError::Request));
3059                } else {
3060                    unreachable!()
3061                }
3062            }
3063            WakeUpReason::NetworkEvent(service::Event::RequestResult {
3064                substream_id,
3065                peer_id,
3066                chain_id,
3067                response: service::RequestResult::CallProof(response),
3068            }) => {
3069                match &response {
3070                    Ok(items) => {
3071                        let decoded = items.decode();
3072                        log!(
3073                            &task.platform,
3074                            Debug,
3075                            "network",
3076                            "call-proof-request-success",
3077                            chain = task.network[chain_id].log_name,
3078                            target = peer_id,
3079                            total_size = BytesDisplay(u64::try_from(decoded.len()).unwrap())
3080                        );
3081                    }
3082                    Err(error) => {
3083                        log!(
3084                            &task.platform,
3085                            Debug,
3086                            "network",
3087                            "call-proof-request-error",
3088                            chain = task.network[chain_id].log_name,
3089                            target = peer_id,
3090                            ?error
3091                        );
3092                    }
3093                }
3094
3095                let _ = task
3096                    .call_proof_requests
3097                    .remove(&substream_id)
3098                    .unwrap()
3099                    .send(response.map_err(CallProofRequestError::Request));
3100            }
3101            WakeUpReason::NetworkEvent(service::Event::RequestResult {
3102                peer_id: requestee_peer_id,
3103                chain_id,
3104                response: service::RequestResult::KademliaFindNode(Ok(nodes)),
3105                ..
3106            }) => {
3107                // Track whether this response taught us anything new. If so, we reset the
3108                // chain's discovery backoff so that the next FindNode round runs at the
3109                // initial 2s interval rather than continuing to back off — Kademlia is
3110                // making progress, walk the DHT eagerly.
3111                let mut any_new_peer = false;
3112                for (peer_id, mut addrs) in nodes {
3113                    // Make sure to not insert too many address for a single peer.
3114                    // While the .
3115                    if addrs.len() >= 10 {
3116                        addrs.truncate(10);
3117                    }
3118
3119                    let mut valid_addrs = Vec::with_capacity(addrs.len());
3120                    for addr in addrs {
3121                        match Multiaddr::from_bytes(addr) {
3122                            Ok(mut a) => {
3123                                if !pop_p2p_if_matches(&mut a, &peer_id) {
3124                                    let reason = DiscoveredAddressDropReason::PeerIdMismatch;
3125                                    log!(
3126                                        &task.platform,
3127                                        Debug,
3128                                        "network",
3129                                        "discovered-address-dropped",
3130                                        chain = &task.network[chain_id].log_name,
3131                                        reason,
3132                                        announced_peer_id = peer_id,
3133                                        addr = &a,
3134                                        obtained_from = requestee_peer_id
3135                                    );
3136                                    task.metrics.discovery_addresses_dropped.inc(reason);
3137                                    continue;
3138                                }
3139                                if platform::address_parse::multiaddr_to_address(&a)
3140                                    .ok()
3141                                    .map_or(false, |addr| {
3142                                        task.platform.supports_connection_type((&addr).into())
3143                                    })
3144                                {
3145                                    valid_addrs.push(a)
3146                                } else {
3147                                    let reason = DiscoveredAddressDropReason::NotSupported;
3148                                    log!(
3149                                        &task.platform,
3150                                        Debug,
3151                                        "network",
3152                                        "discovered-address-dropped",
3153                                        chain = &task.network[chain_id].log_name,
3154                                        reason,
3155                                        peer_id,
3156                                        addr = &a,
3157                                        obtained_from = requestee_peer_id
3158                                    );
3159                                    task.metrics.discovery_addresses_dropped.inc(reason);
3160                                }
3161                            }
3162                            Err((error, addr)) => {
3163                                let reason = DiscoveredAddressDropReason::Invalid;
3164                                log!(
3165                                    &task.platform,
3166                                    Debug,
3167                                    "network",
3168                                    "discovered-address-dropped",
3169                                    chain = &task.network[chain_id].log_name,
3170                                    reason,
3171                                    peer_id,
3172                                    error,
3173                                    addr = hex::encode(&addr),
3174                                    obtained_from = requestee_peer_id
3175                                );
3176                                task.metrics.discovery_addresses_dropped.inc(reason);
3177                            }
3178                        }
3179                    }
3180
3181                    if !valid_addrs.is_empty() {
3182                        // Note that we must call this function before `insert_address`,
3183                        // as documented in `basic_peering_strategy`.
3184                        let insert_outcome =
3185                            task.peering_strategy
3186                                .insert_chain_peer(chain_id, peer_id.clone(), 30); // TODO: constant
3187
3188                        if let basic_peering_strategy::InsertChainPeerResult::Inserted {
3189                            peer_removed,
3190                        } = insert_outcome
3191                        {
3192                            any_new_peer = true;
3193                            if let Some(peer_removed) = peer_removed {
3194                                log!(
3195                                    &task.platform,
3196                                    Debug,
3197                                    "network",
3198                                    "peer-purged-from-address-book",
3199                                    chain = &task.network[chain_id].log_name,
3200                                    peer_id = peer_removed,
3201                                );
3202                            }
3203
3204                            log!(
3205                                &task.platform,
3206                                Debug,
3207                                "network",
3208                                "peer-discovered",
3209                                chain = &task.network[chain_id].log_name,
3210                                peer_id,
3211                                addrs = ?valid_addrs.iter().map(|a| a.to_string()).collect::<Vec<_>>(), // TODO: better formatting?
3212                                obtained_from = requestee_peer_id
3213                            );
3214                        }
3215                    }
3216
3217                    for addr in valid_addrs {
3218                        let _insert_result =
3219                            task.peering_strategy
3220                                .insert_address(&peer_id, addr.into_bytes(), 10); // TODO: constant
3221                        debug_assert!(!matches!(
3222                            _insert_result,
3223                            basic_peering_strategy::InsertAddressResult::UnknownPeer
3224                        ));
3225                    }
3226                }
3227
3228                if any_new_peer {
3229                    task.network[chain_id].next_discovery_period = Duration::from_secs(2);
3230                }
3231            }
3232            WakeUpReason::NetworkEvent(service::Event::RequestResult {
3233                peer_id,
3234                chain_id,
3235                response: service::RequestResult::KademliaFindNode(Err(error)),
3236                ..
3237            }) => {
3238                log!(
3239                    &task.platform,
3240                    Debug,
3241                    "network",
3242                    "discovery-find-node-error",
3243                    chain = &task.network[chain_id].log_name,
3244                    ?error,
3245                    find_node_target = peer_id,
3246                );
3247
3248                // No error is printed if the request fails due to a benign networking error such
3249                // as an unresponsive peer.
3250                match error {
3251                    service::KademliaFindNodeError::RequestFailed(err)
3252                        if !err.is_protocol_error() => {}
3253
3254                    service::KademliaFindNodeError::RequestFailed(
3255                        service::RequestError::Substream(
3256                            connection::established::RequestError::ProtocolNotAvailable,
3257                        ),
3258                    ) => {
3259                        // TODO: remove this warning in a long time
3260                        log!(
3261                            &task.platform,
3262                            Warn,
3263                            "network",
3264                            format!(
3265                                "Problem during discovery on {}: protocol not available. \
3266                                This might indicate that the version of Substrate used by \
3267                                the chain doesn't include \
3268                                <https://github.com/paritytech/substrate/pull/12545>.",
3269                                &task.network[chain_id].log_name
3270                            )
3271                        );
3272                    }
3273                    _ => {
3274                        log!(
3275                            &task.platform,
3276                            Debug,
3277                            "network",
3278                            format!(
3279                                "Problem during discovery on {}: {}",
3280                                &task.network[chain_id].log_name, error
3281                            )
3282                        );
3283                    }
3284                }
3285            }
3286            WakeUpReason::NetworkEvent(service::Event::RequestResult { .. }) => {
3287                // We never start any other kind of requests.
3288                unreachable!()
3289            }
3290            WakeUpReason::NetworkEvent(service::Event::GossipInDesired {
3291                peer_id,
3292                chain_id,
3293                kind: service::GossipKind::ConsensusTransactions,
3294            }) => {
3295                // The networking state machine guarantees that `GossipInDesired`
3296                // can't happen if we are already opening an out slot, which we do
3297                // immediately.
3298                // TODO: add debug_assert! ^
3299                if task
3300                    .network
3301                    .opened_gossip_undesired_by_chain(chain_id)
3302                    .count()
3303                    < 4
3304                {
3305                    log!(
3306                        &task.platform,
3307                        Debug,
3308                        "network",
3309                        "gossip-in-request",
3310                        chain = &task.network[chain_id].log_name,
3311                        peer_id,
3312                        outcome = "accepted"
3313                    );
3314                    task.network
3315                        .gossip_open(
3316                            chain_id,
3317                            &peer_id,
3318                            service::GossipKind::ConsensusTransactions,
3319                        )
3320                        .unwrap();
3321                } else {
3322                    log!(
3323                        &task.platform,
3324                        Debug,
3325                        "network",
3326                        "gossip-in-request",
3327                        chain = &task.network[chain_id].log_name,
3328                        peer_id,
3329                        outcome = "rejected",
3330                    );
3331                    task.network
3332                        .gossip_close(
3333                            chain_id,
3334                            &peer_id,
3335                            service::GossipKind::ConsensusTransactions,
3336                        )
3337                        .unwrap();
3338                }
3339            }
3340            WakeUpReason::NetworkEvent(service::Event::GossipInDesiredCancel { .. }) => {
3341                // Can't happen as we already instantaneously accept or reject gossip in requests.
3342                unreachable!()
3343            }
3344            WakeUpReason::NetworkEvent(service::Event::IdentifyRequestIn {
3345                peer_id,
3346                substream_id,
3347            }) => {
3348                log!(
3349                    &task.platform,
3350                    Debug,
3351                    "network",
3352                    "identify-request-received",
3353                    peer_id,
3354                );
3355                task.network
3356                    .respond_identify(substream_id, &task.identify_agent_version);
3357            }
3358            WakeUpReason::NetworkEvent(service::Event::BlocksRequestIn { .. }) => unreachable!(),
3359            WakeUpReason::NetworkEvent(service::Event::RequestInCancel { .. }) => {
3360                // All incoming requests are immediately answered.
3361                unreachable!()
3362            }
3363            WakeUpReason::NetworkEvent(service::Event::GrandpaNeighborPacket {
3364                chain_id,
3365                peer_id,
3366                state,
3367            }) => {
3368                log!(
3369                    &task.platform,
3370                    Debug,
3371                    "network",
3372                    "grandpa-neighbor-packet-received",
3373                    chain = &task.network[chain_id].log_name,
3374                    peer_id,
3375                    round_number = state.round_number,
3376                    set_id = state.set_id,
3377                    commit_finalized_height = state.commit_finalized_height,
3378                );
3379
3380                task.open_gossip_links
3381                    .get_mut(&(chain_id, peer_id.clone()))
3382                    .unwrap()
3383                    .finalized_block_height = Some(state.commit_finalized_height);
3384
3385                debug_assert!(task.events_pending_send.is_empty());
3386                task.events_pending_send.push_back((
3387                    chain_id,
3388                    Event::GrandpaNeighborPacket {
3389                        peer_id,
3390                        finalized_block_height: state.commit_finalized_height,
3391                    },
3392                ));
3393            }
3394            WakeUpReason::NetworkEvent(service::Event::GrandpaCommitMessage {
3395                chain_id,
3396                peer_id,
3397                message,
3398            }) => {
3399                log!(
3400                    &task.platform,
3401                    Debug,
3402                    "network",
3403                    "grandpa-commit-message-received",
3404                    chain = &task.network[chain_id].log_name,
3405                    peer_id,
3406                    target_block_hash = HashDisplay(message.decode().target_hash),
3407                );
3408
3409                debug_assert!(task.events_pending_send.is_empty());
3410                task.events_pending_send
3411                    .push_back((chain_id, Event::GrandpaCommitMessage { peer_id, message }));
3412            }
3413            WakeUpReason::NetworkEvent(service::Event::StatementsNotification {
3414                chain_id,
3415                peer_id,
3416                statements,
3417            }) => {
3418                debug_assert!(task.statement_event_pending_send.is_none());
3419
3420                if statements.is_empty() {
3421                    continue;
3422                }
3423
3424                task.statement_event_pending_send = Some((
3425                    chain_id,
3426                    StatementEvent::StatementsNotification {
3427                        peer_id,
3428                        statements,
3429                    },
3430                ));
3431            }
3432            WakeUpReason::NetworkEvent(service::Event::StatementProtocolConnected {
3433                peer_id,
3434                chain_id,
3435                version,
3436            }) => {
3437                log!(
3438                    &task.platform,
3439                    Trace,
3440                    "network",
3441                    "statement-protocol-open-success",
3442                    chain = &task.network[chain_id].log_name,
3443                    peer_id,
3444                    ?version,
3445                );
3446
3447                if matches!(version, codec::StatementProtocolVersion::V2) {
3448                    task.v2_statement_peers
3449                        .entry(chain_id)
3450                        .or_insert_with(|| {
3451                            HashSet::with_capacity_and_hasher(16, Default::default())
3452                        })
3453                        .insert(peer_id.clone());
3454                    if let Some(filter) = task.current_affinity_filter.get(&chain_id) {
3455                        if let Err(
3456                            SendTopicAffinityError::NoConnection
3457                            | SendTopicAffinityError::ProtocolV1,
3458                        ) = task.network.send_topic_affinity(&peer_id, chain_id, filter)
3459                        {
3460                            task.v2_statement_peers
3461                                .get_mut(&chain_id)
3462                                .unwrap()
3463                                .remove(&peer_id);
3464                        }
3465                    }
3466                }
3467            }
3468            // TODO: we don't filter outbound statements yet
3469            WakeUpReason::NetworkEvent(service::Event::StatementTopicAffinityReceived {
3470                ..
3471            }) => {}
3472            WakeUpReason::NetworkEvent(service::Event::ProtocolError { peer_id, error }) => {
3473                // TODO: handle properly?
3474                log!(
3475                    &task.platform,
3476                    Warn,
3477                    "network",
3478                    "protocol-error",
3479                    peer_id,
3480                    ?error
3481                );
3482
3483                // TODO: disconnect peer
3484            }
3485            WakeUpReason::CanAssignSlot(peer_id, chain_id) => {
3486                task.peering_strategy.assign_slot(&chain_id, &peer_id);
3487
3488                log!(
3489                    &task.platform,
3490                    Debug,
3491                    "network",
3492                    "slot-assigned",
3493                    chain = &task.network[chain_id].log_name,
3494                    peer_id
3495                );
3496
3497                task.network.gossip_insert_desired(
3498                    chain_id,
3499                    peer_id,
3500                    service::GossipKind::ConsensusTransactions,
3501                );
3502            }
3503            WakeUpReason::CanAssignBitswapSlot(peer_id) => {
3504                task.bitswap_peering_strategy.assign_slot(&peer_id).unwrap();
3505
3506                log!(
3507                    &task.platform,
3508                    Debug,
3509                    "network",
3510                    "bitswap-slot-assigned",
3511                    peer_id
3512                );
3513
3514                task.network.bitswap_insert_desired(peer_id);
3515            }
3516            WakeUpReason::NextRecentConnectionRestore => {
3517                task.num_recent_connection_opening =
3518                    task.num_recent_connection_opening.saturating_sub(1);
3519            }
3520            WakeUpReason::CanStartConnect(expected_peer_id) => {
3521                let Some(multiaddr) = task
3522                    .peering_strategy
3523                    .pick_address_and_add_connection(&expected_peer_id)
3524                else {
3525                    // There is no address for that peer in the address book.
3526                    task.network.gossip_remove_desired_all(
3527                        &expected_peer_id,
3528                        service::GossipKind::ConsensusTransactions,
3529                    );
3530                    let ban_duration = Duration::from_secs(10);
3531                    for (&chain_id, what_happened) in task.peering_strategy.unassign_slots_and_ban(
3532                        &expected_peer_id,
3533                        task.platform.now() + ban_duration,
3534                    ) {
3535                        if matches!(
3536                            what_happened,
3537                            basic_peering_strategy::UnassignSlotsAndBan::Banned { had_slot: true }
3538                        ) {
3539                            log!(
3540                                &task.platform,
3541                                Debug,
3542                                "network",
3543                                "slot-unassigned",
3544                                chain = &task.network[chain_id].log_name,
3545                                peer_id = expected_peer_id,
3546                                ?ban_duration,
3547                                reason = "no-address"
3548                            );
3549                        }
3550                    }
3551                    continue;
3552                };
3553
3554                let multiaddr = match multiaddr::Multiaddr::from_bytes(multiaddr.to_owned()) {
3555                    Ok(a) => a,
3556                    Err((multiaddr::FromBytesError, addr)) => {
3557                        // Address is in an invalid format.
3558                        let _was_in = task
3559                            .peering_strategy
3560                            .decrease_address_connections_and_remove_if_zero(
3561                                &expected_peer_id,
3562                                &addr,
3563                            );
3564                        debug_assert!(_was_in.is_ok());
3565                        continue;
3566                    }
3567                };
3568
3569                let address = address_parse::multiaddr_to_address(&multiaddr)
3570                    .ok()
3571                    .filter(|addr| {
3572                        task.platform.supports_connection_type(match &addr {
3573                            address_parse::AddressOrMultiStreamAddress::Address(addr) => {
3574                                From::from(addr)
3575                            }
3576                            address_parse::AddressOrMultiStreamAddress::MultiStreamAddress(
3577                                addr,
3578                            ) => From::from(addr),
3579                        })
3580                    });
3581
3582                let Some(address) = address else {
3583                    // Address is in an invalid format or isn't supported by the platform.
3584                    let _was_in = task
3585                        .peering_strategy
3586                        .decrease_address_connections_and_remove_if_zero(
3587                            &expected_peer_id,
3588                            multiaddr.as_ref(),
3589                        );
3590                    debug_assert!(_was_in.is_ok());
3591                    continue;
3592                };
3593
3594                // Each connection has its own individual Noise key.
3595                let noise_key = {
3596                    let mut noise_static_key = zeroize::Zeroizing::new([0u8; 32]);
3597                    task.platform.fill_random_bytes(&mut *noise_static_key);
3598                    let mut libp2p_key = zeroize::Zeroizing::new([0u8; 32]);
3599                    task.platform.fill_random_bytes(&mut *libp2p_key);
3600                    connection::NoiseKey::new(&libp2p_key, &noise_static_key)
3601                };
3602
3603                log!(
3604                    &task.platform,
3605                    Debug,
3606                    "network",
3607                    "connection-started",
3608                    expected_peer_id,
3609                    remote_addr = multiaddr,
3610                    local_peer_id =
3611                        peer_id::PublicKey::Ed25519(*noise_key.libp2p_public_ed25519_key())
3612                            .into_peer_id(),
3613                );
3614
3615                task.num_recent_connection_opening += 1;
3616                task.metrics.connections_dialed.inc();
3617
3618                let (coordinator_to_connection_tx, coordinator_to_connection_rx) =
3619                    async_channel::bounded(8);
3620                let task_name = format!("connection-{}", multiaddr);
3621
3622                match address {
3623                    address_parse::AddressOrMultiStreamAddress::Address(address) => {
3624                        // As documented in the `PlatformRef` trait, `connect_stream` must
3625                        // return as soon as possible.
3626                        let connection = task.platform.connect_stream(address).await;
3627
3628                        let (connection_id, connection_task) =
3629                            task.network.add_single_stream_connection(
3630                                task.platform.now(),
3631                                service::SingleStreamHandshakeKind::MultistreamSelectNoiseYamux {
3632                                    is_initiator: true,
3633                                    noise_key: &noise_key,
3634                                },
3635                                multiaddr.clone().into_bytes(),
3636                                Some(expected_peer_id.clone()),
3637                                coordinator_to_connection_tx,
3638                            );
3639
3640                        task.platform.spawn_task(
3641                            task_name.into(),
3642                            tasks::single_stream_connection_task::<TPlat>(
3643                                connection,
3644                                multiaddr.to_string(),
3645                                task.platform.clone(),
3646                                connection_id,
3647                                connection_task,
3648                                coordinator_to_connection_rx,
3649                                task.tasks_messages_tx.clone(),
3650                            ),
3651                        );
3652                    }
3653                    address_parse::AddressOrMultiStreamAddress::MultiStreamAddress(
3654                        platform::MultiStreamAddress::WebRtc {
3655                            ip,
3656                            port,
3657                            remote_certificate_sha256,
3658                        },
3659                    ) => {
3660                        // We need to know the local TLS certificate in order to insert the
3661                        // connection, and as such we need to call `connect_multistream` here.
3662                        // As documented in the `PlatformRef` trait, `connect_multistream` must
3663                        // return as soon as possible.
3664                        let connection = task
3665                            .platform
3666                            .connect_multistream(platform::MultiStreamAddress::WebRtc {
3667                                ip,
3668                                port,
3669                                remote_certificate_sha256,
3670                            })
3671                            .await;
3672
3673                        // Convert the SHA256 hashes into multihashes.
3674                        let local_tls_certificate_multihash = [18u8, 32]
3675                            .into_iter()
3676                            .chain(connection.local_tls_certificate_sha256.into_iter())
3677                            .collect();
3678                        let remote_tls_certificate_multihash = [18u8, 32]
3679                            .into_iter()
3680                            .chain(remote_certificate_sha256.iter().copied())
3681                            .collect();
3682
3683                        let (connection_id, connection_task) =
3684                            task.network.add_multi_stream_connection(
3685                                task.platform.now(),
3686                                service::MultiStreamHandshakeKind::WebRtc {
3687                                    is_initiator: true,
3688                                    local_tls_certificate_multihash,
3689                                    remote_tls_certificate_multihash,
3690                                    noise_key: &noise_key,
3691                                },
3692                                multiaddr.clone().into_bytes(),
3693                                Some(expected_peer_id.clone()),
3694                                coordinator_to_connection_tx,
3695                            );
3696
3697                        task.platform.spawn_task(
3698                            task_name.into(),
3699                            tasks::webrtc_multi_stream_connection_task::<TPlat>(
3700                                connection.connection,
3701                                multiaddr.to_string(),
3702                                task.platform.clone(),
3703                                connection_id,
3704                                connection_task,
3705                                coordinator_to_connection_rx,
3706                                task.tasks_messages_tx.clone(),
3707                            ),
3708                        );
3709                    }
3710                }
3711            }
3712            WakeUpReason::CanOpenGossip(peer_id, chain_id) => {
3713                task.network
3714                    .gossip_open(
3715                        chain_id,
3716                        &peer_id,
3717                        service::GossipKind::ConsensusTransactions,
3718                    )
3719                    .unwrap();
3720
3721                log!(
3722                    &task.platform,
3723                    Debug,
3724                    "network",
3725                    "gossip-open-start",
3726                    chain = &task.network[chain_id].log_name,
3727                    peer_id,
3728                );
3729            }
3730            WakeUpReason::CanOpenBitswap(peer_id) => {
3731                task.network.bitswap_open(&peer_id).unwrap();
3732
3733                log!(
3734                    &task.platform,
3735                    Debug,
3736                    "network",
3737                    "bitswap-open-start",
3738                    peer_id
3739                );
3740            }
3741            WakeUpReason::MessageToConnection {
3742                connection_id,
3743                message,
3744            } => {
3745                // Note that it is critical for the sending to not take too long here, in order to
3746                // not block the process of the network service.
3747                // In particular, if sending the message to the connection is blocked due to
3748                // sending a message on the connection-to-coordinator channel, this will result
3749                // in a deadlock.
3750                // For this reason, the connection task is always ready to immediately accept a
3751                // message on the coordinator-to-connection channel.
3752                let _send_result = task.network[connection_id].send(message).await;
3753                debug_assert!(_send_result.is_ok());
3754            }
3755        }
3756    }
3757}
3758
3759/// Starts find-node requests against `candidates` until `max` have started, each with a fresh
3760/// random target key, and returns the `(request_target, requested_peer_id)` of each.
3761///
3762/// A `kademlia_capable_peers` candidate may have no usable connection (the flag outlives the
3763/// connection it was learned on); such a peer fails with `NoConnection` and is skipped without
3764/// counting towards `max`.
3765fn dispatch_find_node_requests<TChain, TConn, TNow>(
3766    network: &mut service::ChainNetwork<TChain, TConn, TNow>,
3767    randomness: &mut impl rand_chacha::rand_core::RngCore,
3768    chain_id: service::ChainId,
3769    candidates: &[PeerId],
3770    max: usize,
3771) -> Vec<(PeerId, PeerId)>
3772where
3773    TNow: Clone
3774        + core::ops::Add<Duration, Output = TNow>
3775        + core::ops::Sub<TNow, Output = Duration>
3776        + Ord,
3777{
3778    let mut started = Vec::with_capacity(max);
3779
3780    for target in candidates {
3781        if started.len() >= max {
3782            break;
3783        }
3784
3785        let random_peer_id = {
3786            let mut pub_key = [0; 32];
3787            randomness.fill_bytes(&mut pub_key);
3788            PeerId::from_public_key(&peer_id::PublicKey::Ed25519(pub_key))
3789        };
3790
3791        match network.start_kademlia_find_node_request(
3792            target,
3793            chain_id,
3794            &random_peer_id,
3795            Duration::from_secs(20),
3796        ) {
3797            Ok(_) => started.push((target.clone(), random_peer_id)),
3798            Err(service::StartRequestError::NoConnection) => {}
3799        }
3800    }
3801
3802    started
3803}
3804
3805/// Pops a trailing `/p2p/<peer_id>` from `addr` if it matches `expected_peer`. Returns `false`
3806/// (caller must discard the address) on mismatch.
3807fn pop_p2p_if_matches(
3808    addr: &mut smoldot::libp2p::multiaddr::Multiaddr,
3809    expected_peer: &smoldot::libp2p::peer_id::PeerId,
3810) -> bool {
3811    use smoldot::libp2p::multiaddr::Protocol;
3812    match addr.iter().last() {
3813        Some(Protocol::P2p(mh)) => {
3814            if mh.into_bytes() == expected_peer.as_bytes() {
3815                addr.pop();
3816                true
3817            } else {
3818                false
3819            }
3820        }
3821        _ => true,
3822    }
3823}
3824
3825#[cfg(test)]
3826mod tests {
3827    use super::{Role, dispatch_find_node_requests, pop_p2p_if_matches, service};
3828    use core::time::Duration;
3829    use rand_chacha::rand_core::SeedableRng as _;
3830    use smoldot::libp2p::{multiaddr::Multiaddr, peer_id::PeerId};
3831
3832    // Two distinct, valid PeerIds. The first is reused from existing smoldot tests in
3833    // `lib/src/libp2p/multiaddr.rs:629`; the second is the bootnode peer-id observed in the
3834    // test environment that motivated this change.
3835    const PEER_A: &str = "12D3KooWDpJ7As7BWAwRMfu1VU2WCqNjvq387JEYKDBj4kx6nXTN";
3836    const PEER_B: &str = "12D3KooWQk1yQtG1YugyKjiQf6KNk8VjGGAT5xy1FWcnRKN4yXYJ";
3837
3838    fn peer(s: &str) -> PeerId {
3839        PeerId::from_bytes(bs58::decode(s).into_vec().unwrap()).unwrap()
3840    }
3841
3842    #[test]
3843    fn no_suffix_passes_through_unchanged() {
3844        let mut addr: Multiaddr = "/ip4/127.0.0.1/tcp/30333/ws".parse().unwrap();
3845        let before = addr.clone();
3846        assert!(pop_p2p_if_matches(&mut addr, &peer(PEER_A)));
3847        assert_eq!(addr, before);
3848    }
3849
3850    #[test]
3851    fn matching_suffix_is_stripped() {
3852        let mut addr: Multiaddr = format!("/ip4/127.0.0.1/tcp/30333/ws/p2p/{PEER_A}")
3853            .parse()
3854            .unwrap();
3855        assert!(pop_p2p_if_matches(&mut addr, &peer(PEER_A)));
3856        let expected: Multiaddr = "/ip4/127.0.0.1/tcp/30333/ws".parse().unwrap();
3857        assert_eq!(addr, expected);
3858    }
3859
3860    #[test]
3861    fn mismatched_suffix_rejects_and_keeps_addr() {
3862        let original: Multiaddr = format!("/ip4/127.0.0.1/tcp/30333/ws/p2p/{PEER_A}")
3863            .parse()
3864            .unwrap();
3865        let mut addr = original.clone();
3866        assert!(!pop_p2p_if_matches(&mut addr, &peer(PEER_B)));
3867        assert_eq!(addr, original);
3868    }
3869
3870    fn empty_network() -> (service::ChainNetwork<(), (), Duration>, service::ChainId) {
3871        let mut network = service::ChainNetwork::new(service::Config {
3872            connections_capacity: 8,
3873            chains_capacity: 1,
3874            randomness_seed: [0; 32],
3875            handshake_timeout: Duration::from_secs(10),
3876        });
3877        let chain_id = network
3878            .add_chain(service::ChainConfig {
3879                user_data: (),
3880                genesis_hash: [0; 32],
3881                fork_id: None,
3882                block_number_bytes: 4,
3883                grandpa_protocol_config: None,
3884                allow_inbound_block_requests: false,
3885                best_hash: [0; 32],
3886                best_number: 0,
3887                role: Role::Light,
3888                enable_statement_protocol: false,
3889            })
3890            .unwrap();
3891        (network, chain_id)
3892    }
3893
3894    // With no connections every candidate returns `NoConnection`, so all are skipped and no
3895    // request is started. The dispatch loop must not treat that as unreachable.
3896    #[test]
3897    fn dispatch_skips_unreachable_candidates() {
3898        let (mut network, chain_id) = empty_network();
3899        let mut randomness = rand_chacha::ChaCha20Rng::from_seed([7; 32]);
3900
3901        let candidates = [peer(PEER_A), peer(PEER_B)];
3902        let started =
3903            dispatch_find_node_requests(&mut network, &mut randomness, chain_id, &candidates, 3);
3904
3905        assert!(started.is_empty());
3906    }
3907}