network.rs (14595B)
1 use std::{ 2 collections::{BTreeMap, BTreeSet}, 3 net::{IpAddr, SocketAddr}, 4 sync::{Arc, Mutex as StdMutex}, 5 }; 6 7 use anyhow::{Context, Result}; 8 use tokio::{net::TcpListener, sync::mpsc}; 9 10 use crate::app::{ 11 BlockInventory, GossipEnvelope, SharedNode, SharedPeerBook, debug_logging_enabled, now_ms, 12 }; 13 14 use super::{ 15 GossipNetwork, GossipNetworkInner, InboundConnectionLimiter, InboundSessionPermit, 16 InboundSessionRejection, OutboundBatch, P2pMetrics, P2pMetricsCounters, PEER_QUEUE_SIZE, 17 STALE_INBOUND_PEER_RETENTION_MS, accept_loop, is_self_peer_address_for, new_node_id, 18 outbound_session, outbound_supervisor, 19 }; 20 21 impl GossipNetwork { 22 #[cfg(test)] 23 pub(crate) fn new_for_tests(node: SharedNode, peers: SharedPeerBook) -> Self { 24 Self { 25 inner: Arc::new(GossipNetworkInner { 26 node, 27 peers, 28 listen_addr: "127.0.0.1:0".parse().unwrap(), 29 p2p_announce_addr: tokio::sync::Mutex::new(None), 30 node_id: new_node_id(), 31 accept_task: tokio::sync::Mutex::new(None), 32 sessions: tokio::sync::Mutex::new( 33 BTreeMap::<String, mpsc::Sender<OutboundBatch>>::new(), 34 ), 35 inbound_limiter: Arc::new(StdMutex::new(InboundConnectionLimiter::default())), 36 metrics: P2pMetricsCounters::default(), 37 }), 38 } 39 } 40 41 pub async fn start( 42 node: SharedNode, 43 peers: SharedPeerBook, 44 addr: SocketAddr, 45 p2p_announce_addr: Option<SocketAddr>, 46 accept_inbound: bool, 47 ) -> Result<Self> { 48 let network = Self { 49 inner: Arc::new(GossipNetworkInner { 50 node, 51 peers, 52 listen_addr: addr, 53 p2p_announce_addr: tokio::sync::Mutex::new(p2p_announce_addr), 54 node_id: new_node_id(), 55 accept_task: tokio::sync::Mutex::new(None), 56 sessions: tokio::sync::Mutex::new( 57 BTreeMap::<String, mpsc::Sender<OutboundBatch>>::new(), 58 ), 59 inbound_limiter: Arc::new(StdMutex::new(InboundConnectionLimiter::default())), 60 metrics: P2pMetricsCounters::default(), 61 }), 62 }; 63 64 if accept_inbound { 65 network.set_accept_inbound(true).await?; 66 } 67 tokio::spawn(outbound_supervisor(network.clone())); 68 network.ensure_outbound_sessions().await; 69 Ok(network) 70 } 71 72 pub async fn set_accept_inbound(&self, enabled: bool) -> Result<()> { 73 let mut accept_task = self.inner.accept_task.lock().await; 74 if enabled { 75 if accept_task.is_some() { 76 return Ok(()); 77 } 78 let listener = TcpListener::bind(self.inner.listen_addr) 79 .await 80 .with_context(|| format!("binding p2p listener on {}", self.inner.listen_addr))?; 81 *accept_task = Some(tokio::spawn(accept_loop(self.clone(), listener))); 82 } else if let Some(task) = accept_task.take() { 83 task.abort(); 84 } 85 Ok(()) 86 } 87 88 pub async fn accepts_inbound(&self) -> bool { 89 self.inner.accept_task.lock().await.is_some() 90 } 91 92 pub fn listen_addr(&self) -> SocketAddr { 93 self.inner.listen_addr 94 } 95 96 pub async fn set_p2p_announce_addr(&self, addr: Option<SocketAddr>) { 97 *self.inner.p2p_announce_addr.lock().await = addr; 98 } 99 100 pub(super) async fn advertised_addr(&self) -> Option<SocketAddr> { 101 if !self.accepts_inbound().await { 102 return None; 103 } 104 Some((*self.inner.p2p_announce_addr.lock().await).unwrap_or(self.inner.listen_addr)) 105 } 106 107 pub(super) async fn self_filter_addr(&self) -> Option<SocketAddr> { 108 if let Some(addr) = *self.inner.p2p_announce_addr.lock().await { 109 return Some(addr); 110 } 111 self.accepts_inbound() 112 .await 113 .then_some(self.inner.listen_addr) 114 } 115 116 pub(super) async fn is_self_peer(&self, address: &str) -> bool { 117 is_self_peer_address_for( 118 address, 119 self.inner.listen_addr, 120 self.self_filter_addr().await, 121 ) 122 } 123 124 pub fn metrics(&self) -> P2pMetrics { 125 self.inner.metrics.snapshot() 126 } 127 128 pub(super) fn try_acquire_inbound_session( 129 &self, 130 ip: IpAddr, 131 ) -> std::result::Result<InboundSessionPermit, InboundSessionRejection> { 132 self.inner 133 .inbound_limiter 134 .lock() 135 .expect("inbound limiter mutex poisoned") 136 .try_acquire(ip, now_ms())?; 137 Ok(InboundSessionPermit { 138 limiter: Arc::clone(&self.inner.inbound_limiter), 139 ip, 140 }) 141 } 142 143 pub async fn broadcast(&self, envelopes: Vec<GossipEnvelope>) -> Result<()> { 144 let envelopes = self.prepare_gossip(envelopes).await; 145 if envelopes.is_empty() { 146 return Ok(()); 147 } 148 149 let sessions = self.inner.sessions.lock().await.clone(); 150 for (peer, sender) in sessions { 151 if self.inner.peers.lock().await.is_banned(&peer) { 152 continue; 153 } 154 match sender.try_send(envelopes.clone()) { 155 Ok(()) => {} 156 Err(mpsc::error::TrySendError::Full(_)) => { 157 P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_full); 158 } 159 Err(mpsc::error::TrySendError::Closed(_)) => { 160 P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_closed); 161 } 162 } 163 } 164 Ok(()) 165 } 166 167 async fn prepare_gossip(&self, envelopes: Vec<GossipEnvelope>) -> Vec<GossipEnvelope> { 168 let mut blocks = Vec::new(); 169 let mut passthrough = Vec::new(); 170 171 for envelope in envelopes { 172 match envelope { 173 GossipEnvelope::Block(block) => blocks.push(BlockInventory { 174 height: block.height, 175 hash: block.hash, 176 }), 177 GossipEnvelope::Blocks { blocks: batch } => { 178 blocks.extend(batch.into_iter().map(|block| BlockInventory { 179 height: block.height, 180 hash: block.hash, 181 })); 182 } 183 GossipEnvelope::Inventory { blocks: inv_blocks } => blocks.extend(inv_blocks), 184 other => passthrough.push(other), 185 } 186 } 187 188 blocks.sort_by(|left, right| { 189 left.height 190 .cmp(&right.height) 191 .then_with(|| left.hash.cmp(&right.hash)) 192 }); 193 blocks.dedup_by(|left, right| left.hash == right.hash); 194 195 if !blocks.is_empty() { 196 passthrough.push(GossipEnvelope::Inventory { blocks }); 197 } 198 passthrough 199 } 200 201 pub async fn peer_exchange(&self) -> GossipEnvelope { 202 let advertised_addr = self.advertised_addr().await; 203 let self_filter_addr = self.self_filter_addr().await; 204 let self_addr = advertised_addr.map(|addr| addr.to_string()); 205 let peers = self 206 .inner 207 .peers 208 .lock() 209 .await 210 .addresses_except(self_addr.as_deref().unwrap_or("")) 211 .into_iter() 212 .filter(|peer| peer.parse::<SocketAddr>().is_ok()) 213 .filter(|peer| { 214 !is_self_peer_address_for(peer, self.inner.listen_addr, self_filter_addr) 215 }) 216 .collect::<Vec<_>>(); 217 GossipEnvelope::PeerList { 218 peers: self_addr.into_iter().chain(peers.into_iter()).collect(), 219 } 220 } 221 222 pub(super) async fn ensure_outbound_sessions(&self) { 223 self.inner 224 .peers 225 .lock() 226 .await 227 .prune_stale_inbound_peers_at(crate::app::now_ms(), STALE_INBOUND_PEER_RETENTION_MS); 228 let addresses = self 229 .inner 230 .peers 231 .lock() 232 .await 233 .connectable_addresses_at(crate::app::now_ms()); 234 let address_set = addresses.iter().cloned().collect::<BTreeSet<_>>(); 235 let self_filter_addr = self.self_filter_addr().await; 236 let mut sessions = self.inner.sessions.lock().await; 237 sessions.retain(|peer, _| { 238 let keep = address_set.contains(peer) 239 && !is_self_peer_address_for(peer, self.inner.listen_addr, self_filter_addr); 240 if !keep { 241 P2pMetricsCounters::inc(&self.inner.metrics.self_peer_skips); 242 } 243 keep 244 }); 245 for peer in addresses { 246 if is_self_peer_address_for(&peer, self.inner.listen_addr, self_filter_addr) { 247 P2pMetricsCounters::inc(&self.inner.metrics.self_peer_skips); 248 continue; 249 } 250 if sessions.contains_key(&peer) { 251 continue; 252 } 253 254 let (sender, receiver) = mpsc::channel(PEER_QUEUE_SIZE); 255 sessions.insert(peer.clone(), sender); 256 tokio::spawn(outbound_session(self.clone(), peer, receiver)); 257 } 258 } 259 260 pub(super) async fn forward_outbox(&self) { 261 let outbox = self.inner.node.lock().await.drain_outbox(); 262 if let Err(error) = self.broadcast(outbox).await { 263 if debug_logging_enabled() { 264 eprintln!("p2p rebroadcast failed: {error:#}"); 265 } 266 } 267 } 268 } 269 270 #[cfg(test)] 271 mod tests { 272 use std::sync::Arc; 273 274 use crate::{ 275 app::{GossipEnvelope, PeerBook}, 276 domain::Wallet, 277 }; 278 279 use super::super::test_support::{allocations, gossip_network, node}; 280 281 #[tokio::test] 282 async fn peer_exchange_does_not_advertise_self_when_outbound_only() { 283 let alice = Wallet::from_seed("px-private-alice"); 284 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 285 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 286 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 287 "127.0.0.1:9545".to_string(), 288 ]))); 289 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 290 match network.peer_exchange().await { 291 GossipEnvelope::PeerList { peers } => { 292 assert!(!peers.contains(&"127.0.0.1:9544".to_string())); 293 assert!(peers.contains(&"127.0.0.1:9545".to_string())); 294 } 295 other => panic!("expected peer list, got {other:?}"), 296 } 297 } 298 299 #[tokio::test] 300 async fn peer_exchange_omits_hostname_bootstrap_peers() { 301 let alice = Wallet::from_seed("px-hostname-alice"); 302 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 303 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 304 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 305 "iuna.jhx.app:9444".to_string(), 306 "127.0.0.1:9545".to_string(), 307 ]))); 308 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 309 match network.peer_exchange().await { 310 GossipEnvelope::PeerList { peers } => { 311 assert!(!peers.contains(&"iuna.jhx.app:9444".to_string())); 312 assert!(peers.contains(&"127.0.0.1:9545".to_string())); 313 } 314 other => panic!("expected peer list, got {other:?}"), 315 } 316 } 317 318 #[tokio::test] 319 async fn peer_exchange_advertises_discovered_listening_peers() { 320 let alice = Wallet::from_seed("px-discovered-alice"); 321 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 322 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 323 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 324 peers 325 .lock() 326 .await 327 .add_discovered_peer("127.0.0.1:9546".to_string()); 328 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 329 330 match network.peer_exchange().await { 331 GossipEnvelope::PeerList { peers } => { 332 assert!(peers.contains(&"127.0.0.1:9546".to_string())); 333 } 334 other => panic!("expected peer list, got {other:?}"), 335 } 336 } 337 338 #[tokio::test] 339 async fn peer_exchange_advertises_stable_listen_and_known_peers() { 340 let alice = Wallet::from_seed("px-alice"); 341 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 342 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 343 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 344 "127.0.0.1:9545".to_string(), 345 ]))); 346 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 347 network.set_accept_inbound(true).await.unwrap(); 348 349 match network.peer_exchange().await { 350 GossipEnvelope::PeerList { peers } => { 351 assert!(peers.contains(&"127.0.0.1:9544".to_string())); 352 assert!(peers.contains(&"127.0.0.1:9545".to_string())); 353 } 354 other => panic!("expected peer list, got {other:?}"), 355 } 356 network.set_accept_inbound(false).await.unwrap(); 357 } 358 359 #[tokio::test] 360 async fn peer_exchange_filters_announced_self_from_known_peers() { 361 let alice = Wallet::from_seed("px-announced-self-alice"); 362 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 363 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 364 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 365 "8.8.8.8:9444".to_string(), 366 "8.8.4.4:9444".to_string(), 367 ]))); 368 let network = gossip_network( 369 node, 370 peers, 371 "127.0.0.1:0".parse().unwrap(), 372 Some("8.8.8.8:9444".parse().unwrap()), 373 ); 374 network.set_accept_inbound(true).await.unwrap(); 375 376 match network.peer_exchange().await { 377 GossipEnvelope::PeerList { peers } => { 378 assert_eq!( 379 peers.iter().filter(|peer| *peer == "8.8.8.8:9444").count(), 380 1 381 ); 382 assert!(peers.contains(&"8.8.4.4:9444".to_string())); 383 } 384 other => panic!("expected peer list, got {other:?}"), 385 } 386 network.set_accept_inbound(false).await.unwrap(); 387 } 388 }