snarkos_node/bootstrap_client/
network.rs1use crate::{
17 BootstrapClient,
18 bft::{
19 MAX_VALIDATORS_TO_SEND,
20 events::{self, Event},
21 },
22 bootstrap_client::codec::BootstrapClientCodec,
23 network::{ConnectionMode, NodeType, Peer, PeerPoolHandling, Resolver},
24 router::{
25 MAX_PEERS_TO_SEND,
26 messages::{self, Message},
27 },
28 tcp::{ConnectionSide, P2P, Tcp, connections::DisconnectOrigin, protocols::*},
29};
30use snarkvm::prelude::Network;
31
32use indexmap::IndexMap;
33#[cfg(feature = "locktick")]
34use locktick::parking_lot::RwLock;
35#[cfg(not(feature = "locktick"))]
36use parking_lot::RwLock;
37use std::{collections::HashMap, io, net::SocketAddr};
38use tokio::time::sleep;
39use tokio_util::codec::Decoder;
40
41impl<N: Network> P2P for BootstrapClient<N> {
42 fn tcp(&self) -> &Tcp {
43 &self.tcp
44 }
45}
46
47impl<N: Network> PeerPoolHandling<N> for BootstrapClient<N> {
48 const MAXIMUM_POOL_SIZE: usize = 10_000;
49 const OWNER: &'static str = "[Network]";
50 const PEER_SLASHING_COUNT: usize = 200;
51
52 fn is_dev(&self) -> bool {
53 self.dev.is_some()
54 }
55
56 fn trusted_peers_only(&self) -> bool {
57 false
58 }
59
60 fn node_type(&self) -> NodeType {
61 NodeType::BootstrapClient
62 }
63
64 fn peer_pool(&self) -> &RwLock<HashMap<SocketAddr, Peer<N>>> {
65 &self.peer_pool
66 }
67
68 fn resolver(&self) -> &RwLock<Resolver<N>> {
69 &self.resolver
70 }
71}
72
73#[derive(Debug)]
75pub enum MessageOrEvent<N: Network> {
76 Message(Message<N>),
77 Event(Event<N>),
78}
79
80#[async_trait]
81impl<N: Network> OnConnect for BootstrapClient<N> {
82 async fn on_connect(&self, peer_addr: SocketAddr) {
83 if let Some(listener_addr) = self.resolve_to_listener(peer_addr)
86 && let Some(peer) = self.get_connected_peer(listener_addr)
87 && peer.node_type == NodeType::Validator
88 {
89 self.known_validators.write().insert(listener_addr, (peer.aleo_addr, peer.connection_mode));
90 }
91 let tcp = self.tcp().clone();
94 tokio::spawn(async move {
95 sleep(Self::CONNECTION_LIFETIME).await;
96 tcp.disconnect(peer_addr).await;
97 });
98 }
99}
100
101#[async_trait]
102impl<N: Network> Disconnect for BootstrapClient<N> {
103 async fn handle_disconnect(&self, peer_addr: SocketAddr, origin: DisconnectOrigin) {
105 debug!("Physically disconnecting from {peer_addr}; origin: {origin:?}");
106
107 if let Some(listener_addr) = self.resolve_to_listener(peer_addr) {
108 self.downgrade_peer_to_candidate(listener_addr);
109 }
110 }
111}
112
113#[async_trait]
114impl<N: Network> Reading for BootstrapClient<N> {
115 type Codec = BootstrapClientCodec<N>;
116 type Message = <BootstrapClientCodec<N> as Decoder>::Item;
117
118 fn codec(&self, _peer_addr: SocketAddr, _side: ConnectionSide) -> Self::Codec {
121 Default::default()
122 }
123
124 async fn process_message(&self, peer_addr: SocketAddr, message: Self::Message) -> io::Result<()> {
126 let Some(listener_addr) = self.resolve_to_listener(peer_addr) else {
128 return Ok(());
130 };
131
132 match message {
134 MessageOrEvent::Message(Message::PeerRequest(_)) => {
135 debug!("Received a PeerRequest from '{listener_addr}'");
136 let mut peers = self.get_candidate_peers();
137
138 let Some(peer) = self.get_connected_peer(listener_addr) else {
141 return Ok(());
142 };
143 let validators = self.get_validator_addrs().await;
144
145 if peer.node_type == NodeType::Validator {
146 peers.retain(|peer| {
148 validators
149 .get(&peer.listener_addr)
150 .map(|(_, connection_mode)| *connection_mode != ConnectionMode::Gateway)
151 .unwrap_or(true)
152 });
153 } else {
154 peers.retain(|peer| !validators.contains_key(&peer.listener_addr));
156 }
157 peers.truncate(MAX_PEERS_TO_SEND);
158 let peers = peers.into_iter().map(|peer| (peer.listener_addr, None)).collect::<Vec<_>>();
159
160 debug!("Sending {} peer address(es) to '{listener_addr}'", peers.len());
161 let msg = MessageOrEvent::Message(Message::PeerResponse(messages::PeerResponse { peers }));
162 if let Err(err) = self.unicast(peer_addr, msg)?.await {
163 warn!("Couldn't deliver a peer list to '{listener_addr}': {err}; disconnecting");
164 } else {
165 debug!("Disconnecting from '{listener_addr}' - peers provided");
166 }
167
168 self.tcp().disconnect(peer_addr).await;
169 }
170 MessageOrEvent::Event(Event::ValidatorsRequest(_)) => {
171 debug!("Received a ValidatorsRequest from '{listener_addr}'");
172
173 let validators = self.get_validator_addrs().await;
175 let validators = validators
176 .into_iter()
177 .filter_map(|(listener_addr, (aleo_addr, connection_mode))| {
178 (connection_mode == ConnectionMode::Gateway).then_some((listener_addr, aleo_addr))
180 })
181 .take(MAX_VALIDATORS_TO_SEND)
182 .collect::<IndexMap<_, _>>();
183
184 debug!("Sending {} validator address(es) to '{listener_addr}'", validators.len());
185 let msg = MessageOrEvent::Event(Event::ValidatorsResponse(events::ValidatorsResponse { validators }));
186 if let Err(err) = self.unicast(peer_addr, msg)?.await {
187 warn!("Couldn't deliver a peer list to '{listener_addr}': {err}; disconnecting");
188 } else {
189 debug!("Disconnecting from '{listener_addr}' - peers provided");
190 }
191
192 self.tcp().disconnect(peer_addr).await;
193 }
194 msg => {
195 let name = match msg {
196 MessageOrEvent::Message(msg) => msg.name(),
197 MessageOrEvent::Event(msg) => msg.name(),
198 };
199 trace!("Ignoring an unhandled message ({name}) from {listener_addr}");
200 }
201 }
202
203 Ok(())
204 }
205}
206
207#[async_trait]
208impl<N: Network> Writing for BootstrapClient<N> {
209 type Codec = BootstrapClientCodec<N>;
210 type Message = MessageOrEvent<N>;
211
212 fn codec(&self, _addr: SocketAddr, _side: ConnectionSide) -> Self::Codec {
215 Default::default()
216 }
217}