Skip to main content

snarkos_node_network/
peering.rs

1// Copyright (c) 2019-2026 Provable Inc.
2// This file is part of the snarkOS library.
3
4// Licensed under the Apache License, Version 2.0 (the "License");
5// you may not use this file except in compliance with the License.
6// You may obtain a copy of the License at:
7
8// http://www.apache.org/licenses/LICENSE-2.0
9
10// Unless required by applicable law or agreed to in writing, software
11// distributed under the License is distributed on an "AS IS" BASIS,
12// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13// See the License for the specific language governing permissions and
14// limitations under the License.
15
16use crate::{CandidatePeer, ConnectedPeer, ConnectionMode, NodeType, Peer, Resolver};
17
18use snarkos_node_tcp::{ConnectError, P2P, is_bogon_ip, is_unspecified_or_broadcast_ip};
19use snarkvm::prelude::{Address, Network};
20
21use anyhow::Result;
22#[cfg(feature = "locktick")]
23use locktick::parking_lot::RwLock;
24#[cfg(not(feature = "locktick"))]
25use parking_lot::RwLock;
26use std::{
27    cmp,
28    collections::{
29        HashSet,
30        hash_map::{Entry, HashMap},
31    },
32    fs,
33    io::{self, Write},
34    net::{IpAddr, SocketAddr},
35    path::Path,
36    str::FromStr,
37    time::Instant,
38};
39use tokio::task;
40use tracing::*;
41
42/// Application-level errors generated by the peering module.
43/// This is never returned directly, but only as the payload for a `ConnectError`.
44#[derive(Debug)]
45pub enum PeeringError {
46    NoExternalPeersAllowed,
47}
48
49impl snarkos_node_tcp::ApplicationError for PeeringError {}
50
51impl std::fmt::Display for PeeringError {
52    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
53        match self {
54            Self::NoExternalPeersAllowed => write!(f, "no untrusted peers allowed"),
55        }
56    }
57}
58
59pub trait PeerPoolHandling<N: Network>: P2P {
60    const OWNER: &str;
61
62    /// The maximum number of peers permitted to be stored in the peer pool.
63    const MAXIMUM_POOL_SIZE: usize;
64
65    /// The number of candidate peers to be removed from the pool once `MAXIMUM_POOL_SIZE` is reached.
66    /// It must be lower than `MAXIMUM_POOL_SIZE`.
67    const PEER_SLASHING_COUNT: usize;
68
69    /// Returns the mapping of all known peers (connected or otherwise), keyed by their public listener address.
70    fn peer_pool(&self) -> &RwLock<HashMap<SocketAddr, Peer<N>>>;
71
72    /// Returns the resolver for translating between public listener addresses and connected addresses.
73    fn resolver(&self) -> &RwLock<Resolver<N>>;
74
75    /// Returns `true` if the owning node is in development mode.
76    fn is_dev(&self) -> bool;
77
78    /// Returns `true` if the node is in trusted peers only mode.
79    fn trusted_peers_only(&self) -> bool;
80
81    /// Returns the node type.
82    fn node_type(&self) -> NodeType;
83
84    /// Returns the listener address of this node.
85    fn local_ip(&self) -> SocketAddr {
86        self.tcp().listening_addr().expect("The TCP listener is not enabled")
87    }
88
89    /// Returns `true` if the given IP is this node.
90    fn is_local_ip(&self, addr: SocketAddr) -> bool {
91        addr == self.local_ip()
92            || (addr.ip().is_unspecified() || addr.ip().is_loopback()) && addr.port() == self.local_ip().port()
93    }
94
95    /// Returns `true` if the given IP is not this node, is not a bogon address, and is not unspecified.
96    fn is_valid_peer_ip(&self, ip: SocketAddr) -> bool {
97        !self.is_local_ip(ip) && !is_bogon_ip(ip.ip()) && !is_unspecified_or_broadcast_ip(ip.ip())
98    }
99
100    /// Returns the maximum number of connected peers.
101    fn max_connected_peers(&self) -> usize {
102        self.tcp().config().max_connections as usize
103    }
104
105    /// Ensure we can and are allowed to connect to the given listener address of a peer.
106    fn check_connection_attempt(&self, listener_addr: SocketAddr) -> Result<(), ConnectError> {
107        // Ensure the peer IP is not this node.
108        if self.is_local_ip(listener_addr) {
109            return Err(ConnectError::SelfConnect { address: listener_addr });
110        }
111        // Ensure the node does not surpass the maximum number of peer connections.
112        if self.number_of_connected_peers() >= self.max_connected_peers() {
113            return Err(ConnectError::MaximumConnectionsReached { limit: self.max_connected_peers() as u16 });
114        }
115        // Ensure the node is not already connected to this peer.
116        if self.is_connected(listener_addr) {
117            return Err(ConnectError::AlreadyConnected { address: listener_addr });
118        }
119        // Ensure the node is not already connecting to this peer.
120        if self.is_connecting(listener_addr) {
121            return Err(ConnectError::AlreadyConnecting { address: listener_addr });
122        }
123        // Ensure the peer IP is not banned.
124        if self.is_ip_banned(listener_addr.ip()) {
125            return Err(ConnectError::BannedIp { ip: listener_addr.ip() });
126        }
127        // If the node is in trusted peers only mode, ensure the peer is trusted.
128        if self.trusted_peers_only() && !self.is_trusted(listener_addr) {
129            return Err(ConnectError::application(PeeringError::NoExternalPeersAllowed));
130        }
131
132        Ok(())
133    }
134
135    /// Attempts to connect to the given peer's listener address.
136    ///
137    /// Returns an earlier error, if, for example, we are already connected to the peer.
138    /// Otherwise, it returns a handle to the tokio tasks that sets up the connection.
139    ///
140    /// # Concurrency
141    /// Only one task may call this function for a given listener address at a time.
142    fn connect(&self, listener_addr: SocketAddr) -> Result<task::JoinHandle<Result<(), ConnectError>>, ConnectError> {
143        // Return early if the attempt is against the protocol rules.
144        self.check_connection_attempt(listener_addr)?;
145
146        // Update the last connection attempt time for the peer.
147        if let Some(Peer::Candidate(peer)) = self.peer_pool().write().get_mut(&listener_addr) {
148            peer.last_connection_attempt = Some(Instant::now());
149            peer.total_connection_attempts += 1;
150        } else {
151            warn!("{} No candidate peer entry exists for '{listener_addr:?}' while connecting.", Self::OWNER);
152        }
153
154        let tcp = self.tcp().clone();
155        Ok(tokio::spawn(async move {
156            debug!("{} Connecting to {listener_addr}...", Self::OWNER);
157            tcp.connect(listener_addr).await
158        }))
159    }
160
161    /// Disconnects from the given peer IP, if the peer is connected. The returned boolean
162    /// indicates whether the peer was actually disconnected from, or if this was a noop.
163    fn disconnect(&self, listener_addr: SocketAddr) -> task::JoinHandle<bool> {
164        if let Some(connected_addr) = self.resolve_to_ambiguous(listener_addr) {
165            let tcp = self.tcp().clone();
166            tokio::spawn(async move { tcp.disconnect(connected_addr).await })
167        } else {
168            tokio::spawn(async { false })
169        }
170    }
171
172    /// Downgrades a connected peer to candidate status.
173    ///
174    /// Returns true if the peer was fully connected.
175    fn downgrade_peer_to_candidate(&self, listener_addr: SocketAddr) -> bool {
176        let mut peer_pool = self.peer_pool().write();
177        let Some(peer) = peer_pool.get_mut(&listener_addr) else {
178            trace!("{} Downgrade peer to candidate failed - peer not found", Self::OWNER);
179            return false;
180        };
181
182        if let Peer::Connected(conn_peer) = peer {
183            // Exception: the BootstrapClient only has a single Resolver,
184            // so it may only map a validator's Aleo address once, for its
185            // Gateway-mode connection. This also means that the Router-mode
186            // connection may not remove that mapping.
187            let aleo_addr = if self.node_type() == NodeType::BootstrapClient
188                && conn_peer.connection_mode == ConnectionMode::Router
189            {
190                None
191            } else {
192                Some(conn_peer.aleo_addr)
193            };
194            self.resolver().write().remove_peer(conn_peer.connected_addr, aleo_addr);
195            peer.downgrade_to_candidate(listener_addr);
196            true
197        } else {
198            peer.downgrade_to_candidate(listener_addr);
199            false
200        }
201    }
202
203    /// Adds new candidate peers to the peer pool, ensuring their validity and following the
204    /// limit on the number of peers in the pool. The listener addresses may be paired with
205    /// the last known block height of the associated peer.
206    fn insert_candidate_peers(&self, mut listener_addrs: Vec<(SocketAddr, Option<u32>)>) {
207        let trusted_peers = self.trusted_peers();
208
209        // Hold a write guard from now on, so as not to accidentally slash multiple times
210        // based on multiple batches of candidate peers, and to not overwrite any entries.
211        let mut peer_pool = self.peer_pool().write();
212
213        // Perform filtering to ensure candidate validity. Also count how many entries are updates.
214        let mut num_updates: usize = 0;
215        listener_addrs.retain(|&(addr, height)| {
216            !self.is_ip_banned(addr.ip())
217                && if self.is_dev() { !is_bogon_ip(addr.ip()) } else { self.is_valid_peer_ip(addr) }
218                && peer_pool
219                    .get(&addr)
220                    .map(|peer| peer.is_candidate() && height.is_some())
221                    .inspect(|is_valid_update| {
222                        if *is_valid_update {
223                            num_updates += 1
224                        }
225                    })
226                    .unwrap_or(true)
227        });
228
229        // If we've managed to filter out every entry, there's nothing to do.
230        if listener_addrs.is_empty() {
231            return;
232        }
233
234        // If we're about to exceed the peer pool size limit, apply candidate slashing.
235        if peer_pool.len() + listener_addrs.len() - num_updates >= Self::MAXIMUM_POOL_SIZE
236            && Self::PEER_SLASHING_COUNT != 0
237        {
238            // Collect the addresses of prospect peers.
239            let mut peers_to_slash = peer_pool
240                .iter()
241                .filter_map(|(addr, peer)| {
242                    (matches!(peer, Peer::Candidate(_)) && !trusted_peers.contains(addr)).then_some(*addr)
243                })
244                .collect::<Vec<_>>();
245
246            // Get the low-level peer stats.
247            let known_peers = self.tcp().known_peers().snapshot();
248
249            // Sort the list of candidate peers by failure count (descending) and timestamp (ascending).
250            let default_value = (0, Instant::now());
251            peers_to_slash.sort_unstable_by_key(|addr| {
252                let (num_failures, last_seen) = known_peers
253                    .get(&addr.ip())
254                    .map(|stats| (stats.failures(), stats.timestamp()))
255                    .unwrap_or(default_value);
256                (cmp::Reverse(num_failures), last_seen)
257            });
258
259            // Retain the candidate peers with the most failures and oldest timestamps.
260            peers_to_slash.truncate(Self::PEER_SLASHING_COUNT);
261
262            // Remove the peers to slash from the pool.
263            peer_pool.retain(|addr, _| !peers_to_slash.contains(addr));
264
265            // Remove the peers to slash from the low-level list of known peers.
266            self.tcp().known_peers().batch_remove(peers_to_slash.iter().map(|addr| addr.ip()));
267        }
268
269        // Make sure that we won't breach the pool size limit in case the slashing didn't suffice.
270        listener_addrs.truncate(Self::MAXIMUM_POOL_SIZE.saturating_sub(peer_pool.len()));
271
272        // If we've managed to truncate to 0, exit.
273        if listener_addrs.is_empty() {
274            return;
275        }
276
277        // Insert or update the applicable candidate peers.
278        for (addr, height) in listener_addrs {
279            match peer_pool.entry(addr) {
280                Entry::Vacant(entry) => {
281                    entry.insert(Peer::new_candidate(addr, false));
282                }
283                Entry::Occupied(mut entry) => {
284                    if let Peer::Candidate(peer) = entry.get_mut() {
285                        peer.last_height_seen = height;
286                    }
287                }
288            }
289        }
290    }
291
292    /// Completely removes an entry from the peer pool.
293    fn remove_peer(&self, listener_addr: SocketAddr) {
294        self.peer_pool().write().remove(&listener_addr);
295    }
296
297    /// Returns the connected peer address from the listener IP address.
298    fn resolve_to_ambiguous(&self, listener_addr: SocketAddr) -> Option<SocketAddr> {
299        if let Some(Peer::Connected(peer)) = self.peer_pool().read().get(&listener_addr) {
300            Some(peer.connected_addr)
301        } else {
302            None
303        }
304    }
305
306    /// Returns the connected peer aleo address from the listener IP address.
307    fn resolve_to_aleo_addr(&self, listener_addr: SocketAddr) -> Option<Address<N>> {
308        if let Some(Peer::Connected(peer)) = self.peer_pool().read().get(&listener_addr) {
309            Some(peer.aleo_addr)
310        } else {
311            None
312        }
313    }
314
315    /// Returns `true` if the node is connecting to the given peer's listener address.
316    fn is_connecting(&self, listener_addr: SocketAddr) -> bool {
317        self.peer_pool().read().get(&listener_addr).is_some_and(|peer| peer.is_connecting())
318    }
319
320    /// Returns `true` if the node is connected to the given peer listener address.
321    fn is_connected(&self, listener_addr: SocketAddr) -> bool {
322        self.peer_pool().read().get(&listener_addr).is_some_and(|peer| peer.is_connected())
323    }
324
325    /// Returns `true` if the node is connected to the given Aleo address.
326    fn is_connected_address(&self, aleo_address: Address<N>) -> bool {
327        // The resolver only contains data on connected peers.
328        self.resolver().read().get_peer_ip_for_address(aleo_address).is_some()
329    }
330
331    /// Returns `true` if the node is connected or connecting to the given peer listener address.
332    fn is_connecting_or_connected(&self, listener_addr: SocketAddr) -> bool {
333        self.peer_pool().read().get(&listener_addr).is_some_and(|peer| peer.is_connecting() || peer.is_connected())
334    }
335
336    /// Returns `true` if the given listener address is trusted.
337    fn is_trusted(&self, listener_addr: SocketAddr) -> bool {
338        self.peer_pool().read().get(&listener_addr).is_some_and(|peer| peer.is_trusted())
339    }
340
341    /// Returns the number of all peers.
342    fn number_of_peers(&self) -> usize {
343        self.peer_pool().read().len()
344    }
345
346    /// Returns the number of connected peers.
347    fn number_of_connected_peers(&self) -> usize {
348        self.peer_pool().read().values().filter(|peer| peer.is_connected()).count()
349    }
350
351    /// Returns the number of connected validators.
352    #[cfg(feature = "metrics")]
353    fn number_of_connected_validators(&self) -> Option<usize> {
354        Some(
355            self.peer_pool()
356                .try_read()?
357                .values()
358                .filter(|peer| peer.as_connected().is_some_and(|peer| peer.is_validator()))
359                .count(),
360        )
361    }
362
363    /// Returns the number of connecting peers.
364    #[cfg(feature = "metrics")]
365    fn number_of_connecting_peers(&self) -> Option<usize> {
366        Some(self.peer_pool().try_read()?.values().filter(|peer| peer.is_connecting()).count())
367    }
368
369    /// Returns the number of candidate peers.
370    fn number_of_candidate_peers(&self) -> usize {
371        self.peer_pool().read().values().filter(|peer| matches!(peer, Peer::Candidate(_))).count()
372    }
373
374    /// Returns the connected peer given the peer IP, if it exists.
375    fn get_connected_peer(&self, listener_addr: SocketAddr) -> Option<ConnectedPeer<N>> {
376        if let Some(Peer::Connected(peer)) = self.peer_pool().read().get(&listener_addr) {
377            Some(peer.clone())
378        } else {
379            None
380        }
381    }
382
383    /// Updates the connected peer - if it exists -  given the peer IP and a closure.
384    /// The returned status indicates whether the update was successful, i.e. the peer had existed.
385    fn update_connected_peer<F: FnMut(&mut ConnectedPeer<N>)>(
386        &self,
387        listener_addr: &SocketAddr,
388        mut update_fn: F,
389    ) -> bool {
390        if let Some(Peer::Connected(peer)) = self.peer_pool().write().get_mut(listener_addr) {
391            update_fn(peer);
392            true
393        } else {
394            false
395        }
396    }
397
398    /// Returns the list of all peers (connected, connecting, and candidate).
399    fn get_peers(&self) -> Vec<Peer<N>> {
400        self.peer_pool().read().values().cloned().collect()
401    }
402
403    /// Returns all connected peers.
404    fn get_connected_peers(&self) -> Vec<ConnectedPeer<N>> {
405        self.filter_connected_peers(|_| true)
406    }
407
408    /// Returns an optionally bounded list of all connected peers sorted by their
409    /// block height (highest first) and failure count (lowest first).
410    fn get_best_connected_peers(&self, max_entries: Option<usize>) -> Vec<ConnectedPeer<N>> {
411        // Get a snapshot of the currently connected peers.
412        let mut peers = self.get_connected_peers();
413        // Get the low-level peer stats.
414        let known_peers = self.tcp().known_peers().snapshot();
415
416        // Sort the prospect peers.
417        peers.sort_unstable_by_key(|peer| {
418            if let Some(peer_stats) = known_peers.get(&peer.listener_addr.ip()) {
419                // Prioritize greatest height, then lowest failure count.
420                (cmp::Reverse(peer.last_height_seen), peer_stats.failures())
421            } else {
422                // Unreachable; use an else-compatible dummy.
423                (cmp::Reverse(peer.last_height_seen), 0)
424            }
425        });
426        if let Some(max) = max_entries {
427            peers.truncate(max);
428        }
429
430        peers
431    }
432
433    /// Returns all connected peers that satisify the given predicate.
434    fn filter_connected_peers<P: FnMut(&ConnectedPeer<N>) -> bool>(&self, mut predicate: P) -> Vec<ConnectedPeer<N>> {
435        self.peer_pool()
436            .read()
437            .values()
438            .filter_map(|p| {
439                if let Peer::Connected(peer) = p
440                    && predicate(peer)
441                {
442                    Some(peer)
443                } else {
444                    None
445                }
446            })
447            .cloned()
448            .collect()
449    }
450
451    /// Returns the list of connected peers.
452    fn connected_peers(&self) -> Vec<SocketAddr> {
453        self.peer_pool().read().iter().filter_map(|(addr, peer)| peer.is_connected().then_some(*addr)).collect()
454    }
455
456    /// Returns the list of trusted peers.
457    fn trusted_peers(&self) -> Vec<SocketAddr> {
458        self.peer_pool().read().iter().filter_map(|(addr, peer)| peer.is_trusted().then_some(*addr)).collect()
459    }
460
461    /// Returns the list of candidate peers.
462    fn get_candidate_peers(&self) -> Vec<CandidatePeer<N>> {
463        self.peer_pool()
464            .read()
465            .values()
466            .filter_map(|peer| if let Peer::Candidate(peer) = peer { Some(peer.clone()) } else { None })
467            .collect()
468    }
469
470    /// Returns the list of trusted candidate peers.
471    fn get_trusted_candidate_peers(&self) -> Vec<CandidatePeer<N>> {
472        self.peer_pool()
473            .read()
474            .values()
475            .filter_map(|peer| {
476                if let Peer::Candidate(peer) = peer
477                    && peer.trusted
478                {
479                    Some(peer.clone())
480                } else {
481                    None
482                }
483            })
484            .collect()
485    }
486
487    /// Loads any previously cached peer addresses so they can be introduced as initial
488    /// candidate peers to connect to.
489    fn load_cached_peers(path: &Path) -> Result<Vec<SocketAddr>> {
490        let peers = match fs::read_to_string(path) {
491            Ok(cached_peers_str) => {
492                let mut cached_peers = Vec::new();
493                for peer_addr_str in cached_peers_str.lines() {
494                    match SocketAddr::from_str(peer_addr_str) {
495                        Ok(addr) => cached_peers.push(addr),
496                        Err(error) => warn!("Couldn't parse the cached peer address '{peer_addr_str}': {error}"),
497                    }
498                }
499                cached_peers
500            }
501            Err(error) if error.kind() == io::ErrorKind::NotFound => {
502                // Not an issue - the cache may not exist yet.
503                Vec::new()
504            }
505            Err(error) => {
506                warn!("{} Couldn't load cached peers at {}: {error}", Self::OWNER, path.display());
507                Vec::new()
508            }
509        };
510
511        Ok(peers)
512    }
513
514    /// Preserve the peers who have the greatest known block heights, and the lowest
515    /// number of registered network failures.
516    ///
517    /// # Arguments
518    /// * `path` - The path to the file to save the peers to.
519    /// * `max_entries` - The maximum number of peers to save (if there are more, the extra ones are truncated).
520    /// * `store_ports` - Whether to store the ports of the peers, or just the IP addresses.
521    fn save_best_peers(&self, path: &Path, max_entries: Option<usize>, store_ports: bool) -> Result<()> {
522        // Collect all prospect peers.
523        let mut peers = self.get_peers();
524
525        // Get the low-level peer stats.
526        let known_peers = self.tcp().known_peers().snapshot();
527
528        // Sort the list of peers.
529        peers.sort_unstable_by_key(|peer| {
530            if let Some(peer_stats) = known_peers.get(&peer.listener_addr().ip()) {
531                // Prioritize greatest height, then lowest failure count.
532                (cmp::Reverse(peer.last_height_seen()), peer_stats.failures())
533            } else {
534                // Unreachable; use an else-compatible dummy.
535                (cmp::Reverse(peer.last_height_seen()), 0)
536            }
537        });
538        if let Some(max) = max_entries {
539            peers.truncate(max);
540        }
541
542        // Dump the connected and deduplicated peers to a file.
543        let addrs: HashSet<_> = peers
544            .iter()
545            .map(
546                |peer| {
547                    if store_ports { peer.listener_addr().to_string() } else { peer.listener_addr().ip().to_string() }
548                },
549            )
550            .collect();
551
552        let mut file = fs::File::create(path)?;
553        for addr in addrs {
554            writeln!(file, "{addr}")?;
555        }
556
557        Ok(())
558    }
559
560    // Introduces a new connecting peer into the peer pool if unknown, or promotes
561    // a known candidate peer to a connecting one. The returned boolean indicates
562    // whether the peer has been added/promoted, or rejected due to already being
563    // shaken hands with or connected.
564    fn add_connecting_peer(&self, listener_addr: SocketAddr) -> Result<(), ConnectError> {
565        match self.peer_pool().write().entry(listener_addr) {
566            Entry::Vacant(entry) => {
567                entry.insert(Peer::new_connecting(listener_addr, false));
568                Ok(())
569            }
570            Entry::Occupied(mut entry) => match entry.get() {
571                peer @ Peer::Candidate(_) => {
572                    entry.insert(Peer::new_connecting(listener_addr, peer.is_trusted()));
573                    Ok(())
574                }
575                Peer::Connecting(_) => Err(ConnectError::AlreadyConnecting { address: listener_addr }),
576                Peer::Connected(_) => Err(ConnectError::AlreadyConnected { address: listener_addr }),
577            },
578        }
579    }
580
581    /// Temporarily IP-ban and disconnect from the peer with the given listener address and an
582    /// optional reason for the ban. This also removes the peer from the candidate pool.
583    fn ip_ban_peer(&self, listener_addr: SocketAddr, reason: Option<&str>) {
584        // Ignore IP-banning if we are in dev mode.
585        if self.is_dev() {
586            return;
587        }
588
589        let ip = listener_addr.ip();
590        debug!("IP-banning {ip}{}", reason.map(|r| format!(" reason: {r}")).unwrap_or_default());
591
592        // Insert/update the low-level IP ban list.
593        self.tcp().banned_peers().update_ip_ban(ip);
594
595        // Disconnect from the peer.
596        self.disconnect(listener_addr);
597        // Remove the peer from the pool.
598        self.remove_peer(listener_addr);
599    }
600
601    /// Check whether the given IP address is currently banned.
602    fn is_ip_banned(&self, ip: IpAddr) -> bool {
603        self.tcp().banned_peers().is_ip_banned(&ip)
604    }
605
606    /// Insert or update a banned IP.
607    fn update_ip_ban(&self, ip: IpAddr) {
608        self.tcp().banned_peers().update_ip_ban(ip);
609    }
610}
611
612#[cfg(test)]
613mod tests {
614    use super::*;
615    use crate::Peer;
616    use snarkos_node_tcp::{Config, P2P, Tcp};
617    use snarkvm::{prelude::Rng, utilities::TestRng};
618
619    use std::{collections::HashMap, net::SocketAddr, time::Instant};
620
621    type CurrentNetwork = snarkvm::prelude::MainnetV0;
622
623    struct MockPeerPool<N: Network> {
624        tcp: Tcp,
625        peer_pool: RwLock<HashMap<SocketAddr, Peer<N>>>,
626        resolver: RwLock<Resolver<N>>,
627    }
628
629    impl<N: Network> MockPeerPool<N> {
630        fn new() -> Self {
631            let config = Config { listener_ip: None, ..Default::default() };
632            Self { tcp: Tcp::new(config), peer_pool: Default::default(), resolver: Default::default() }
633        }
634    }
635
636    impl<N: Network> P2P for MockPeerPool<N> {
637        fn tcp(&self) -> &Tcp {
638            &self.tcp
639        }
640    }
641
642    impl<N: Network> PeerPoolHandling<N> for MockPeerPool<N> {
643        const MAXIMUM_POOL_SIZE: usize = 100;
644        const OWNER: &str = "MockPeerPool";
645        const PEER_SLASHING_COUNT: usize = 10;
646
647        fn peer_pool(&self) -> &RwLock<HashMap<SocketAddr, Peer<N>>> {
648            &self.peer_pool
649        }
650
651        fn resolver(&self) -> &RwLock<Resolver<N>> {
652            &self.resolver
653        }
654
655        fn is_dev(&self) -> bool {
656            false
657        }
658
659        fn trusted_peers_only(&self) -> bool {
660            false
661        }
662
663        fn node_type(&self) -> NodeType {
664            NodeType::Client
665        }
666    }
667
668    fn make_connected_peer(port: u16, node_type: NodeType, rng: &mut TestRng) -> (SocketAddr, Peer<CurrentNetwork>) {
669        use snarkvm::prelude::Address;
670        let listener_addr = SocketAddr::from(([127, 0, 0, 1], port));
671        let connected_addr = SocketAddr::from(([127, 0, 0, 1], port + 10000));
672        let now = Instant::now();
673        let peer = Peer::Connected(ConnectedPeer {
674            listener_addr,
675            connected_addr,
676            connection_mode: ConnectionMode::Router,
677            trusted: false,
678            aleo_addr: Address::<CurrentNetwork>::new(rng.random()),
679            node_type,
680            version: 1,
681            snarkos_sha: None,
682            last_height_seen: None,
683            first_seen: now,
684            last_seen: now,
685        });
686        (listener_addr, peer)
687    }
688
689    #[test]
690    fn test_peer_state_transitions() {
691        use snarkvm::prelude::Address;
692
693        let pool = MockPeerPool::<CurrentNetwork>::new();
694        let mut rng = TestRng::default();
695
696        let listener_addr = SocketAddr::from(([192, 0, 2, 1], 4000));
697        let connected_addr = SocketAddr::from(([192, 0, 2, 1], 14000));
698        let aleo_addr = Address::<CurrentNetwork>::new(rng.random());
699
700        // Step 1: insert as a candidate.
701        pool.peer_pool().write().insert(listener_addr, Peer::new_candidate(listener_addr, false));
702
703        assert_eq!(pool.number_of_candidate_peers(), 1);
704        assert_eq!(pool.number_of_connecting_peers(), Some(0));
705        assert_eq!(pool.number_of_connected_peers(), 0);
706        assert!(!pool.is_connecting(listener_addr));
707        assert!(!pool.is_connected(listener_addr));
708
709        // Step 2: promote to connecting.
710        assert!(pool.add_connecting_peer(listener_addr).is_ok());
711
712        assert_eq!(pool.number_of_candidate_peers(), 0);
713        assert_eq!(pool.number_of_connecting_peers(), Some(1));
714        assert_eq!(pool.number_of_connected_peers(), 0);
715        assert!(pool.is_connecting(listener_addr));
716        assert!(!pool.is_connected(listener_addr));
717
718        // Step 3: complete the handshake — upgrade to connected.
719        pool.peer_pool().write().get_mut(&listener_addr).unwrap().upgrade_to_connected(
720            connected_addr,
721            listener_addr.port(),
722            aleo_addr,
723            NodeType::Validator,
724            1,
725            None,
726            ConnectionMode::Router,
727        );
728
729        assert_eq!(pool.number_of_candidate_peers(), 0);
730        assert_eq!(pool.number_of_connecting_peers(), Some(0));
731        assert_eq!(pool.number_of_connected_peers(), 1);
732        assert!(!pool.is_connecting(listener_addr));
733        assert!(pool.is_connected(listener_addr));
734        assert_eq!(pool.number_of_connected_validators(), Some(1));
735
736        // Verify the connected peer's fields.
737        let connected = pool.get_connected_peer(listener_addr).expect("peer should be connected");
738        assert_eq!(connected.listener_addr, listener_addr);
739        assert_eq!(connected.connected_addr, connected_addr);
740        assert_eq!(connected.aleo_addr, aleo_addr);
741        assert_eq!(connected.node_type, NodeType::Validator);
742    }
743
744    #[test]
745    fn test_number_of_connected_validators() {
746        let pool = MockPeerPool::<CurrentNetwork>::new();
747        let mut rng = TestRng::default();
748
749        // Empty pool: no validators.
750        assert_eq!(pool.number_of_connected_validators(), Some(0));
751
752        // Insert 2 validators and 1 client.
753        let (addr1, peer1) = make_connected_peer(3000, NodeType::Validator, &mut rng);
754        let (addr2, peer2) = make_connected_peer(3001, NodeType::Validator, &mut rng);
755        let (addr3, peer3) = make_connected_peer(3002, NodeType::Client, &mut rng);
756        {
757            let mut pool_write = pool.peer_pool().write();
758            pool_write.insert(addr1, peer1);
759            pool_write.insert(addr2, peer2);
760            pool_write.insert(addr3, peer3);
761        }
762
763        assert_eq!(pool.number_of_connected_validators(), Some(2));
764        assert_eq!(pool.number_of_connected_peers(), 3);
765
766        // A candidate peer should not be counted as a validator.
767        let candidate_addr = SocketAddr::from(([127, 0, 0, 1], 3003));
768        pool.peer_pool().write().insert(candidate_addr, Peer::new_candidate(candidate_addr, false));
769
770        assert_eq!(pool.number_of_connected_validators(), Some(2));
771        assert_eq!(pool.number_of_connected_peers(), 3);
772    }
773}