snarkos_node_router/
heartbeat.rs1use crate::{
17 CandidatePeer,
18 ConnectedPeer,
19 NodeType,
20 Outbound,
21 PeerPoolHandling,
22 Router,
23 bootstrap_peers,
24 messages::{DisconnectReason, Message, PeerRequest},
25};
26
27use snarkos_node_tcp::{ConnectError, P2P};
28
29use snarkvm::prelude::Network;
30
31use colored::Colorize;
32use futures::future::join_all;
33use rand::{SeedableRng, prelude::IteratorRandom};
34use rand_chacha::ChaChaRng;
35use std::time::Duration;
36use tokio::task::JoinError;
37
38pub const fn max(a: usize, b: usize) -> usize {
41 match a > b {
42 true => a,
43 false => b,
44 }
45}
46
47#[async_trait]
48pub trait Heartbeat<N: Network>: Outbound<N> {
49 const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(25);
51 const MINIMUM_NUMBER_OF_PEERS: usize = 3;
53 const MINIMUM_TIME_BETWEEN_CONNECTION_ATTEMPTS: Duration = Duration::from_secs(10);
55 const STARTUP_GRACE_PERIOD: Duration = Duration::from_secs(60);
57 const MEDIAN_NUMBER_OF_PEERS: usize = max(Self::MAXIMUM_NUMBER_OF_PEERS / 2, Self::MINIMUM_NUMBER_OF_PEERS);
59 const MAXIMUM_NUMBER_OF_PEERS: usize = 21;
61 const MAXIMUM_NUMBER_OF_PROVERS: usize = Self::MAXIMUM_NUMBER_OF_PEERS / 4;
63 const IP_BAN_TIME_IN_SECS: u64 = 300;
65
66 async fn heartbeat(&self) {
68 self.safety_check_minimum_number_of_peers();
69 self.log_connected_peers();
70
71 self.remove_oldest_connected_peer();
73 self.handle_connected_peers().await;
75 self.handle_bootstrap_peers().await;
77 self.handle_trusted_peers().await;
79 self.handle_puzzle_request();
81 self.handle_banned_ips();
83 self.clear_stale_peers();
85 }
86
87 fn safety_check_minimum_number_of_peers(&self) {
90 assert!(Self::MINIMUM_NUMBER_OF_PEERS >= 1, "The minimum number of peers must be at least 1.");
92 assert!(Self::MINIMUM_NUMBER_OF_PEERS <= Self::MAXIMUM_NUMBER_OF_PEERS);
93 assert!(Self::MINIMUM_NUMBER_OF_PEERS <= Self::MEDIAN_NUMBER_OF_PEERS);
94 assert!(Self::MEDIAN_NUMBER_OF_PEERS <= Self::MAXIMUM_NUMBER_OF_PEERS);
95 assert!(Self::MAXIMUM_NUMBER_OF_PROVERS <= Self::MAXIMUM_NUMBER_OF_PEERS);
96 }
97
98 fn log_connected_peers(&self) {
100 let connected_peers = self.router().connected_peers();
102 let connected_peers_fmt = format!("{connected_peers:?}").dimmed();
103 match connected_peers.len() {
104 0 => {
105 if self.router().tcp().uptime() > Self::STARTUP_GRACE_PERIOD {
107 warn!("No connected peers")
108 }
109 }
110 1 => debug!("Connected to 1 peer: {connected_peers_fmt}"),
111 num_connected => debug!("Connected to {num_connected} peers {connected_peers_fmt}"),
112 }
113 }
114
115 fn get_removable_peers(&self) -> Vec<ConnectedPeer<N>> {
123 let is_block_synced = self.is_block_synced();
125
126 let mut peers = self.router().filter_connected_peers(|peer| {
131 !peer.trusted
132 && peer.node_type != NodeType::BootstrapClient
133 && !self.router().cache.contains_inbound_block_request(&peer.listener_addr) && (is_block_synced || self.router().cache.num_outbound_block_requests(&peer.listener_addr) == 0) });
136 peers.sort_by_key(|peer| peer.last_seen);
137
138 peers
139 }
140
141 fn remove_oldest_connected_peer(&self) {
145 if self.router().trusted_peers_only() {
147 return;
148 }
149
150 if self.router().number_of_connected_peers() <= Self::MINIMUM_NUMBER_OF_PEERS {
152 return;
153 }
154
155 if let Some(oldest) = self.get_removable_peers().first().map(|peer| peer.listener_addr) {
159 info!("Disconnecting from '{oldest}' (periodic refresh of peers)");
160 let _ = self.router().send(oldest, Message::Disconnect(DisconnectReason::PeerRefresh.into()));
161 self.router().disconnect(oldest);
162 }
163 }
164
165 #[inline]
168 fn log_if_connect_error(result: Result<Result<(), ConnectError>, JoinError>, context: &str) {
169 match result {
170 Ok(Ok(())) => {}
172 Ok(Err(err @ ConnectError::AlreadyConnecting { .. }))
173 | Ok(Err(err @ ConnectError::AlreadyConnected { .. })) => {
174 debug!("{context}: {err}");
176 }
177 Ok(Err(err)) => warn!("{context}: {err}"),
179 Err(err) => error!("{context}: {err}"),
181 }
182 }
183
184 async fn handle_connected_peers(&self) {
186 let rng = &mut ChaChaRng::from_rng(&mut rand::rng());
188
189 let num_connected = self.router().number_of_connected_peers();
191 let num_connected_provers = self.router().filter_connected_peers(|peer| peer.node_type.is_prover()).len();
193
194 let (max_peers, max_provers) = (Self::MAXIMUM_NUMBER_OF_PEERS, Self::MAXIMUM_NUMBER_OF_PROVERS);
196
197 let num_surplus_peers = num_connected.saturating_sub(max_peers);
199 let num_surplus_provers = num_connected_provers.saturating_sub(max_provers);
201 let num_remaining_provers = num_connected_provers.saturating_sub(num_surplus_provers);
203 let num_surplus_clients_validators = num_surplus_peers.saturating_sub(num_remaining_provers);
205
206 if num_surplus_provers > 0 || num_surplus_clients_validators > 0 {
207 debug!(
208 "Exceeded maximum number of connected peers, disconnecting from ({num_surplus_provers} + {num_surplus_clients_validators}) peers"
209 );
210
211 let provers_to_disconnect = self
213 .router()
214 .filter_connected_peers(|peer| peer.node_type.is_prover() && !peer.trusted)
215 .into_iter()
216 .sample(rng, num_surplus_provers);
217
218 let peers_to_disconnect = self
220 .get_removable_peers()
221 .into_iter()
222 .filter(|peer| !peer.node_type.is_prover()) .take(num_surplus_clients_validators);
224
225 for peer in peers_to_disconnect.chain(provers_to_disconnect) {
227 if self.router().node_type().is_prover() {
229 continue;
230 }
231
232 let peer_addr = peer.listener_addr;
233 info!("Disconnecting from '{peer_addr}' (exceeded maximum connections)");
234 self.router().send(peer_addr, Message::Disconnect(DisconnectReason::TooManyPeers.into()));
235 self.router().disconnect(peer_addr);
237 }
238 }
239
240 let num_connected = self.router().number_of_connected_peers();
242 let num_deficient = Self::MEDIAN_NUMBER_OF_PEERS.saturating_sub(num_connected);
244
245 if num_deficient > 0 {
246 let rng = &mut ChaChaRng::from_rng(&mut rand::rng());
248
249 let own_height = self.router().ledger.latest_block_height();
252 let (higher_peers, other_peers): (Vec<_>, Vec<_>) = self
253 .router()
254 .get_candidate_peers()
255 .into_iter()
256 .partition(|peer| peer.last_height_seen.map(|h| h > own_height).unwrap_or(false));
257 let num_higher_peers = num_deficient.div_ceil(2).min(higher_peers.len());
259
260 let higher_peers = higher_peers.into_iter().sample(rng, num_higher_peers);
261 let other_peers = other_peers.into_iter().sample(rng, num_deficient.saturating_sub(num_higher_peers));
262
263 self.try_connect_to_peers(higher_peers.into_iter().chain(other_peers)).await;
265
266 if !self.router().trusted_peers_only() {
267 for peer_ip in self.router().connected_peers().into_iter().sample(rng, 3) {
269 self.router().send(peer_ip, Message::PeerRequest(PeerRequest));
270 }
271 }
272 }
273 }
274
275 async fn handle_bootstrap_peers(&self) {
277 if self.router().trusted_peers_only() {
279 return;
280 }
281 let mut candidate_bootstrap = Vec::new();
283 let connected_bootstrap =
284 self.router().filter_connected_peers(|peer| peer.node_type == NodeType::BootstrapClient);
285 for bootstrap_ip in bootstrap_peers::<N>(self.router().is_dev()) {
286 if !connected_bootstrap.iter().any(|peer| peer.listener_addr == bootstrap_ip) {
287 candidate_bootstrap.push(bootstrap_ip);
288 }
289 }
290 if connected_bootstrap.is_empty() {
292 let rng = &mut ChaChaRng::from_rng(&mut rand::rng());
294 if let Some(peer_ip) = candidate_bootstrap.into_iter().choose(rng) {
296 match self.router().connect(peer_ip) {
297 Ok(hdl) => {
298 Self::log_if_connect_error(
299 hdl.await,
300 &format!("Could not connect to bootstrap peer at '{peer_ip:?}'"),
301 );
302 }
303 Err(ConnectError::AlreadyConnected { .. }) | Err(ConnectError::AlreadyConnecting { .. }) => {}
304 Err(err) => warn!("Could not initiate connection to bootstrap peer at '{peer_ip:?}' - {err}"),
305 }
306 }
307 }
308 let num_surplus = connected_bootstrap.len().saturating_sub(1);
310 if num_surplus > 0 {
311 let rng = &mut ChaChaRng::from_rng(&mut rand::rng());
313 for peer in connected_bootstrap.into_iter().sample(rng, num_surplus) {
315 info!("Disconnecting from '{}' (exceeded maximum bootstrap)", peer.listener_addr);
316 self.router().send(peer.listener_addr, Message::Disconnect(DisconnectReason::TooManyPeers.into()));
317 self.router().disconnect(peer.listener_addr);
319 }
320 }
321 }
322
323 async fn try_connect_to_peers(&self, peers: impl Iterator<Item = CandidatePeer<N>> + Send + 'static) {
327 let (peer_info, hdls): (Vec<_>, Vec<_>) = peers
328 .filter_map(|peer| {
329 let peer_type = if peer.trusted { "trusted peer" } else { "peer" };
330
331 if let Some(last_connection_attempt) = peer.last_connection_attempt
334 && last_connection_attempt.elapsed() < Self::MINIMUM_TIME_BETWEEN_CONNECTION_ATTEMPTS
335 {
336 return None;
337 }
338
339 let addr = peer.listener_addr;
341 let attempt_no = peer.total_connection_attempts + 1;
342
343 debug!("(Re-)connecting to {peer_type} '{addr}' (attempt #{attempt_no})");
345 match self.router().connect(addr) {
346 Ok(hdl) => Some(((addr, attempt_no, peer_type), hdl)),
347 Err(ConnectError::AlreadyConnected { .. }) | Err(ConnectError::AlreadyConnecting { .. }) => None,
348 Err(err) => {
349 warn!("Could not initiate connection to {peer_type} at '{addr}' - {err}");
350 None
351 }
352 }
353 })
354 .unzip();
355
356 for ((peer_addr, attempt_no, peer_type), result) in peer_info.into_iter().zip(join_all(hdls).await) {
358 Self::log_if_connect_error(
359 result,
360 &format!("Could not connect to {peer_type} at '{peer_addr}' (attempt #{attempt_no})"),
361 );
362 }
363 }
364
365 async fn handle_trusted_peers(&self) {
367 self.try_connect_to_peers(self.router().get_trusted_candidate_peers().into_iter()).await;
368 }
369
370 fn handle_puzzle_request(&self) {
372 }
374
375 fn handle_banned_ips(&self) {
377 self.router().tcp().banned_peers().remove_old_bans(Self::IP_BAN_TIME_IN_SECS);
378 }
379
380 fn clear_stale_peers(&self) {
381 self.router().cache().clear_stale_entries(
382 Router::<N>::CONNECTION_ATTEMPTS_SINCE_SECS,
383 Router::<N>::MESSAGE_LIMIT_TIME_FRAME_IN_SECS,
384 );
385 }
386}