1use 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
83pub struct Config<TPlat> {
85 pub platform: TPlat,
87
88 pub identify_agent_version: String,
90
91 pub chains_capacity: usize,
93
94 pub connections_open_pool_size: u32,
98
99 pub connections_open_pool_restore_delay: Duration,
104}
105
106pub struct ConfigChain {
112 pub log_name: String,
114
115 pub num_out_slots: usize,
118
119 pub genesis_block_hash: [u8; 32],
125
126 pub best_block: (u64, [u8; 32]),
129
130 pub fork_id: Option<String>,
133
134 pub block_number_bytes: usize,
136
137 pub grandpa_protocol_finalized_block_height: Option<u64>,
140
141 pub enable_statement_protocol: bool,
143
144 pub metrics: Arc<metrics::ChainMetrics>,
146}
147
148pub struct NetworkService<TPlat: PlatformRef> {
149 messages_tx: async_channel::Sender<ToBackground<TPlat>>,
151
152 platform: TPlat,
154
155 metrics: Arc<metrics::NetworkMetrics>,
157}
158
159impl<TPlat: PlatformRef> NetworkService<TPlat> {
160 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 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 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, 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, },
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 pub fn metrics(&self) -> Arc<metrics::NetworkMetrics> {
259 self.metrics.clone()
260 }
261
262 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 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 _keep_alive_messages_tx: async_channel::Sender<ToBackground<TPlat>>,
324
325 messages_tx: async_channel::Sender<ToBackgroundChain>,
327
328 metrics: Arc<metrics::ChainMetrics>,
330
331 platform: TPlat,
333}
334
335#[derive(Debug, Copy, Clone, PartialEq, Eq)]
337pub enum BanSeverity {
338 Low,
339 High,
340}
341
342#[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 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 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 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 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 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 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 pub async fn storage_proof_request(
544 self: Arc<Self>,
545 target: PeerId, 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()) .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 pub async fn call_proof_request(
581 self: Arc<Self>,
582 target: PeerId, 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()) .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 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 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(), result: tx,
667 })
668 .await
669 .unwrap();
670
671 rx.await.unwrap()
672 }
673
674 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(), scale_encoded_header: scale_encoded_header.to_vec(), is_best,
688 result: tx,
689 })
690 .await
691 .unwrap();
692
693 rx.await.unwrap()
694 }
695
696 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 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 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 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 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 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 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#[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 GrandpaCommitMessage {
853 peer_id: PeerId,
854 message: service::EncodedGrandpaCommitMessage,
855 },
856}
857
858#[derive(Debug, Clone)]
862pub enum BitswapEvent {
863 BitswapMessage {
864 peer_id: PeerId,
865 message: service::EncodedBitswapMessage,
866 },
867}
868
869#[derive(Debug, Clone)]
871pub enum StatementEvent {
872 StatementsNotification {
874 peer_id: PeerId,
875 statements: Vec<([u8; 32], codec::Statement)>,
876 },
877}
878
879#[derive(Debug, derive_more::Display, derive_more::Error)]
881pub enum BlocksRequestError {
882 NoConnection,
884 #[display("{_0}")]
886 Request(service::BlocksRequestError),
887}
888
889#[derive(Debug, derive_more::Display, derive_more::Error)]
891pub enum WarpSyncRequestError {
892 NoConnection,
894 #[display("{_0}")]
896 Request(service::GrandpaWarpSyncRequestError),
897}
898
899#[derive(Debug, derive_more::Display, derive_more::Error, Clone)]
901pub enum StorageProofRequestError {
902 NoConnection,
904 RequestTooLarge,
906 #[display("{_0}")]
908 Request(service::StorageProofRequestError),
909}
910
911#[derive(Debug, derive_more::Display, derive_more::Error, Clone)]
913pub enum CallProofRequestError {
914 NoConnection,
916 RequestTooLarge,
918 #[display("{_0}")]
920 Request(service::CallProofRequestError),
921}
922
923impl CallProofRequestError {
924 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#[derive(Debug, derive_more::Display, derive_more::Error, Clone)]
937pub enum ChildStorageProofRequestError {
938 NoConnection,
940 RequestTooLarge,
942 #[display("{_0}")]
944 Request(service::StorageProofRequestError),
945}
946
947impl ChildStorageProofRequestError {
948 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#[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 PeerIdMismatch,
970 NotSupported,
972 Invalid,
974}
975
976struct 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 StartBlocksRequest {
1008 target: PeerId, config: codec::BlocksRequestConfig,
1010 timeout: Duration,
1011 result: oneshot::Sender<Result<Vec<codec::BlockData>, BlocksRequestError>>,
1012 },
1013 StartWarpSyncRequest {
1015 target: PeerId,
1016 begin_hash: [u8; 32],
1017 timeout: Duration,
1018 result:
1019 oneshot::Sender<Result<service::EncodedGrandpaWarpSyncResponse, WarpSyncRequestError>>,
1020 },
1021 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 StartCallProofRequest {
1030 target: PeerId, config: codec::CallProofRequestConfig<'static, vec::IntoIter<Vec<u8>>>,
1032 timeout: Duration,
1033 result: oneshot::Sender<Result<service::EncodedMerkleProof, CallProofRequestError>>,
1034 },
1035 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 platform: TPlat,
1090
1091 metrics: Arc<metrics::NetworkMetrics>,
1093
1094 randomness: rand_chacha::ChaCha20Rng,
1096
1097 identify_agent_version: String,
1099
1100 tasks_messages_tx:
1102 async_channel::Sender<(service::ConnectionId, service::ConnectionToCoordinator)>,
1103
1104 tasks_messages_rx: Pin<
1106 Box<async_channel::Receiver<(service::ConnectionId, service::ConnectionToCoordinator)>>,
1107 >,
1108
1109 network: service::ChainNetwork<
1111 Chain<TPlat>,
1112 async_channel::Sender<service::CoordinatorToConnection>,
1113 TPlat::Instant,
1114 >,
1115
1116 peering_strategy: basic_peering_strategy::BasicPeeringStrategy<ChainId, TPlat::Instant>,
1118
1119 bitswap_peering_strategy: bitswap_peering_strategy::BitswapPeeringStrategy<TPlat::Instant>,
1121
1122 connections_open_pool_size: u32,
1124
1125 connections_open_pool_restore_delay: Duration,
1127
1128 num_recent_connection_opening: u32,
1132
1133 next_recent_connection_restore: Option<Pin<Box<TPlat::Delay>>>,
1135
1136 open_gossip_links: BTreeMap<(ChainId, PeerId), OpenGossipLinkState>,
1139
1140 chains_ever_gossip_connected: HashSet<ChainId, fnv::FnvBuildHasher>,
1143
1144 v2_statement_peers: HashMap<ChainId, HashSet<PeerId, fnv::FnvBuildHasher>, fnv::FnvBuildHasher>,
1146
1147 current_affinity_filter: HashMap<ChainId, AffinityFilter, fnv::FnvBuildHasher>,
1149
1150 important_nodes: HashMap<ChainId, HashSet<PeerId, fnv::FnvBuildHasher>, fnv::FnvBuildHasher>,
1154
1155 events_pending_send: VecDeque<(ChainId, Event)>,
1161
1162 bitswap_event_pending_send: Option<BitswapEvent>,
1164
1165 bitswap_connected_peers: usize,
1170
1171 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 pending_new_subscriptions: Vec<(ChainId, async_channel::Sender<Event>)>,
1184
1185 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 pending_new_bitswap_subscriptions: Vec<async_channel::Sender<BitswapEvent>>,
1203
1204 statement_event_pending_send: Option<(ChainId, StatementEvent)>,
1207
1208 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 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 chains_by_next_discovery: BTreeMap<(TPlat::Instant, ChainId), Pin<Box<TPlat::Delay>>>,
1259}
1260
1261struct Chain<TPlat: PlatformRef> {
1262 log_name: String,
1263
1264 num_references: NonZero<usize>,
1266
1267 block_number_bytes: usize,
1270
1271 num_out_slots: usize,
1273
1274 metrics: Arc<metrics::ChainMetrics>,
1276
1277 next_discovery_when: TPlat::Instant,
1279
1280 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 finalized_block_height: Option<u64>,
1292}
1293
1294async fn background_task<TPlat: PlatformRef>(mut task: BackgroundTask<TPlat>) {
1295 loop {
1296 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 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 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 return;
1559 }
1560 WakeUpReason::Message(ToBackground::AddChain {
1561 messages_rx,
1562 config,
1563 }) => {
1564 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 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 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 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 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 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 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 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 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 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 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 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 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 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<_>>(); let succeeded_peers = results
2256 .iter()
2257 .filter_map(|(peer, r)| r.is_ok().then(|| (*peer).clone()))
2258 .collect::<Vec<_>>();
2259
2260 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 Err(SendBitswapMessageError::QueueFull)
2270 } else {
2271 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 task.peering_strategy
2352 .insert_chain_peer(chain_id, peer_id.clone(), 30); for addr in addrs {
2355 let _ =
2356 task.peering_strategy
2357 .insert_address(&peer_id, addr.into_bytes(), 10);
2358 }
2360 }
2361 }
2362 WakeUpReason::MessageForChain(
2363 chain_id,
2364 ToBackgroundChain::DiscoveredNodes { result },
2365 ) => {
2366 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 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 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(); 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 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 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 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(); 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 let ban_duration = Duration::from_millis(5500);
2736
2737 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 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 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 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 let mut any_new_peer = false;
3112 for (peer_id, mut addrs) in nodes {
3113 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 let insert_outcome =
3185 task.peering_strategy
3186 .insert_chain_peer(chain_id, peer_id.clone(), 30); 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<_>>(), 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); 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 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 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 unreachable!()
3289 }
3290 WakeUpReason::NetworkEvent(service::Event::GossipInDesired {
3291 peer_id,
3292 chain_id,
3293 kind: service::GossipKind::ConsensusTransactions,
3294 }) => {
3295 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 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 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 WakeUpReason::NetworkEvent(service::Event::StatementTopicAffinityReceived {
3470 ..
3471 }) => {}
3472 WakeUpReason::NetworkEvent(service::Event::ProtocolError { peer_id, error }) => {
3473 log!(
3475 &task.platform,
3476 Warn,
3477 "network",
3478 "protocol-error",
3479 peer_id,
3480 ?error
3481 );
3482
3483 }
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 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 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 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 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 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 let connection = task
3665 .platform
3666 .connect_multistream(platform::MultiStreamAddress::WebRtc {
3667 ip,
3668 port,
3669 remote_certificate_sha256,
3670 })
3671 .await;
3672
3673 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 let _send_result = task.network[connection_id].send(message).await;
3753 debug_assert!(_send_result.is_ok());
3754 }
3755 }
3756 }
3757}
3758
3759fn 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
3805fn 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 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 #[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}