Skip to main content

snarkos_node_router/
heartbeat.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::{
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
38/// A helper function to compute the maximum of two numbers.
39/// See Rust issue 92391: https://github.com/rust-lang/rust/issues/92391.
40pub 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    /// The duration in seconds to sleep in between heartbeat executions.
50    const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(25);
51    /// The minimum number of peers required to maintain connections with.
52    const MINIMUM_NUMBER_OF_PEERS: usize = 3;
53    /// The minimum time between connection attempts to a peer.
54    const MINIMUM_TIME_BETWEEN_CONNECTION_ATTEMPTS: Duration = Duration::from_secs(10);
55    /// The time we consider the node to be starting up and avoid certain warnings such as "No connected peers".
56    const STARTUP_GRACE_PERIOD: Duration = Duration::from_secs(60);
57    /// The median number of peers to maintain connections with.
58    const MEDIAN_NUMBER_OF_PEERS: usize = max(Self::MAXIMUM_NUMBER_OF_PEERS / 2, Self::MINIMUM_NUMBER_OF_PEERS);
59    /// The maximum number of peers permitted to maintain connections with.
60    const MAXIMUM_NUMBER_OF_PEERS: usize = 21;
61    /// The maximum number of provers to maintain connections with.
62    const MAXIMUM_NUMBER_OF_PROVERS: usize = Self::MAXIMUM_NUMBER_OF_PEERS / 4;
63    /// The amount of time an IP address is prohibited from connecting.
64    const IP_BAN_TIME_IN_SECS: u64 = 300;
65
66    /// Handles the heartbeat request.
67    async fn heartbeat(&self) {
68        self.safety_check_minimum_number_of_peers();
69        self.log_connected_peers();
70
71        // Remove the oldest connected peer.
72        self.remove_oldest_connected_peer();
73        // Keep the number of connected peers within the allowed range.
74        self.handle_connected_peers().await;
75        // Keep the bootstrap peers within the allowed range.
76        self.handle_bootstrap_peers().await;
77        // Keep the trusted peers connected.
78        self.handle_trusted_peers().await;
79        // Keep the puzzle request up to date.
80        self.handle_puzzle_request();
81        // Unban any addresses whose ban time has expired.
82        self.handle_banned_ips();
83        // Clear stale rate-limit cache entries.
84        self.clear_stale_peers();
85    }
86
87    /// TODO (howardwu): Consider checking minimum number of validators, to exclude clients and provers.
88    /// This function performs safety checks on the setting for the minimum number of peers.
89    fn safety_check_minimum_number_of_peers(&self) {
90        // Perform basic sanity checks on the configuration for the number of peers.
91        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    /// This function logs the connected peers.
99    fn log_connected_peers(&self) {
100        // Log the connected peers.
101        let connected_peers = self.router().connected_peers();
102        let connected_peers_fmt = format!("{connected_peers:?}").dimmed();
103        match connected_peers.len() {
104            0 => {
105                // Only log a warning if the node has been running for a while.
106                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    /// Returns a sorted vector of network addresses of all removable connected peers
116    /// where the first entry has the lowest priority and the last one the highest.
117    ///
118    /// Rules:
119    ///     - Trusted peers and bootstrap nodes are not removable.
120    ///     - Peers that we are currently syncing with are not removable.
121    ///     - Connections that have not been seen in a while are considered lower priority.
122    fn get_removable_peers(&self) -> Vec<ConnectedPeer<N>> {
123        // Are we synced already? (cache this here, so it does not need to be recomputed)
124        let is_block_synced = self.is_block_synced();
125
126        // Sort by priority, where lowest priority will be at the beginning
127        // of the vector.
128        // Note, that this gives equal priority to clients and provers, which
129        // we might want to change in the future.
130        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) // This peer is currently syncing from us.
134                && (is_block_synced || self.router().cache.num_outbound_block_requests(&peer.listener_addr) == 0) // We are currently syncing from this peer.
135        });
136        peers.sort_by_key(|peer| peer.last_seen);
137
138        peers
139    }
140
141    /// This function removes the peer that we have not heard from the longest,
142    /// to keep the connections fresh.
143    /// It only triggers if the router is above the minimum number of connected peers.
144    fn remove_oldest_connected_peer(&self) {
145        // Skip if the node is not requesting peers.
146        if self.router().trusted_peers_only() {
147            return;
148        }
149
150        // Skip if the router is at or below the minimum number of connected peers.
151        if self.router().number_of_connected_peers() <= Self::MINIMUM_NUMBER_OF_PEERS {
152            return;
153        }
154
155        // Disconnect from the oldest connected peer, which is the first entry in the list
156        // of removable peers.
157        // Do nothing, if the list is empty.
158        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    /// Logs a message with the error and `context` if the connection attempt failed,
166    /// and sets the log level based on the severity of the error.
167    #[inline]
168    fn log_if_connect_error(result: Result<Result<(), ConnectError>, JoinError>, context: &str) {
169        match result {
170            // Success!
171            Ok(Ok(())) => {}
172            Ok(Err(err @ ConnectError::AlreadyConnecting { .. }))
173            | Ok(Err(err @ ConnectError::AlreadyConnected { .. })) => {
174                // Log benign errors at a lower level.
175                debug!("{context}: {err}");
176            }
177            // Print regular connection errors (such as "connection refused" as warnings)
178            Ok(Err(err)) => warn!("{context}: {err}"),
179            // Print join errors as error, as they most likely indicate a crash.
180            Err(err) => error!("{context}: {err}"),
181        }
182    }
183
184    /// This function keeps the number of connected peers within the allowed range.
185    async fn handle_connected_peers(&self) {
186        // Initialize an RNG.
187        let rng = &mut ChaChaRng::from_rng(&mut rand::rng());
188
189        // Obtain the number of connected peers.
190        let num_connected = self.router().number_of_connected_peers();
191        // Obtain the number of connected provers.
192        let num_connected_provers = self.router().filter_connected_peers(|peer| peer.node_type.is_prover()).len();
193
194        // Determine the maximum number of peers and provers to keep.
195        let (max_peers, max_provers) = (Self::MAXIMUM_NUMBER_OF_PEERS, Self::MAXIMUM_NUMBER_OF_PROVERS);
196
197        // Compute the number of surplus peers.
198        let num_surplus_peers = num_connected.saturating_sub(max_peers);
199        // Compute the number of surplus provers.
200        let num_surplus_provers = num_connected_provers.saturating_sub(max_provers);
201        // Compute the number of provers remaining connected.
202        let num_remaining_provers = num_connected_provers.saturating_sub(num_surplus_provers);
203        // Compute the number of surplus clients and validators.
204        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            // Determine the provers to disconnect from.
212            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            // Determine the clients and validators to disconnect from.
219            let peers_to_disconnect = self
220                .get_removable_peers()
221                .into_iter()
222                .filter(|peer| !peer.node_type.is_prover()) // remove provers as those are handled separately
223                .take(num_surplus_clients_validators);
224
225            // Proceed to send disconnect requests to these peers.
226            for peer in peers_to_disconnect.chain(provers_to_disconnect) {
227                // TODO (howardwu): Remove this after specializing this function.
228                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                // Disconnect from this peer.
236                self.router().disconnect(peer_addr);
237            }
238        }
239
240        // Obtain the number of connected peers.
241        let num_connected = self.router().number_of_connected_peers();
242        // Compute the number of deficit peers.
243        let num_deficient = Self::MEDIAN_NUMBER_OF_PEERS.saturating_sub(num_connected);
244
245        if num_deficient > 0 {
246            // Initialize an RNG.
247            let rng = &mut ChaChaRng::from_rng(&mut rand::rng());
248
249            // Attempt to connect to more peers, separately choosing from those at a greater block
250            // height, and those whose height is lower or unknown to us.
251            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            // We may not know of half of `num_deficient` candidates; account for it using `min`.
258            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            // Initiate connection attempts and wait for them to complete.
264            self.try_connect_to_peers(higher_peers.into_iter().chain(other_peers)).await;
265
266            if !self.router().trusted_peers_only() {
267                // Request more peers from the connected peers.
268                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    /// This function keeps the number of bootstrap peers within the allowed range.
276    async fn handle_bootstrap_peers(&self) {
277        // Return early if we are in trusted peers only mode.
278        if self.router().trusted_peers_only() {
279            return;
280        }
281        // Split the bootstrap peers into connected and candidate lists.
282        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 there are not enough connected bootstrap peers, connect to more.
291        if connected_bootstrap.is_empty() {
292            // Initialize an RNG.
293            let rng = &mut ChaChaRng::from_rng(&mut rand::rng());
294            // Attempt to connect to a random bootstrap peer.
295            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        // Determine if the node is connected to more bootstrap peers than allowed.
309        let num_surplus = connected_bootstrap.len().saturating_sub(1);
310        if num_surplus > 0 {
311            // Initialize an RNG.
312            let rng = &mut ChaChaRng::from_rng(&mut rand::rng());
313            // Proceed to send disconnect requests to these bootstrap peers.
314            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                // Disconnect from this peer.
318                self.router().disconnect(peer.listener_addr);
319            }
320        }
321    }
322
323    /// Helper function that attempts to connect the given peers.
324    ///
325    /// Used by [`Self::handle_trusted_peers`] and [`Self::handle_connected_peers`].
326    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                // Do not attempt to reconnect too frequently.
332                // TODO (kaimast): Consider increasing the minimum time based on the number of failed attempts.
333                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                // Get the peers address.
340                let addr = peer.listener_addr;
341                let attempt_no = peer.total_connection_attempts + 1;
342
343                // Start connection attempt.
344                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        // Wait for all the connection attempts to complete.
357        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    /// This function attempts to connect to any disconnected trusted peers.
366    async fn handle_trusted_peers(&self) {
367        self.try_connect_to_peers(self.router().get_trusted_candidate_peers().into_iter()).await;
368    }
369
370    /// This function updates the puzzle if network has updated.
371    fn handle_puzzle_request(&self) {
372        // No-op
373    }
374
375    // Remove addresses whose ban time has expired.
376    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}