snarkos_node_router/
lib.rs1#![forbid(unsafe_code)]
17
18#[macro_use]
19extern crate async_trait;
20#[macro_use]
21extern crate tracing;
22
23#[cfg(feature = "metrics")]
24extern crate snarkos_node_metrics as metrics;
25
26pub use snarkos_node_router_messages as messages;
27use snarkos_utilities::NodeDataDir;
28
29mod handshake;
30
31mod heartbeat;
32pub use heartbeat::*;
33
34mod helpers;
35pub use helpers::*;
36
37mod inbound;
38pub use inbound::*;
39
40mod outbound;
41pub use outbound::*;
42
43mod routing;
44pub use routing::*;
45
46mod writing;
47
48use crate::messages::{BlockRequest, Message, MessageCodec};
49
50use snarkos_account::Account;
51use snarkos_node_bft_ledger_service::LedgerService;
52use snarkos_node_network::{
53 CandidatePeer,
54 ConnectedPeer,
55 ConnectionMode,
56 NodeType,
57 Peer,
58 PeerPoolHandling,
59 Resolver,
60 bootstrap_peers,
61};
62use snarkos_node_sync_communication_service::CommunicationService;
63use snarkos_node_tcp::{Config, ConnectionSide, Tcp};
64
65use snarkvm::prelude::{Address, Network, PrivateKey, ViewKey};
66
67use anyhow::Result;
68#[cfg(feature = "locktick")]
69use locktick::parking_lot::{Mutex, RwLock};
70#[cfg(not(feature = "locktick"))]
71use parking_lot::{Mutex, RwLock};
72use std::{collections::HashMap, future::Future, io, net::SocketAddr, ops::Deref, sync::Arc};
73use tokio::task::JoinHandle;
74
75pub const DEFAULT_NODE_PORT: u16 = 4130;
77
78#[derive(Clone)]
82pub struct Router<N: Network>(Arc<InnerRouter<N>>);
83
84impl<N: Network> Deref for Router<N> {
85 type Target = Arc<InnerRouter<N>>;
86
87 fn deref(&self) -> &Self::Target {
88 &self.0
89 }
90}
91
92impl<N: Network> PeerPoolHandling<N> for Router<N> {
93 const MAXIMUM_POOL_SIZE: usize = 10_000;
94 const OWNER: &str = "[Router]";
95 const PEER_SLASHING_COUNT: usize = 200;
96
97 fn peer_pool(&self) -> &RwLock<HashMap<SocketAddr, Peer<N>>> {
98 &self.peer_pool
99 }
100
101 fn resolver(&self) -> &RwLock<Resolver<N>> {
102 &self.resolver
103 }
104
105 fn is_dev(&self) -> bool {
106 self.is_dev
107 }
108
109 fn trusted_peers_only(&self) -> bool {
110 self.trusted_peers_only
111 }
112
113 fn node_type(&self) -> NodeType {
114 self.node_type
115 }
116}
117
118pub struct InnerRouter<N: Network> {
119 tcp: Tcp,
121 node_type: NodeType,
123 account: Account<N>,
125 ledger: Arc<dyn LedgerService<N>>,
127 cache: Cache<N>,
129 resolver: RwLock<Resolver<N>>,
131 peer_pool: RwLock<HashMap<SocketAddr, Peer<N>>>,
133 handles: Mutex<Vec<JoinHandle<()>>>,
135 trusted_peers_only: bool,
137 node_data_dir: NodeDataDir,
139 is_dev: bool,
141}
142
143impl<N: Network> Router<N> {
144 const CONNECTION_ATTEMPTS_SINCE_SECS: i64 = 10;
146 #[cfg(not(feature = "test"))]
148 const MAX_CONNECTION_ATTEMPTS: usize = 10;
149 const MESSAGE_LIMIT_TIME_FRAME_IN_SECS: i64 = 5;
151}
152
153impl<N: Network> Router<N> {
154 #[allow(clippy::too_many_arguments)]
156 pub async fn new(
157 node_ip: SocketAddr,
158 node_type: NodeType,
159 account: Account<N>,
160 ledger: Arc<dyn LedgerService<N>>,
161 trusted_peers: &[SocketAddr],
162 max_peers: u16,
163 trusted_peers_only: bool,
164 node_data_dir: NodeDataDir,
165 is_dev: bool,
166 ) -> Result<Self> {
167 let tcp = Tcp::new(Config::new(node_ip, max_peers));
169
170 let mut initial_peers = HashMap::new();
172
173 if !trusted_peers_only {
175 let cached_peers = Self::load_cached_peers(&node_data_dir.router_peer_cache_path())?;
176 for addr in cached_peers {
177 initial_peers.insert(addr, Peer::new_candidate(addr, false));
178 }
179 }
180
181 initial_peers.extend(trusted_peers.iter().copied().map(|addr| (addr, Peer::new_candidate(addr, true))));
184
185 Ok(Self(Arc::new(InnerRouter {
187 tcp,
188 node_type,
189 account,
190 ledger,
191 cache: Default::default(),
192 resolver: Default::default(),
193 peer_pool: RwLock::new(initial_peers),
194 handles: Default::default(),
195 trusted_peers_only,
196 node_data_dir,
197 is_dev,
198 })))
199 }
200}
201
202impl<N: Network> Router<N> {
203 pub fn is_valid_message_version(&self, message_version: u32) -> bool {
205 let lowest_accepted_message_version = match self.node_type {
209 NodeType::Prover | NodeType::BootstrapClient => Message::<N>::latest_message_version(),
212 NodeType::Validator | NodeType::Client => {
214 Message::<N>::lowest_accepted_message_version(self.ledger.latest_block_height())
215 }
216 };
217
218 message_version >= lowest_accepted_message_version
220 }
221
222 pub fn private_key(&self) -> &PrivateKey<N> {
224 self.account.private_key()
225 }
226
227 pub fn view_key(&self) -> &ViewKey<N> {
229 self.account.view_key()
230 }
231
232 pub fn address(&self) -> Address<N> {
234 self.account.address()
235 }
236
237 pub fn cache(&self) -> &Cache<N> {
239 &self.cache
240 }
241
242 pub fn ledger(&self) -> &Arc<dyn LedgerService<N>> {
244 &self.ledger
245 }
246
247 pub fn trusted_peers_only(&self) -> bool {
249 self.trusted_peers_only
250 }
251
252 pub fn resolve_to_listener(&self, connected_addr: SocketAddr) -> Option<SocketAddr> {
254 self.resolver.read().get_listener(connected_addr)
255 }
256
257 pub fn connected_metrics(&self) -> Vec<(SocketAddr, NodeType)> {
259 self.get_connected_peers().iter().map(|peer| (peer.listener_addr, peer.node_type)).collect()
260 }
261
262 #[cfg(feature = "metrics")]
263 pub fn update_metrics(&self) {
264 metrics::gauge(metrics::router::CONNECTED, self.number_of_connected_peers() as f64);
265 metrics::gauge(metrics::router::CANDIDATE, self.number_of_candidate_peers() as f64);
266 }
267
268 pub fn update_last_seen_for_connected_peer(&self, peer_ip: SocketAddr) {
269 if let Some(peer) = self.peer_pool.write().get_mut(&peer_ip) {
270 peer.update_last_seen();
271 }
272 }
273
274 pub fn spawn<T: Future<Output = ()> + Send + 'static>(&self, future: T) {
276 self.handles.lock().push(tokio::spawn(future));
277 }
278
279 pub async fn shut_down(&self) {
281 info!("Shutting down the router...");
282 if let Err(e) =
284 self.save_best_peers(&self.node_data_dir.router_peer_cache_path(), Some(MAX_PEERS_TO_SEND), true)
285 {
286 warn!("Failed to persist best peers to disk: {e}");
287 }
288 self.handles.lock().iter().for_each(|handle| handle.abort());
290 self.tcp.shut_down().await;
292 }
293}
294
295#[async_trait]
296impl<N: Network> CommunicationService for Router<N> {
297 type Message = Message<N>;
299
300 fn prepare_block_request(start_height: u32, end_height: u32) -> Self::Message {
302 debug_assert!(start_height < end_height, "Invalid block request format");
303 Message::BlockRequest(BlockRequest { start_height, end_height })
304 }
305
306 async fn send(
312 &self,
313 peer_ip: SocketAddr,
314 message: Self::Message,
315 ) -> Option<tokio::sync::oneshot::Receiver<io::Result<()>>> {
316 self.send(peer_ip, message)
317 }
318}