Skip to main content

reth_network/
manager.rs

1//! High level network management.
2//!
3//! The [`NetworkManager`] contains the state of the network as a whole. It controls how connections
4//! are handled and keeps track of connections to peers.
5//!
6//! ## Capabilities
7//!
8//! The network manages peers depending on their announced capabilities via their `RLPx` sessions. Most importantly the [Ethereum Wire Protocol](https://github.com/ethereum/devp2p/blob/master/caps/eth.md)(`eth`).
9//!
10//! ## Overview
11//!
12//! The [`NetworkManager`] is responsible for advancing the state of the `network`. The `network` is
13//! made up of peer-to-peer connections between nodes that are available on the same network.
14//! Responsible for peer discovery is ethereum's discovery protocol (discv4, discv5). If the address
15//! (IP+port) of our node is published via discovery, remote peers can initiate inbound connections
16//! to the local node. Once a (tcp) connection is established, both peers start to authenticate a [RLPx session](https://github.com/ethereum/devp2p/blob/master/rlpx.md) via a handshake. If the handshake was successful, both peers announce their capabilities and are now ready to exchange sub-protocol messages via the `RLPx` session.
17
18use crate::{
19    budget::{DEFAULT_BUDGET_TRY_DRAIN_NETWORK_HANDLE_CHANNEL, DEFAULT_BUDGET_TRY_DRAIN_SWARM},
20    config::NetworkConfig,
21    discovery::Discovery,
22    error::{NetworkError, ServiceKind},
23    eth_requests::IncomingEthRequest,
24    import::{BlockImport, BlockImportEvent, BlockImportOutcome, BlockValidation, NewBlockEvent},
25    listener::ConnectionListener,
26    message::{NewBlockMessage, PeerMessage},
27    metrics::{
28        BackedOffPeersMetrics, ClosedSessionsMetrics, DirectionalDisconnectMetrics, NetworkMetrics,
29        PendingSessionFailureMetrics,
30    },
31    network::{NetworkHandle, NetworkHandleMessage},
32    peers::{BackoffReason, PeersManager},
33    poll_nested_stream_with_budget,
34    protocol::IntoRlpxSubProtocol,
35    required_block_filter::RequiredBlockFilter,
36    session::SessionManager,
37    state::NetworkState,
38    swarm::{Swarm, SwarmEvent},
39    transactions::NetworkTransactionEvent,
40    FetchClient, NetworkBuilder,
41};
42use futures::{Future, StreamExt};
43use parking_lot::Mutex;
44use reth_chainspec::EnrForkIdEntry;
45use reth_eth_wire::{DisconnectReason, EthNetworkPrimitives, NetworkPrimitives};
46use reth_fs_util::{self as fs, FsPathError};
47use reth_metrics::common::mpsc::MemoryBoundedSender;
48use reth_network_api::{
49    events::{PeerEvent, SessionInfo},
50    test_utils::PeersHandle,
51    EthProtocolInfo, NetworkEvent, NetworkStatus, PeerInfo, PeerRequest,
52};
53use reth_network_peers::{NodeRecord, PeerId};
54use reth_network_types::ReputationChangeKind;
55use reth_storage_api::BlockNumReader;
56use reth_tasks::shutdown::GracefulShutdown;
57use reth_tokio_util::EventSender;
58use secp256k1::SecretKey;
59use std::{
60    net::SocketAddr,
61    path::Path,
62    pin::Pin,
63    sync::{
64        atomic::{AtomicU64, AtomicUsize, Ordering},
65        Arc,
66    },
67    task::{Context, Poll},
68    time::{Duration, Instant},
69};
70use tokio::sync::mpsc::{self, error::TrySendError};
71use tokio_stream::wrappers::UnboundedReceiverStream;
72use tracing::{debug, error, trace, warn};
73
74#[cfg_attr(doc, aquamarine::aquamarine)]
75// TODO: Inlined diagram due to a bug in aquamarine library, should become an include when it's
76// fixed. See https://github.com/mersinvald/aquamarine/issues/50
77// include_mmd!("docs/mermaid/network-manager.mmd")
78/// Manages the _entire_ state of the network.
79///
80/// This is an endless [`Future`] that consistently drives the state of the entire network forward.
81///
82/// The [`NetworkManager`] is the container type for all parts involved with advancing the network.
83///
84/// ```mermaid
85/// graph TB
86///   handle(NetworkHandle)
87///   events(NetworkEvents)
88///   transactions(Transactions Task)
89///   ethrequest(ETH Request Task)
90///   discovery(Discovery Task)
91///   subgraph NetworkManager
92///     direction LR
93///     subgraph Swarm
94///         direction TB
95///         B1[(Session Manager)]
96///         B2[(Connection Listener)]
97///         B3[(Network State)]
98///     end
99///  end
100///  handle <--> |request response channel| NetworkManager
101///  NetworkManager --> |Network events| events
102///  transactions <--> |transactions| NetworkManager
103///  ethrequest <--> |ETH request handing| NetworkManager
104///  discovery --> |Discovered peers| NetworkManager
105/// ```
106#[derive(Debug)]
107#[must_use = "The NetworkManager does nothing unless polled"]
108pub struct NetworkManager<N: NetworkPrimitives = EthNetworkPrimitives> {
109    /// The type that manages the actual network part, which includes connections.
110    swarm: Swarm<N>,
111    /// Underlying network handle that can be shared.
112    handle: NetworkHandle<N>,
113    /// Receiver half of the command channel set up between this type and the [`NetworkHandle`]
114    from_handle_rx: UnboundedReceiverStream<NetworkHandleMessage<N>>,
115    /// Handles block imports according to the `eth` protocol.
116    block_import: Box<dyn BlockImport<N::NewBlockPayload>>,
117    /// Sender for high level network events.
118    event_sender: EventSender<NetworkEvent<PeerRequest<N>>>,
119    /// Sender half to send events to the
120    /// [`TransactionsManager`](crate::transactions::TransactionsManager) task, if configured.
121    to_transactions_manager: Option<MemoryBoundedSender<NetworkTransactionEvent<N>>>,
122    /// Sender half to send events to the
123    /// [`EthRequestHandler`](crate::eth_requests::EthRequestHandler) task, if configured.
124    ///
125    /// The channel that originally receives and bundles all requests from all sessions is already
126    /// bounded. However, since handling an eth request is more I/O intensive than delegating
127    /// them from the bounded channel to the eth-request channel, it is possible that this
128    /// builds up if the node is flooded with requests.
129    ///
130    /// Even though nonmalicious requests are relatively cheap, it's possible to craft
131    /// body requests with bogus data up until the allowed max message size limit.
132    /// Thus, we use a bounded channel here to avoid unbounded build up if the node is flooded with
133    /// requests. This channel size is set at
134    /// [`ETH_REQUEST_CHANNEL_CAPACITY`](crate::builder::ETH_REQUEST_CHANNEL_CAPACITY)
135    to_eth_request_handler: Option<mpsc::Sender<IncomingEthRequest<N>>>,
136    /// Tracks the number of active session (connected peers).
137    ///
138    /// This is updated via internal events and shared via `Arc` with the [`NetworkHandle`]
139    /// Updated by the `NetworkWorker` and loaded by the `NetworkService`.
140    num_active_peers: Arc<AtomicUsize>,
141    /// Metrics for the Network
142    metrics: NetworkMetrics,
143    /// Disconnect metrics for the Network, split by connection direction.
144    disconnect_metrics: DirectionalDisconnectMetrics,
145    /// Closed sessions metrics, split by direction.
146    closed_sessions_metrics: ClosedSessionsMetrics,
147    /// Pending session failure metrics, split by direction.
148    pending_session_failure_metrics: PendingSessionFailureMetrics,
149    /// Backed off peers metrics, split by reason.
150    backed_off_peers_metrics: BackedOffPeersMetrics,
151}
152
153impl NetworkManager {
154    /// Creates the manager of a new network with [`EthNetworkPrimitives`] types.
155    ///
156    /// ```no_run
157    /// # async fn f() {
158    /// use reth_chainspec::MAINNET;
159    /// use reth_network::{NetworkConfig, NetworkManager};
160    /// use reth_tasks::Runtime;
161    /// let config = NetworkConfig::builder_with_rng_secret_key(Runtime::test())
162    ///     .build_with_noop_provider(MAINNET.clone());
163    /// let manager = NetworkManager::eth(config).await;
164    /// # }
165    /// ```
166    pub async fn eth<C: BlockNumReader + 'static>(
167        config: NetworkConfig<C, EthNetworkPrimitives>,
168    ) -> Result<Self, NetworkError> {
169        Self::new(config).await
170    }
171}
172
173impl<N: NetworkPrimitives> NetworkManager<N> {
174    /// Sets the dedicated channel for events intended for the
175    /// [`TransactionsManager`](crate::transactions::TransactionsManager).
176    pub fn with_transactions(
177        mut self,
178        tx: MemoryBoundedSender<NetworkTransactionEvent<N>>,
179    ) -> Self {
180        self.set_transactions(tx);
181        self
182    }
183
184    /// Sets the dedicated channel for events intended for the
185    /// [`TransactionsManager`](crate::transactions::TransactionsManager).
186    pub fn set_transactions(&mut self, tx: MemoryBoundedSender<NetworkTransactionEvent<N>>) {
187        self.to_transactions_manager = Some(tx);
188    }
189
190    /// Sets the dedicated channel for events intended for the
191    /// [`EthRequestHandler`](crate::eth_requests::EthRequestHandler).
192    pub fn with_eth_request_handler(mut self, tx: mpsc::Sender<IncomingEthRequest<N>>) -> Self {
193        self.set_eth_request_handler(tx);
194        self
195    }
196
197    /// Sets the dedicated channel for events intended for the
198    /// [`EthRequestHandler`](crate::eth_requests::EthRequestHandler).
199    pub fn set_eth_request_handler(&mut self, tx: mpsc::Sender<IncomingEthRequest<N>>) {
200        self.to_eth_request_handler = Some(tx);
201    }
202
203    /// Adds an additional protocol handler to the `RLPx` sub-protocol list.
204    pub fn add_rlpx_sub_protocol(&mut self, protocol: impl IntoRlpxSubProtocol) {
205        self.swarm.add_rlpx_sub_protocol(protocol)
206    }
207
208    /// Returns the [`NetworkHandle`] that can be cloned and shared.
209    ///
210    /// The [`NetworkHandle`] can be used to interact with this [`NetworkManager`]
211    pub const fn handle(&self) -> &NetworkHandle<N> {
212        &self.handle
213    }
214
215    /// Returns the secret key used for authenticating sessions.
216    pub const fn secret_key(&self) -> SecretKey {
217        self.swarm.sessions().secret_key()
218    }
219
220    #[inline]
221    fn update_poll_metrics(&self, start: Instant, poll_durations: NetworkManagerPollDurations) {
222        let metrics = &self.metrics;
223
224        let NetworkManagerPollDurations { acc_network_handle, acc_swarm } = poll_durations;
225
226        // update metrics for whole poll function
227        metrics.duration_poll_network_manager.set(start.elapsed().as_secs_f64());
228        // update poll metrics for nested items
229        metrics.acc_duration_poll_network_handle.set(acc_network_handle.as_secs_f64());
230        metrics.acc_duration_poll_swarm.set(acc_swarm.as_secs_f64());
231    }
232
233    /// Creates the manager of a new network.
234    ///
235    /// The [`NetworkManager`] is an endless future that needs to be polled in order to advance the
236    /// state of the entire network.
237    pub async fn new<C: BlockNumReader + 'static>(
238        config: NetworkConfig<C, N>,
239    ) -> Result<Self, NetworkError> {
240        let NetworkConfig {
241            client,
242            secret_key,
243            discovery_v4_addr,
244            mut discovery_v4_config,
245            mut discovery_v5_config,
246            listener_addr,
247            peers_config,
248            sessions_config,
249            chain_id,
250            block_import,
251            network_mode,
252            boot_nodes,
253            executor,
254            hello_message,
255            status,
256            fork_filter,
257            dns_discovery_config,
258            extra_protocols,
259            tx_gossip_disabled,
260            transactions_manager_config: _,
261            nat,
262            handshake,
263            eth_max_message_size,
264            required_block_hashes,
265        } = config;
266
267        let peers_manager = PeersManager::new(peers_config);
268        let peers_handle = peers_manager.handle();
269
270        let incoming = ConnectionListener::bind(listener_addr).await.map_err(|err| {
271            NetworkError::from_io_error(err, ServiceKind::Listener(listener_addr))
272        })?;
273
274        // retrieve the tcp address of the socket
275        let listener_addr = incoming.local_address();
276
277        // resolve boot nodes
278        let resolved_boot_nodes =
279            futures::future::try_join_all(boot_nodes.iter().map(|record| record.resolve())).await?;
280
281        if let Some(disc_config) = discovery_v4_config.as_mut() {
282            // merge configured boot nodes
283            disc_config.bootstrap_nodes.extend(resolved_boot_nodes.clone());
284            // add the forkid entry for EIP-868, but wrap it in an `EnrForkIdEntry` for proper
285            // encoding
286            disc_config.add_eip868_pair("eth", EnrForkIdEntry::from(status.forkid));
287        }
288
289        if let Some(discv5) = discovery_v5_config.as_mut() {
290            // merge configured boot nodes
291            discv5.extend_unsigned_boot_nodes(resolved_boot_nodes)
292        }
293
294        let discovery = Discovery::new(
295            listener_addr,
296            discovery_v4_addr,
297            secret_key,
298            discovery_v4_config,
299            discovery_v5_config,
300            dns_discovery_config,
301            nat.clone(),
302        )
303        .await?;
304        // need to retrieve the addr here since provided port could be `0`
305        let local_peer_id = discovery.local_id();
306        let discv4 = discovery.discv4();
307        let discv5 = discovery.discv5();
308
309        let num_active_peers = Arc::new(AtomicUsize::new(0));
310
311        let sessions = SessionManager::new(
312            secret_key,
313            sessions_config,
314            executor,
315            status,
316            hello_message,
317            fork_filter,
318            extra_protocols,
319            handshake,
320            eth_max_message_size,
321            network_mode.is_stake(),
322        );
323
324        let state = NetworkState::new(
325            crate::state::BlockNumReader::new(client),
326            discovery,
327            peers_manager,
328            Arc::clone(&num_active_peers),
329        );
330
331        let swarm = Swarm::new(incoming, sessions, state);
332
333        let (to_manager_tx, from_handle_rx) = mpsc::unbounded_channel();
334
335        let event_sender: EventSender<NetworkEvent<PeerRequest<N>>> = Default::default();
336
337        let handle = NetworkHandle::new(
338            Arc::clone(&num_active_peers),
339            Arc::new(Mutex::new(listener_addr)),
340            to_manager_tx,
341            secret_key,
342            local_peer_id,
343            peers_handle,
344            network_mode,
345            Arc::new(AtomicU64::new(chain_id)),
346            tx_gossip_disabled,
347            discv4,
348            discv5,
349            event_sender.clone(),
350            nat,
351        );
352
353        // Spawn required block peer filter if configured
354        if !required_block_hashes.is_empty() {
355            let filter = RequiredBlockFilter::new(handle.clone(), required_block_hashes);
356            filter.spawn();
357        }
358
359        Ok(Self {
360            swarm,
361            handle,
362            from_handle_rx: UnboundedReceiverStream::new(from_handle_rx),
363            block_import,
364            event_sender,
365            to_transactions_manager: None,
366            to_eth_request_handler: None,
367            num_active_peers,
368            metrics: Default::default(),
369            disconnect_metrics: Default::default(),
370            closed_sessions_metrics: Default::default(),
371            pending_session_failure_metrics: Default::default(),
372            backed_off_peers_metrics: Default::default(),
373        })
374    }
375
376    /// Create a new [`NetworkManager`] instance and start a [`NetworkBuilder`] to configure all
377    /// components of the network
378    ///
379    /// ```
380    /// use reth_network::{
381    ///     config::rng_secret_key, EthNetworkPrimitives, NetworkConfig, NetworkManager,
382    /// };
383    /// use reth_network_peers::mainnet_nodes;
384    /// use reth_storage_api::noop::NoopProvider;
385    /// use reth_tasks::Runtime;
386    /// use reth_transaction_pool::TransactionPool;
387    /// async fn launch<Pool: TransactionPool>(pool: Pool) {
388    ///     // This block provider implementation is used for testing purposes.
389    ///     let client = NoopProvider::default();
390    ///
391    ///     // The key that's used for encrypting sessions and to identify our node.
392    ///     let local_key = rng_secret_key();
393    ///
394    ///     let config = NetworkConfig::<_, EthNetworkPrimitives>::builder(local_key, Runtime::test())
395    ///         .boot_nodes(mainnet_nodes())
396    ///         .build(client.clone());
397    ///     let transactions_manager_config = config.transactions_manager_config.clone();
398    ///
399    ///     // create the network instance
400    ///     let (handle, network, transactions, request_handler) = NetworkManager::builder(config)
401    ///         .await
402    ///         .unwrap()
403    ///         .transactions(pool, transactions_manager_config)
404    ///         .request_handler(client)
405    ///         .split_with_handle();
406    /// }
407    /// ```
408    pub async fn builder<C: BlockNumReader + 'static>(
409        config: NetworkConfig<C, N>,
410    ) -> Result<NetworkBuilder<(), (), N>, NetworkError> {
411        let network = Self::new(config).await?;
412        Ok(network.into_builder())
413    }
414
415    /// Create a [`NetworkBuilder`] to configure all components of the network
416    pub const fn into_builder(self) -> NetworkBuilder<(), (), N> {
417        NetworkBuilder { network: self, transactions: (), request_handler: () }
418    }
419
420    /// Returns the [`SocketAddr`] that listens for incoming tcp connections.
421    pub const fn local_addr(&self) -> SocketAddr {
422        self.swarm.listener().local_address()
423    }
424
425    /// How many peers we're currently connected to.
426    pub fn num_connected_peers(&self) -> usize {
427        self.swarm.state().num_active_peers()
428    }
429
430    /// Returns the [`PeerId`] used in the network.
431    pub fn peer_id(&self) -> &PeerId {
432        self.handle.peer_id()
433    }
434
435    /// Returns an iterator over all peers in the peer set.
436    pub fn all_peers(&self) -> impl Iterator<Item = NodeRecord> + '_ {
437        self.swarm.peers().iter_peers()
438    }
439
440    /// Returns the number of peers in the peer set.
441    pub fn num_known_peers(&self) -> usize {
442        self.swarm.peers().num_known_peers()
443    }
444
445    /// Returns a new [`PeersHandle`] that can be cloned and shared.
446    ///
447    /// The [`PeersHandle`] can be used to interact with the network's peer set.
448    pub fn peers_handle(&self) -> PeersHandle {
449        self.swarm.peers().handle()
450    }
451
452    /// Collect the peers from the [`NetworkManager`] and write them to the given
453    /// `persistent_peers_file`.
454    ///
455    /// Only persists peers that are not currently backed off or banned. Includes metadata like
456    /// peer kind, fork ID, and reputation.
457    pub fn write_peers_to_file(&self, persistent_peers_file: &Path) -> Result<(), FsPathError> {
458        let peers = self.swarm.peers().persistable_peers().collect::<Vec<_>>();
459        persistent_peers_file.parent().map(fs::create_dir_all).transpose()?;
460        reth_fs_util::write_json_file(persistent_peers_file, &peers)?;
461        Ok(())
462    }
463
464    /// Returns a new [`FetchClient`] that can be cloned and shared.
465    ///
466    /// The [`FetchClient`] is the entrypoint for sending requests to the network, including
467    /// `snap/2` requests via its [`SnapClient`](reth_network_p2p::snap::client::SnapClient) impl.
468    pub fn fetch_client(&self) -> FetchClient<N> {
469        self.swarm.state().fetch_client()
470    }
471
472    /// Returns the current [`NetworkStatus`] for the local node.
473    pub fn status(&self) -> NetworkStatus {
474        let sessions = self.swarm.sessions();
475        let status = sessions.status();
476        let hello_message = sessions.hello_message();
477
478        #[expect(deprecated)]
479        NetworkStatus {
480            client_version: hello_message.client_version,
481            protocol_version: hello_message.protocol_version as u64,
482            eth_protocol_info: EthProtocolInfo {
483                difficulty: None,
484                head: status.blockhash,
485                network: status.chain.id(),
486                genesis: status.genesis,
487                config: Default::default(),
488            },
489            capabilities: hello_message
490                .protocols
491                .into_iter()
492                .map(|protocol| protocol.cap)
493                .collect(),
494        }
495    }
496
497    /// Sends an event to the [`TransactionsManager`](crate::transactions::TransactionsManager) if
498    /// configured.
499    fn notify_tx_manager(&self, event: NetworkTransactionEvent<N>) {
500        if let Some(ref tx) = self.to_transactions_manager &&
501            let Err(e) = tx.try_send(event)
502        {
503            match e {
504                TrySendError::Full(_) => {
505                    trace!(target: "net", "Transaction events channel at capacity, dropping event");
506                    self.metrics.total_dropped_tx_events_at_full_capacity.increment(1);
507                }
508                TrySendError::Closed(_) => {}
509            }
510        }
511    }
512
513    /// Sends an event to the [`EthRequestManager`](crate::eth_requests::EthRequestHandler) if
514    /// configured.
515    fn delegate_eth_request(&self, event: IncomingEthRequest<N>) {
516        if let Some(ref reqs) = self.to_eth_request_handler {
517            let _ = reqs.try_send(event).map_err(|e| {
518                if let TrySendError::Full(_) = e {
519                    debug!(target:"net", "EthRequestHandler channel is full!");
520                    self.metrics.total_dropped_eth_requests_at_full_capacity.increment(1);
521                }
522            });
523        }
524    }
525
526    /// Handle an incoming request from the peer
527    fn on_eth_request(&self, peer_id: PeerId, req: PeerRequest<N>) {
528        match req {
529            PeerRequest::GetBlockHeaders { request, response } => {
530                self.delegate_eth_request(IncomingEthRequest::GetBlockHeaders {
531                    peer_id,
532                    request,
533                    response,
534                })
535            }
536            PeerRequest::GetBlockBodies { request, response } => {
537                self.delegate_eth_request(IncomingEthRequest::GetBlockBodies {
538                    peer_id,
539                    request,
540                    response,
541                })
542            }
543            PeerRequest::GetNodeData { request, response } => {
544                self.delegate_eth_request(IncomingEthRequest::GetNodeData {
545                    peer_id,
546                    request,
547                    response,
548                })
549            }
550            PeerRequest::GetReceipts { request, response } => {
551                self.delegate_eth_request(IncomingEthRequest::GetReceipts {
552                    peer_id,
553                    request,
554                    response,
555                })
556            }
557            PeerRequest::GetReceipts69 { request, response } => {
558                self.delegate_eth_request(IncomingEthRequest::GetReceipts69 {
559                    peer_id,
560                    request,
561                    response,
562                })
563            }
564            PeerRequest::GetReceipts70 { request, response } => {
565                self.delegate_eth_request(IncomingEthRequest::GetReceipts70 {
566                    peer_id,
567                    request,
568                    response,
569                })
570            }
571            PeerRequest::GetBlockAccessLists { request, response } => {
572                self.delegate_eth_request(IncomingEthRequest::GetBlockAccessLists {
573                    peer_id,
574                    request,
575                    response,
576                })
577            }
578            PeerRequest::GetCells { request, response } => self
579                .delegate_eth_request(IncomingEthRequest::GetCells { peer_id, request, response }),
580            PeerRequest::GetPooledTransactions { request, response } => {
581                self.notify_tx_manager(NetworkTransactionEvent::GetPooledTransactions {
582                    peer_id,
583                    request,
584                    response,
585                });
586            }
587            PeerRequest::GetSnap { request, response } => self
588                .delegate_eth_request(IncomingEthRequest::GetSnap { peer_id, request, response }),
589        }
590    }
591
592    /// Invoked after a `NewBlock` message from the peer was validated
593    fn on_block_import_result(&mut self, event: BlockImportEvent<N::NewBlockPayload>) {
594        match event {
595            BlockImportEvent::Announcement(validation) => match validation {
596                BlockValidation::ValidHeader { block } => {
597                    self.swarm.state_mut().announce_new_block(block);
598                }
599                BlockValidation::ValidBlock { block } => {
600                    self.swarm.state_mut().announce_new_block_hash(block);
601                }
602            },
603            BlockImportEvent::Outcome(outcome) => {
604                let BlockImportOutcome { peer, result } = outcome;
605                match result {
606                    Ok(validated_block) => match validated_block {
607                        BlockValidation::ValidHeader { block } => {
608                            self.swarm.state_mut().update_peer_block(
609                                &peer,
610                                block.hash,
611                                block.number(),
612                            );
613                            self.swarm.state_mut().announce_new_block(block);
614                        }
615                        BlockValidation::ValidBlock { block } => {
616                            self.swarm.state_mut().announce_new_block_hash(block);
617                        }
618                    },
619                    Err(_err) => {
620                        self.swarm
621                            .state_mut()
622                            .peers_mut()
623                            .apply_reputation_change(&peer, ReputationChangeKind::BadBlock);
624                    }
625                }
626            }
627        }
628    }
629
630    /// Enforces [EIP-3675](https://eips.ethereum.org/EIPS/eip-3675#devp2p) consensus rules for the network protocol
631    ///
632    /// Depending on the mode of the network:
633    ///    - disconnect peer if in POS
634    ///    - execute the closure if in POW
635    fn within_pow_or_disconnect<F>(&mut self, peer_id: PeerId, only_pow: F)
636    where
637        F: FnOnce(&mut Self),
638    {
639        // reject message in POS
640        if self.handle.mode().is_stake() {
641            // connections to peers which send invalid messages should be terminated
642            self.swarm
643                .sessions_mut()
644                .disconnect(peer_id, Some(DisconnectReason::SubprotocolSpecific));
645        } else {
646            only_pow(self);
647        }
648    }
649
650    /// Handles a received Message from the peer's session.
651    fn on_peer_message(&mut self, peer_id: PeerId, msg: PeerMessage<N>) {
652        match msg {
653            PeerMessage::NewBlockHashes(hashes) => {
654                self.within_pow_or_disconnect(peer_id, |this| {
655                    // update peer's state, to track what blocks this peer has seen
656                    this.swarm.state_mut().on_new_block_hashes(peer_id, hashes.to_vec());
657                    // start block import process for the hashes
658                    this.block_import.on_new_block(peer_id, NewBlockEvent::Hashes(hashes));
659                })
660            }
661            PeerMessage::NewBlock(block) => {
662                self.within_pow_or_disconnect(peer_id, move |this| {
663                    this.swarm.state_mut().on_new_block(peer_id, block.hash);
664                    // start block import process
665                    this.block_import.on_new_block(peer_id, NewBlockEvent::Block(block));
666                });
667            }
668            PeerMessage::PooledTransactions(msg) => {
669                self.notify_tx_manager(NetworkTransactionEvent::IncomingPooledTransactionHashes {
670                    peer_id,
671                    msg,
672                });
673            }
674            PeerMessage::EthRequest(req) => {
675                self.on_eth_request(peer_id, req);
676            }
677            PeerMessage::ReceivedTransaction(msg) => {
678                self.notify_tx_manager(NetworkTransactionEvent::IncomingTransactions {
679                    peer_id,
680                    msg,
681                });
682            }
683            PeerMessage::SendTransactions(_) | PeerMessage::SendBroadcastPoolTransactions(_) => {
684                unreachable!("Not emitted by session")
685            }
686            PeerMessage::BlockRangeUpdated(_) => {}
687            PeerMessage::Other(other) => {
688                debug!(target: "net", message_id=%other.id, "Ignoring unsupported message");
689            }
690        }
691    }
692
693    /// Handler for received messages from a handle
694    fn on_handle_message(&mut self, msg: NetworkHandleMessage<N>) {
695        match msg {
696            NetworkHandleMessage::DiscoveryListener(tx) => {
697                self.swarm.state_mut().discovery_mut().add_listener(tx);
698            }
699            NetworkHandleMessage::AnnounceBlock(block, hash) => {
700                if self.handle.mode().is_stake() {
701                    // See [EIP-3675](https://eips.ethereum.org/EIPS/eip-3675#devp2p)
702                    warn!(target: "net", "Peer performed block propagation, but it is not supported in proof of stake (EIP-3675)");
703                    return
704                }
705                let msg = NewBlockMessage { hash, block: Arc::new(block) };
706                self.swarm.state_mut().announce_new_block(msg);
707            }
708            NetworkHandleMessage::EthRequest { peer_id, request } => {
709                self.swarm.sessions_mut().send_message(&peer_id, PeerMessage::EthRequest(request))
710            }
711            NetworkHandleMessage::SendTransaction { peer_id, msg } => {
712                self.swarm.sessions_mut().send_message(&peer_id, PeerMessage::SendTransactions(msg))
713            }
714            NetworkHandleMessage::SendBroadcastPoolTransactions { peer_id, msg } => self
715                .swarm
716                .sessions_mut()
717                .send_message(&peer_id, PeerMessage::SendBroadcastPoolTransactions(msg)),
718            NetworkHandleMessage::SendPooledTransactionHashes { peer_id, msg } => self
719                .swarm
720                .sessions_mut()
721                .send_message(&peer_id, PeerMessage::PooledTransactions(msg)),
722            NetworkHandleMessage::AddTrustedPeerId(peer_id) => {
723                self.swarm.state_mut().add_trusted_peer_id(peer_id);
724            }
725            NetworkHandleMessage::AddTrustedPeerNode(trusted_peer) => {
726                if !self.swarm.is_shutting_down() {
727                    self.swarm.state_mut().add_trusted_peer_node(trusted_peer);
728                }
729            }
730            NetworkHandleMessage::AddPeerAddress(peer, kind, addr) => {
731                // only add peer if we are not shutting down
732                if !self.swarm.is_shutting_down() {
733                    self.swarm.state_mut().add_peer_kind(peer, kind, addr);
734                }
735            }
736            NetworkHandleMessage::RemovePeer(peer_id, kind) => {
737                self.swarm.state_mut().remove_peer_kind(peer_id, kind);
738            }
739            NetworkHandleMessage::DisconnectPeer(peer_id, reason) => {
740                self.swarm.sessions_mut().disconnect(peer_id, reason);
741            }
742            NetworkHandleMessage::BanPeer(peer_id) => {
743                self.swarm.peers_mut().ban_peer_by_admin(peer_id);
744            }
745            NetworkHandleMessage::UnbanPeer(peer_id) => {
746                self.swarm.peers_mut().unban_peer_by_admin(peer_id);
747            }
748            NetworkHandleMessage::ConnectPeer(peer_id, kind, addr) => {
749                self.swarm.state_mut().add_and_connect(peer_id, kind, addr);
750            }
751            NetworkHandleMessage::SetNetworkState(net_state) => {
752                // Sets network connection state between Active and Hibernate.
753                // If hibernate stops the node to fill new outbound
754                // connections, this is beneficial for sync stages that do not require a network
755                // connection.
756                self.swarm.on_network_state_change(net_state);
757            }
758
759            NetworkHandleMessage::Shutdown(tx) => {
760                self.perform_network_shutdown();
761                let _ = tx.send(());
762            }
763            NetworkHandleMessage::ReputationChange(peer_id, kind) => {
764                self.swarm.peers_mut().apply_reputation_change(&peer_id, kind);
765            }
766            NetworkHandleMessage::GetReputationById(peer_id, tx) => {
767                let _ = tx.send(self.swarm.peers().get_reputation(&peer_id));
768            }
769            NetworkHandleMessage::FetchClient(tx) => {
770                let _ = tx.send(self.fetch_client());
771            }
772            NetworkHandleMessage::GetStatus(tx) => {
773                let _ = tx.send(self.status());
774            }
775            NetworkHandleMessage::StatusUpdate { head } => {
776                if let Some(transition) = self.swarm.sessions_mut().on_status_update(head) {
777                    self.swarm.state_mut().update_fork_id(transition.current);
778                }
779            }
780            NetworkHandleMessage::SetForkFilter { fork_filter } => {
781                let fork_id = self.swarm.sessions_mut().set_fork_filter(fork_filter);
782                self.swarm.state_mut().update_fork_id(fork_id);
783            }
784            NetworkHandleMessage::GetPeerInfos(tx) => {
785                let _ = tx.send(self.get_peer_infos());
786            }
787            NetworkHandleMessage::GetPeerInfoById(peer_id, tx) => {
788                let _ = tx.send(self.get_peer_info_by_id(peer_id));
789            }
790            NetworkHandleMessage::GetPeerInfosByIds(peer_ids, tx) => {
791                let _ = tx.send(self.get_peer_infos_by_ids(peer_ids));
792            }
793            NetworkHandleMessage::GetPeerInfosByPeerKind(kind, tx) => {
794                let peer_ids = self.swarm.peers().peers_by_kind(kind);
795                let _ = tx.send(self.get_peer_infos_by_ids(peer_ids));
796            }
797            NetworkHandleMessage::AddRlpxSubProtocol(proto) => self.add_rlpx_sub_protocol(proto),
798            NetworkHandleMessage::GetTransactionsHandle(tx) => {
799                if let Some(ref tx_inner) = self.to_transactions_manager {
800                    let _ = tx_inner.try_send(NetworkTransactionEvent::GetTransactionsHandle(tx));
801                } else {
802                    let _ = tx.send(None);
803                }
804            }
805            NetworkHandleMessage::InternalBlockRangeUpdate(block_range_update) => {
806                self.swarm.sessions_mut().update_advertised_block_range(block_range_update);
807            }
808            NetworkHandleMessage::EthMessage { peer_id, message } => {
809                self.swarm.sessions_mut().send_message(&peer_id, message)
810            }
811        }
812    }
813
814    fn on_swarm_event(&mut self, event: SwarmEvent<N>) {
815        // handle event
816        match event {
817            SwarmEvent::ValidMessage { peer_id, message } => self.on_peer_message(peer_id, message),
818            SwarmEvent::TcpListenerClosed { remote_addr } => {
819                trace!(target: "net", ?remote_addr, "TCP listener closed.");
820            }
821            SwarmEvent::TcpListenerError(err) => {
822                trace!(target: "net", %err, "TCP connection error.");
823            }
824            SwarmEvent::IncomingTcpConnection { remote_addr, session_id } => {
825                trace!(target: "net", ?session_id, ?remote_addr, "Incoming connection");
826                self.metrics.total_incoming_connections.increment(1);
827                self.metrics
828                    .incoming_connections
829                    .set(self.swarm.peers().num_inbound_connections() as f64);
830            }
831            SwarmEvent::OutgoingTcpConnection { remote_addr, peer_id } => {
832                trace!(target: "net", ?remote_addr, ?peer_id, "Starting outbound connection.");
833                self.metrics.total_outgoing_connections.increment(1);
834                self.update_pending_connection_metrics()
835            }
836            SwarmEvent::SessionEstablished {
837                peer_id,
838                remote_addr,
839                client_version,
840                capabilities,
841                version,
842                messages,
843                status,
844                direction,
845            } => {
846                let total_active = self.num_active_peers.fetch_add(1, Ordering::Relaxed) + 1;
847                self.metrics.connected_peers.set(total_active as f64);
848                debug!(
849                    target: "net",
850                    ?remote_addr,
851                    %client_version,
852                    ?peer_id,
853                    ?total_active,
854                    kind=%direction,
855                    peer_enode=%NodeRecord::new(remote_addr, peer_id),
856                    "Session established"
857                );
858
859                if direction.is_incoming() {
860                    self.swarm
861                        .state_mut()
862                        .peers_mut()
863                        .on_incoming_session_established(peer_id, remote_addr);
864                }
865
866                if direction.is_outgoing() {
867                    self.swarm.peers_mut().on_active_outgoing_established(peer_id);
868                }
869
870                self.update_active_connection_metrics();
871
872                let peer_kind = self
873                    .swarm
874                    .state()
875                    .peers()
876                    .peer_by_id(peer_id)
877                    .map(|(_, kind)| kind)
878                    .unwrap_or_default();
879                let session_info = SessionInfo {
880                    peer_id,
881                    remote_addr,
882                    client_version,
883                    capabilities,
884                    status,
885                    version,
886                    peer_kind,
887                };
888
889                self.event_sender
890                    .notify(NetworkEvent::ActivePeerSession { info: session_info, messages });
891            }
892            SwarmEvent::PeerAdded(peer_id) => {
893                trace!(target: "net", ?peer_id, "Peer added");
894                self.event_sender.notify(NetworkEvent::Peer(PeerEvent::PeerAdded(peer_id)));
895                self.metrics.tracked_peers.set(self.swarm.peers().num_known_peers() as f64);
896            }
897            SwarmEvent::PeerRemoved(peer_id) => {
898                trace!(target: "net", ?peer_id, "Peer dropped");
899                self.event_sender.notify(NetworkEvent::Peer(PeerEvent::PeerRemoved(peer_id)));
900                self.metrics.tracked_peers.set(self.swarm.peers().num_known_peers() as f64);
901            }
902            SwarmEvent::SessionClosed { peer_id, remote_addr, error } => {
903                let total_active = self.num_active_peers.fetch_sub(1, Ordering::Relaxed) - 1;
904                self.metrics.connected_peers.set(total_active as f64);
905                trace!(
906                    target: "net",
907                    ?remote_addr,
908                    ?peer_id,
909                    ?total_active,
910                    ?error,
911                    "Session disconnected"
912                );
913
914                // Capture direction before state is reset to Idle
915                let is_inbound = self.swarm.peers().is_inbound_peer(&peer_id);
916
917                let reason = if let Some(ref err) = error {
918                    // If the connection was closed due to an error, we report
919                    // the peer
920                    self.swarm.peers_mut().on_active_session_dropped(&remote_addr, &peer_id, err);
921                    self.backed_off_peers_metrics.increment_for_reason(
922                        BackoffReason::from_disconnect(err.as_disconnected()),
923                    );
924                    err.as_disconnected()
925                } else {
926                    // Gracefully disconnected
927                    self.swarm.peers_mut().on_active_session_gracefully_closed(peer_id);
928                    self.backed_off_peers_metrics
929                        .increment_for_reason(BackoffReason::GracefulClose);
930                    None
931                };
932                self.closed_sessions_metrics.active.increment(1);
933                self.update_active_connection_metrics();
934
935                if let Some(reason) = reason {
936                    if is_inbound {
937                        self.disconnect_metrics.increment_inbound(reason);
938                    } else {
939                        self.disconnect_metrics.increment_outbound(reason);
940                    }
941                }
942                self.metrics.backed_off_peers.set(self.swarm.peers().num_backed_off_peers() as f64);
943                self.event_sender
944                    .notify(NetworkEvent::Peer(PeerEvent::SessionClosed { peer_id, reason }));
945            }
946            SwarmEvent::IncomingPendingSessionClosed { remote_addr, error } => {
947                trace!(
948                    target: "net",
949                    ?remote_addr,
950                    ?error,
951                    "Incoming pending session failed"
952                );
953
954                if let Some(ref err) = error {
955                    self.swarm
956                        .state_mut()
957                        .peers_mut()
958                        .on_incoming_pending_session_dropped(remote_addr, err);
959                    self.pending_session_failure_metrics.inbound.increment(1);
960                    if let Some(reason) = err.as_disconnected() {
961                        self.disconnect_metrics.increment_inbound(reason);
962                    }
963                } else {
964                    self.swarm
965                        .state_mut()
966                        .peers_mut()
967                        .on_incoming_pending_session_gracefully_closed();
968                }
969                self.closed_sessions_metrics.incoming_pending.increment(1);
970                self.metrics
971                    .incoming_connections
972                    .set(self.swarm.peers().num_inbound_connections() as f64);
973            }
974            SwarmEvent::OutgoingPendingSessionClosed { remote_addr, peer_id, error } => {
975                trace!(
976                    target: "net",
977                    ?remote_addr,
978                    ?peer_id,
979                    ?error,
980                    "Outgoing pending session failed"
981                );
982
983                if let Some(ref err) = error {
984                    self.swarm.peers_mut().on_outgoing_pending_session_dropped(
985                        &remote_addr,
986                        &peer_id,
987                        err,
988                    );
989                    self.pending_session_failure_metrics.outbound.increment(1);
990                    self.backed_off_peers_metrics.increment_for_reason(
991                        BackoffReason::from_disconnect(err.as_disconnected()),
992                    );
993                    if let Some(reason) = err.as_disconnected() {
994                        self.disconnect_metrics.increment_outbound(reason);
995                    }
996                } else {
997                    self.swarm
998                        .state_mut()
999                        .peers_mut()
1000                        .on_outgoing_pending_session_gracefully_closed(&peer_id);
1001                }
1002                self.closed_sessions_metrics.outgoing_pending.increment(1);
1003                self.update_pending_connection_metrics();
1004                self.metrics.backed_off_peers.set(self.swarm.peers().num_backed_off_peers() as f64);
1005            }
1006            SwarmEvent::OutgoingConnectionError { remote_addr, peer_id, error } => {
1007                trace!(
1008                    target: "net",
1009                    ?remote_addr,
1010                    ?peer_id,
1011                    %error,
1012                    "Outgoing connection error"
1013                );
1014
1015                self.swarm.peers_mut().on_outgoing_connection_failure(
1016                    &remote_addr,
1017                    &peer_id,
1018                    &error,
1019                );
1020
1021                self.backed_off_peers_metrics.increment_for_reason(BackoffReason::ConnectionError);
1022                self.metrics.backed_off_peers.set(self.swarm.peers().num_backed_off_peers() as f64);
1023                self.update_pending_connection_metrics();
1024            }
1025            SwarmEvent::BadMessage { peer_id } => {
1026                self.swarm
1027                    .state_mut()
1028                    .peers_mut()
1029                    .apply_reputation_change(&peer_id, ReputationChangeKind::BadMessage);
1030                self.metrics.invalid_messages_received.increment(1);
1031            }
1032            SwarmEvent::ProtocolBreach { peer_id } => {
1033                self.swarm
1034                    .state_mut()
1035                    .peers_mut()
1036                    .apply_reputation_change(&peer_id, ReputationChangeKind::BadProtocol);
1037            }
1038        }
1039    }
1040
1041    /// Returns [`PeerInfo`] for all connected peers
1042    fn get_peer_infos(&self) -> Vec<PeerInfo> {
1043        self.swarm
1044            .sessions()
1045            .active_sessions()
1046            .iter()
1047            .filter_map(|(&peer_id, session)| {
1048                self.swarm
1049                    .state()
1050                    .peers()
1051                    .peer_by_id(peer_id)
1052                    .map(|(record, kind)| session.peer_info(&record, kind))
1053            })
1054            .collect()
1055    }
1056
1057    /// Returns [`PeerInfo`] for a given peer.
1058    ///
1059    /// Returns `None` if there's no active session to the peer.
1060    fn get_peer_info_by_id(&self, peer_id: PeerId) -> Option<PeerInfo> {
1061        self.swarm.sessions().active_sessions().get(&peer_id).and_then(|session| {
1062            self.swarm
1063                .state()
1064                .peers()
1065                .peer_by_id(peer_id)
1066                .map(|(record, kind)| session.peer_info(&record, kind))
1067        })
1068    }
1069
1070    /// Returns [`PeerInfo`] for a given peers.
1071    ///
1072    /// Ignore the non-active peer.
1073    fn get_peer_infos_by_ids(&self, peer_ids: impl IntoIterator<Item = PeerId>) -> Vec<PeerInfo> {
1074        peer_ids.into_iter().filter_map(|peer_id| self.get_peer_info_by_id(peer_id)).collect()
1075    }
1076
1077    /// Updates the metrics for active,established connections
1078    #[inline]
1079    fn update_active_connection_metrics(&self) {
1080        self.metrics.incoming_connections.set(self.swarm.peers().num_inbound_connections() as f64);
1081        self.metrics.outgoing_connections.set(self.swarm.peers().num_outbound_connections() as f64);
1082    }
1083
1084    /// Updates the metrics for pending connections
1085    #[inline]
1086    fn update_pending_connection_metrics(&self) {
1087        self.metrics
1088            .pending_outgoing_connections
1089            .set(self.swarm.peers().num_pending_outbound_connections() as f64);
1090        self.metrics
1091            .total_pending_connections
1092            .set(self.swarm.sessions().num_pending_connections() as f64);
1093    }
1094
1095    /// Drives the [`NetworkManager`] future until a [`GracefulShutdown`] signal is received.
1096    ///
1097    /// This invokes the given function `shutdown_hook` while holding the graceful shutdown guard.
1098    pub async fn run_until_graceful_shutdown<F, R>(
1099        mut self,
1100        shutdown: GracefulShutdown,
1101        shutdown_hook: F,
1102    ) -> R
1103    where
1104        F: FnOnce(Self) -> R,
1105    {
1106        let mut graceful_guard = None;
1107        tokio::select! {
1108            _ = &mut self => {},
1109            guard = shutdown => {
1110                graceful_guard = Some(guard);
1111            },
1112        }
1113
1114        self.perform_network_shutdown();
1115        let res = shutdown_hook(self);
1116        drop(graceful_guard);
1117        res
1118    }
1119
1120    /// Performs a graceful network shutdown by stopping new connections from being accepted while
1121    /// draining current and pending connections.
1122    fn perform_network_shutdown(&mut self) {
1123        // Set connection status to `Shutdown`. Stops node from accepting
1124        // new incoming connections as well as sending connection requests to newly
1125        // discovered nodes.
1126        self.swarm.on_shutdown_requested();
1127        // Disconnect all active connections
1128        self.swarm.sessions_mut().disconnect_all(Some(DisconnectReason::ClientQuitting));
1129        // drop pending connections
1130        self.swarm.sessions_mut().disconnect_all_pending();
1131    }
1132}
1133
1134impl<N: NetworkPrimitives> Future for NetworkManager<N> {
1135    type Output = ();
1136
1137    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
1138        let start = Instant::now();
1139        let mut poll_durations = NetworkManagerPollDurations::default();
1140
1141        let this = self.get_mut();
1142
1143        // poll new block imports (expected to be a noop for POS)
1144        while let Poll::Ready(outcome) = this.block_import.poll(cx) {
1145            this.on_block_import_result(outcome);
1146        }
1147
1148        // These loops drive the entire state of network and does a lot of work. Under heavy load
1149        // (many messages/events), data may arrive faster than it can be processed (incoming
1150        // messages/requests -> events), and it is possible that more data has already arrived by
1151        // the time an internal event is processed. Which could turn this loop into a busy loop.
1152        // Without yielding back to the executor, it can starve other tasks waiting on that
1153        // executor to execute them, or drive underlying resources To prevent this, we
1154        // preemptively return control when the `budget` is exhausted. The value itself is chosen
1155        // somewhat arbitrarily, it is high enough so the swarm can make meaningful progress but
1156        // low enough that this loop does not starve other tasks for too long. If the budget is
1157        // exhausted we manually yield back control to the (coop) scheduler. This manual yield
1158        // point should prevent situations where polling appears to be frozen. See also
1159        // <https://tokio.rs/blog/2020-04-preemption> And tokio's docs on cooperative scheduling
1160        // <https://docs.rs/tokio/latest/tokio/task/#cooperative-scheduling>
1161        //
1162        // Testing has shown that this loop naturally reaches the pending state within 1-5
1163        // iterations in << 100µs in most cases. On average it requires ~50µs, which is inside the
1164        // range of what's recommended as rule of thumb.
1165        // <https://ryhl.io/blog/async-what-is-blocking/>
1166
1167        // process incoming messages from a handle (`TransactionsManager` has one)
1168        //
1169        // will only be closed if the channel was deliberately closed since we always have an
1170        // instance of `NetworkHandle`
1171        let start_network_handle = Instant::now();
1172        let maybe_more_handle_messages = poll_nested_stream_with_budget!(
1173            "net",
1174            "Network message channel",
1175            DEFAULT_BUDGET_TRY_DRAIN_NETWORK_HANDLE_CHANNEL,
1176            this.from_handle_rx.poll_next_unpin(cx),
1177            |msg| this.on_handle_message(msg),
1178            error!("Network channel closed");
1179        );
1180        poll_durations.acc_network_handle = start_network_handle.elapsed();
1181
1182        // process incoming messages from the network
1183        let maybe_more_swarm_events = poll_nested_stream_with_budget!(
1184            "net",
1185            "Swarm events stream",
1186            DEFAULT_BUDGET_TRY_DRAIN_SWARM,
1187            this.swarm.poll_next_unpin(cx),
1188            |event| this.on_swarm_event(event),
1189        );
1190        poll_durations.acc_swarm =
1191            start_network_handle.elapsed() - poll_durations.acc_network_handle;
1192
1193        // all streams are fully drained and import futures pending
1194        if maybe_more_handle_messages || maybe_more_swarm_events {
1195            // make sure we're woken up again
1196            cx.waker().wake_by_ref();
1197            return Poll::Pending
1198        }
1199
1200        this.update_poll_metrics(start, poll_durations);
1201
1202        Poll::Pending
1203    }
1204}
1205
1206#[derive(Debug, Default)]
1207struct NetworkManagerPollDurations {
1208    acc_network_handle: Duration,
1209    acc_swarm: Duration,
1210}