1use 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#[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 const MAXIMUM_POOL_SIZE: usize;
64
65 const PEER_SLASHING_COUNT: usize;
68
69 fn peer_pool(&self) -> &RwLock<HashMap<SocketAddr, Peer<N>>>;
71
72 fn resolver(&self) -> &RwLock<Resolver<N>>;
74
75 fn is_dev(&self) -> bool;
77
78 fn trusted_peers_only(&self) -> bool;
80
81 fn node_type(&self) -> NodeType;
83
84 fn local_ip(&self) -> SocketAddr {
86 self.tcp().listening_addr().expect("The TCP listener is not enabled")
87 }
88
89 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 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 fn max_connected_peers(&self) -> usize {
102 self.tcp().config().max_connections as usize
103 }
104
105 fn check_connection_attempt(&self, listener_addr: SocketAddr) -> Result<(), ConnectError> {
107 if self.is_local_ip(listener_addr) {
109 return Err(ConnectError::SelfConnect { address: listener_addr });
110 }
111 if self.number_of_connected_peers() >= self.max_connected_peers() {
113 return Err(ConnectError::MaximumConnectionsReached { limit: self.max_connected_peers() as u16 });
114 }
115 if self.is_connected(listener_addr) {
117 return Err(ConnectError::AlreadyConnected { address: listener_addr });
118 }
119 if self.is_connecting(listener_addr) {
121 return Err(ConnectError::AlreadyConnecting { address: listener_addr });
122 }
123 if self.is_ip_banned(listener_addr.ip()) {
125 return Err(ConnectError::BannedIp { ip: listener_addr.ip() });
126 }
127 if self.trusted_peers_only() && !self.is_trusted(listener_addr) {
129 return Err(ConnectError::application(PeeringError::NoExternalPeersAllowed));
130 }
131
132 Ok(())
133 }
134
135 fn connect(&self, listener_addr: SocketAddr) -> Result<task::JoinHandle<Result<(), ConnectError>>, ConnectError> {
143 self.check_connection_attempt(listener_addr)?;
145
146 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 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 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 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 fn insert_candidate_peers(&self, mut listener_addrs: Vec<(SocketAddr, Option<u32>)>) {
207 let trusted_peers = self.trusted_peers();
208
209 let mut peer_pool = self.peer_pool().write();
212
213 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 listener_addrs.is_empty() {
231 return;
232 }
233
234 if peer_pool.len() + listener_addrs.len() - num_updates >= Self::MAXIMUM_POOL_SIZE
236 && Self::PEER_SLASHING_COUNT != 0
237 {
238 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 let known_peers = self.tcp().known_peers().snapshot();
248
249 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 peers_to_slash.truncate(Self::PEER_SLASHING_COUNT);
261
262 peer_pool.retain(|addr, _| !peers_to_slash.contains(addr));
264
265 self.tcp().known_peers().batch_remove(peers_to_slash.iter().map(|addr| addr.ip()));
267 }
268
269 listener_addrs.truncate(Self::MAXIMUM_POOL_SIZE.saturating_sub(peer_pool.len()));
271
272 if listener_addrs.is_empty() {
274 return;
275 }
276
277 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 fn remove_peer(&self, listener_addr: SocketAddr) {
294 self.peer_pool().write().remove(&listener_addr);
295 }
296
297 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 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 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 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 fn is_connected_address(&self, aleo_address: Address<N>) -> bool {
327 self.resolver().read().get_peer_ip_for_address(aleo_address).is_some()
329 }
330
331 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 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 fn number_of_peers(&self) -> usize {
343 self.peer_pool().read().len()
344 }
345
346 fn number_of_connected_peers(&self) -> usize {
348 self.peer_pool().read().values().filter(|peer| peer.is_connected()).count()
349 }
350
351 #[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 #[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 fn number_of_candidate_peers(&self) -> usize {
371 self.peer_pool().read().values().filter(|peer| matches!(peer, Peer::Candidate(_))).count()
372 }
373
374 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 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 fn get_peers(&self) -> Vec<Peer<N>> {
400 self.peer_pool().read().values().cloned().collect()
401 }
402
403 fn get_connected_peers(&self) -> Vec<ConnectedPeer<N>> {
405 self.filter_connected_peers(|_| true)
406 }
407
408 fn get_best_connected_peers(&self, max_entries: Option<usize>) -> Vec<ConnectedPeer<N>> {
411 let mut peers = self.get_connected_peers();
413 let known_peers = self.tcp().known_peers().snapshot();
415
416 peers.sort_unstable_by_key(|peer| {
418 if let Some(peer_stats) = known_peers.get(&peer.listener_addr.ip()) {
419 (cmp::Reverse(peer.last_height_seen), peer_stats.failures())
421 } else {
422 (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 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 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 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 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 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 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 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 fn save_best_peers(&self, path: &Path, max_entries: Option<usize>, store_ports: bool) -> Result<()> {
522 let mut peers = self.get_peers();
524
525 let known_peers = self.tcp().known_peers().snapshot();
527
528 peers.sort_unstable_by_key(|peer| {
530 if let Some(peer_stats) = known_peers.get(&peer.listener_addr().ip()) {
531 (cmp::Reverse(peer.last_height_seen()), peer_stats.failures())
533 } else {
534 (cmp::Reverse(peer.last_height_seen()), 0)
536 }
537 });
538 if let Some(max) = max_entries {
539 peers.truncate(max);
540 }
541
542 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 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 fn ip_ban_peer(&self, listener_addr: SocketAddr, reason: Option<&str>) {
584 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 self.tcp().banned_peers().update_ip_ban(ip);
594
595 self.disconnect(listener_addr);
597 self.remove_peer(listener_addr);
599 }
600
601 fn is_ip_banned(&self, ip: IpAddr) -> bool {
603 self.tcp().banned_peers().is_ip_banned(&ip)
604 }
605
606 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 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 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 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 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 assert_eq!(pool.number_of_connected_validators(), Some(0));
751
752 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 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}