1use 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#[derive(Debug)]
51pub struct PeersManager {
52 peers: HashMap<PeerId, Peer, FbBuildHasher<64>>,
54 trusted_peer_ids: HashSet<PeerId, FbBuildHasher<64>>,
59 trusted_peers_resolver: TrustedPeersResolver,
62 manager_tx: mpsc::UnboundedSender<PeerCommand>,
64 handle_rx: UnboundedReceiverStream<PeerCommand>,
66 queued_actions: VecDeque<PeerAction>,
68 refill_slots_interval: Interval,
70 reputation_weights: ReputationChangeWeights,
72 connection_info: ConnectionInfo,
74 ban_list: BanList,
76 backed_off_peers: HashMap<PeerId, std::time::Instant, FbBuildHasher<64>>,
78 release_interval: Interval,
80 ban_duration: Duration,
82 backoff_durations: PeerBackoffDurations,
85 trusted_nodes_only: bool,
88 last_tick: Instant,
90 max_backoff_count: u8,
92 net_connection_state: NetworkConnectionState,
94 incoming_ip_throttle_duration: Duration,
96 ip_filter: reth_net_banlist::IpFilter,
98 enforce_enr_fork_id: bool,
101 peer_rotation_sleep: Option<Pin<Box<Sleep>>>,
104 peer_rotation_mean: Option<Duration>,
106}
107
108impl PeersManager {
109 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 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 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), ),
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 pub(crate) fn handle(&self) -> PeersHandle {
216 PeersHandle::new(self.manager_tx.clone())
217 }
218
219 pub(crate) const fn enforce_enr_fork_id(&self) -> bool {
221 self.enforce_enr_fork_id
222 }
223
224 #[inline]
226 pub(crate) fn num_known_peers(&self) -> usize {
227 self.peers.len()
228 }
229
230 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 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 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 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 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 #[inline]
291 pub(crate) const fn num_inbound_connections(&self) -> usize {
292 self.connection_info.num_inbound
293 }
294
295 #[inline]
297 pub(crate) const fn num_outbound_connections(&self) -> usize {
298 self.connection_info.num_outbound
299 }
300
301 #[inline]
303 pub(crate) const fn num_pending_outbound_connections(&self) -> usize {
304 self.connection_info.num_pending_out
305 }
306
307 #[inline]
309 pub(crate) fn num_backed_off_peers(&self) -> usize {
310 self.backed_off_peers.len()
311 }
312
313 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 pub(crate) fn on_incoming_pending_session(
322 &mut self,
323 addr: IpAddr,
324 ) -> Result<(), InboundConnectionError> {
325 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 if !self.connection_info.has_in_capacity() {
337 if self.trusted_peer_ids.is_empty() {
338 return Err(InboundConnectionError::ExceedsCapacity)
341 }
342
343 let num_idle_trusted_peers = self.num_idle_trusted_peers();
347 if num_idle_trusted_peers <= self.trusted_peer_ids.len() {
348 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 return Err(InboundConnectionError::ExceedsCapacity)
359 }
360
361 if !self.connection_info.has_in_pending_capacity() {
363 return Err(InboundConnectionError::ExceedsCapacity)
364 }
365
366 self.throttle_incoming_ip(addr);
368
369 self.connection_info.inc_pending_in();
370 Ok(())
371 }
372
373 pub(crate) const fn on_incoming_pending_session_rejected_internally(&mut self) {
376 self.connection_info.decr_pending_in();
377 }
378
379 pub(crate) const fn on_incoming_pending_session_gracefully_closed(&mut self) {
381 self.connection_info.decr_pending_in()
382 }
383
384 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 pub(crate) fn on_incoming_session_established(&mut self, peer_id: PeerId, addr: SocketAddr) {
409 self.connection_info.decr_pending_in();
410
411 if self.ban_list.is_banned_peer(&peer_id) {
414 self.queued_actions.push_back(PeerAction::DisconnectBannedIncoming { peer_id });
415 return
416 }
417
418 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 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 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 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 self.connection_info.inc_in();
460
461 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 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 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 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 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 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 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 fn tick(&mut self) {
518 let now = Instant::now();
519 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 for peer in self.peers.iter_mut().filter(|(_, peer)| peer.state.is_connected()) {
528 if peer.1.reputation < DEFAULT_REPUTATION {
531 peer.1.reputation += secs_since_last_tick;
532 }
533 }
534 }
535
536 pub(crate) fn get_reputation(&self, peer_id: &PeerId) -> Option<i32> {
538 self.peers.get(peer_id).map(|peer| peer.reputation)
539 }
540
541 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 if matches!(
564 rep,
565 ReputationChangeKind::Dropped |
566 ReputationChangeKind::BadAnnouncement |
567 ReputationChangeKind::Timeout |
568 ReputationChangeKind::AlreadySeenTransaction
569 ) {
570 return
571 }
572
573 if reputation_change < MAX_TRUSTED_PEER_REPUTATION_CHANGE {
575 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 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 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 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 entry.remove();
633 self.queued_actions.push_back(PeerAction::PeerRemoved(peer_id));
634 } else {
635 let peer = entry.get_mut();
636 peer.severe_backoff_counter = 0;
640 peer.state = PeerConnectionState::Idle;
641 peer.mark_disconnected();
642
643 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 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 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 pub(crate) fn on_outgoing_connection_failure(
688 &mut self,
689 remote_addr: &SocketAddr,
690 peer_id: &PeerId,
691 err: &io::Error,
692 ) {
693 if let Some(peer) = self.peers.get(peer_id) {
697 if peer.state.is_incoming() {
698 return
700 }
701
702 if peer.is_trusted() && is_connection_failed_reputation(peer.reputation) {
703 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 if let Entry::Occupied(mut entry) = self.peers.entry(*peer_id) {
726 self.connection_info.decr_state(entry.get().state);
727 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 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 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 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 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 backoff_until = Some(backoff_time);
777 }
778 } else {
779 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 remove_peer = true;
795 }
796 }
797
798 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 self.backoff_peer_until(*peer_id, backoff_until);
806 }
807 }
808
809 self.fill_outbound_slots();
810 }
811
812 pub(crate) const fn on_already_connected(&mut self, direction: Direction) {
818 match direction {
819 Direction::Incoming => {
820 self.connection_info.decr_pending_in();
822 }
823 Direction::Outgoing(_) => {
824 }
827 }
828 }
829
830 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 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 #[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 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 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 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 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 self.trusted_peer_ids.insert(peer_id);
914 }
915 }
916
917 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 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 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 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 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 #[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 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 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 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 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 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 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 maybe_better.1.is_trusted() || maybe_better.1.is_static() {
1098 return Some((*maybe_better.0, maybe_better.1))
1099 }
1100
1101 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 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 fn fill_outbound_slots(&mut self) {
1168 self.tick();
1169
1170 if !self.net_connection_state.is_active() {
1171 return
1173 }
1174
1175 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 pub const fn on_network_state_change(&mut self, state: NetworkConnectionState) {
1220 self.net_connection_state = state;
1221 }
1222
1223 pub const fn connection_state(&self) -> &NetworkConnectionState {
1225 &self.net_connection_state
1226 }
1227
1228 pub const fn on_shutdown(&mut self) {
1230 self.net_connection_state = NetworkConnectionState::ShuttingDown;
1231 }
1232
1233 pub fn poll(&mut self, cx: &mut Context<'_>) -> Poll<PeerAction> {
1238 loop {
1239 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 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 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#[derive(Debug, Clone, PartialEq, Eq, Default)]
1325pub struct ConnectionInfo {
1326 num_outbound: usize,
1328 num_pending_out: usize,
1330 num_inbound: usize,
1332 num_pending_in: usize,
1334 config: ConnectionsConfig,
1336}
1337
1338impl ConnectionInfo {
1341 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 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 const fn is_outbound_at_capacity(&self) -> bool {
1354 self.num_outbound >= self.config.max_outbound
1355 }
1356
1357 const fn is_inbound_at_capacity(&self) -> bool {
1359 self.num_inbound >= self.config.max_inbound
1360 }
1361
1362 const fn has_in_capacity(&self) -> bool {
1364 self.num_inbound < self.config.max_inbound
1365 }
1366
1367 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#[derive(Debug)]
1416pub enum PeerAction {
1417 Connect {
1419 peer_id: PeerId,
1421 remote_addr: SocketAddr,
1423 },
1424 Disconnect {
1426 peer_id: PeerId,
1428 reason: Option<DisconnectReason>,
1430 },
1431 DisconnectBannedIncoming {
1434 peer_id: PeerId,
1436 },
1437 DisconnectUntrustedIncoming {
1439 peer_id: PeerId,
1441 },
1442 DiscoveryBanPeerId {
1444 peer_id: PeerId,
1446 ip_addr: IpAddr,
1448 },
1449 DiscoveryBanIp {
1451 ip_addr: IpAddr,
1453 },
1454 BanPeer {
1456 peer_id: PeerId,
1458 },
1459 UnBanPeer {
1461 peer_id: PeerId,
1463 },
1464 PeerAdded(PeerId),
1466 PeerRemoved(PeerId),
1468}
1469
1470#[derive(Debug, Error, PartialEq, Eq)]
1472pub enum InboundConnectionError {
1473 IpBanned,
1475 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1487pub enum BackoffReason {
1488 TooManyPeers,
1490 GracefulClose,
1492 ConnectionError,
1494}
1495
1496impl BackoffReason {
1497 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
1506fn 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 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 peer_struct.severe_backoff_counter = 1;
1895
1896 let now = std::time::Instant::now();
1897
1898 peer_struct.severe_backoff_counter += 1;
1900 let backoff_time = peers
1902 .backoff_durations
1903 .backoff_until(BackoffKind::High, peer_struct.severe_backoff_counter);
1904
1905 let backoff_duration = std::time::Duration::new(30 * 60, 0);
1907
1908 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 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 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 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 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 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 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 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 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 let trusted = PeerId::random();
2306 peers.add_trusted_peer_id(trusted);
2307
2308 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 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 assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2373 assert_eq!(peers.connection_info.num_pending_in, 1);
2374
2375 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 assert!(peers.on_incoming_pending_session(socket_addr.ip()).is_ok());
2386 assert_eq!(peers.connection_info.num_pending_in, 1);
2387
2388 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 assert!(!p.is_banned());
2431 }
2432
2433 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 match event!(peers) {
2553 PeerAction::BanPeer { peer_id } => {
2554 assert_eq!(peer_id, peer);
2555 }
2556 err => unreachable!("{err:?}"),
2557 }
2558
2559 assert!(peers.peers.contains_key(&peer));
2561
2562 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 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 assert_eq!(peer_manager.connection_info.num_pending_in, 1);
2684 peer_manager.on_incoming_session_established(given_peer_id, socket_addr);
2685 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 assert_eq!(peers.connection_info.num_pending_in, 1);
2831 peers.on_incoming_session_established(basic_peer, basic_sock);
2832 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 assert_eq!(peers.connection_info.num_pending_in, 1);
2865 peers.on_incoming_session_established(basic_peer, basic_sock);
2866 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 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 peer_manager.last_tick -= Duration::from_secs(1);
2936 peer_manager.tick();
2937
2938 assert_eq!(peer_manager.peers.get_mut(&peer_id).unwrap().reputation, DEFAULT_REPUTATION);
2940
2941 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 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 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 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 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 peers.on_active_outgoing_established(peer_id);
3022 assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::Out);
3023
3024 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 assert!(peers.backed_off_peers.contains_key(&peer_id));
3033
3034 poll_fn(|cx| {
3036 assert!(peers.poll(cx).is_pending());
3037 Poll::Ready(())
3038 })
3039 .await;
3040
3041 assert!(peers.backed_off_peers.contains_key(&peer_id));
3043 assert!(peers.peers.get(&peer_id).unwrap().backed_off);
3044
3045 tokio::time::sleep(throttle_duration).await;
3047
3048 match event!(peers) {
3050 PeerAction::Connect { peer_id: id, .. } => assert_eq!(id, peer_id),
3051 _ => unreachable!(),
3052 }
3053
3054 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 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 peers.on_active_outgoing_established(peer_id);
3084 assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::Out);
3085
3086 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 assert!(peers.on_incoming_pending_session(addr.ip()).is_ok());
3097 assert_eq!(peers.connection_info.num_pending_in, 1);
3098
3099 peers.on_incoming_session_established(peer_id, addr);
3101
3102 assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::In);
3104 assert_eq!(peers.connection_info.num_inbound, 1);
3105
3106 assert!(peers.backed_off_peers.contains_key(&peer_id));
3108 assert!(peers.peers.get(&peer_id).unwrap().backed_off);
3109
3110 poll_fn(|cx| {
3112 assert!(peers.poll(cx).is_pending());
3113 Poll::Ready(())
3114 })
3115 .await;
3116
3117 assert_eq!(peers.peers.get(&peer_id).unwrap().state, PeerConnectionState::In);
3119
3120 tokio::time::sleep(throttle_duration).await;
3122
3123 poll_fn(|cx| {
3124 let _ = peers.poll(cx);
3125 Poll::Ready(())
3126 })
3127 .await;
3128
3129 assert!(!peers.backed_off_peers.contains_key(&peer_id));
3131 assert!(!peers.peers.get(&peer_id).unwrap().backed_off);
3132
3133 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 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 peer_manager.fill_outbound_slots();
3245
3246 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 for peer_id in &num_pendingout_states {
3263 peer_manager.on_active_outgoing_established(*peer_id);
3264 }
3265
3266 for peer_id in &num_pendingout_states {
3268 assert_eq!(peer_manager.peers.get(peer_id).unwrap().state, PeerConnectionState::Out);
3269 }
3270
3271 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 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 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 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!(peers.on_incoming_pending_session(new_addr.ip()).is_ok());
3383 assert_eq!(peers.connection_info.num_pending_in, 1);
3384
3385 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 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 let allowed_ip: IpAddr = "192.168.1.100".parse().unwrap();
3450 assert!(peers.on_incoming_pending_session(allowed_ip).is_ok());
3451
3452 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 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 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 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 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 let allowed_ip: IpAddr = "2001:db8::1".parse().unwrap();
3493 assert!(peers.on_incoming_pending_session(allowed_ip).is_ok());
3494
3495 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 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 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 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 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 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 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 #[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 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 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}