Skip to main content

reth_network/
peers.rs

1//! Peer related implementations
2
3use crate::{
4    error::SessionError,
5    session::{Direction, PendingSessionHandshakeError},
6    swarm::NetworkConnectionState,
7    trusted_peers_resolver::TrustedPeersResolver,
8};
9use alloy_primitives::map::{hash_map::Entry, FbBuildHasher, HashMap, HashSet};
10use futures::StreamExt;
11
12use rand::Rng;
13use reth_eth_wire::{errors::EthStreamError, DisconnectReason};
14use reth_ethereum_forks::ForkId;
15use reth_net_banlist::BanList;
16use reth_network_api::test_utils::{PeerCommand, PeersHandle};
17use reth_network_peers::{NodeRecord, PeerId, TrustedPeer};
18use reth_network_types::{
19    is_connection_failed_reputation,
20    peers::{
21        config::{PeerBackoffDurations, PEER_ROTATION_MIN_UPTIME},
22        reputation::{DEFAULT_REPUTATION, MAX_TRUSTED_PEER_REPUTATION_CHANGE},
23    },
24    ConnectionsConfig, Peer, PeerAddr, PeerConnectionState, PeerKind, PeersConfig,
25    PersistedPeerInfo, ReputationChangeKind, ReputationChangeOutcome, ReputationChangeWeights,
26};
27use std::{
28    collections::VecDeque,
29    fmt::Display,
30    io,
31    net::{IpAddr, SocketAddr},
32    pin::Pin,
33    task::{Context, Poll},
34    time::Duration,
35};
36use thiserror::Error;
37use tokio::{
38    sync::mpsc,
39    time::{Instant, Interval, Sleep},
40};
41use tokio_stream::wrappers::UnboundedReceiverStream;
42use tracing::{trace, warn};
43
44/// Maintains the state of _all_ the peers known to the network.
45///
46/// This is supposed to be owned by the network itself, but can be reached via the [`PeersHandle`].
47/// From this type, connections to peers are established or disconnected, see [`PeerAction`].
48///
49/// The [`PeersManager`] will be notified on peer related changes
50#[derive(Debug)]
51pub struct PeersManager {
52    /// All peers known to the network
53    peers: HashMap<PeerId, Peer, FbBuildHasher<64>>,
54    /// The set of trusted peer ids.
55    ///
56    /// This tracks peer ids that are considered trusted, but for which we don't necessarily have
57    /// an address: [`Self::add_trusted_peer_id`]
58    trusted_peer_ids: HashSet<PeerId, FbBuildHasher<64>>,
59    /// A resolver used to periodically resolve DNS names for trusted peers. This updates the
60    /// peer's address when the DNS records change.
61    trusted_peers_resolver: TrustedPeersResolver,
62    /// Copy of the sender half, so new [`PeersHandle`] can be created on demand.
63    manager_tx: mpsc::UnboundedSender<PeerCommand>,
64    /// Receiver half of the command channel.
65    handle_rx: UnboundedReceiverStream<PeerCommand>,
66    /// Buffered actions until the manager is polled.
67    queued_actions: VecDeque<PeerAction>,
68    /// Interval for triggering connections if there are free slots.
69    refill_slots_interval: Interval,
70    /// How to weigh reputation changes
71    reputation_weights: ReputationChangeWeights,
72    /// Tracks current slot stats.
73    connection_info: ConnectionInfo,
74    /// Tracks unwanted ips/peer ids.
75    ban_list: BanList,
76    /// Tracks currently backed off peers.
77    backed_off_peers: HashMap<PeerId, std::time::Instant, FbBuildHasher<64>>,
78    /// Interval at which to check for peers to unban and release from the backoff map.
79    release_interval: Interval,
80    /// How long to ban bad peers.
81    ban_duration: Duration,
82    /// How long peers to which we could not connect for non-fatal reasons, e.g.
83    /// [`DisconnectReason::TooManyPeers`], are put in time out.
84    backoff_durations: PeerBackoffDurations,
85    /// If non-trusted peers should be connected to, or the connection from non-trusted
86    /// incoming peers should be accepted.
87    trusted_nodes_only: bool,
88    /// Timestamp of the last time [`Self::tick`] was called.
89    last_tick: Instant,
90    /// Maximum number of backoff attempts before we give up on a peer and dropping.
91    max_backoff_count: u8,
92    /// Tracks the connection state of the node
93    net_connection_state: NetworkConnectionState,
94    /// How long to temporarily ban ip on an incoming connection attempt.
95    incoming_ip_throttle_duration: Duration,
96    /// IP address filter for restricting network connections to specific IP ranges.
97    ip_filter: reth_net_banlist::IpFilter,
98    /// If true, discovered peers without a confirmed ENR fork ID will not be added until their
99    /// fork ID is verified via EIP-868.
100    enforce_enr_fork_id: bool,
101    /// One-shot sleep that fires when it's time to rotate a peer; reset with jitter after each
102    /// fire. `None` when rotation is disabled.
103    peer_rotation_sleep: Option<Pin<Box<Sleep>>>,
104    /// Mean duration for computing jittered rotation intervals. `None` when rotation is disabled.
105    peer_rotation_mean: Option<Duration>,
106}
107
108impl PeersManager {
109    /// Create a new instance with the given config
110    pub fn new(config: PeersConfig) -> Self {
111        let PeersConfig {
112            refill_slots_interval,
113            connection_info,
114            reputation_weights,
115            ban_list,
116            ban_duration,
117            backoff_durations,
118            trusted_nodes,
119            trusted_nodes_only,
120            trusted_nodes_resolution_interval,
121            basic_nodes,
122            persisted_peers,
123            max_backoff_count,
124            incoming_ip_throttle_duration,
125            ip_filter,
126            enforce_enr_fork_id,
127            peer_rotation_interval,
128        } = config;
129        let (manager_tx, handle_rx) = mpsc::unbounded_channel();
130        let now = Instant::now();
131
132        // We use half of the interval to decrease the max duration to `150%` in worst case
133        let unban_interval = ban_duration.min(backoff_durations.low) / 2;
134
135        let mut peers: HashMap<PeerId, Peer, FbBuildHasher<64>> = HashMap::with_capacity_and_hasher(
136            trusted_nodes.len() + basic_nodes.len() + persisted_peers.len(),
137            Default::default(),
138        );
139        let mut trusted_peer_ids: HashSet<PeerId, FbBuildHasher<64>> =
140            HashSet::with_capacity_and_hasher(trusted_nodes.len(), Default::default());
141
142        for trusted_peer in &trusted_nodes {
143            match trusted_peer.resolve_blocking() {
144                Ok(NodeRecord { address, tcp_port, udp_port, id }) => {
145                    trusted_peer_ids.insert(id);
146                    peers.entry(id).or_insert_with(|| {
147                        Peer::trusted(PeerAddr::new_with_ports(address, tcp_port, Some(udp_port)))
148                    });
149                }
150                Err(err) => {
151                    warn!(target: "net::peers", ?err, "Failed to resolve trusted peer");
152                }
153            }
154        }
155
156        for PersistedPeerInfo { record, kind, fork_id, reputation } in persisted_peers {
157            // When enforce_enr_fork_id is enabled, skip persisted peers that don't have a
158            // confirmed fork ID. These were likely accumulated from a different network during
159            // a prior run without the flag.
160            if enforce_enr_fork_id && fork_id.is_none() {
161                continue
162            }
163            let NodeRecord { address, tcp_port, udp_port, id } = record;
164            peers.entry(id).or_insert_with(|| {
165                let mut peer = Peer::with_kind(
166                    PeerAddr::new_with_ports(address, tcp_port, Some(udp_port)),
167                    kind,
168                );
169                peer.fork_id = fork_id.map(Box::new);
170                peer.reputation = reputation;
171                peer
172            });
173        }
174
175        for NodeRecord { address, tcp_port, udp_port, id } in basic_nodes {
176            peers.entry(id).or_insert_with(|| {
177                Peer::new(PeerAddr::new_with_ports(address, tcp_port, Some(udp_port)))
178            });
179        }
180
181        trace!(target: "net::peers", trusted_peers=?trusted_peer_ids, "Initialized peers manager");
182
183        Self {
184            peers,
185            trusted_peer_ids,
186            trusted_peers_resolver: TrustedPeersResolver::new(
187                trusted_nodes,
188                tokio::time::interval(trusted_nodes_resolution_interval), // 1 hour
189            ),
190            manager_tx,
191            handle_rx: UnboundedReceiverStream::new(handle_rx),
192            queued_actions: Default::default(),
193            reputation_weights,
194            refill_slots_interval: tokio::time::interval(refill_slots_interval),
195            release_interval: tokio::time::interval_at(now + unban_interval, unban_interval),
196            connection_info: ConnectionInfo::new(connection_info),
197            ban_list,
198            backed_off_peers: Default::default(),
199            ban_duration,
200            backoff_durations,
201            trusted_nodes_only,
202            last_tick: Instant::now(),
203            max_backoff_count,
204            net_connection_state: NetworkConnectionState::default(),
205            incoming_ip_throttle_duration,
206            ip_filter,
207            enforce_enr_fork_id,
208            peer_rotation_sleep: peer_rotation_interval
209                .map(|mean| Box::pin(tokio::time::sleep(jitter_rotation_interval(mean)))),
210            peer_rotation_mean: peer_rotation_interval,
211        }
212    }
213
214    /// Returns a new [`PeersHandle`] that can send commands to this type.
215    pub(crate) fn handle(&self) -> PeersHandle {
216        PeersHandle::new(self.manager_tx.clone())
217    }
218
219    /// Returns `true` if discovered peers must have a confirmed ENR fork ID before being added.
220    pub(crate) const fn enforce_enr_fork_id(&self) -> bool {
221        self.enforce_enr_fork_id
222    }
223
224    /// Returns the number of peers in the peer set
225    #[inline]
226    pub(crate) fn num_known_peers(&self) -> usize {
227        self.peers.len()
228    }
229
230    /// Returns an iterator over all peers as [`NodeRecord`]s.
231    pub(crate) fn iter_peers(&self) -> impl Iterator<Item = NodeRecord> + '_ {
232        self.peers.iter().map(|(peer_id, v)| {
233            NodeRecord::new_with_ports(
234                v.addr.tcp().ip(),
235                v.addr.tcp().port(),
236                v.addr.udp().map(|addr| addr.port()),
237                *peer_id,
238            )
239        })
240    }
241
242    /// Returns an iterator over peers suitable for persisting to disk.
243    ///
244    /// Filters out backed-off and banned peers, and includes metadata like kind, fork ID, and
245    /// reputation.
246    pub(crate) fn persistable_peers(&self) -> impl Iterator<Item = PersistedPeerInfo> + '_ {
247        self.peers.iter().filter(|(_, peer)| !peer.is_backed_off() && !peer.is_banned()).map(
248            |(peer_id, peer)| PersistedPeerInfo {
249                record: NodeRecord::new_with_ports(
250                    peer.addr.tcp().ip(),
251                    peer.addr.tcp().port(),
252                    peer.addr.udp().map(|addr| addr.port()),
253                    *peer_id,
254                ),
255                kind: peer.kind,
256                fork_id: peer.fork_id.as_deref().copied(),
257                reputation: peer.reputation,
258            },
259        )
260    }
261
262    /// Returns the `NodeRecord` and `PeerKind` for the given peer id
263    pub(crate) fn peer_by_id(&self, peer_id: PeerId) -> Option<(NodeRecord, PeerKind)> {
264        self.peers.get(&peer_id).map(|v| {
265            (
266                NodeRecord::new_with_ports(
267                    v.addr.tcp().ip(),
268                    v.addr.tcp().port(),
269                    v.addr.udp().map(|addr| addr.port()),
270                    peer_id,
271                ),
272                v.kind,
273            )
274        })
275    }
276
277    /// Returns `true` if the given peer is connected via an inbound session.
278    pub(crate) fn is_inbound_peer(&self, peer_id: &PeerId) -> bool {
279        self.peers.get(peer_id).is_some_and(|p| {
280            matches!(p.state, PeerConnectionState::In | PeerConnectionState::DisconnectingIn)
281        })
282    }
283
284    /// Returns an iterator over all peer ids for peers with the given kind
285    pub(crate) fn peers_by_kind(&self, kind: PeerKind) -> impl Iterator<Item = PeerId> + '_ {
286        self.peers.iter().filter_map(move |(peer_id, peer)| (peer.kind == kind).then_some(*peer_id))
287    }
288
289    /// Returns the number of currently active inbound connections.
290    #[inline]
291    pub(crate) const fn num_inbound_connections(&self) -> usize {
292        self.connection_info.num_inbound
293    }
294
295    /// Returns the number of currently __active__ outbound connections.
296    #[inline]
297    pub(crate) const fn num_outbound_connections(&self) -> usize {
298        self.connection_info.num_outbound
299    }
300
301    /// Returns the number of currently pending outbound connections.
302    #[inline]
303    pub(crate) const fn num_pending_outbound_connections(&self) -> usize {
304        self.connection_info.num_pending_out
305    }
306
307    /// Returns the number of currently backed off peers.
308    #[inline]
309    pub(crate) fn num_backed_off_peers(&self) -> usize {
310        self.backed_off_peers.len()
311    }
312
313    /// Returns the number of idle trusted peers.
314    fn num_idle_trusted_peers(&self) -> usize {
315        self.peers.iter().filter(|(_, peer)| peer.kind.is_trusted() && peer.state.is_idle()).count()
316    }
317
318    /// Invoked when a new _incoming_ tcp connection is accepted.
319    ///
320    /// returns an error if the inbound ip address is on the ban list
321    pub(crate) fn on_incoming_pending_session(
322        &mut self,
323        addr: IpAddr,
324    ) -> Result<(), InboundConnectionError> {
325        // Check if the IP is in the allowed ranges (netrestrict)
326        if !self.ip_filter.is_allowed(&addr) {
327            trace!(target: "net", ?addr, "Rejecting connection from IP not in allowed ranges");
328            return Err(InboundConnectionError::IpBanned)
329        }
330
331        if self.ban_list.is_banned_ip(&addr) {
332            return Err(InboundConnectionError::IpBanned)
333        }
334
335        // check if we even have slots for a new incoming connection
336        if !self.connection_info.has_in_capacity() {
337            if self.trusted_peer_ids.is_empty() {
338                // if we don't have any incoming slots and no trusted peers, we don't accept any new
339                // connections
340                return Err(InboundConnectionError::ExceedsCapacity)
341            }
342
343            // there's an edge case here where no incoming connections besides from trusted peers
344            // are allowed (max_inbound == 0), in which case we still need to allow new pending
345            // incoming connections until all trusted peers are connected.
346            let num_idle_trusted_peers = self.num_idle_trusted_peers();
347            if num_idle_trusted_peers <= self.trusted_peer_ids.len() {
348                // we still want to limit concurrent pending connections
349                let max_inbound =
350                    self.trusted_peer_ids.len().max(self.connection_info.config.max_inbound);
351                if self.connection_info.num_pending_in < max_inbound {
352                    self.connection_info.inc_pending_in();
353                    return Ok(())
354                }
355            }
356
357            // all trusted peers are either connected or connecting
358            return Err(InboundConnectionError::ExceedsCapacity)
359        }
360
361        // also cap the incoming connections we can process at once
362        if !self.connection_info.has_in_pending_capacity() {
363            return Err(InboundConnectionError::ExceedsCapacity)
364        }
365
366        // apply the rate limit
367        self.throttle_incoming_ip(addr);
368
369        self.connection_info.inc_pending_in();
370        Ok(())
371    }
372
373    /// Invoked when a previous call to [`Self::on_incoming_pending_session`] succeeded but it was
374    /// rejected.
375    pub(crate) const fn on_incoming_pending_session_rejected_internally(&mut self) {
376        self.connection_info.decr_pending_in();
377    }
378
379    /// Invoked when a pending session was closed.
380    pub(crate) const fn on_incoming_pending_session_gracefully_closed(&mut self) {
381        self.connection_info.decr_pending_in()
382    }
383
384    /// Invoked when a pending session was closed.
385    pub(crate) fn on_incoming_pending_session_dropped(
386        &mut self,
387        remote_addr: SocketAddr,
388        err: &PendingSessionHandshakeError,
389    ) {
390        if err.is_fatal_protocol_error() {
391            self.ban_ip(remote_addr.ip());
392
393            if err.merits_discovery_ban() {
394                self.queued_actions
395                    .push_back(PeerAction::DiscoveryBanIp { ip_addr: remote_addr.ip() })
396            }
397        }
398
399        self.connection_info.decr_pending_in();
400    }
401
402    /// Called when a new _incoming_ active session was established to the given peer.
403    ///
404    /// This will update the state of the peer if not yet tracked.
405    ///
406    /// If the reputation of the peer is below the `BANNED_REPUTATION` threshold, a disconnect will
407    /// be scheduled.
408    pub(crate) fn on_incoming_session_established(&mut self, peer_id: PeerId, addr: SocketAddr) {
409        self.connection_info.decr_pending_in();
410
411        // we only need to check the peer id here as the ip address will have been checked at
412        // on_incoming_pending_session. We also check if the peer is in the backoff list here.
413        if self.ban_list.is_banned_peer(&peer_id) {
414            self.queued_actions.push_back(PeerAction::DisconnectBannedIncoming { peer_id });
415            return
416        }
417
418        // check if the peer is trustable or not
419        let mut is_trusted = self.trusted_peer_ids.contains(&peer_id);
420        if self.trusted_nodes_only && !is_trusted {
421            self.queued_actions.push_back(PeerAction::DisconnectUntrustedIncoming { peer_id });
422            return
423        }
424
425        // start a new tick, so the peer is not immediately rewarded for the time since last tick
426        self.tick();
427
428        match self.peers.entry(peer_id) {
429            Entry::Occupied(mut entry) => {
430                let peer = entry.get_mut();
431                if peer.is_banned() {
432                    self.queued_actions.push_back(PeerAction::DisconnectBannedIncoming { peer_id });
433                    return
434                }
435                // it might be the case that we're also trying to connect to this peer at the same
436                // time, so we need to adjust the state here
437                if peer.state.is_pending_out() {
438                    self.connection_info.decr_state(peer.state);
439                }
440
441                peer.state = PeerConnectionState::In;
442                peer.mark_connected();
443
444                is_trusted = is_trusted || peer.is_trusted();
445            }
446            Entry::Vacant(entry) => {
447                // peer is missing in the table, we add it but mark it as to be removed after
448                // disconnect, because we only know the outgoing port
449                let mut peer = Peer::with_state(PeerAddr::from_tcp(addr), PeerConnectionState::In);
450                peer.mark_connected();
451                peer.remove_after_disconnect = true;
452                entry.insert(peer);
453                self.queued_actions.push_back(PeerAction::PeerAdded(peer_id));
454            }
455        }
456
457        let has_in_capacity = self.connection_info.has_in_capacity();
458        // increment new incoming connection
459        self.connection_info.inc_in();
460
461        // disconnect the peer if we don't have capacity for more inbound connections
462        if !is_trusted && !has_in_capacity {
463            self.queued_actions.push_back(PeerAction::Disconnect {
464                peer_id,
465                reason: Some(DisconnectReason::TooManyPeers),
466            });
467        }
468    }
469
470    /// Bans the peer temporarily with the configured ban timeout
471    fn ban_peer(&mut self, peer_id: PeerId) {
472        let ban_duration = if let Some(peer) = self.peers.get(&peer_id) &&
473            (peer.is_trusted() || peer.is_static())
474        {
475            // For misbehaving trusted or static peers, we provide a bit more leeway when
476            // penalizing them.
477            self.backoff_durations.low / 2
478        } else {
479            self.ban_duration
480        };
481
482        self.ban_list.ban_peer_until(peer_id, std::time::Instant::now() + ban_duration);
483        self.queued_actions.push_back(PeerAction::BanPeer { peer_id });
484    }
485
486    /// Bans the IP temporarily with the configured ban timeout
487    fn ban_ip(&mut self, ip: IpAddr) {
488        self.ban_list.ban_ip_until(ip, std::time::Instant::now() + self.ban_duration);
489    }
490
491    /// Bans the IP temporarily to rate limit inbound connection attempts per IP.
492    fn throttle_incoming_ip(&mut self, ip: IpAddr) {
493        self.ban_list
494            .ban_ip_until(ip, std::time::Instant::now() + self.incoming_ip_throttle_duration);
495    }
496
497    /// Temporarily puts the peer in timeout by inserting it into the backedoff peers set
498    fn backoff_peer_until(&mut self, peer_id: PeerId, until: std::time::Instant) {
499        trace!(target: "net::peers", ?peer_id, "backing off");
500
501        if let Some(peer) = self.peers.get_mut(&peer_id) {
502            peer.backed_off = true;
503            self.backed_off_peers.insert(peer_id, until);
504        }
505    }
506
507    /// Unbans the peer
508    fn unban_peer(&mut self, peer_id: PeerId) {
509        self.ban_list.unban_peer(&peer_id);
510        self.queued_actions.push_back(PeerAction::UnBanPeer { peer_id });
511    }
512
513    /// Tick function to update reputation of all connected peers.
514    /// Peers are rewarded with reputation increases for the time they are connected since the last
515    /// tick. This is to prevent peers from being disconnected eventually due to slashed
516    /// reputation because of some bad messages (most likely transaction related)
517    fn tick(&mut self) {
518        let now = Instant::now();
519        // Determine the number of seconds since the last tick.
520        // Ensuring that now is always greater than last_tick to account for issues with system
521        // time.
522        let secs_since_last_tick =
523            if self.last_tick > now { 0 } else { (now - self.last_tick).as_secs() as i32 };
524        self.last_tick = now;
525
526        // update reputation via seconds connected
527        for peer in self.peers.iter_mut().filter(|(_, peer)| peer.state.is_connected()) {
528            // update reputation via seconds connected, but keep the target _around_ the default
529            // reputation.
530            if peer.1.reputation < DEFAULT_REPUTATION {
531                peer.1.reputation += secs_since_last_tick;
532            }
533        }
534    }
535
536    /// Returns the tracked reputation for a peer.
537    pub(crate) fn get_reputation(&self, peer_id: &PeerId) -> Option<i32> {
538        self.peers.get(peer_id).map(|peer| peer.reputation)
539    }
540
541    /// Apply the corresponding reputation change to the given peer.
542    ///
543    /// If the peer is a trusted peer, it will be exempt from reputation slashing for certain
544    /// reputation changes that can be attributed to network conditions. If the peer is a
545    /// trusted peer, it will also be less strict with the reputation slashing.
546    pub(crate) fn apply_reputation_change(&mut self, peer_id: &PeerId, rep: ReputationChangeKind) {
547        trace!(target: "net::peers", ?peer_id, reputation=?rep, "applying reputation change");
548
549        let reputation_change = if rep.is_reset() {
550            None
551        } else {
552            let reputation_change = self.reputation_weights.change(rep).as_i32();
553            if reputation_change == 0 {
554                return
555            }
556            Some(reputation_change)
557        };
558
559        let outcome = if let Some(peer) = self.peers.get_mut(peer_id) {
560            if let Some(mut reputation_change) = reputation_change {
561                if peer.is_trusted() || peer.is_static() {
562                    // exempt trusted and static peers from reputation slashing for
563                    if matches!(
564                        rep,
565                        ReputationChangeKind::Dropped |
566                            ReputationChangeKind::BadAnnouncement |
567                            ReputationChangeKind::Timeout |
568                            ReputationChangeKind::AlreadySeenTransaction
569                    ) {
570                        return
571                    }
572
573                    // also be less strict with the reputation slashing for trusted peers
574                    if reputation_change < MAX_TRUSTED_PEER_REPUTATION_CHANGE {
575                        // this caps the reputation change to the maximum allowed for trusted peers
576                        reputation_change = MAX_TRUSTED_PEER_REPUTATION_CHANGE;
577                    }
578                }
579                peer.apply_reputation(reputation_change, rep)
580            } else {
581                peer.reset_reputation()
582            }
583        } else {
584            return
585        };
586
587        match outcome {
588            ReputationChangeOutcome::None => {}
589            ReputationChangeOutcome::Ban => {
590                self.ban_peer(*peer_id);
591            }
592            ReputationChangeOutcome::Unban => self.unban_peer(*peer_id),
593            ReputationChangeOutcome::DisconnectAndBan => {
594                self.queued_actions.push_back(PeerAction::Disconnect {
595                    peer_id: *peer_id,
596                    reason: Some(DisconnectReason::DisconnectRequested),
597                });
598                self.ban_peer(*peer_id);
599            }
600        }
601    }
602
603    /// Gracefully disconnected a pending _outgoing_ session
604    pub(crate) fn on_outgoing_pending_session_gracefully_closed(&mut self, peer_id: &PeerId) {
605        if let Some(peer) = self.peers.get_mut(peer_id) {
606            self.connection_info.decr_state(peer.state);
607            peer.state = PeerConnectionState::Idle;
608            peer.mark_disconnected();
609        }
610    }
611
612    /// Invoked when an _outgoing_ pending session was closed during authentication or the
613    /// handshake.
614    pub(crate) fn on_outgoing_pending_session_dropped(
615        &mut self,
616        remote_addr: &SocketAddr,
617        peer_id: &PeerId,
618        err: &PendingSessionHandshakeError,
619    ) {
620        self.on_connection_failure(remote_addr, peer_id, err, ReputationChangeKind::FailedToConnect)
621    }
622
623    /// Gracefully disconnected an active session
624    pub(crate) fn on_active_session_gracefully_closed(&mut self, peer_id: PeerId) {
625        match self.peers.entry(peer_id) {
626            Entry::Occupied(mut entry) => {
627                trace!(target: "net::peers", ?peer_id, direction=?entry.get().state, "active session gracefully closed");
628                self.connection_info.decr_state(entry.get().state);
629
630                if entry.get().remove_after_disconnect && !entry.get().is_trusted() {
631                    // this peer should be removed from the set
632                    entry.remove();
633                    self.queued_actions.push_back(PeerAction::PeerRemoved(peer_id));
634                } else {
635                    let peer = entry.get_mut();
636                    // reset the peer's state
637                    // we reset the backoff counter since we're able to establish a successful
638                    // session to that peer
639                    peer.severe_backoff_counter = 0;
640                    peer.state = PeerConnectionState::Idle;
641                    peer.mark_disconnected();
642
643                    // but we're backing off slightly to avoid dialing the peer again right away, to
644                    // give the remote time to also properly register the closed session and clean
645                    // up and to avoid any issues with ip throttling on the remote in case this
646                    // session was terminated right away.
647                    peer.backed_off = true;
648                    self.backed_off_peers.insert(
649                        peer_id,
650                        std::time::Instant::now() + self.incoming_ip_throttle_duration,
651                    );
652                    trace!(target: "net::peers", ?peer_id, kind=?peer.kind, duration=?self.incoming_ip_throttle_duration, "backing off on gracefully closed session");
653                }
654            }
655            Entry::Vacant(_) => return,
656        }
657
658        self.fill_outbound_slots();
659    }
660
661    /// Called when a _pending_ outbound connection is successful.
662    pub(crate) fn on_active_outgoing_established(&mut self, peer_id: PeerId) {
663        if let Some(peer) = self.peers.get_mut(&peer_id) {
664            trace!(target: "net::peers", ?peer_id, "established active outgoing connection");
665            self.connection_info.decr_state(peer.state);
666            self.connection_info.inc_out();
667            peer.state = PeerConnectionState::Out;
668            peer.mark_connected();
669        }
670    }
671
672    /// Called when an _active_ session to a peer was forcefully dropped due to an error.
673    ///
674    /// Depending on whether the error is fatal, the peer will be removed from the peer set
675    /// otherwise its reputation is slashed.
676    pub(crate) fn on_active_session_dropped(
677        &mut self,
678        remote_addr: &SocketAddr,
679        peer_id: &PeerId,
680        err: &EthStreamError,
681    ) {
682        self.on_connection_failure(remote_addr, peer_id, err, ReputationChangeKind::Dropped)
683    }
684
685    /// Called when an attempt to create an _outgoing_ pending session failed while setting up a tcp
686    /// connection.
687    pub(crate) fn on_outgoing_connection_failure(
688        &mut self,
689        remote_addr: &SocketAddr,
690        peer_id: &PeerId,
691        err: &io::Error,
692    ) {
693        // there's a race condition where we accepted an incoming connection while we were trying to
694        // connect to the same peer at the same time. if the outgoing connection failed
695        // after the incoming connection was accepted, we can ignore this error
696        if let Some(peer) = self.peers.get(peer_id) {
697            if peer.state.is_incoming() {
698                // we already have an active connection to the peer, so we can ignore this error
699                return
700            }
701
702            if peer.is_trusted() && is_connection_failed_reputation(peer.reputation) {
703                // trigger resolution task for trusted peer since multiple connection failures
704                // occurred
705                self.trusted_peers_resolver.interval.reset_immediately();
706            }
707        }
708
709        self.on_connection_failure(remote_addr, peer_id, err, ReputationChangeKind::FailedToConnect)
710    }
711
712    fn on_connection_failure(
713        &mut self,
714        remote_addr: &SocketAddr,
715        peer_id: &PeerId,
716        err: impl SessionError,
717        reputation_change: ReputationChangeKind,
718    ) {
719        trace!(target: "net::peers", ?remote_addr, ?peer_id, %err, "handling failed connection");
720
721        if err.is_fatal_protocol_error() {
722            trace!(target: "net::peers", ?remote_addr, ?peer_id, %err, "fatal connection error");
723            // remove the peer to which we can't establish a connection due to protocol related
724            // issues.
725            if let Entry::Occupied(mut entry) = self.peers.entry(*peer_id) {
726                self.connection_info.decr_state(entry.get().state);
727                // only remove if the peer is not trusted
728                if entry.get().is_trusted() {
729                    let peer = entry.get_mut();
730                    peer.state = PeerConnectionState::Idle;
731                    peer.mark_disconnected();
732                } else {
733                    entry.remove();
734                    self.queued_actions.push_back(PeerAction::PeerRemoved(*peer_id));
735                    // If the error is caused by a peer that should be banned from discovery
736                    if err.merits_discovery_ban() {
737                        self.queued_actions.push_back(PeerAction::DiscoveryBanPeerId {
738                            peer_id: *peer_id,
739                            ip_addr: remote_addr.ip(),
740                        })
741                    }
742                }
743            }
744
745            // ban the peer
746            self.ban_peer(*peer_id);
747        } else {
748            let mut backoff_until = None;
749            let mut remove_peer = false;
750
751            if let Some(peer) = self.peers.get_mut(peer_id) {
752                if let Some(kind) = err.should_backoff() {
753                    if peer.is_trusted() || peer.is_static() {
754                        // provide a bit more leeway for trusted peers and use a lower backoff so
755                        // that we keep re-trying them after backing off shortly, but we should at
756                        // least backoff for the low duration to not violate the ip based inbound
757                        // connection throttle that peer has in place, because this peer might not
758                        // have us registered as a trusted peer.
759                        let backoff = self.backoff_durations.low;
760                        backoff_until = Some(std::time::Instant::now() + backoff);
761                        trace!(target: "net::peers", ?peer_id, ?backoff, "backing off trusted peer");
762                    } else {
763                        // Increment peer.backoff_counter
764                        if kind.is_severe() {
765                            peer.severe_backoff_counter =
766                                peer.severe_backoff_counter.saturating_add(1);
767                        }
768                        trace!(target: "net::peers", ?peer_id, ?kind, severe_backoff_counter=peer.severe_backoff_counter, "backing off basic peer");
769
770                        let backoff_time =
771                            self.backoff_durations.backoff_until(kind, peer.severe_backoff_counter);
772
773                        // The peer has signaled that it is currently unable to process any more
774                        // connections, so we will hold off on attempting any new connections for a
775                        // while
776                        backoff_until = Some(backoff_time);
777                    }
778                } else {
779                    // If the error was not a backoff error, we reduce the peer's reputation
780                    let reputation_change = self.reputation_weights.change(reputation_change);
781                    peer.reputation = peer.reputation.saturating_add(reputation_change.as_i32());
782                };
783
784                self.connection_info.decr_state(peer.state);
785                peer.state = PeerConnectionState::Idle;
786                peer.mark_disconnected();
787
788                if peer.severe_backoff_counter > self.max_backoff_count &&
789                    !peer.is_trusted() &&
790                    !peer.is_static()
791                {
792                    // mark peer for removal if it has been backoff too many times and is _not_
793                    // trusted or static
794                    remove_peer = true;
795                }
796            }
797
798            // remove peer if it has been marked for removal
799            if remove_peer {
800                trace!(target: "net", ?peer_id, "removed peer after exceeding backoff counter");
801                let (peer_id, _) = self.peers.remove_entry(peer_id).expect("peer must exist");
802                self.queued_actions.push_back(PeerAction::PeerRemoved(peer_id));
803            } else if let Some(backoff_until) = backoff_until {
804                // otherwise, backoff the peer if marked as such
805                self.backoff_peer_until(*peer_id, backoff_until);
806            }
807        }
808
809        self.fill_outbound_slots();
810    }
811
812    /// Invoked if a pending session was disconnected because there's already a connection to the
813    /// peer.
814    ///
815    /// If the session was an outgoing connection, this means that the peer initiated a connection
816    /// to us at the same time and this connection is already established.
817    pub(crate) const fn on_already_connected(&mut self, direction: Direction) {
818        match direction {
819            Direction::Incoming => {
820                // need to decrement the ingoing counter
821                self.connection_info.decr_pending_in();
822            }
823            Direction::Outgoing(_) => {
824                // cleanup is handled when the incoming active session is activated in
825                // `on_incoming_session_established`
826            }
827        }
828    }
829
830    /// Called for a newly discovered peer.
831    ///
832    /// If the peer already exists, then the address, kind and `fork_id` will be updated.
833    pub(crate) fn add_peer(&mut self, peer_id: PeerId, addr: PeerAddr, fork_id: Option<ForkId>) {
834        self.add_peer_kind(peer_id, None, addr, fork_id)
835    }
836
837    /// Marks the given peer as trusted.
838    pub(crate) fn add_trusted_peer_id(&mut self, peer_id: PeerId) {
839        self.trusted_peer_ids.insert(peer_id);
840        if let Some(peer) = self.peers.get_mut(&peer_id) {
841            peer.kind = PeerKind::Trusted;
842        }
843    }
844
845    /// Called for a newly discovered trusted peer.
846    ///
847    /// If the peer already exists, then the address and kind will be updated.
848    #[cfg_attr(not(test), expect(dead_code))]
849    pub(crate) fn add_trusted_peer(&mut self, peer_id: PeerId, addr: PeerAddr) {
850        self.add_peer_kind(peer_id, Some(PeerKind::Trusted), addr, None)
851    }
852
853    /// Adds a trusted peer that may use a hostname instead of an IP address.
854    pub(crate) fn add_trusted_peer_node(&mut self, trusted: TrustedPeer) {
855        let peer_id = trusted.id;
856        self.trusted_peer_ids.insert(peer_id);
857        self.trusted_peers_resolver.remove(peer_id);
858        self.trusted_peers_resolver.trusted_peers.push(trusted);
859        self.trusted_peers_resolver.interval.reset_immediately();
860    }
861
862    /// Called for a newly discovered peer.
863    ///
864    /// If the peer already exists, then the address, kind and `fork_id` will be updated.
865    /// If the peer exists and a [`PeerKind`] is provided then the peer's kind is updated
866    pub(crate) fn add_peer_kind(
867        &mut self,
868        peer_id: PeerId,
869        kind: Option<PeerKind>,
870        addr: PeerAddr,
871        fork_id: Option<ForkId>,
872    ) {
873        let ip_addr = addr.tcp().ip();
874
875        // Check if the IP is in the allowed ranges (netrestrict)
876        if !self.ip_filter.is_allowed(&ip_addr) {
877            trace!(target: "net", ?peer_id, ?ip_addr, "Skipping peer from IP not in allowed ranges");
878            return
879        }
880
881        if self.ban_list.is_banned(&peer_id, &ip_addr) {
882            return
883        }
884
885        match self.peers.entry(peer_id) {
886            Entry::Occupied(mut entry) => {
887                let peer = entry.get_mut();
888                peer.fork_id = fork_id.map(Box::new);
889                peer.addr = addr;
890
891                if let Some(kind) = kind {
892                    peer.kind = kind;
893                }
894
895                if peer.state.is_incoming() {
896                    // now that we have an actual discovered address, for that peer and not just the
897                    // ip of the incoming connection, we don't need to remove the peer after
898                    // disconnecting, See `on_incoming_session_established`
899                    peer.remove_after_disconnect = false;
900                }
901            }
902            Entry::Vacant(entry) => {
903                trace!(target: "net::peers", ?peer_id, addr=?addr.tcp(), "discovered new node");
904                let mut peer = Peer::with_kind(addr, kind.unwrap_or(PeerKind::Basic));
905                peer.fork_id = fork_id.map(Box::new);
906                entry.insert(peer);
907                self.queued_actions.push_back(PeerAction::PeerAdded(peer_id));
908            }
909        }
910
911        if kind.is_some_and(|kind| kind.is_trusted()) {
912            // also track the peer in the peer id set
913            self.trusted_peer_ids.insert(peer_id);
914        }
915    }
916
917    /// Called for a peer that was explicitly requested, e.g. via `admin_addPeer`.
918    ///
919    /// Same as [`Self::add_peer_kind`], but the peer is dialed right away if there's an outbound
920    /// slot instead of on the next refill tick like discovered peers.
921    pub(crate) fn add_requested_peer(
922        &mut self,
923        peer_id: PeerId,
924        kind: Option<PeerKind>,
925        addr: PeerAddr,
926    ) {
927        self.add_peer_kind(peer_id, kind, addr, None);
928        self.fill_outbound_slots();
929    }
930
931    /// Removes the tracked node from the set.
932    pub(crate) fn remove_peer(&mut self, peer_id: PeerId) {
933        let Entry::Occupied(entry) = self.peers.entry(peer_id) else { return };
934        if entry.get().is_trusted() {
935            return
936        }
937        let mut peer = entry.remove();
938
939        trace!(target: "net::peers", ?peer_id, "remove discovered node");
940        self.queued_actions.push_back(PeerAction::PeerRemoved(peer_id));
941
942        if peer.state.is_connected() {
943            trace!(target: "net::peers", ?peer_id, "disconnecting on remove from discovery");
944            // we terminate the active session here, but only remove the peer after the session
945            // was disconnected, this prevents the case where the session is scheduled for
946            // disconnect but the node is immediately rediscovered, See also
947            // [`Self::on_disconnected()`]
948            peer.remove_after_disconnect = true;
949            peer.state.disconnect();
950            self.peers.insert(peer_id, peer);
951            self.queued_actions.push_back(PeerAction::Disconnect {
952                peer_id,
953                reason: Some(DisconnectReason::DisconnectRequested),
954            })
955        }
956    }
957
958    /// Bans the peer indefinitely and removes it from the peer set.
959    ///
960    /// This follows [`Self::remove_peer`] and does not override trusted status. Remove trusted
961    /// peers from the trusted set before banning them.
962    pub(crate) fn ban_peer_by_admin(&mut self, peer_id: PeerId) {
963        if self.trusted_peer_ids.contains(&peer_id) ||
964            self.peers.get(&peer_id).is_some_and(Peer::is_trusted)
965        {
966            return
967        }
968
969        self.remove_peer(peer_id);
970        self.ban_list.ban_peer(peer_id);
971    }
972
973    /// Removes the peer from the ban list and resets its reputation.
974    pub(crate) fn unban_peer_by_admin(&mut self, peer_id: PeerId) {
975        self.ban_list.unban_peer(&peer_id);
976        if let Some(peer) = self.peers.get_mut(&peer_id) &&
977            peer.is_banned()
978        {
979            peer.unban();
980            self.queued_actions.push_back(PeerAction::UnBanPeer { peer_id });
981        }
982    }
983
984    /// Connect to the given peer. NOTE: if the maximum number of outbound sessions is reached,
985    /// this won't do anything. See `reth_network::SessionManager::dial_outbound`.
986    #[cfg_attr(not(test), expect(dead_code))]
987    pub(crate) fn add_and_connect(
988        &mut self,
989        peer_id: PeerId,
990        addr: PeerAddr,
991        fork_id: Option<ForkId>,
992    ) {
993        self.add_and_connect_kind(peer_id, PeerKind::Basic, addr, fork_id)
994    }
995
996    /// Connects a peer and its address with the given kind.
997    ///
998    /// Note: This is invoked on demand via an external command received by the manager
999    pub(crate) fn add_and_connect_kind(
1000        &mut self,
1001        peer_id: PeerId,
1002        kind: PeerKind,
1003        addr: PeerAddr,
1004        fork_id: Option<ForkId>,
1005    ) {
1006        let ip_addr = addr.tcp().ip();
1007
1008        // Check if the IP is in the allowed ranges (netrestrict)
1009        if !self.ip_filter.is_allowed(&ip_addr) {
1010            trace!(target: "net", ?peer_id, ?ip_addr, "Skipping outbound connection to IP not in allowed ranges");
1011            return
1012        }
1013
1014        if self.ban_list.is_banned(&peer_id, &ip_addr) {
1015            return
1016        }
1017
1018        match self.peers.entry(peer_id) {
1019            Entry::Occupied(mut entry) => {
1020                let peer = entry.get_mut();
1021                peer.kind = kind;
1022                peer.fork_id = fork_id.map(Box::new);
1023                peer.addr = addr;
1024
1025                if peer.state == PeerConnectionState::Idle {
1026                    // Try connecting again.
1027                    peer.state = PeerConnectionState::PendingOut;
1028                    self.connection_info.inc_pending_out();
1029                    self.queued_actions
1030                        .push_back(PeerAction::Connect { peer_id, remote_addr: addr.tcp() });
1031                }
1032            }
1033            Entry::Vacant(entry) => {
1034                trace!(target: "net::peers", ?peer_id, addr=?addr.tcp(), "connects new node");
1035                let mut peer = Peer::with_kind(addr, kind);
1036                peer.state = PeerConnectionState::PendingOut;
1037                peer.fork_id = fork_id.map(Box::new);
1038                entry.insert(peer);
1039                self.connection_info.inc_pending_out();
1040                self.queued_actions
1041                    .push_back(PeerAction::Connect { peer_id, remote_addr: addr.tcp() });
1042            }
1043        }
1044
1045        if kind.is_trusted() {
1046            self.trusted_peer_ids.insert(peer_id);
1047        }
1048    }
1049
1050    /// Removes the tracked node from the trusted set.
1051    pub(crate) fn remove_peer_from_trusted_set(&mut self, peer_id: PeerId) {
1052        self.trusted_peers_resolver.remove(peer_id);
1053
1054        let Entry::Occupied(mut entry) = self.peers.entry(peer_id) else {
1055            self.trusted_peer_ids.remove(&peer_id);
1056            return
1057        };
1058        if !entry.get().is_trusted() {
1059            return
1060        }
1061
1062        let peer = entry.get_mut();
1063        peer.kind = PeerKind::Basic;
1064
1065        self.trusted_peer_ids.remove(&peer_id);
1066    }
1067
1068    /// Returns the best idle peer to connect to.
1069    ///
1070    /// Peers that are `trusted` or `static`, see [`PeerKind`], are prioritized as long as they're
1071    /// not currently marked as banned or backed off.
1072    ///
1073    /// Among remaining peers, the one with the highest reputation is selected. When reputation is
1074    /// equal, a peer with a discovered `fork_id` is preferred since it indicates a compatible fork.
1075    ///
1076    /// If `trusted_nodes_only` is enabled, see [`PeersConfig`], then this will only consider
1077    /// `trusted` peers.
1078    ///
1079    /// Returns `None` if no peer is available.
1080    fn best_unconnected(&mut self) -> Option<(PeerId, &mut Peer)> {
1081        let mut unconnected = self.peers.iter_mut().filter(|(_, peer)| {
1082            !peer.is_backed_off() &&
1083                !peer.is_banned() &&
1084                peer.state.is_unconnected() &&
1085                (!self.trusted_nodes_only || peer.is_trusted())
1086        });
1087
1088        // keep track of the best peer, if there's one
1089        let mut best_peer = unconnected.next()?;
1090
1091        if best_peer.1.is_trusted() || best_peer.1.is_static() {
1092            return Some((*best_peer.0, best_peer.1))
1093        }
1094
1095        for maybe_better in unconnected {
1096            // if the peer is trusted or static, return it immediately
1097            if maybe_better.1.is_trusted() || maybe_better.1.is_static() {
1098                return Some((*maybe_better.0, maybe_better.1))
1099            }
1100
1101            // prefer higher reputation, break ties by fork_id presence
1102            match maybe_better.1.reputation.cmp(&best_peer.1.reputation) {
1103                std::cmp::Ordering::Greater => best_peer = maybe_better,
1104                std::cmp::Ordering::Equal
1105                    if maybe_better.1.fork_id.is_some() && best_peer.1.fork_id.is_none() =>
1106                {
1107                    best_peer = maybe_better
1108                }
1109                _ => {}
1110            }
1111        }
1112        Some((*best_peer.0, best_peer.1))
1113    }
1114
1115    /// Disconnects one eligible peer (inbound or outbound) to open a slot for new nodes,
1116    /// mirroring Geth's peer dropper logic:
1117    /// <https://github.com/ethereum/go-ethereum/blob/c5c75977ab55e4d7ea6147cc0e221b588e5e3754/eth/dropper.go#L106-L138>
1118    ///
1119    /// A peer is eligible when its pool (inbound or outbound) is at capacity, the peer is not
1120    /// trusted or static, and it has been connected for at least [`PEER_ROTATION_MIN_UPTIME`].
1121    fn try_rotate_peer(&mut self) {
1122        let outbound_at_capacity = self.connection_info.is_outbound_at_capacity();
1123        let inbound_at_capacity = self.connection_info.is_inbound_at_capacity();
1124
1125        if !outbound_at_capacity && !inbound_at_capacity {
1126            return
1127        }
1128
1129        let now = std::time::Instant::now();
1130
1131        let candidates = self
1132            .peers
1133            .iter()
1134            .filter_map(|(peer_id, peer)| {
1135                let eligible = match peer.state {
1136                    PeerConnectionState::Out => outbound_at_capacity,
1137                    PeerConnectionState::In => inbound_at_capacity,
1138                    _ => false,
1139                };
1140                (eligible &&
1141                    !peer.is_trusted() &&
1142                    !peer.is_static() &&
1143                    !self.trusted_peer_ids.contains(peer_id) &&
1144                    peer.connected_for_at_least(now, PEER_ROTATION_MIN_UPTIME))
1145                .then_some(*peer_id)
1146            })
1147            .collect::<Vec<_>>();
1148
1149        if candidates.is_empty() {
1150            return
1151        }
1152        let peer_id = candidates[rand::rng().random_range(0..candidates.len())];
1153
1154        trace!(target: "net::peers", ?peer_id, "rotating peer to open slot for new nodes");
1155
1156        self.queued_actions.push_back(PeerAction::Disconnect {
1157            peer_id,
1158            reason: Some(DisconnectReason::UselessPeer),
1159        });
1160    }
1161
1162    /// If there's capacity for new outbound connections, this will queue new
1163    /// [`PeerAction::Connect`] actions.
1164    ///
1165    /// New connections are only initiated, if slots are available and appropriate peers are
1166    /// available.
1167    fn fill_outbound_slots(&mut self) {
1168        self.tick();
1169
1170        if !self.net_connection_state.is_active() {
1171            // nothing to fill
1172            return
1173        }
1174
1175        // as long as there are slots available fill them with the best peers
1176        while self.connection_info.has_out_capacity() {
1177            let action = {
1178                let (peer_id, peer) = match self.best_unconnected() {
1179                    Some(peer) => peer,
1180                    _ => break,
1181                };
1182
1183                trace!(target: "net::peers", ?peer_id, addr=?peer.addr, "schedule outbound connection");
1184
1185                peer.state = PeerConnectionState::PendingOut;
1186                PeerAction::Connect { peer_id, remote_addr: peer.addr.tcp() }
1187            };
1188
1189            self.connection_info.inc_pending_out();
1190
1191            self.queued_actions.push_back(action);
1192        }
1193    }
1194
1195    fn on_resolved_peer(&mut self, peer_id: PeerId, new_record: NodeRecord) {
1196        if !self.trusted_peer_ids.contains(&peer_id) {
1197            trace!(target: "net::peers", ?peer_id, "Ignoring resolved trusted peer after removal");
1198            return
1199        }
1200
1201        let new_addr = PeerAddr::new_with_ports(
1202            new_record.address,
1203            new_record.tcp_port,
1204            Some(new_record.udp_port),
1205        );
1206
1207        if let Some(peer) = self.peers.get_mut(&peer_id) {
1208            if peer.addr != new_addr {
1209                peer.addr = new_addr;
1210                trace!(target: "net::peers", ?peer_id, addr=?peer.addr, "Updated resolved trusted peer address");
1211            }
1212        } else {
1213            trace!(target: "net::peers", ?peer_id, ?new_addr, "Adding trusted peer after first successful resolution");
1214            self.add_peer_kind(peer_id, Some(PeerKind::Trusted), new_addr, None);
1215        }
1216    }
1217
1218    /// Keeps track of network state changes.
1219    pub const fn on_network_state_change(&mut self, state: NetworkConnectionState) {
1220        self.net_connection_state = state;
1221    }
1222
1223    /// Returns the current network connection state.
1224    pub const fn connection_state(&self) -> &NetworkConnectionState {
1225        &self.net_connection_state
1226    }
1227
1228    /// Sets `net_connection_state` to `ShuttingDown`.
1229    pub const fn on_shutdown(&mut self) {
1230        self.net_connection_state = NetworkConnectionState::ShuttingDown;
1231    }
1232
1233    /// Advances the state.
1234    ///
1235    /// Event hooks invoked externally may trigger a new [`PeerAction`] that are buffered until
1236    /// [`PeersManager`] is polled.
1237    pub fn poll(&mut self, cx: &mut Context<'_>) -> Poll<PeerAction> {
1238        loop {
1239            // drain buffered actions
1240            if let Some(action) = self.queued_actions.pop_front() {
1241                return Poll::Ready(action)
1242            }
1243
1244            while let Poll::Ready(Some(cmd)) = self.handle_rx.poll_next_unpin(cx) {
1245                match cmd {
1246                    PeerCommand::Add(peer_id, addr) => {
1247                        self.add_requested_peer(peer_id, None, PeerAddr::from_tcp(addr));
1248                    }
1249                    PeerCommand::Remove(peer) => self.remove_peer(peer),
1250                    PeerCommand::ReputationChange(peer_id, rep) => {
1251                        self.apply_reputation_change(&peer_id, rep)
1252                    }
1253                    PeerCommand::GetPeer(peer, tx) => {
1254                        let _ = tx.send(self.peers.get(&peer).cloned());
1255                    }
1256                    PeerCommand::GetPeers(tx) => {
1257                        let _ = tx.send(self.iter_peers().collect());
1258                    }
1259                }
1260            }
1261
1262            if self.release_interval.poll_tick(cx).is_ready() {
1263                let now = std::time::Instant::now();
1264                let (_, unbanned_peers) = self.ban_list.evict(now);
1265
1266                for peer_id in unbanned_peers {
1267                    if let Some(peer) = self.peers.get_mut(&peer_id) {
1268                        peer.unban();
1269                        self.queued_actions.push_back(PeerAction::UnBanPeer { peer_id });
1270                    }
1271                }
1272
1273                // clear the backoff list of expired backoffs, and mark the relevant peers as
1274                // ready to be dialed
1275                self.backed_off_peers.retain(|peer_id, until| {
1276                    if now > *until {
1277                        if let Some(peer) = self.peers.get_mut(peer_id) {
1278                            peer.backed_off = false;
1279                        }
1280                        return false
1281                    }
1282                    true
1283                })
1284            }
1285
1286            while self.refill_slots_interval.poll_tick(cx).is_ready() {
1287                self.fill_outbound_slots();
1288            }
1289
1290            if let Poll::Ready((peer_id, new_record)) = self.trusted_peers_resolver.poll(cx) {
1291                self.on_resolved_peer(peer_id, new_record);
1292            }
1293
1294            // Poll the jittered rotation timer and rotate one eligible peer when ready.
1295            let rotation_ready = self
1296                .peer_rotation_sleep
1297                .as_mut()
1298                .is_some_and(|sleep| sleep.as_mut().poll(cx).is_ready());
1299            if rotation_ready {
1300                self.try_rotate_peer();
1301                if let Some((mean, sleep)) =
1302                    self.peer_rotation_mean.zip(self.peer_rotation_sleep.as_mut())
1303                {
1304                    sleep
1305                        .as_mut()
1306                        .reset(tokio::time::Instant::now() + jitter_rotation_interval(mean));
1307                }
1308            }
1309
1310            if self.queued_actions.is_empty() {
1311                return Poll::Pending
1312            }
1313        }
1314    }
1315}
1316
1317impl Default for PeersManager {
1318    fn default() -> Self {
1319        Self::new(Default::default())
1320    }
1321}
1322
1323/// Tracks stats about connected nodes
1324#[derive(Debug, Clone, PartialEq, Eq, Default)]
1325pub struct ConnectionInfo {
1326    /// Counter for currently occupied slots for active outbound connections.
1327    num_outbound: usize,
1328    /// Counter for pending outbound connections.
1329    num_pending_out: usize,
1330    /// Counter for currently occupied slots for active inbound connections.
1331    num_inbound: usize,
1332    /// Counter for pending inbound connections.
1333    num_pending_in: usize,
1334    /// Restrictions on number of connections.
1335    config: ConnectionsConfig,
1336}
1337
1338// === impl ConnectionInfo ===
1339
1340impl ConnectionInfo {
1341    /// Returns a new [`ConnectionInfo`] with the given config.
1342    const fn new(config: ConnectionsConfig) -> Self {
1343        Self { config, num_outbound: 0, num_pending_out: 0, num_inbound: 0, num_pending_in: 0 }
1344    }
1345
1346    ///  Returns `true` if there's still capacity to perform an outgoing connection.
1347    const fn has_out_capacity(&self) -> bool {
1348        self.num_pending_out < self.config.max_concurrent_outbound_dials &&
1349            self.num_outbound < self.config.max_outbound
1350    }
1351
1352    /// Returns `true` if all active outbound slots are occupied (ignoring pending dials).
1353    const fn is_outbound_at_capacity(&self) -> bool {
1354        self.num_outbound >= self.config.max_outbound
1355    }
1356
1357    /// Returns `true` if all active inbound slots are occupied.
1358    const fn is_inbound_at_capacity(&self) -> bool {
1359        self.num_inbound >= self.config.max_inbound
1360    }
1361
1362    ///  Returns `true` if there's still capacity to accept a new incoming connection.
1363    const fn has_in_capacity(&self) -> bool {
1364        self.num_inbound < self.config.max_inbound
1365    }
1366
1367    /// Returns `true` if we can handle an additional incoming pending connection.
1368    const fn has_in_pending_capacity(&self) -> bool {
1369        self.num_pending_in < self.config.max_inbound
1370    }
1371
1372    const fn decr_state(&mut self, state: PeerConnectionState) {
1373        match state {
1374            PeerConnectionState::Idle => {}
1375            PeerConnectionState::DisconnectingIn | PeerConnectionState::In => self.decr_in(),
1376            PeerConnectionState::DisconnectingOut | PeerConnectionState::Out => self.decr_out(),
1377            PeerConnectionState::PendingOut => self.decr_pending_out(),
1378        }
1379    }
1380
1381    const fn decr_out(&mut self) {
1382        self.num_outbound -= 1;
1383    }
1384
1385    const fn inc_out(&mut self) {
1386        self.num_outbound += 1;
1387    }
1388
1389    const fn inc_pending_out(&mut self) {
1390        self.num_pending_out += 1;
1391    }
1392
1393    const fn inc_in(&mut self) {
1394        self.num_inbound += 1;
1395    }
1396
1397    const fn inc_pending_in(&mut self) {
1398        self.num_pending_in += 1;
1399    }
1400
1401    const fn decr_in(&mut self) {
1402        self.num_inbound -= 1;
1403    }
1404
1405    const fn decr_pending_out(&mut self) {
1406        self.num_pending_out -= 1;
1407    }
1408
1409    const fn decr_pending_in(&mut self) {
1410        self.num_pending_in -= 1;
1411    }
1412}
1413
1414/// Actions the peer manager can trigger.
1415#[derive(Debug)]
1416pub enum PeerAction {
1417    /// Start a new connection to a peer.
1418    Connect {
1419        /// The peer to connect to.
1420        peer_id: PeerId,
1421        /// Where to reach the node
1422        remote_addr: SocketAddr,
1423    },
1424    /// Disconnect an existing connection.
1425    Disconnect {
1426        /// The peer ID of the established connection.
1427        peer_id: PeerId,
1428        /// An optional reason for the disconnect.
1429        reason: Option<DisconnectReason>,
1430    },
1431    /// Disconnect an existing incoming connection, because the peers reputation is below the
1432    /// banned threshold or is on the [`BanList`]
1433    DisconnectBannedIncoming {
1434        /// The peer ID of the established connection.
1435        peer_id: PeerId,
1436    },
1437    /// Disconnect an untrusted incoming connection when trust-node-only is enabled.
1438    DisconnectUntrustedIncoming {
1439        /// The peer ID.
1440        peer_id: PeerId,
1441    },
1442    /// Ban the peer in discovery.
1443    DiscoveryBanPeerId {
1444        /// The peer ID.
1445        peer_id: PeerId,
1446        /// The IP address.
1447        ip_addr: IpAddr,
1448    },
1449    /// Ban the IP in discovery.
1450    DiscoveryBanIp {
1451        /// The IP address.
1452        ip_addr: IpAddr,
1453    },
1454    /// Ban the peer temporarily
1455    BanPeer {
1456        /// The peer ID.
1457        peer_id: PeerId,
1458    },
1459    /// Unban the peer temporarily
1460    UnBanPeer {
1461        /// The peer ID.
1462        peer_id: PeerId,
1463    },
1464    /// Emit peerAdded event
1465    PeerAdded(PeerId),
1466    /// Emit peerRemoved event
1467    PeerRemoved(PeerId),
1468}
1469
1470/// Error thrown when an incoming connection is rejected right away
1471#[derive(Debug, Error, PartialEq, Eq)]
1472pub enum InboundConnectionError {
1473    /// The remote's ip address is banned
1474    IpBanned,
1475    /// No capacity for new inbound connections
1476    ExceedsCapacity,
1477}
1478
1479impl Display for InboundConnectionError {
1480    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1481        write!(f, "{self:?}")
1482    }
1483}
1484
1485/// The reason a peer was backed off.
1486#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1487pub enum BackoffReason {
1488    /// The remote peer responded with `TooManyPeers` (0x04).
1489    TooManyPeers,
1490    /// The session was gracefully closed and we're backing off briefly.
1491    GracefulClose,
1492    /// A connection or protocol-level error occurred.
1493    ConnectionError,
1494}
1495
1496impl BackoffReason {
1497    /// Derives the backoff reason from an optional [`DisconnectReason`].
1498    pub const fn from_disconnect(reason: Option<DisconnectReason>) -> Self {
1499        match reason {
1500            Some(DisconnectReason::TooManyPeers) => Self::TooManyPeers,
1501            _ => Self::ConnectionError,
1502        }
1503    }
1504}
1505
1506/// Returns a random duration uniformly distributed in `[mean * 3/5, mean * 7/5]`.
1507///
1508/// With the default 5-minute mean this gives `[3 min, 7 min]`, matching Geth's peer dropper:
1509/// <https://github.com/ethereum/go-ethereum/blob/c5c75977ab55e4d7ea6147cc0e221b588e5e3754/eth/dropper.go#L32-L36>
1510fn jitter_rotation_interval(mean: Duration) -> Duration {
1511    let min_nanos = (mean * 3 / 5).as_nanos() as u64;
1512    let max_nanos = (mean * 7 / 5).as_nanos() as u64;
1513    Duration::from_nanos(rand::rng().random_range(min_nanos..=max_nanos))
1514}
1515
1516#[cfg(test)]
1517mod tests {
1518    use alloy_primitives::B512;
1519    use reth_eth_wire::{
1520        errors::{EthHandshakeError, EthStreamError, P2PHandshakeError, P2PStreamError},
1521        DisconnectReason,
1522    };
1523    use reth_ethereum_forks::{ForkHash, ForkId};
1524    use reth_net_banlist::BanList;
1525    use reth_network_api::Direction;
1526    use reth_network_peers::{NodeRecord, PeerId, TrustedPeer};
1527    use reth_network_types::{
1528        peers::reputation::DEFAULT_REPUTATION, BackoffKind, Peer, ReputationChangeKind,
1529    };
1530    use std::{
1531        future::{poll_fn, Future},
1532        io,
1533        net::{IpAddr, Ipv4Addr, SocketAddr},
1534        pin::Pin,
1535        task::{Context, Poll, Waker},
1536        time::Duration,
1537    };
1538    use url::Host;
1539
1540    use super::PeersManager;
1541    use crate::{
1542        error::SessionError,
1543        peers::{
1544            ConnectionInfo, InboundConnectionError, PeerAction, PeerAddr, PeerBackoffDurations,
1545            PeerConnectionState,
1546        },
1547        session::PendingSessionHandshakeError,
1548        PeersConfig,
1549    };
1550
1551    struct PeerActionFuture<'a> {
1552        peers: &'a mut PeersManager,
1553    }
1554
1555    impl Future for PeerActionFuture<'_> {
1556        type Output = PeerAction;
1557
1558        fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
1559            self.get_mut().peers.poll(cx)
1560        }
1561    }
1562
1563    macro_rules! event {
1564        ($peers:expr) => {
1565            PeerActionFuture { peers: &mut $peers }.await
1566        };
1567    }
1568
1569    fn set_connected_at(
1570        peers: &mut PeersManager,
1571        peer_id: PeerId,
1572        connected_at: std::time::Instant,
1573    ) {
1574        peers.peers.get_mut(&peer_id).expect("peer exists").connected_at = Some(connected_at);
1575    }
1576
1577    /// Polls the manager once without waiting for an action.
1578    fn poll_now(peers: &mut PeersManager) -> Poll<PeerAction> {
1579        peers.poll(&mut Context::from_waker(Waker::noop()))
1580    }
1581
1582    #[tokio::test]
1583    async fn test_insert() {
1584        let peer = PeerId::random();
1585        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1586        let mut peers = PeersManager::default();
1587        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
1588
1589        match event!(peers) {
1590            PeerAction::PeerAdded(peer_id) => {
1591                assert_eq!(peer_id, peer);
1592            }
1593            _ => unreachable!(),
1594        }
1595        match event!(peers) {
1596            PeerAction::Connect { peer_id, remote_addr } => {
1597                assert_eq!(peer_id, peer);
1598                assert_eq!(remote_addr, socket_addr);
1599            }
1600            _ => unreachable!(),
1601        }
1602
1603        let (record, _) = peers.peer_by_id(peer).unwrap();
1604        assert_eq!(record.tcp_addr(), socket_addr);
1605        assert_eq!(record.udp_addr(), socket_addr);
1606    }
1607
1608    #[tokio::test]
1609    async fn test_insert_udp() {
1610        let peer = PeerId::random();
1611        let tcp_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1612        let udp_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
1613        let mut peers = PeersManager::default();
1614        peers.add_peer(peer, PeerAddr::new(tcp_addr, Some(udp_addr)), None);
1615
1616        match event!(peers) {
1617            PeerAction::PeerAdded(peer_id) => {
1618                assert_eq!(peer_id, peer);
1619            }
1620            _ => unreachable!(),
1621        }
1622        match event!(peers) {
1623            PeerAction::Connect { peer_id, remote_addr } => {
1624                assert_eq!(peer_id, peer);
1625                assert_eq!(remote_addr, tcp_addr);
1626            }
1627            _ => unreachable!(),
1628        }
1629
1630        let (record, _) = peers.peer_by_id(peer).unwrap();
1631        assert_eq!(record.tcp_addr(), tcp_addr);
1632        assert_eq!(record.udp_addr(), udp_addr);
1633    }
1634
1635    #[tokio::test]
1636    async fn test_ban() {
1637        let peer = PeerId::random();
1638        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1639        let mut peers = PeersManager::default();
1640        peers.ban_peer(peer);
1641        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
1642
1643        match event!(peers) {
1644            PeerAction::BanPeer { peer_id } => {
1645                assert_eq!(peer_id, peer);
1646            }
1647            _ => unreachable!(),
1648        }
1649
1650        poll_fn(|cx| {
1651            assert!(peers.poll(cx).is_pending());
1652            Poll::Ready(())
1653        })
1654        .await;
1655    }
1656
1657    #[tokio::test]
1658    async fn test_unban() {
1659        let peer = PeerId::random();
1660        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1661        let mut peers = PeersManager::default();
1662        peers.ban_peer(peer);
1663        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
1664
1665        match event!(peers) {
1666            PeerAction::BanPeer { peer_id } => {
1667                assert_eq!(peer_id, peer);
1668            }
1669            _ => unreachable!(),
1670        }
1671
1672        poll_fn(|cx| {
1673            assert!(peers.poll(cx).is_pending());
1674            Poll::Ready(())
1675        })
1676        .await;
1677
1678        peers.unban_peer(peer);
1679
1680        match event!(peers) {
1681            PeerAction::UnBanPeer { peer_id } => {
1682                assert_eq!(peer_id, peer);
1683            }
1684            _ => unreachable!(),
1685        }
1686
1687        poll_fn(|cx| {
1688            assert!(peers.poll(cx).is_pending());
1689            Poll::Ready(())
1690        })
1691        .await;
1692    }
1693
1694    #[tokio::test]
1695    async fn test_admin_ban_removes_peer_and_bans_indefinitely() {
1696        let peer = PeerId::random();
1697        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1698        let mut peers = PeersManager::default();
1699        peers.peers.insert(peer, Peer::new(PeerAddr::from_tcp(socket_addr)));
1700
1701        peers.ban_peer_by_admin(peer);
1702
1703        assert!(peers.ban_list.is_banned_peer(&peer));
1704        assert!(peers.peer_by_id(peer).is_none());
1705
1706        match peers.queued_actions.pop_front() {
1707            Some(PeerAction::PeerRemoved(peer_id)) => assert_eq!(peer_id, peer),
1708            other => panic!("unexpected action: {other:?}"),
1709        }
1710
1711        let (_, unbanned_peers) =
1712            peers.ban_list.evict(std::time::Instant::now() + Duration::from_secs(1));
1713        assert!(unbanned_peers.is_empty());
1714    }
1715
1716    #[tokio::test]
1717    async fn test_admin_ban_does_not_override_trusted_peer() {
1718        let peer = PeerId::random();
1719        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1720        let mut peers = PeersManager::default();
1721        peers.peers.insert(peer, Peer::trusted(PeerAddr::from_tcp(socket_addr)));
1722
1723        peers.ban_peer_by_admin(peer);
1724
1725        assert!(!peers.ban_list.is_banned_peer(&peer));
1726        assert!(peers.peer_by_id(peer).is_some());
1727        assert!(peers.queued_actions.is_empty());
1728    }
1729
1730    #[tokio::test]
1731    async fn test_admin_unban_resets_reputation() {
1732        let peer = PeerId::random();
1733        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1734        let mut peers = PeersManager::default();
1735        let mut peer_info = Peer::new(PeerAddr::from_tcp(socket_addr));
1736        peer_info.reputation = i32::MIN;
1737        peers.peers.insert(peer, peer_info);
1738        peers.ban_list.ban_peer(peer);
1739        assert!(peers.peers.get(&peer).is_some_and(Peer::is_banned));
1740
1741        peers.unban_peer_by_admin(peer);
1742
1743        assert!(!peers.ban_list.is_banned_peer(&peer));
1744        assert!(!peers.peers.get(&peer).is_some_and(Peer::is_banned));
1745        assert!(matches!(
1746            peers.queued_actions.pop_front(),
1747            Some(PeerAction::UnBanPeer { peer_id }) if peer_id == peer
1748        ));
1749    }
1750
1751    #[tokio::test]
1752    async fn test_backoff_on_busy() {
1753        let peer = PeerId::random();
1754        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1755
1756        let mut peers = PeersManager::new(PeersConfig::test());
1757        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
1758
1759        match event!(peers) {
1760            PeerAction::PeerAdded(peer_id) => {
1761                assert_eq!(peer_id, peer);
1762            }
1763            _ => unreachable!(),
1764        }
1765        match event!(peers) {
1766            PeerAction::Connect { peer_id, .. } => {
1767                assert_eq!(peer_id, peer);
1768            }
1769            _ => unreachable!(),
1770        }
1771
1772        poll_fn(|cx| {
1773            assert!(peers.poll(cx).is_pending());
1774            Poll::Ready(())
1775        })
1776        .await;
1777
1778        peers.on_active_session_dropped(
1779            &socket_addr,
1780            &peer,
1781            &EthStreamError::P2PStreamError(P2PStreamError::Disconnected(
1782                DisconnectReason::TooManyPeers,
1783            )),
1784        );
1785
1786        poll_fn(|cx| {
1787            assert!(peers.poll(cx).is_pending());
1788            Poll::Ready(())
1789        })
1790        .await;
1791
1792        assert!(peers.backed_off_peers.contains_key(&peer));
1793        assert!(peers.peers.get(&peer).unwrap().is_backed_off());
1794
1795        tokio::time::sleep(peers.backoff_durations.low).await;
1796
1797        match event!(peers) {
1798            PeerAction::Connect { peer_id, .. } => {
1799                assert_eq!(peer_id, peer);
1800            }
1801            _ => unreachable!(),
1802        }
1803
1804        assert!(!peers.backed_off_peers.contains_key(&peer));
1805        assert!(!peers.peers.get(&peer).unwrap().is_backed_off());
1806    }
1807
1808    #[tokio::test]
1809    async fn test_backoff_on_no_response() {
1810        let peer = PeerId::random();
1811        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1812
1813        let backoff_durations = PeerBackoffDurations::test();
1814        let config = PeersConfig { backoff_durations, ..PeersConfig::test() };
1815        let mut peers = PeersManager::new(config);
1816        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
1817
1818        match event!(peers) {
1819            PeerAction::PeerAdded(peer_id) => {
1820                assert_eq!(peer_id, peer);
1821            }
1822            _ => unreachable!(),
1823        }
1824        match event!(peers) {
1825            PeerAction::Connect { peer_id, .. } => {
1826                assert_eq!(peer_id, peer);
1827            }
1828            _ => unreachable!(),
1829        }
1830
1831        poll_fn(|cx| {
1832            assert!(peers.poll(cx).is_pending());
1833            Poll::Ready(())
1834        })
1835        .await;
1836
1837        peers.on_outgoing_pending_session_dropped(
1838            &socket_addr,
1839            &peer,
1840            &PendingSessionHandshakeError::Eth(EthStreamError::EthHandshakeError(
1841                EthHandshakeError::NoResponse,
1842            )),
1843        );
1844
1845        poll_fn(|cx| {
1846            assert!(peers.poll(cx).is_pending());
1847            Poll::Ready(())
1848        })
1849        .await;
1850
1851        assert!(peers.backed_off_peers.contains_key(&peer));
1852        assert!(peers.peers.get(&peer).unwrap().is_backed_off());
1853
1854        tokio::time::sleep(backoff_durations.high).await;
1855
1856        match event!(peers) {
1857            PeerAction::Connect { peer_id, .. } => {
1858                assert_eq!(peer_id, peer);
1859            }
1860            _ => unreachable!(),
1861        }
1862
1863        assert!(!peers.backed_off_peers.contains_key(&peer));
1864        assert!(!peers.peers.get(&peer).unwrap().is_backed_off());
1865    }
1866
1867    #[tokio::test]
1868    async fn test_low_backoff() {
1869        let peer = PeerId::random();
1870        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1871        let config = PeersConfig::test();
1872        let mut peers = PeersManager::new(config);
1873        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
1874        let peer_struct = peers.peers.get_mut(&peer).unwrap();
1875
1876        let backoff_timestamp = peers
1877            .backoff_durations
1878            .backoff_until(BackoffKind::Low, peer_struct.severe_backoff_counter);
1879
1880        let expected = std::time::Instant::now() + peers.backoff_durations.low;
1881        assert!(backoff_timestamp <= expected);
1882    }
1883
1884    #[tokio::test]
1885    async fn test_multiple_backoff_calculations() {
1886        let peer = PeerId::random();
1887        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1888        let config = PeersConfig::default();
1889        let mut peers = PeersManager::new(config);
1890        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
1891        let peer_struct = peers.peers.get_mut(&peer).unwrap();
1892
1893        // Simulate a peer that was already backed off once
1894        peer_struct.severe_backoff_counter = 1;
1895
1896        let now = std::time::Instant::now();
1897
1898        // Simulate the increment that happens in on_connection_failure
1899        peer_struct.severe_backoff_counter += 1;
1900        // Get official backoff time
1901        let backoff_time = peers
1902            .backoff_durations
1903            .backoff_until(BackoffKind::High, peer_struct.severe_backoff_counter);
1904
1905        // Duration of the backoff should be 2 * 15 minutes = 30 minutes
1906        let backoff_duration = std::time::Duration::new(30 * 60, 0);
1907
1908        // We can't use assert_eq! since there is a very small diff in the nano secs
1909        // Usually it is 1800s != 1799.9999996s
1910        assert!(backoff_time.duration_since(now) > backoff_duration);
1911    }
1912
1913    #[tokio::test]
1914    async fn test_ban_on_active_drop() {
1915        let peer = PeerId::random();
1916        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1917        let mut peers = PeersManager::default();
1918        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
1919
1920        match event!(peers) {
1921            PeerAction::PeerAdded(peer_id) => {
1922                assert_eq!(peer_id, peer);
1923            }
1924            _ => unreachable!(),
1925        }
1926        match event!(peers) {
1927            PeerAction::Connect { peer_id, .. } => {
1928                assert_eq!(peer_id, peer);
1929            }
1930            _ => unreachable!(),
1931        }
1932
1933        poll_fn(|cx| {
1934            assert!(peers.poll(cx).is_pending());
1935            Poll::Ready(())
1936        })
1937        .await;
1938
1939        peers.on_active_session_dropped(
1940            &socket_addr,
1941            &peer,
1942            &EthStreamError::P2PStreamError(P2PStreamError::Disconnected(
1943                DisconnectReason::UselessPeer,
1944            )),
1945        );
1946
1947        match event!(peers) {
1948            PeerAction::PeerRemoved(peer_id) => {
1949                assert_eq!(peer_id, peer);
1950            }
1951            _ => unreachable!(),
1952        }
1953        match event!(peers) {
1954            PeerAction::BanPeer { peer_id } => {
1955                assert_eq!(peer_id, peer);
1956            }
1957            _ => unreachable!(),
1958        }
1959
1960        poll_fn(|cx| {
1961            assert!(peers.poll(cx).is_pending());
1962            Poll::Ready(())
1963        })
1964        .await;
1965
1966        assert!(!peers.peers.contains_key(&peer));
1967    }
1968
1969    #[tokio::test]
1970    async fn test_remove_on_max_backoff_count() {
1971        let peer = PeerId::random();
1972        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
1973        let config = PeersConfig::test();
1974        let mut peers = PeersManager::new(config.clone());
1975        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
1976        let peer_struct = peers.peers.get_mut(&peer).unwrap();
1977
1978        // Simulate a peer that was already backed off once
1979        peer_struct.severe_backoff_counter = config.max_backoff_count;
1980
1981        match event!(peers) {
1982            PeerAction::PeerAdded(peer_id) => {
1983                assert_eq!(peer_id, peer);
1984            }
1985            _ => unreachable!(),
1986        }
1987        match event!(peers) {
1988            PeerAction::Connect { peer_id, .. } => {
1989                assert_eq!(peer_id, peer);
1990            }
1991            _ => unreachable!(),
1992        }
1993
1994        poll_fn(|cx| {
1995            assert!(peers.poll(cx).is_pending());
1996            Poll::Ready(())
1997        })
1998        .await;
1999
2000        peers.on_outgoing_pending_session_dropped(
2001            &socket_addr,
2002            &peer,
2003            &PendingSessionHandshakeError::Eth(
2004                io::Error::new(io::ErrorKind::ConnectionRefused, "peer unreachable").into(),
2005            ),
2006        );
2007
2008        match event!(peers) {
2009            PeerAction::PeerRemoved(peer_id) => {
2010                assert_eq!(peer_id, peer);
2011            }
2012            _ => unreachable!(),
2013        }
2014
2015        poll_fn(|cx| {
2016            assert!(peers.poll(cx).is_pending());
2017            Poll::Ready(())
2018        })
2019        .await;
2020
2021        assert!(!peers.peers.contains_key(&peer));
2022    }
2023
2024    #[tokio::test]
2025    async fn test_ban_on_pending_drop() {
2026        let peer = PeerId::random();
2027        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2028        let mut peers = PeersManager::default();
2029        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
2030
2031        match event!(peers) {
2032            PeerAction::PeerAdded(peer_id) => {
2033                assert_eq!(peer_id, peer);
2034            }
2035            _ => unreachable!(),
2036        }
2037        match event!(peers) {
2038            PeerAction::Connect { peer_id, .. } => {
2039                assert_eq!(peer_id, peer);
2040            }
2041            _ => unreachable!(),
2042        }
2043
2044        poll_fn(|cx| {
2045            assert!(peers.poll(cx).is_pending());
2046            Poll::Ready(())
2047        })
2048        .await;
2049
2050        peers.on_outgoing_pending_session_dropped(
2051            &socket_addr,
2052            &peer,
2053            &PendingSessionHandshakeError::Eth(EthStreamError::P2PStreamError(
2054                P2PStreamError::Disconnected(DisconnectReason::UselessPeer),
2055            )),
2056        );
2057
2058        match event!(peers) {
2059            PeerAction::PeerRemoved(peer_id) => {
2060                assert_eq!(peer_id, peer);
2061            }
2062            _ => unreachable!(),
2063        }
2064        match event!(peers) {
2065            PeerAction::BanPeer { peer_id } => {
2066                assert_eq!(peer_id, peer);
2067            }
2068            _ => unreachable!(),
2069        }
2070
2071        poll_fn(|cx| {
2072            assert!(peers.poll(cx).is_pending());
2073            Poll::Ready(())
2074        })
2075        .await;
2076
2077        assert!(!peers.peers.contains_key(&peer));
2078    }
2079
2080    #[tokio::test]
2081    async fn test_internally_closed_incoming() {
2082        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2083        let mut peers = PeersManager::default();
2084
2085        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2086        assert_eq!(peers.connection_info.num_pending_in, 1);
2087        peers.on_incoming_pending_session_rejected_internally();
2088        assert_eq!(peers.connection_info.num_pending_in, 0);
2089    }
2090
2091    #[tokio::test]
2092    async fn test_reject_incoming_at_pending_capacity() {
2093        let mut peers = PeersManager::default();
2094
2095        for count in 1..=peers.connection_info.config.max_inbound {
2096            let socket_addr =
2097                SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, count as u8)), 8008);
2098            assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2099            assert_eq!(peers.connection_info.num_pending_in, count);
2100        }
2101        assert!(peers.connection_info.has_in_capacity());
2102        assert!(!peers.connection_info.has_in_pending_capacity());
2103
2104        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 100)), 8008);
2105        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_err());
2106    }
2107
2108    #[tokio::test]
2109    async fn test_reject_incoming_at_pending_capacity_trusted_peers() {
2110        let mut peers = PeersManager::new(PeersConfig::test().with_max_inbound(2));
2111        let trusted = PeerId::random();
2112        peers.add_trusted_peer_id(trusted);
2113
2114        // connect the trusted peer
2115        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 0)), 8008);
2116        assert!(peers.on_incoming_pending_session(addr.ip()).is_ok());
2117        peers.on_incoming_session_established(trusted, addr);
2118
2119        match event!(peers) {
2120            PeerAction::PeerAdded(id) => {
2121                assert_eq!(id, trusted);
2122            }
2123            _ => unreachable!(),
2124        }
2125
2126        // saturate the remaining inbound slots with untrusted peers
2127        let mut connected_untrusted_peer_ids = Vec::new();
2128        for i in 0..(peers.connection_info.config.max_inbound - 1) {
2129            let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, (i + 1) as u8)), 8008);
2130            assert!(peers.on_incoming_pending_session(addr.ip()).is_ok());
2131            let peer_id = PeerId::random();
2132            peers.on_incoming_session_established(peer_id, addr);
2133            connected_untrusted_peer_ids.push(peer_id);
2134
2135            match event!(peers) {
2136                PeerAction::PeerAdded(id) => {
2137                    assert_eq!(id, peer_id);
2138                }
2139                _ => unreachable!(),
2140            }
2141        }
2142
2143        let mut pending_addrs = Vec::new();
2144
2145        // saturate available slots
2146        for i in 0..2 {
2147            let socket_addr =
2148                SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, (i + 10) as u8)), 8008);
2149            assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2150
2151            pending_addrs.push(socket_addr);
2152        }
2153
2154        assert_eq!(peers.connection_info.num_pending_in, 2);
2155
2156        // try to handle additional incoming connections at capacity
2157        for i in 0..2 {
2158            let socket_addr =
2159                SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, (i + 20) as u8)), 8008);
2160            assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_err());
2161        }
2162
2163        let err = PendingSessionHandshakeError::Eth(EthStreamError::P2PStreamError(
2164            P2PStreamError::HandshakeError(P2PHandshakeError::Disconnected(
2165                DisconnectReason::UselessPeer,
2166            )),
2167        ));
2168
2169        // Remove all pending peers
2170        for pending_addr in pending_addrs {
2171            peers.on_incoming_pending_session_dropped(pending_addr, &err);
2172        }
2173
2174        println!("num_pending_in: {}", peers.connection_info.num_pending_in);
2175
2176        println!(
2177            "num_inbound: {}, has_in_capacity: {}",
2178            peers.connection_info.num_inbound,
2179            peers.connection_info.has_in_capacity()
2180        );
2181
2182        // disconnect a connected peer
2183        peers.on_active_session_gracefully_closed(connected_untrusted_peer_ids[0]);
2184
2185        println!(
2186            "num_inbound: {}, has_in_capacity: {}",
2187            peers.connection_info.num_inbound,
2188            peers.connection_info.has_in_capacity()
2189        );
2190
2191        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 99)), 8008);
2192        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2193    }
2194
2195    #[tokio::test]
2196    async fn test_closed_incoming() {
2197        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2198        let mut peers = PeersManager::default();
2199
2200        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2201        assert_eq!(peers.connection_info.num_pending_in, 1);
2202        peers.on_incoming_pending_session_gracefully_closed();
2203        assert_eq!(peers.connection_info.num_pending_in, 0);
2204    }
2205
2206    #[tokio::test]
2207    async fn test_dropped_incoming() {
2208        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(1, 0, 1, 2)), 8008);
2209        let ban_duration = Duration::from_millis(500);
2210        let config = PeersConfig { ban_duration, ..PeersConfig::test() };
2211        let mut peers = PeersManager::new(config);
2212
2213        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2214        assert_eq!(peers.connection_info.num_pending_in, 1);
2215        let err = PendingSessionHandshakeError::Eth(EthStreamError::P2PStreamError(
2216            P2PStreamError::HandshakeError(P2PHandshakeError::Disconnected(
2217                DisconnectReason::UselessPeer,
2218            )),
2219        ));
2220
2221        peers.on_incoming_pending_session_dropped(socket_addr, &err);
2222        assert_eq!(peers.connection_info.num_pending_in, 0);
2223        assert!(peers.ban_list.is_banned_ip(&socket_addr.ip()));
2224
2225        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_err());
2226
2227        // unbanned after timeout
2228        tokio::time::sleep(ban_duration).await;
2229
2230        poll_fn(|cx| {
2231            let _ = peers.poll(cx);
2232            Poll::Ready(())
2233        })
2234        .await;
2235
2236        assert!(!peers.ban_list.is_banned_ip(&socket_addr.ip()));
2237        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2238    }
2239
2240    #[tokio::test]
2241    async fn test_reputation_change_connected() {
2242        let peer = PeerId::random();
2243        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2244        let mut peers = PeersManager::default();
2245        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
2246
2247        match event!(peers) {
2248            PeerAction::PeerAdded(peer_id) => {
2249                assert_eq!(peer_id, peer);
2250            }
2251            _ => unreachable!(),
2252        }
2253        match event!(peers) {
2254            PeerAction::Connect { peer_id, remote_addr } => {
2255                assert_eq!(peer_id, peer);
2256                assert_eq!(remote_addr, socket_addr);
2257            }
2258            _ => unreachable!(),
2259        }
2260
2261        let p = peers.peers.get_mut(&peer).unwrap();
2262        assert_eq!(p.state, PeerConnectionState::PendingOut);
2263
2264        peers.apply_reputation_change(&peer, ReputationChangeKind::BadProtocol);
2265
2266        let p = peers.peers.get(&peer).unwrap();
2267        assert_eq!(p.state, PeerConnectionState::PendingOut);
2268        assert!(p.is_banned());
2269
2270        peers.on_active_session_gracefully_closed(peer);
2271
2272        let p = peers.peers.get(&peer).unwrap();
2273        assert_eq!(p.state, PeerConnectionState::Idle);
2274        assert!(p.is_banned());
2275
2276        match event!(peers) {
2277            PeerAction::Disconnect { peer_id, .. } => {
2278                assert_eq!(peer_id, peer);
2279            }
2280            _ => unreachable!(),
2281        }
2282    }
2283
2284    #[tokio::test]
2285    async fn retain_trusted_status() {
2286        let _socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 99)), 8008);
2287        let trusted = PeerId::random();
2288        let mut peers =
2289            PeersManager::new(PeersConfig::test().with_trusted_nodes(vec![TrustedPeer {
2290                host: Host::Ipv4(Ipv4Addr::new(127, 0, 1, 2)),
2291                tcp_port: 8008,
2292                udp_port: 8008,
2293                id: trusted,
2294            }]));
2295        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2296        peers.add_peer(trusted, PeerAddr::from_tcp(socket_addr), None);
2297        assert!(peers.peers.get(&trusted).unwrap().is_trusted());
2298    }
2299
2300    #[tokio::test]
2301    async fn accept_incoming_trusted_unknown_peer_address() {
2302        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 99)), 8008);
2303        let mut peers = PeersManager::new(PeersConfig::test().with_max_inbound(2));
2304        // try to connect trusted peer
2305        let trusted = PeerId::random();
2306        peers.add_trusted_peer_id(trusted);
2307
2308        // saturate the inbound slots
2309        for i in 0..peers.connection_info.config.max_inbound {
2310            let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, i as u8)), 8008);
2311            assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2312            let peer_id = PeerId::random();
2313            peers.on_incoming_session_established(peer_id, addr);
2314
2315            match event!(peers) {
2316                PeerAction::PeerAdded(id) => {
2317                    assert_eq!(id, peer_id);
2318                }
2319                _ => unreachable!(),
2320            }
2321        }
2322
2323        // try to connect untrusted peer
2324        let untrusted = PeerId::random();
2325        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 99)), 8008);
2326        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2327        peers.on_incoming_session_established(untrusted, socket_addr);
2328
2329        match event!(peers) {
2330            PeerAction::PeerAdded(id) => {
2331                assert_eq!(id, untrusted);
2332            }
2333            _ => unreachable!(),
2334        }
2335
2336        match event!(peers) {
2337            PeerAction::Disconnect { peer_id, reason } => {
2338                assert_eq!(peer_id, untrusted);
2339                assert_eq!(reason, Some(DisconnectReason::TooManyPeers));
2340            }
2341            _ => unreachable!(),
2342        }
2343
2344        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 100)), 8008);
2345        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2346        peers.on_incoming_session_established(trusted, socket_addr);
2347
2348        match event!(peers) {
2349            PeerAction::PeerAdded(id) => {
2350                assert_eq!(id, trusted);
2351            }
2352            _ => unreachable!(),
2353        }
2354
2355        poll_fn(|cx| {
2356            assert!(peers.poll(cx).is_pending());
2357            Poll::Ready(())
2358        })
2359        .await;
2360
2361        let peer = peers.peers.get(&trusted).unwrap();
2362        assert_eq!(peer.state, PeerConnectionState::In);
2363    }
2364
2365    #[tokio::test]
2366    async fn test_already_connected() {
2367        let peer = PeerId::random();
2368        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2369        let mut peers = PeersManager::default();
2370
2371        // Attempt to establish an incoming session, expecting `num_pending_in` to increase by 1
2372        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2373        assert_eq!(peers.connection_info.num_pending_in, 1);
2374
2375        // Establish a session with the peer, expecting the peer to be added and the `num_inbound`
2376        // to increase by 1
2377        peers.on_incoming_session_established(peer, socket_addr);
2378        let p = peers.peers.get_mut(&peer).expect("peer not found");
2379        assert_eq!(p.addr.tcp(), socket_addr);
2380        assert_eq!(peers.connection_info.num_pending_in, 0);
2381        assert_eq!(peers.connection_info.num_inbound, 1);
2382
2383        // Attempt to establish another incoming session, expecting the `num_pending_in` to increase
2384        // by 1
2385        assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2386        assert_eq!(peers.connection_info.num_pending_in, 1);
2387
2388        // Simulate a rejection due to an already established connection, expecting the
2389        // `num_pending_in` to decrease by 1. The peer should remain connected and the `num_inbound`
2390        // should not be changed.
2391        peers.on_already_connected(Direction::Incoming);
2392
2393        let p = peers.peers.get_mut(&peer).expect("peer not found");
2394        assert_eq!(p.addr.tcp(), socket_addr);
2395        assert_eq!(peers.connection_info.num_pending_in, 0);
2396        assert_eq!(peers.connection_info.num_inbound, 1);
2397    }
2398
2399    #[tokio::test]
2400    async fn test_reputation_change_trusted_peer() {
2401        let peer = PeerId::random();
2402        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2403        let mut peers = PeersManager::default();
2404        peers.add_trusted_peer(peer, PeerAddr::from_tcp(socket_addr));
2405
2406        match event!(peers) {
2407            PeerAction::PeerAdded(peer_id) => {
2408                assert_eq!(peer_id, peer);
2409            }
2410            _ => unreachable!(),
2411        }
2412        match event!(peers) {
2413            PeerAction::Connect { peer_id, remote_addr } => {
2414                assert_eq!(peer_id, peer);
2415                assert_eq!(remote_addr, socket_addr);
2416            }
2417            _ => unreachable!(),
2418        }
2419
2420        assert_eq!(peers.peers.get_mut(&peer).unwrap().state, PeerConnectionState::PendingOut);
2421        peers.on_active_outgoing_established(peer);
2422        assert_eq!(peers.peers.get_mut(&peer).unwrap().state, PeerConnectionState::Out);
2423
2424        peers.apply_reputation_change(&peer, ReputationChangeKind::BadMessage);
2425
2426        {
2427            let p = peers.peers.get(&peer).unwrap();
2428            assert_eq!(p.state, PeerConnectionState::Out);
2429            // not banned yet
2430            assert!(!p.is_banned());
2431        }
2432
2433        // ensure peer is banned eventually
2434        loop {
2435            peers.apply_reputation_change(&peer, ReputationChangeKind::BadMessage);
2436
2437            let p = peers.peers.get(&peer).unwrap();
2438            if p.is_banned() {
2439                break
2440            }
2441        }
2442
2443        match event!(peers) {
2444            PeerAction::Disconnect { peer_id, .. } => {
2445                assert_eq!(peer_id, peer);
2446            }
2447            _ => unreachable!(),
2448        }
2449    }
2450
2451    #[tokio::test]
2452    async fn test_reputation_management() {
2453        let peer = PeerId::random();
2454        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2455        let mut peers = PeersManager::default();
2456        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
2457        assert_eq!(peers.get_reputation(&peer), Some(0));
2458
2459        peers.apply_reputation_change(&peer, ReputationChangeKind::Other(1024));
2460        assert_eq!(peers.get_reputation(&peer), Some(1024));
2461
2462        peers.apply_reputation_change(&peer, ReputationChangeKind::Reset);
2463        assert_eq!(peers.get_reputation(&peer), Some(0));
2464    }
2465
2466    #[tokio::test]
2467    async fn test_remove_discovered_active() {
2468        let peer = PeerId::random();
2469        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2470        let mut peers = PeersManager::default();
2471        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
2472
2473        match event!(peers) {
2474            PeerAction::PeerAdded(peer_id) => {
2475                assert_eq!(peer_id, peer);
2476            }
2477            _ => unreachable!(),
2478        }
2479        match event!(peers) {
2480            PeerAction::Connect { peer_id, remote_addr } => {
2481                assert_eq!(peer_id, peer);
2482                assert_eq!(remote_addr, socket_addr);
2483            }
2484            _ => unreachable!(),
2485        }
2486
2487        let p = peers.peers.get(&peer).unwrap();
2488        assert_eq!(p.state, PeerConnectionState::PendingOut);
2489
2490        peers.remove_peer(peer);
2491
2492        match event!(peers) {
2493            PeerAction::PeerRemoved(peer_id) => {
2494                assert_eq!(peer_id, peer);
2495            }
2496            _ => unreachable!(),
2497        }
2498        match event!(peers) {
2499            PeerAction::Disconnect { peer_id, .. } => {
2500                assert_eq!(peer_id, peer);
2501            }
2502            _ => unreachable!(),
2503        }
2504
2505        let p = peers.peers.get(&peer).unwrap();
2506        assert_eq!(p.state, PeerConnectionState::PendingOut);
2507
2508        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
2509        let p = peers.peers.get(&peer).unwrap();
2510        assert_eq!(p.state, PeerConnectionState::PendingOut);
2511
2512        peers.on_active_session_gracefully_closed(peer);
2513        assert!(!peers.peers.contains_key(&peer));
2514    }
2515
2516    #[tokio::test]
2517    async fn test_fatal_outgoing_connection_error_trusted() {
2518        let peer = PeerId::random();
2519        let config = PeersConfig::test()
2520            .with_trusted_nodes(vec![TrustedPeer {
2521                host: Host::Ipv4(Ipv4Addr::new(127, 0, 1, 2)),
2522                tcp_port: 8008,
2523                udp_port: 8008,
2524                id: peer,
2525            }])
2526            .with_trusted_nodes_only(true);
2527        let mut peers = PeersManager::new(config);
2528        let socket_addr = peers.peers.get(&peer).unwrap().addr.tcp();
2529
2530        match event!(peers) {
2531            PeerAction::Connect { peer_id, remote_addr } => {
2532                assert_eq!(peer_id, peer);
2533                assert_eq!(remote_addr, socket_addr);
2534            }
2535            _ => unreachable!(),
2536        }
2537
2538        let p = peers.peers.get(&peer).unwrap();
2539        assert_eq!(p.state, PeerConnectionState::PendingOut);
2540
2541        assert_eq!(peers.num_outbound_connections(), 0);
2542
2543        let err = PendingSessionHandshakeError::Eth(EthStreamError::EthHandshakeError(
2544            EthHandshakeError::NonStatusMessageInHandshake,
2545        ));
2546        assert!(err.is_fatal_protocol_error());
2547
2548        peers.on_outgoing_pending_session_dropped(&socket_addr, &peer, &err);
2549        assert_eq!(peers.num_outbound_connections(), 0);
2550
2551        // try tmp ban peer
2552        match event!(peers) {
2553            PeerAction::BanPeer { peer_id } => {
2554                assert_eq!(peer_id, peer);
2555            }
2556            err => unreachable!("{err:?}"),
2557        }
2558
2559        // ensure we still have trusted peer
2560        assert!(peers.peers.contains_key(&peer));
2561
2562        // await for the ban to expire
2563        tokio::time::sleep(peers.backoff_durations.medium).await;
2564
2565        match event!(peers) {
2566            PeerAction::Connect { peer_id, remote_addr } => {
2567                assert_eq!(peer_id, peer);
2568                assert_eq!(remote_addr, socket_addr);
2569            }
2570            err => unreachable!("{err:?}"),
2571        }
2572    }
2573
2574    #[tokio::test]
2575    async fn test_outgoing_connection_error() {
2576        let peer = PeerId::random();
2577        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2578        let mut peers = PeersManager::default();
2579        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
2580
2581        match event!(peers) {
2582            PeerAction::PeerAdded(peer_id) => {
2583                assert_eq!(peer_id, peer);
2584            }
2585            _ => unreachable!(),
2586        }
2587        match event!(peers) {
2588            PeerAction::Connect { peer_id, remote_addr } => {
2589                assert_eq!(peer_id, peer);
2590                assert_eq!(remote_addr, socket_addr);
2591            }
2592            _ => unreachable!(),
2593        }
2594
2595        let p = peers.peers.get(&peer).unwrap();
2596        assert_eq!(p.state, PeerConnectionState::PendingOut);
2597
2598        assert_eq!(peers.num_outbound_connections(), 0);
2599
2600        peers.on_outgoing_connection_failure(
2601            &socket_addr,
2602            &peer,
2603            &io::Error::new(io::ErrorKind::ConnectionRefused, ""),
2604        );
2605
2606        assert_eq!(peers.num_outbound_connections(), 0);
2607    }
2608
2609    #[tokio::test]
2610    async fn test_outgoing_connection_gracefully_closed() {
2611        let peer = PeerId::random();
2612        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2613        let mut peers = PeersManager::default();
2614        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
2615
2616        match event!(peers) {
2617            PeerAction::PeerAdded(peer_id) => {
2618                assert_eq!(peer_id, peer);
2619            }
2620            _ => unreachable!(),
2621        }
2622        match event!(peers) {
2623            PeerAction::Connect { peer_id, remote_addr } => {
2624                assert_eq!(peer_id, peer);
2625                assert_eq!(remote_addr, socket_addr);
2626            }
2627            _ => unreachable!(),
2628        }
2629
2630        let p = peers.peers.get(&peer).unwrap();
2631        assert_eq!(p.state, PeerConnectionState::PendingOut);
2632
2633        assert_eq!(peers.num_outbound_connections(), 0);
2634
2635        peers.on_outgoing_pending_session_gracefully_closed(&peer);
2636
2637        assert_eq!(peers.num_outbound_connections(), 0);
2638        assert_eq!(peers.connection_info.num_pending_out, 0);
2639    }
2640
2641    #[tokio::test]
2642    async fn test_discovery_ban_list() {
2643        let ip = IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2));
2644        let socket_addr = SocketAddr::new(ip, 8008);
2645        let ban_list = BanList::new(vec![], vec![ip]);
2646        let config = PeersConfig::default().with_ban_list(ban_list);
2647        let mut peer_manager = PeersManager::new(config);
2648        peer_manager.add_peer(B512::default(), PeerAddr::from_tcp(socket_addr), None);
2649
2650        assert!(peer_manager.peers.is_empty());
2651    }
2652
2653    #[tokio::test]
2654    async fn test_on_pending_ban_list() {
2655        let ip = IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2));
2656        let socket_addr = SocketAddr::new(ip, 8008);
2657        let ban_list = BanList::new(vec![], vec![ip]);
2658        let config = PeersConfig::test().with_ban_list(ban_list);
2659        let mut peer_manager = PeersManager::new(config);
2660        let a = peer_manager.on_incoming_pending_session(socket_addr.ip());
2661        // because we have no active peers this should be fine for testings
2662        match a {
2663            Ok(_) => panic!(),
2664            Err(err) => match err {
2665                InboundConnectionError::IpBanned => {
2666                    assert_eq!(peer_manager.connection_info.num_pending_in, 0)
2667                }
2668                _ => unreachable!(),
2669            },
2670        }
2671    }
2672
2673    #[tokio::test]
2674    async fn test_on_active_inbound_ban_list() {
2675        let ip = IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2));
2676        let socket_addr = SocketAddr::new(ip, 8008);
2677        let given_peer_id = PeerId::random();
2678        let ban_list = BanList::new(vec![given_peer_id], vec![]);
2679        let config = PeersConfig::test().with_ban_list(ban_list);
2680        let mut peer_manager = PeersManager::new(config);
2681        assert!(peer_manager.on_incoming_pending_session(socket_addr.ip()).is_ok());
2682        // non-trusted nodes should also increase pending_in
2683        assert_eq!(peer_manager.connection_info.num_pending_in, 1);
2684        peer_manager.on_incoming_session_established(given_peer_id, socket_addr);
2685        // after the connection is established, the peer should be removed, the num_pending_in
2686        // should be decreased, and the num_inbound should not be increased
2687        assert_eq!(peer_manager.connection_info.num_pending_in, 0);
2688        assert_eq!(peer_manager.connection_info.num_inbound, 0);
2689
2690        let Some(PeerAction::DisconnectBannedIncoming { peer_id }) =
2691            peer_manager.queued_actions.pop_front()
2692        else {
2693            panic!()
2694        };
2695
2696        assert_eq!(peer_id, given_peer_id)
2697    }
2698
2699    #[test]
2700    fn test_connection_limits() {
2701        let mut info = ConnectionInfo::default();
2702        info.inc_in();
2703        assert_eq!(info.num_inbound, 1);
2704        assert_eq!(info.num_outbound, 0);
2705        assert!(info.has_in_capacity());
2706
2707        info.decr_in();
2708        assert_eq!(info.num_inbound, 0);
2709        assert_eq!(info.num_outbound, 0);
2710
2711        info.inc_out();
2712        assert_eq!(info.num_inbound, 0);
2713        assert_eq!(info.num_outbound, 1);
2714        assert!(info.has_out_capacity());
2715
2716        info.decr_out();
2717        assert_eq!(info.num_inbound, 0);
2718        assert_eq!(info.num_outbound, 0);
2719    }
2720
2721    #[test]
2722    fn test_connection_peer_state() {
2723        let mut info = ConnectionInfo::default();
2724        info.inc_in();
2725
2726        info.decr_state(PeerConnectionState::In);
2727        assert_eq!(info.num_inbound, 0);
2728        assert_eq!(info.num_outbound, 0);
2729
2730        info.inc_out();
2731
2732        info.decr_state(PeerConnectionState::Out);
2733        assert_eq!(info.num_inbound, 0);
2734        assert_eq!(info.num_outbound, 0);
2735    }
2736
2737    #[tokio::test]
2738    async fn test_trusted_peers_are_prioritized() {
2739        let trusted_peer = PeerId::random();
2740        let trusted_sock = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2741        let config = PeersConfig::test().with_trusted_nodes(vec![TrustedPeer {
2742            host: Host::Ipv4(Ipv4Addr::new(127, 0, 1, 2)),
2743            tcp_port: 8008,
2744            udp_port: 8008,
2745            id: trusted_peer,
2746        }]);
2747        let mut peers = PeersManager::new(config);
2748
2749        let basic_peer = PeerId::random();
2750        let basic_sock = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
2751        peers.add_peer(basic_peer, PeerAddr::from_tcp(basic_sock), None);
2752
2753        match event!(peers) {
2754            PeerAction::PeerAdded(peer_id) => {
2755                assert_eq!(peer_id, basic_peer);
2756            }
2757            _ => unreachable!(),
2758        }
2759        match event!(peers) {
2760            PeerAction::Connect { peer_id, remote_addr } => {
2761                assert_eq!(peer_id, trusted_peer);
2762                assert_eq!(remote_addr, trusted_sock);
2763            }
2764            _ => unreachable!(),
2765        }
2766        match event!(peers) {
2767            PeerAction::Connect { peer_id, remote_addr } => {
2768                assert_eq!(peer_id, basic_peer);
2769                assert_eq!(remote_addr, basic_sock);
2770            }
2771            _ => unreachable!(),
2772        }
2773    }
2774
2775    #[tokio::test]
2776    async fn test_connect_trusted_nodes_only() {
2777        let trusted_peer = PeerId::random();
2778        let trusted_sock = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
2779        let config = PeersConfig::test()
2780            .with_trusted_nodes(vec![TrustedPeer {
2781                host: Host::Ipv4(Ipv4Addr::new(127, 0, 1, 2)),
2782                tcp_port: 8008,
2783                udp_port: 8008,
2784                id: trusted_peer,
2785            }])
2786            .with_trusted_nodes_only(true);
2787        let mut peers = PeersManager::new(config);
2788
2789        let basic_peer = PeerId::random();
2790        let basic_sock = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
2791        peers.add_peer(basic_peer, PeerAddr::from_tcp(basic_sock), None);
2792
2793        match event!(peers) {
2794            PeerAction::PeerAdded(peer_id) => {
2795                assert_eq!(peer_id, basic_peer);
2796            }
2797            _ => unreachable!(),
2798        }
2799        match event!(peers) {
2800            PeerAction::Connect { peer_id, remote_addr } => {
2801                assert_eq!(peer_id, trusted_peer);
2802                assert_eq!(remote_addr, trusted_sock);
2803            }
2804            _ => unreachable!(),
2805        }
2806        poll_fn(|cx| {
2807            assert!(peers.poll(cx).is_pending());
2808            Poll::Ready(())
2809        })
2810        .await;
2811    }
2812
2813    #[tokio::test]
2814    async fn test_incoming_with_trusted_nodes_only() {
2815        let trusted_peer = PeerId::random();
2816        let config = PeersConfig::test()
2817            .with_trusted_nodes(vec![TrustedPeer {
2818                host: Host::Ipv4(Ipv4Addr::new(127, 0, 1, 2)),
2819                tcp_port: 8008,
2820                udp_port: 8008,
2821                id: trusted_peer,
2822            }])
2823            .with_trusted_nodes_only(true);
2824        let mut peers = PeersManager::new(config);
2825
2826        let basic_peer = PeerId::random();
2827        let basic_sock = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
2828        assert!(peers.on_incoming_pending_session(basic_sock.ip()).is_ok());
2829        // non-trusted nodes should also increase pending_in
2830        assert_eq!(peers.connection_info.num_pending_in, 1);
2831        peers.on_incoming_session_established(basic_peer, basic_sock);
2832        // after the connection is established, the peer should be removed, the num_pending_in
2833        // should be decreased, and the num_inbound mut not be increased
2834        assert_eq!(peers.connection_info.num_pending_in, 0);
2835        assert_eq!(peers.connection_info.num_inbound, 0);
2836
2837        let Some(PeerAction::DisconnectUntrustedIncoming { peer_id }) =
2838            peers.queued_actions.pop_front()
2839        else {
2840            panic!()
2841        };
2842        assert_eq!(basic_peer, peer_id);
2843        assert!(!peers.peers.contains_key(&basic_peer));
2844    }
2845
2846    #[tokio::test]
2847    async fn test_incoming_without_trusted_nodes_only() {
2848        let trusted_peer = PeerId::random();
2849        let config = PeersConfig::test()
2850            .with_trusted_nodes(vec![TrustedPeer {
2851                host: Host::Ipv4(Ipv4Addr::new(127, 0, 1, 2)),
2852                tcp_port: 8008,
2853                udp_port: 8008,
2854                id: trusted_peer,
2855            }])
2856            .with_trusted_nodes_only(false);
2857        let mut peers = PeersManager::new(config);
2858
2859        let basic_peer = PeerId::random();
2860        let basic_sock = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
2861        assert!(peers.on_incoming_pending_session(basic_sock.ip()).is_ok());
2862
2863        // non-trusted nodes should also increase pending_in
2864        assert_eq!(peers.connection_info.num_pending_in, 1);
2865        peers.on_incoming_session_established(basic_peer, basic_sock);
2866        // after the connection is established, the peer should be removed, the num_pending_in
2867        // should be decreased, and the num_inbound must be increased
2868        assert_eq!(peers.connection_info.num_pending_in, 0);
2869        assert_eq!(peers.connection_info.num_inbound, 1);
2870        assert!(peers.peers.contains_key(&basic_peer));
2871    }
2872
2873    #[tokio::test]
2874    async fn test_incoming_at_capacity() {
2875        let mut config = PeersConfig::test();
2876        config.connection_info.max_inbound = 1;
2877        let mut peers = PeersManager::new(config);
2878
2879        let peer = PeerId::random();
2880        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
2881        assert!(peers.on_incoming_pending_session(addr.ip()).is_ok());
2882
2883        peers.on_incoming_session_established(peer, addr);
2884
2885        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
2886        assert_eq!(
2887            peers.on_incoming_pending_session(addr.ip()).unwrap_err(),
2888            InboundConnectionError::ExceedsCapacity
2889        );
2890    }
2891
2892    #[tokio::test]
2893    async fn test_incoming_rate_limit() {
2894        let config = PeersConfig {
2895            incoming_ip_throttle_duration: Duration::from_millis(100),
2896            ..PeersConfig::test()
2897        };
2898        let mut peers = PeersManager::new(config);
2899
2900        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(168, 0, 1, 2)), 8009);
2901        assert!(peers.on_incoming_pending_session(addr.ip()).is_ok());
2902        assert_eq!(
2903            peers.on_incoming_pending_session(addr.ip()).unwrap_err(),
2904            InboundConnectionError::IpBanned
2905        );
2906
2907        peers.release_interval.reset_immediately();
2908        tokio::time::sleep(peers.incoming_ip_throttle_duration).await;
2909
2910        // await unban
2911        poll_fn(|cx| loop {
2912            if peers.poll(cx).is_pending() {
2913                return Poll::Ready(());
2914            }
2915        })
2916        .await;
2917
2918        assert!(peers.on_incoming_pending_session(addr.ip()).is_ok());
2919        assert_eq!(
2920            peers.on_incoming_pending_session(addr.ip()).unwrap_err(),
2921            InboundConnectionError::IpBanned
2922        );
2923    }
2924
2925    #[tokio::test]
2926    async fn test_tick() {
2927        let ip = IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2));
2928        let socket_addr = SocketAddr::new(ip, 8008);
2929        let config = PeersConfig::test();
2930        let mut peer_manager = PeersManager::new(config);
2931        let peer_id = PeerId::random();
2932        peer_manager.add_peer(peer_id, PeerAddr::from_tcp(socket_addr), None);
2933
2934        // Age the last tick directly so each reputation update sees a full elapsed second.
2935        peer_manager.last_tick -= Duration::from_secs(1);
2936        peer_manager.tick();
2937
2938        // still unconnected
2939        assert_eq!(peer_manager.peers.get_mut(&peer_id).unwrap().reputation, DEFAULT_REPUTATION);
2940
2941        // mark as connected
2942        peer_manager.peers.get_mut(&peer_id).unwrap().state = PeerConnectionState::Out;
2943
2944        peer_manager.last_tick -= Duration::from_secs(1);
2945        peer_manager.tick();
2946
2947        // still at default reputation
2948        assert_eq!(peer_manager.peers.get_mut(&peer_id).unwrap().reputation, DEFAULT_REPUTATION);
2949
2950        peer_manager.peers.get_mut(&peer_id).unwrap().reputation -= 1;
2951
2952        peer_manager.last_tick -= Duration::from_secs(1);
2953        peer_manager.tick();
2954
2955        // tick applied
2956        assert!(peer_manager.peers.get_mut(&peer_id).unwrap().reputation >= DEFAULT_REPUTATION);
2957    }
2958
2959    #[tokio::test]
2960    async fn test_remove_incoming_after_disconnect() {
2961        let peer_id = PeerId::random();
2962        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
2963        let mut peers = PeersManager::default();
2964
2965        peers.on_incoming_pending_session(addr.ip()).unwrap();
2966        peers.on_incoming_session_established(peer_id, addr);
2967        let peer = peers.peers.get(&peer_id).unwrap();
2968        assert_eq!(peer.state, PeerConnectionState::In);
2969        assert!(peer.remove_after_disconnect);
2970
2971        peers.on_active_session_gracefully_closed(peer_id);
2972        assert!(!peers.peers.contains_key(&peer_id))
2973    }
2974
2975    #[tokio::test]
2976    async fn test_keep_incoming_after_disconnect_if_discovered() {
2977        let peer_id = PeerId::random();
2978        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
2979        let mut peers = PeersManager::default();
2980
2981        peers.on_incoming_pending_session(addr.ip()).unwrap();
2982        peers.on_incoming_session_established(peer_id, addr);
2983        let peer = peers.peers.get(&peer_id).unwrap();
2984        assert_eq!(peer.state, PeerConnectionState::In);
2985        assert!(peer.remove_after_disconnect);
2986
2987        // trigger discovery manually while the peer is still connected
2988        peers.add_peer(peer_id, PeerAddr::from_tcp(addr), None);
2989
2990        peers.on_active_session_gracefully_closed(peer_id);
2991
2992        let peer = peers.peers.get(&peer_id).unwrap();
2993        assert_eq!(peer.state, PeerConnectionState::Idle);
2994        assert!(!peer.remove_after_disconnect);
2995    }
2996
2997    #[tokio::test]
2998    async fn test_peer_reconnect_after_graceful_close_respects_throttle() {
2999        let throttle_duration = Duration::from_millis(100);
3000        let config =
3001            PeersConfig { incoming_ip_throttle_duration: throttle_duration, ..PeersConfig::test() };
3002        let mut peers = PeersManager::new(config);
3003
3004        let peer_id = PeerId::random();
3005        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
3006
3007        // Add as regular peer
3008        peers.add_peer(peer_id, PeerAddr::from_tcp(addr), None);
3009
3010        match event!(peers) {
3011            PeerAction::PeerAdded(id) => assert_eq!(id, peer_id),
3012            _ => unreachable!(),
3013        }
3014
3015        match event!(peers) {
3016            PeerAction::Connect { .. } => {}
3017            _ => unreachable!(),
3018        }
3019
3020        // Simulate outbound connection established
3021        peers.on_active_outgoing_established(peer_id);
3022        assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::Out);
3023
3024        // Gracefully close the session
3025        peers.on_active_session_gracefully_closed(peer_id);
3026
3027        let peer = peers.peers.get(&peer_id).unwrap();
3028        assert_eq!(peer.state, PeerConnectionState::Idle);
3029        assert!(peer.backed_off);
3030
3031        // Verify the peer is in the backed_off_peers set
3032        assert!(peers.backed_off_peers.contains_key(&peer_id));
3033
3034        // Immediately try to poll - should not trigger any actions yet
3035        poll_fn(|cx| {
3036            assert!(peers.poll(cx).is_pending());
3037            Poll::Ready(())
3038        })
3039        .await;
3040
3041        // Peer should still be backed off
3042        assert!(peers.backed_off_peers.contains_key(&peer_id));
3043        assert!(peers.peers.get(&peer_id).unwrap().backed_off);
3044
3045        // Sleep for the throttle duration
3046        tokio::time::sleep(throttle_duration).await;
3047
3048        // After throttle duration, event! will poll until we get a Connect action
3049        match event!(peers) {
3050            PeerAction::Connect { peer_id: id, .. } => assert_eq!(id, peer_id),
3051            _ => unreachable!(),
3052        }
3053
3054        // After connection is initiated, peer should no longer be backed off
3055        assert!(!peers.backed_off_peers.contains_key(&peer_id));
3056        assert!(!peers.peers.get(&peer_id).unwrap().backed_off);
3057    }
3058
3059    #[tokio::test]
3060    async fn test_backed_off_peer_can_accept_incoming_connection() {
3061        let throttle_duration = Duration::from_millis(100);
3062        let config =
3063            PeersConfig { incoming_ip_throttle_duration: throttle_duration, ..PeersConfig::test() };
3064        let mut peers = PeersManager::new(config);
3065
3066        let peer_id = PeerId::random();
3067        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
3068
3069        // Add as regular peer
3070        peers.add_peer(peer_id, PeerAddr::from_tcp(addr), None);
3071
3072        match event!(peers) {
3073            PeerAction::PeerAdded(id) => assert_eq!(id, peer_id),
3074            _ => unreachable!(),
3075        }
3076
3077        match event!(peers) {
3078            PeerAction::Connect { .. } => {}
3079            _ => unreachable!(),
3080        }
3081
3082        // Simulate outbound connection established
3083        peers.on_active_outgoing_established(peer_id);
3084        assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::Out);
3085
3086        // Gracefully close the session - this will back off the peer
3087        peers.on_active_session_gracefully_closed(peer_id);
3088
3089        let peer = peers.peers.get(&peer_id).unwrap();
3090        assert_eq!(peer.state, PeerConnectionState::Idle);
3091        assert!(peer.backed_off);
3092        assert!(peers.backed_off_peers.contains_key(&peer_id));
3093
3094        // Now simulate an incoming connection from the backed-off peer
3095        // First, handle the incoming pending session
3096        assert!(peers.on_incoming_pending_session(addr.ip()).is_ok());
3097        assert_eq!(peers.connection_info.num_pending_in, 1);
3098
3099        // Establish the incoming session
3100        peers.on_incoming_session_established(peer_id, addr);
3101
3102        // Peer should have been added to incoming connections
3103        assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::In);
3104        assert_eq!(peers.connection_info.num_inbound, 1);
3105
3106        // Peer should still be backed off for outbound connections
3107        assert!(peers.backed_off_peers.contains_key(&peer_id));
3108        assert!(peers.peers.get(&peer_id).unwrap().backed_off);
3109
3110        // Verify we don't try to reconnect outbound while peer is backed off
3111        poll_fn(|cx| {
3112            assert!(peers.poll(cx).is_pending());
3113            Poll::Ready(())
3114        })
3115        .await;
3116
3117        // No outbound connection should be attempted while backed off
3118        assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::In);
3119
3120        // After throttle duration, the backoff should be cleared
3121        tokio::time::sleep(throttle_duration).await;
3122
3123        poll_fn(|cx| {
3124            let _ = peers.poll(cx);
3125            Poll::Ready(())
3126        })
3127        .await;
3128
3129        // Backoff should be cleared now
3130        assert!(!peers.backed_off_peers.contains_key(&peer_id));
3131        assert!(!peers.peers.get(&peer_id).unwrap().backed_off);
3132
3133        // Peer should still be in incoming state
3134        assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::In);
3135    }
3136
3137    #[tokio::test]
3138    async fn test_incoming_outgoing_already_connected() {
3139        let peer_id = PeerId::random();
3140        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
3141        let mut peers = PeersManager::default();
3142
3143        peers.on_incoming_pending_session(addr.ip()).unwrap();
3144        peers.add_peer(peer_id, PeerAddr::from_tcp(addr), None);
3145
3146        match event!(peers) {
3147            PeerAction::PeerAdded(_) => {}
3148            _ => unreachable!(),
3149        }
3150
3151        match event!(peers) {
3152            PeerAction::Connect { .. } => {}
3153            _ => unreachable!(),
3154        }
3155
3156        peers.on_incoming_session_established(peer_id, addr);
3157        peers.on_already_connected(Direction::Outgoing(peer_id));
3158        assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::In);
3159        assert_eq!(peers.connection_info.num_inbound, 1);
3160        assert_eq!(peers.connection_info.num_pending_out, 0);
3161        assert_eq!(peers.connection_info.num_pending_in, 0);
3162        assert_eq!(peers.connection_info.num_outbound, 0);
3163    }
3164
3165    #[tokio::test]
3166    async fn test_already_connected_incoming_outgoing_connection_error() {
3167        let peer_id = PeerId::random();
3168        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8009);
3169        let mut peers = PeersManager::default();
3170
3171        peers.on_incoming_pending_session(addr.ip()).unwrap();
3172        peers.add_peer(peer_id, PeerAddr::from_tcp(addr), None);
3173
3174        match event!(peers) {
3175            PeerAction::PeerAdded(_) => {}
3176            _ => unreachable!(),
3177        }
3178
3179        match event!(peers) {
3180            PeerAction::Connect { .. } => {}
3181            _ => unreachable!(),
3182        }
3183
3184        peers.on_incoming_session_established(peer_id, addr);
3185
3186        peers.on_outgoing_connection_failure(
3187            &addr,
3188            &peer_id,
3189            &io::Error::new(io::ErrorKind::ConnectionRefused, ""),
3190        );
3191        assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::In);
3192        assert_eq!(peers.connection_info.num_inbound, 1);
3193        assert_eq!(peers.connection_info.num_pending_out, 0);
3194        assert_eq!(peers.connection_info.num_pending_in, 0);
3195        assert_eq!(peers.connection_info.num_outbound, 0);
3196    }
3197
3198    #[tokio::test]
3199    async fn test_max_concurrent_dials() {
3200        let config = PeersConfig::default();
3201        let mut peer_manager = PeersManager::new(config);
3202        let ip = IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2));
3203        let peer_addr = PeerAddr::from_tcp(SocketAddr::new(ip, 8008));
3204        for _ in 0..peer_manager.connection_info.config.max_concurrent_outbound_dials * 2 {
3205            peer_manager.add_peer(PeerId::random(), peer_addr, None);
3206        }
3207
3208        peer_manager.fill_outbound_slots();
3209        let dials = peer_manager
3210            .queued_actions
3211            .iter()
3212            .filter(|ev| matches!(ev, PeerAction::Connect { .. }))
3213            .count();
3214        assert_eq!(dials, peer_manager.connection_info.config.max_concurrent_outbound_dials);
3215    }
3216
3217    #[tokio::test]
3218    async fn test_max_num_of_pending_dials() {
3219        let config = PeersConfig::default();
3220        let mut peer_manager = PeersManager::new(config);
3221        let ip = IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2));
3222        let peer_addr = PeerAddr::from_tcp(SocketAddr::new(ip, 8008));
3223
3224        // add more peers than allowed
3225        for _ in 0..peer_manager.connection_info.config.max_concurrent_outbound_dials * 2 {
3226            peer_manager.add_peer(PeerId::random(), peer_addr, None);
3227        }
3228
3229        for _ in 0..peer_manager.connection_info.config.max_concurrent_outbound_dials * 2 {
3230            match event!(peer_manager) {
3231                PeerAction::PeerAdded(_) => {}
3232                _ => unreachable!(),
3233            }
3234        }
3235
3236        for _ in 0..peer_manager.connection_info.config.max_concurrent_outbound_dials {
3237            match event!(peer_manager) {
3238                PeerAction::Connect { .. } => {}
3239                _ => unreachable!(),
3240            }
3241        }
3242
3243        // generate 'Connect' actions
3244        peer_manager.fill_outbound_slots();
3245
3246        // all dialed connections should be in 'PendingOut' state
3247        let dials = peer_manager.connection_info.num_pending_out;
3248        assert_eq!(dials, peer_manager.connection_info.config.max_concurrent_outbound_dials);
3249
3250        let num_pendingout_states = peer_manager
3251            .peers
3252            .iter()
3253            .filter(|(_, peer)| peer.state == PeerConnectionState::PendingOut)
3254            .map(|(peer_id, _)| *peer_id)
3255            .collect::<Vec<PeerId>>();
3256        assert_eq!(
3257            num_pendingout_states.len(),
3258            peer_manager.connection_info.config.max_concurrent_outbound_dials
3259        );
3260
3261        // establish dialed connections
3262        for peer_id in &num_pendingout_states {
3263            peer_manager.on_active_outgoing_established(*peer_id);
3264        }
3265
3266        // all dialed connections should now be in 'Out' state
3267        for peer_id in &num_pendingout_states {
3268            assert_eq!(peer_manager.peers.get(peer_id).unwrap().state, PeerConnectionState::Out);
3269        }
3270
3271        // no more pending outbound connections
3272        assert_eq!(peer_manager.connection_info.num_pending_out, 0);
3273    }
3274
3275    #[tokio::test]
3276    async fn test_connect() {
3277        let peer = PeerId::random();
3278        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
3279        let mut peers = PeersManager::default();
3280        peers.add_and_connect(peer, PeerAddr::from_tcp(socket_addr), None);
3281        assert_eq!(peers.peers.get(&peer).unwrap().state, PeerConnectionState::PendingOut);
3282
3283        match event!(peers) {
3284            PeerAction::Connect { peer_id, remote_addr } => {
3285                assert_eq!(peer_id, peer);
3286                assert_eq!(remote_addr, socket_addr);
3287            }
3288            _ => unreachable!(),
3289        }
3290
3291        let (record, _) = peers.peer_by_id(peer).unwrap();
3292        assert_eq!(record.tcp_addr(), socket_addr);
3293        assert_eq!(record.udp_addr(), socket_addr);
3294
3295        // connect again
3296        peers.add_and_connect(peer, PeerAddr::from_tcp(socket_addr), None);
3297
3298        let (record, _) = peers.peer_by_id(peer).unwrap();
3299        assert_eq!(record.tcp_addr(), socket_addr);
3300        assert_eq!(record.udp_addr(), socket_addr);
3301    }
3302
3303    #[tokio::test]
3304    async fn test_incoming_connection_from_banned() {
3305        let peer = PeerId::random();
3306        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
3307        let config = PeersConfig::test().with_max_inbound(3);
3308        let mut peers = PeersManager::new(config);
3309        peers.add_peer(peer, PeerAddr::from_tcp(socket_addr), None);
3310
3311        match event!(peers) {
3312            PeerAction::PeerAdded(peer_id) => {
3313                assert_eq!(peer_id, peer);
3314            }
3315            _ => unreachable!(),
3316        }
3317        match event!(peers) {
3318            PeerAction::Connect { peer_id, .. } => {
3319                assert_eq!(peer_id, peer);
3320            }
3321            _ => unreachable!(),
3322        }
3323
3324        poll_fn(|cx| {
3325            assert!(peers.poll(cx).is_pending());
3326            Poll::Ready(())
3327        })
3328        .await;
3329
3330        // simulate new connection drops with error
3331        loop {
3332            peers.on_active_session_dropped(
3333                &socket_addr,
3334                &peer,
3335                &EthStreamError::InvalidMessage(reth_eth_wire::message::MessageError::Invalid(
3336                    reth_eth_wire::EthVersion::Eth68,
3337                    reth_eth_wire::EthMessageID::Status,
3338                )),
3339            );
3340
3341            if peers.peers.get(&peer).unwrap().is_banned() {
3342                break;
3343            }
3344
3345            assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
3346            peers.on_incoming_session_established(peer, socket_addr);
3347
3348            match event!(peers) {
3349                PeerAction::Connect { peer_id, .. } => {
3350                    assert_eq!(peer_id, peer);
3351                }
3352                _ => unreachable!(),
3353            }
3354        }
3355
3356        assert!(peers.peers.get(&peer).unwrap().is_banned());
3357
3358        // fill all incoming slots
3359        for _ in 0..peers.connection_info.config.max_inbound {
3360            assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
3361            peers.on_incoming_session_established(peer, socket_addr);
3362
3363            match event!(peers) {
3364                PeerAction::DisconnectBannedIncoming { peer_id } => {
3365                    assert_eq!(peer_id, peer);
3366                }
3367                _ => unreachable!(),
3368            }
3369        }
3370
3371        poll_fn(|cx| {
3372            assert!(peers.poll(cx).is_pending());
3373            Poll::Ready(())
3374        })
3375        .await;
3376
3377        assert_eq!(peers.connection_info.num_inbound, 0);
3378
3379        let new_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 3)), 8008);
3380
3381        // Assert we can still accept new connections
3382        assert!(peers.on_incoming_pending_session(new_addr.ip()).is_ok());
3383        assert_eq!(peers.connection_info.num_pending_in, 1);
3384
3385        // the triggered DisconnectBannedIncoming will result in dropped connections, assert that
3386        // connection info is updated via the peer's state which would be a noop here since the
3387        // banned peer's state is idle
3388        peers.on_active_session_gracefully_closed(peer);
3389        assert_eq!(peers.connection_info.num_inbound, 0);
3390    }
3391
3392    #[tokio::test]
3393    async fn test_add_pending_connect() {
3394        let peer = PeerId::random();
3395        let socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
3396        let mut peers = PeersManager::default();
3397        peers.add_and_connect(peer, PeerAddr::from_tcp(socket_addr), None);
3398        assert_eq!(peers.peers.get(&peer).unwrap().state, PeerConnectionState::PendingOut);
3399        assert_eq!(peers.connection_info.num_pending_out, 1);
3400    }
3401
3402    #[tokio::test]
3403    async fn test_dns_updates_peer_address() {
3404        let peer_id = PeerId::random();
3405        let initial_socket = SocketAddr::new("1.1.1.1".parse::<IpAddr>().unwrap(), 8008);
3406        let updated_ip = "2.2.2.2".parse::<IpAddr>().unwrap();
3407
3408        let trusted = TrustedPeer {
3409            host: url::Host::Ipv4("2.2.2.2".parse().unwrap()),
3410            tcp_port: 8008,
3411            udp_port: 8008,
3412            id: peer_id,
3413        };
3414
3415        let config = PeersConfig::test().with_trusted_nodes(vec![trusted.clone()]);
3416        let mut manager = PeersManager::new(config);
3417        manager
3418            .trusted_peers_resolver
3419            .set_interval(tokio::time::interval(Duration::from_millis(1)));
3420
3421        manager.peers.insert(
3422            peer_id,
3423            Peer::trusted(PeerAddr::new_with_ports(initial_socket.ip(), 8008, Some(8008))),
3424        );
3425
3426        for _ in 0..100 {
3427            let _ = event!(manager);
3428            if manager.peers.get(&peer_id).unwrap().addr.tcp().ip() == updated_ip {
3429                break;
3430            }
3431            tokio::time::sleep(Duration::from_millis(10)).await;
3432        }
3433
3434        let updated_peer = manager.peers.get(&peer_id).unwrap();
3435        assert_eq!(updated_peer.addr.tcp().ip(), updated_ip);
3436    }
3437
3438    #[tokio::test]
3439    async fn test_ip_filter_blocks_inbound_connection() {
3440        use reth_net_banlist::IpFilter;
3441        use std::net::IpAddr;
3442
3443        // Create a filter that only allows 192.168.0.0/16
3444        let ip_filter = IpFilter::from_cidr_string("192.168.0.0/16").unwrap();
3445        let config = PeersConfig::test().with_ip_filter(ip_filter);
3446        let mut peers = PeersManager::new(config);
3447
3448        // Try to connect from an allowed IP
3449        let allowed_ip: IpAddr = "192.168.1.100".parse().unwrap();
3450        assert!(peers.on_incoming_pending_session(allowed_ip).is_ok());
3451
3452        // Try to connect from a disallowed IP
3453        let disallowed_ip: IpAddr = "10.0.0.1".parse().unwrap();
3454        assert!(peers.on_incoming_pending_session(disallowed_ip).is_err());
3455    }
3456
3457    #[tokio::test]
3458    async fn test_ip_filter_blocks_outbound_connection() {
3459        use reth_net_banlist::IpFilter;
3460        use std::net::SocketAddr;
3461
3462        // Create a filter that only allows 192.168.0.0/16
3463        let ip_filter = IpFilter::from_cidr_string("192.168.0.0/16").unwrap();
3464        let config = PeersConfig::test().with_ip_filter(ip_filter);
3465        let mut peers = PeersManager::new(config);
3466
3467        let peer_id = PeerId::new([1; 64]);
3468
3469        // Try to add a peer with an allowed IP
3470        let allowed_addr: SocketAddr = "192.168.1.100:30303".parse().unwrap();
3471        peers.add_peer(peer_id, PeerAddr::from_tcp(allowed_addr), None);
3472        assert!(peers.peers.contains_key(&peer_id));
3473
3474        // Try to add a peer with a disallowed IP
3475        let peer_id2 = PeerId::new([2; 64]);
3476        let disallowed_addr: SocketAddr = "10.0.0.1:30303".parse().unwrap();
3477        peers.add_peer(peer_id2, PeerAddr::from_tcp(disallowed_addr), None);
3478        assert!(!peers.peers.contains_key(&peer_id2));
3479    }
3480
3481    #[tokio::test]
3482    async fn test_ip_filter_ipv6() {
3483        use reth_net_banlist::IpFilter;
3484        use std::net::IpAddr;
3485
3486        // Create a filter that only allows IPv6 range 2001:db8::/32
3487        let ip_filter = IpFilter::from_cidr_string("2001:db8::/32").unwrap();
3488        let config = PeersConfig::test().with_ip_filter(ip_filter);
3489        let mut peers = PeersManager::new(config);
3490
3491        // Try to connect from an allowed IPv6 address
3492        let allowed_ip: IpAddr = "2001:db8::1".parse().unwrap();
3493        assert!(peers.on_incoming_pending_session(allowed_ip).is_ok());
3494
3495        // Try to connect from a disallowed IPv6 address
3496        let disallowed_ip: IpAddr = "2001:db9::1".parse().unwrap();
3497        assert!(peers.on_incoming_pending_session(disallowed_ip).is_err());
3498    }
3499
3500    #[tokio::test]
3501    async fn test_ip_filter_multiple_ranges() {
3502        use reth_net_banlist::IpFilter;
3503        use std::net::IpAddr;
3504
3505        // Create a filter that allows multiple ranges
3506        let ip_filter = IpFilter::from_cidr_string("192.168.0.0/16,10.0.0.0/8").unwrap();
3507        let config = PeersConfig::test().with_ip_filter(ip_filter);
3508        let mut peers = PeersManager::new(config);
3509
3510        // Try IPs from both allowed ranges
3511        let ip1: IpAddr = "192.168.1.1".parse().unwrap();
3512        let ip2: IpAddr = "10.5.10.20".parse().unwrap();
3513        assert!(peers.on_incoming_pending_session(ip1).is_ok());
3514        assert!(peers.on_incoming_pending_session(ip2).is_ok());
3515
3516        // Try IP from disallowed range
3517        let disallowed_ip: IpAddr = "172.16.0.1".parse().unwrap();
3518        assert!(peers.on_incoming_pending_session(disallowed_ip).is_err());
3519    }
3520
3521    #[tokio::test]
3522    async fn test_ip_filter_no_restriction() {
3523        use reth_net_banlist::IpFilter;
3524        use std::net::IpAddr;
3525
3526        // Create a filter with no restrictions (allow all)
3527        let ip_filter = IpFilter::allow_all();
3528        let config = PeersConfig::test().with_ip_filter(ip_filter);
3529        let mut peers = PeersManager::new(config);
3530
3531        // All IPs should be allowed
3532        let ip1: IpAddr = "192.168.1.1".parse().unwrap();
3533        let ip2: IpAddr = "10.0.0.1".parse().unwrap();
3534        let ip3: IpAddr = "8.8.8.8".parse().unwrap();
3535        assert!(peers.on_incoming_pending_session(ip1).is_ok());
3536        assert!(peers.on_incoming_pending_session(ip2).is_ok());
3537        assert!(peers.on_incoming_pending_session(ip3).is_ok());
3538    }
3539
3540    #[tokio::test]
3541    async fn test_best_unconnected_prefers_fork_id_as_tiebreaker() {
3542        let mut peers = PeersManager::default();
3543        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1)), 8008);
3544
3545        let fork_id = ForkId { hash: ForkHash([0xaa, 0xbb, 0xcc, 0xdd]), next: 0 };
3546
3547        // add two peers with equal reputation, only one has a fork_id
3548        let no_fork = PeerId::random();
3549        peers.add_peer(no_fork, PeerAddr::from_tcp(addr), None);
3550
3551        let with_fork = PeerId::random();
3552        peers.add_peer(with_fork, PeerAddr::from_tcp(addr), None);
3553        peers.peers.get_mut(&with_fork).unwrap().fork_id = Some(Box::new(fork_id));
3554
3555        let (best_id, _) = peers.best_unconnected().unwrap();
3556        assert_eq!(best_id, with_fork, "fork_id should break tie when reputation is equal");
3557    }
3558
3559    #[tokio::test]
3560    async fn test_add_trusted_peer_node_resolves_absent_peer() {
3561        let peer_id = PeerId::random();
3562        let trusted = TrustedPeer {
3563            host: url::Host::Domain("example.invalid".to_string()),
3564            tcp_port: 30303,
3565            udp_port: 30303,
3566            id: peer_id,
3567        };
3568
3569        let mut manager = PeersManager::default();
3570        manager.add_trusted_peer_node(trusted);
3571
3572        assert!(manager.trusted_peer_ids.contains(&peer_id));
3573        assert_eq!(manager.trusted_peers_resolver.trusted_peers.len(), 1);
3574        assert!(!manager.peers.contains_key(&peer_id));
3575
3576        let resolved = NodeRecord {
3577            address: "10.0.0.1".parse::<IpAddr>().unwrap(),
3578            tcp_port: 30303,
3579            udp_port: 30303,
3580            id: peer_id,
3581        };
3582        manager.on_resolved_peer(peer_id, resolved);
3583
3584        let peer = manager.peers.get(&peer_id).expect("peer should be added after resolution");
3585        assert!(peer.kind.is_trusted());
3586        assert_eq!(peer.addr.tcp().ip(), "10.0.0.1".parse::<IpAddr>().unwrap());
3587    }
3588
3589    #[tokio::test]
3590    async fn test_removed_trusted_peer_ignores_late_resolution() {
3591        let peer_id = PeerId::random();
3592        let trusted = TrustedPeer {
3593            host: url::Host::Domain("example.invalid".to_string()),
3594            tcp_port: 30303,
3595            udp_port: 30303,
3596            id: peer_id,
3597        };
3598
3599        let mut manager = PeersManager::default();
3600        manager.add_trusted_peer_node(trusted);
3601        manager.remove_peer_from_trusted_set(peer_id);
3602
3603        let resolved = NodeRecord {
3604            address: "10.0.0.1".parse::<IpAddr>().unwrap(),
3605            tcp_port: 30303,
3606            udp_port: 30303,
3607            id: peer_id,
3608        };
3609        manager.on_resolved_peer(peer_id, resolved);
3610
3611        assert!(!manager.peers.contains_key(&peer_id));
3612        assert!(!manager.trusted_peer_ids.contains(&peer_id));
3613    }
3614
3615    #[tokio::test]
3616    async fn test_rotation_disconnects_eligible_peer() {
3617        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
3618        let peer = PeerId::random();
3619        let mut peers = PeersManager::new(
3620            PeersConfig::test()
3621                .with_max_outbound(1)
3622                .with_peer_rotation_interval(Some(Duration::from_millis(50))),
3623        );
3624
3625        peers.add_peer(peer, PeerAddr::from_tcp(addr), None);
3626        match event!(peers) {
3627            PeerAction::PeerAdded(_) => {}
3628            _ => unreachable!(),
3629        }
3630        match event!(peers) {
3631            PeerAction::Connect { .. } => {}
3632            _ => unreachable!(),
3633        }
3634        peers.on_active_outgoing_established(peer);
3635
3636        set_connected_at(
3637            &mut peers,
3638            peer,
3639            std::time::Instant::now() - Duration::from_secs(11 * 60),
3640        );
3641
3642        tokio::time::sleep(Duration::from_millis(100)).await;
3643
3644        match event!(peers) {
3645            PeerAction::Disconnect { peer_id, reason } => {
3646                assert_eq!(peer_id, peer);
3647                assert_eq!(reason, Some(DisconnectReason::UselessPeer));
3648            }
3649            _ => unreachable!(),
3650        }
3651    }
3652
3653    #[tokio::test]
3654    async fn test_rotation_then_remove_uses_active_session_removal() {
3655        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 7)), 8008);
3656        let peer = PeerId::random();
3657        let mut peers = PeersManager::new(
3658            PeersConfig::test()
3659                .with_max_outbound(1)
3660                .with_peer_rotation_interval(Some(Duration::from_millis(50))),
3661        );
3662
3663        peers.add_peer(peer, PeerAddr::from_tcp(addr), None);
3664        match event!(peers) {
3665            PeerAction::PeerAdded(_) => {}
3666            _ => unreachable!(),
3667        }
3668        match event!(peers) {
3669            PeerAction::Connect { .. } => {}
3670            _ => unreachable!(),
3671        }
3672        peers.on_active_outgoing_established(peer);
3673        set_connected_at(
3674            &mut peers,
3675            peer,
3676            std::time::Instant::now() - Duration::from_secs(11 * 60),
3677        );
3678
3679        tokio::time::sleep(Duration::from_millis(100)).await;
3680
3681        match event!(peers) {
3682            PeerAction::Disconnect { peer_id, reason } => {
3683                assert_eq!(peer_id, peer);
3684                assert_eq!(reason, Some(DisconnectReason::UselessPeer));
3685            }
3686            _ => unreachable!(),
3687        }
3688        assert_eq!(peers.peers.get(&peer).unwrap().state, PeerConnectionState::Out);
3689        assert_eq!(peers.connection_info.num_outbound, 1);
3690
3691        peers.remove_peer(peer);
3692        match event!(peers) {
3693            PeerAction::PeerRemoved(peer_id) => assert_eq!(peer_id, peer),
3694            _ => unreachable!(),
3695        }
3696        match event!(peers) {
3697            PeerAction::Disconnect { peer_id, reason } => {
3698                assert_eq!(peer_id, peer);
3699                assert_eq!(reason, Some(DisconnectReason::DisconnectRequested));
3700            }
3701            _ => unreachable!(),
3702        }
3703
3704        let p = peers.peers.get(&peer).unwrap();
3705        assert_eq!(p.state, PeerConnectionState::DisconnectingOut);
3706        assert!(p.remove_after_disconnect);
3707        assert_eq!(peers.connection_info.num_outbound, 1);
3708
3709        peers.on_active_session_gracefully_closed(peer);
3710        assert_eq!(peers.connection_info.num_outbound, 0);
3711        assert!(!peers.peers.contains_key(&peer));
3712    }
3713
3714    #[tokio::test]
3715    async fn test_rotation_skips_trusted_peers() {
3716        let _addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 3)), 8008);
3717        let peer = PeerId::random();
3718        let trusted = TrustedPeer {
3719            host: Host::Ipv4(Ipv4Addr::new(127, 0, 1, 3)),
3720            tcp_port: 8008,
3721            udp_port: 8008,
3722            id: peer,
3723        };
3724        let mut peers = PeersManager::new(
3725            PeersConfig::test()
3726                .with_max_outbound(1)
3727                .with_trusted_nodes(vec![trusted])
3728                .with_peer_rotation_interval(Some(Duration::from_millis(50))),
3729        );
3730
3731        match event!(peers) {
3732            PeerAction::Connect { .. } => {}
3733            _ => unreachable!(),
3734        }
3735        peers.on_active_outgoing_established(peer);
3736        set_connected_at(
3737            &mut peers,
3738            peer,
3739            std::time::Instant::now() - Duration::from_secs(11 * 60),
3740        );
3741
3742        tokio::time::sleep(Duration::from_millis(100)).await;
3743
3744        poll_fn(|cx| {
3745            assert!(peers.poll(cx).is_pending(), "trusted peer must not be rotated");
3746            Poll::Ready(())
3747        })
3748        .await;
3749    }
3750
3751    #[tokio::test]
3752    async fn test_rotation_skips_recently_connected_peers() {
3753        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 4)), 8008);
3754        let peer = PeerId::random();
3755        let mut peers = PeersManager::new(
3756            PeersConfig::test()
3757                .with_max_outbound(1)
3758                .with_peer_rotation_interval(Some(Duration::from_millis(50))),
3759        );
3760
3761        peers.add_peer(peer, PeerAddr::from_tcp(addr), None);
3762        match event!(peers) {
3763            PeerAction::PeerAdded(_) => {}
3764            _ => unreachable!(),
3765        }
3766        match event!(peers) {
3767            PeerAction::Connect { .. } => {}
3768            _ => unreachable!(),
3769        }
3770        peers.on_active_outgoing_established(peer);
3771
3772        tokio::time::sleep(Duration::from_millis(100)).await;
3773
3774        poll_fn(|cx| {
3775            assert!(peers.poll(cx).is_pending(), "recently connected peer must not be rotated");
3776            Poll::Ready(())
3777        })
3778        .await;
3779    }
3780
3781    #[tokio::test]
3782    async fn test_rotation_skips_when_slots_available() {
3783        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 5)), 8008);
3784        let peer = PeerId::random();
3785        let mut peers = PeersManager::new(
3786            PeersConfig::test()
3787                .with_max_outbound(2)
3788                .with_peer_rotation_interval(Some(Duration::from_millis(50))),
3789        );
3790
3791        peers.add_peer(peer, PeerAddr::from_tcp(addr), None);
3792        match event!(peers) {
3793            PeerAction::PeerAdded(_) => {}
3794            _ => unreachable!(),
3795        }
3796        match event!(peers) {
3797            PeerAction::Connect { .. } => {}
3798            _ => unreachable!(),
3799        }
3800        peers.on_active_outgoing_established(peer);
3801        set_connected_at(
3802            &mut peers,
3803            peer,
3804            std::time::Instant::now() - Duration::from_secs(11 * 60),
3805        );
3806
3807        tokio::time::sleep(Duration::from_millis(100)).await;
3808
3809        poll_fn(|cx| {
3810            assert!(
3811                peers.poll(cx).is_pending(),
3812                "rotation must not fire when outbound slots are available"
3813            );
3814            Poll::Ready(())
3815        })
3816        .await;
3817    }
3818
3819    #[tokio::test]
3820    async fn test_rotation_disconnects_inbound_peer() {
3821        let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 6)), 8008);
3822        let peer = PeerId::random();
3823        let mut peers = PeersManager::new(
3824            PeersConfig::test()
3825                .with_max_inbound(1)
3826                .with_peer_rotation_interval(Some(Duration::from_millis(50))),
3827        );
3828
3829        assert!(peers.on_incoming_pending_session(addr.ip()).is_ok());
3830        peers.on_incoming_session_established(peer, addr);
3831
3832        match event!(peers) {
3833            PeerAction::PeerAdded(_) => {}
3834            _ => unreachable!(),
3835        }
3836
3837        set_connected_at(
3838            &mut peers,
3839            peer,
3840            std::time::Instant::now() - Duration::from_secs(11 * 60),
3841        );
3842
3843        tokio::time::sleep(Duration::from_millis(100)).await;
3844
3845        match event!(peers) {
3846            PeerAction::Disconnect { peer_id, reason } => {
3847                assert_eq!(peer_id, peer);
3848                assert_eq!(reason, Some(DisconnectReason::UselessPeer));
3849            }
3850            _ => unreachable!(),
3851        }
3852    }
3853
3854    // Ensures explicitly added peers are dialed right away, while discovered peers wait for the
3855    // next refill tick.
3856    #[tokio::test]
3857    async fn test_requested_peer_is_dialed_before_refill_tick() {
3858        let mut peers = PeersManager::new(PeersConfig {
3859            refill_slots_interval: Duration::from_secs(60),
3860            ..PeersConfig::test()
3861        });
3862        // consume the initial refill tick
3863        assert!(poll_now(&mut peers).is_pending());
3864
3865        let discovered = PeerId::random();
3866        let discovered_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 1)), 8008);
3867        peers.add_peer(discovered, PeerAddr::from_tcp(discovered_addr), None);
3868        assert!(
3869            matches!(poll_now(&mut peers), Poll::Ready(PeerAction::PeerAdded(id)) if id == discovered)
3870        );
3871        assert!(poll_now(&mut peers).is_pending());
3872
3873        let requested = PeerId::random();
3874        let requested_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 2)), 8008);
3875        peers.handle().add_peer(requested, requested_addr);
3876        assert!(
3877            matches!(poll_now(&mut peers), Poll::Ready(PeerAction::PeerAdded(id)) if id == requested)
3878        );
3879
3880        // the requested peer fills all free outbound slots, including one for the discovered peer
3881        let mut dialed = Vec::new();
3882        while let Poll::Ready(PeerAction::Connect { peer_id, .. }) = poll_now(&mut peers) {
3883            dialed.push(peer_id);
3884        }
3885        dialed.sort();
3886        let mut expected = vec![discovered, requested];
3887        expected.sort();
3888        assert_eq!(dialed, expected);
3889    }
3890}