network.rs (28428B)
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::{ 9 net::TcpListener, 10 sync::{mpsc, watch}, 11 }; 12 13 use crate::app::{ 14 BlockInventory, GossipEnvelope, SharedNode, SharedPeerBook, debug_logging_enabled, now_ms, 15 }; 16 17 use super::{ 18 ChainValidationCoordinator, ChainValidationGuard, GossipNetwork, GossipNetworkInner, 19 GossipSession, INBOUND_SESSION_PREFIX, InboundConnectionLimiter, InboundSessionPermit, 20 InboundSessionRejection, MAX_DISCOVERED_OUTBOUND_DIALS_PER_CYCLE, MAX_GOSSIP_LINE_BYTES, 21 MAX_OUTBOUND_BATCH_BYTES, OutboundBatch, P2pMetrics, P2pMetricsCounters, PEER_QUEUE_BYTES, 22 PEER_QUEUE_SIZE, STALE_DISCOVERED_PEER_RETENTION_MS, STALE_INBOUND_PEER_RETENTION_MS, 23 accept_loop, is_self_peer_address_for, new_node_id, outbound_session, outbound_supervisor, 24 }; 25 26 impl GossipNetwork { 27 pub async fn start( 28 node: SharedNode, 29 peers: SharedPeerBook, 30 addr: SocketAddr, 31 p2p_announce_addr: Option<SocketAddr>, 32 accept_inbound: bool, 33 ) -> Result<Self> { 34 let network = Self { 35 inner: Arc::new(GossipNetworkInner { 36 node, 37 peers, 38 listen_addr: addr, 39 p2p_announce_addr: tokio::sync::Mutex::new(p2p_announce_addr), 40 node_id: new_node_id(), 41 accept_task: tokio::sync::Mutex::new(None), 42 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 43 inbound_limiter: Arc::new(StdMutex::new(InboundConnectionLimiter::default())), 44 metrics: P2pMetricsCounters::default(), 45 sync_progress: StdMutex::new(super::SyncProgressState::default()), 46 chain_validation: Arc::new(ChainValidationCoordinator::default()), 47 }), 48 }; 49 50 if accept_inbound { 51 network.set_accept_inbound(true).await?; 52 } 53 tokio::spawn(outbound_supervisor(network.clone())); 54 network.ensure_outbound_sessions().await; 55 Ok(network) 56 } 57 58 pub async fn set_accept_inbound(&self, enabled: bool) -> Result<()> { 59 let mut accept_task = self.inner.accept_task.lock().await; 60 if enabled { 61 if accept_task.is_some() { 62 return Ok(()); 63 } 64 let listener = TcpListener::bind(self.inner.listen_addr) 65 .await 66 .with_context(|| format!("binding p2p listener on {}", self.inner.listen_addr))?; 67 *accept_task = Some(tokio::spawn(accept_loop(self.clone(), listener))); 68 } else if let Some(task) = accept_task.take() { 69 task.abort(); 70 } 71 Ok(()) 72 } 73 74 pub async fn accepts_inbound(&self) -> bool { 75 self.inner.accept_task.lock().await.is_some() 76 } 77 78 pub fn listen_addr(&self) -> SocketAddr { 79 self.inner.listen_addr 80 } 81 82 pub async fn set_p2p_announce_addr(&self, addr: Option<SocketAddr>) { 83 *self.inner.p2p_announce_addr.lock().await = addr; 84 } 85 86 pub(super) async fn advertised_addr(&self) -> Option<SocketAddr> { 87 if !self.accepts_inbound().await { 88 return None; 89 } 90 Some((*self.inner.p2p_announce_addr.lock().await).unwrap_or(self.inner.listen_addr)) 91 } 92 93 pub(super) async fn self_filter_addr(&self) -> Option<SocketAddr> { 94 if let Some(addr) = *self.inner.p2p_announce_addr.lock().await { 95 return Some(addr); 96 } 97 self.accepts_inbound() 98 .await 99 .then_some(self.inner.listen_addr) 100 } 101 102 pub(super) async fn is_self_peer(&self, address: &str) -> bool { 103 is_self_peer_address_for( 104 address, 105 self.inner.listen_addr, 106 self.self_filter_addr().await, 107 ) 108 } 109 110 pub fn metrics(&self) -> P2pMetrics { 111 self.inner.metrics.snapshot() 112 } 113 114 pub fn sync_progress(&self) -> Option<super::SyncProgress> { 115 self.inner 116 .sync_progress 117 .lock() 118 .expect("sync progress mutex poisoned") 119 .active 120 .values() 121 .copied() 122 .max_by_key(|progress| (progress.target_height, progress.validated_height)) 123 } 124 125 pub(super) fn begin_sync_progress( 126 &self, 127 start_height: u64, 128 target_height: u64, 129 ) -> super::SyncProgressGuard { 130 let mut state = self 131 .inner 132 .sync_progress 133 .lock() 134 .expect("sync progress mutex poisoned"); 135 state.next_id = state.next_id.wrapping_add(1); 136 let id = state.next_id; 137 state.active.insert( 138 id, 139 super::SyncProgress { 140 start_height, 141 validated_height: start_height, 142 target_height: target_height.max(start_height), 143 }, 144 ); 145 state.last_activity = Some(std::time::Instant::now()); 146 super::SyncProgressGuard { 147 network: self.clone(), 148 id, 149 generation: state.generation, 150 } 151 } 152 153 pub(crate) fn invalidate_sync_progress(&self) { 154 let mut state = self 155 .inner 156 .sync_progress 157 .lock() 158 .expect("sync progress mutex poisoned"); 159 state.generation = state.generation.wrapping_add(1); 160 state.active.clear(); 161 } 162 163 pub(super) fn sync_generation(&self) -> u64 { 164 self.inner 165 .sync_progress 166 .lock() 167 .expect("sync progress mutex poisoned") 168 .generation 169 } 170 171 pub(super) fn sync_generation_is_current(&self, generation: u64) -> bool { 172 self.sync_generation() == generation 173 } 174 175 pub(super) fn update_sync_progress(&self, id: u64, validated_height: u64) { 176 let mut state = self 177 .inner 178 .sync_progress 179 .lock() 180 .expect("sync progress mutex poisoned"); 181 if let Some(progress) = state.active.get_mut(&id) { 182 progress.validated_height = 183 validated_height.clamp(progress.start_height, progress.target_height); 184 state.last_activity = Some(std::time::Instant::now()); 185 } 186 } 187 188 pub(super) fn finish_sync_progress(&self, id: u64) { 189 let mut state = self 190 .inner 191 .sync_progress 192 .lock() 193 .expect("sync progress mutex poisoned"); 194 state.active.remove(&id); 195 state.last_activity = Some(std::time::Instant::now()); 196 } 197 198 pub fn chain_sync_active_or_recent(&self, quiet_period: std::time::Duration) -> bool { 199 let state = self 200 .inner 201 .sync_progress 202 .lock() 203 .expect("sync progress mutex poisoned"); 204 if !state.active.is_empty() { 205 return true; 206 } 207 state 208 .last_activity 209 .is_some_and(|last_activity| last_activity.elapsed() <= quiet_period) 210 } 211 212 pub(super) fn try_acquire_inbound_session( 213 &self, 214 ip: IpAddr, 215 ) -> std::result::Result<InboundSessionPermit, InboundSessionRejection> { 216 self.inner 217 .inbound_limiter 218 .lock() 219 .expect("inbound limiter mutex poisoned") 220 .try_acquire(ip, now_ms())?; 221 Ok(InboundSessionPermit { 222 limiter: Arc::clone(&self.inner.inbound_limiter), 223 ip, 224 }) 225 } 226 227 pub(super) async fn claim_chain_validation(&self, key: String) -> Option<ChainValidationGuard> { 228 loop { 229 let active = self 230 .inner 231 .chain_validation 232 .active 233 .lock() 234 .expect("chain validation mutex poisoned") 235 .get(&key) 236 .map(watch::Sender::subscribe); 237 if let Some(mut active) = active { 238 while !*active.borrow() { 239 if active.changed().await.is_err() { 240 break; 241 } 242 } 243 continue; 244 } 245 246 let permit = Arc::clone(&self.inner.chain_validation.permits) 247 .acquire_owned() 248 .await 249 .ok()?; 250 let duplicate = { 251 let mut active = self 252 .inner 253 .chain_validation 254 .active 255 .lock() 256 .expect("chain validation mutex poisoned"); 257 if let Some(sender) = active.get(&key) { 258 Some(sender.subscribe()) 259 } else { 260 let (sender, _) = watch::channel(false); 261 active.insert(key.clone(), sender); 262 None 263 } 264 }; 265 if let Some(mut receiver) = duplicate { 266 drop(permit); 267 while !*receiver.borrow() { 268 if receiver.changed().await.is_err() { 269 break; 270 } 271 } 272 continue; 273 } 274 return Some(ChainValidationGuard { 275 coordinator: Arc::clone(&self.inner.chain_validation), 276 key, 277 _permit: permit, 278 }); 279 } 280 } 281 282 pub async fn broadcast(&self, envelopes: Vec<GossipEnvelope>) -> Result<()> { 283 let envelopes = self.prepare_gossip(envelopes).await; 284 if envelopes.is_empty() { 285 return Ok(()); 286 } 287 288 let batches = byte_bounded_gossip_batches(envelopes)?; 289 let sessions = self.inner.sessions.lock().await.clone(); 290 let mut disconnect = Vec::new(); 291 for (session_id, session) in sessions { 292 if self.inner.peers.lock().await.is_banned(&session.peer) { 293 continue; 294 } 295 for (envelopes, encoded_bytes) in &batches { 296 let permit = 297 match Arc::clone(&session.queue_bytes).try_acquire_many_owned(*encoded_bytes) { 298 Ok(permit) => permit, 299 Err(_) => { 300 P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_full); 301 if session_id.starts_with(INBOUND_SESSION_PREFIX) { 302 let _ = session.shutdown.send(true); 303 disconnect.push((session_id.clone(), session.sender.clone())); 304 } 305 break; 306 } 307 }; 308 let batch = OutboundBatch { 309 envelopes: Arc::clone(envelopes), 310 _queued_bytes: permit, 311 }; 312 match session.sender.try_send(batch) { 313 Ok(()) => {} 314 Err(mpsc::error::TrySendError::Full(_)) => { 315 P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_full); 316 if session_id.starts_with(INBOUND_SESSION_PREFIX) { 317 let _ = session.shutdown.send(true); 318 disconnect.push((session_id.clone(), session.sender.clone())); 319 } 320 break; 321 } 322 Err(mpsc::error::TrySendError::Closed(_)) => { 323 P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_closed); 324 let _ = session.shutdown.send(true); 325 disconnect.push((session_id.clone(), session.sender.clone())); 326 break; 327 } 328 } 329 } 330 } 331 if !disconnect.is_empty() { 332 let mut sessions = self.inner.sessions.lock().await; 333 for (session_id, sender) in disconnect { 334 if sessions 335 .get(&session_id) 336 .is_some_and(|session| session.sender.same_channel(&sender)) 337 { 338 sessions.remove(&session_id); 339 } 340 } 341 } 342 Ok(()) 343 } 344 345 async fn prepare_gossip(&self, envelopes: Vec<GossipEnvelope>) -> Vec<GossipEnvelope> { 346 let mut blocks = Vec::new(); 347 let mut passthrough = Vec::new(); 348 349 for envelope in envelopes { 350 match envelope { 351 GossipEnvelope::Block(block) => blocks.push(BlockInventory { 352 height: block.height, 353 hash: block.hash, 354 }), 355 GossipEnvelope::Blocks { blocks: batch } => { 356 blocks.extend(batch.into_iter().map(|block| BlockInventory { 357 height: block.height, 358 hash: block.hash, 359 })); 360 } 361 GossipEnvelope::Inventory { blocks: inv_blocks } => blocks.extend(inv_blocks), 362 other => passthrough.push(other), 363 } 364 } 365 366 blocks.sort_by(|left, right| { 367 left.height 368 .cmp(&right.height) 369 .then_with(|| left.hash.cmp(&right.hash)) 370 }); 371 blocks.dedup_by(|left, right| left.hash == right.hash); 372 373 if !blocks.is_empty() { 374 passthrough.push(GossipEnvelope::Inventory { blocks }); 375 } 376 passthrough 377 } 378 379 pub async fn peer_exchange(&self) -> GossipEnvelope { 380 let advertised_addr = self.advertised_addr().await; 381 let self_filter_addr = self.self_filter_addr().await; 382 let self_addr = advertised_addr.map(|addr| addr.to_string()); 383 let peers = self 384 .inner 385 .peers 386 .lock() 387 .await 388 .addresses_except(self_addr.as_deref().unwrap_or("")) 389 .into_iter() 390 .filter(|peer| peer.parse::<SocketAddr>().is_ok()) 391 .filter(|peer| { 392 !is_self_peer_address_for(peer, self.inner.listen_addr, self_filter_addr) 393 }) 394 .collect::<Vec<_>>(); 395 GossipEnvelope::PeerList { 396 peers: self_addr.into_iter().chain(peers.into_iter()).collect(), 397 } 398 } 399 400 pub(super) async fn ensure_outbound_sessions(&self) { 401 self.inner.peers.lock().await.prune_stale_peers_at( 402 crate::app::now_ms(), 403 STALE_INBOUND_PEER_RETENTION_MS, 404 STALE_DISCOVERED_PEER_RETENTION_MS, 405 ); 406 let max_discovered_outbound_dials = { 407 let node = self.inner.node.lock().await; 408 if node.ledger().is_setup_placeholder() { 409 0 410 } else { 411 MAX_DISCOVERED_OUTBOUND_DIALS_PER_CYCLE 412 } 413 }; 414 let addresses = self 415 .inner 416 .peers 417 .lock() 418 .await 419 .outbound_session_candidates_at(crate::app::now_ms(), max_discovered_outbound_dials); 420 let address_set = addresses.iter().cloned().collect::<BTreeSet<_>>(); 421 let self_filter_addr = self.self_filter_addr().await; 422 let mut sessions = self.inner.sessions.lock().await; 423 sessions.retain(|peer, _| { 424 let keep = peer.starts_with(INBOUND_SESSION_PREFIX) 425 || (address_set.contains(peer) 426 && !is_self_peer_address_for(peer, self.inner.listen_addr, self_filter_addr)); 427 if !keep { 428 P2pMetricsCounters::inc(&self.inner.metrics.self_peer_skips); 429 } 430 keep 431 }); 432 for peer in addresses { 433 if is_self_peer_address_for(&peer, self.inner.listen_addr, self_filter_addr) { 434 P2pMetricsCounters::inc(&self.inner.metrics.self_peer_skips); 435 continue; 436 } 437 if sessions.contains_key(&peer) { 438 continue; 439 } 440 441 let (sender, receiver) = mpsc::channel(PEER_QUEUE_SIZE); 442 let (shutdown, shutdown_receiver) = watch::channel(false); 443 let queue_bytes = Arc::new(tokio::sync::Semaphore::new(PEER_QUEUE_BYTES)); 444 sessions.insert( 445 peer.clone(), 446 GossipSession { 447 peer: peer.clone(), 448 sender, 449 shutdown, 450 queue_bytes, 451 }, 452 ); 453 tokio::spawn(outbound_session( 454 self.clone(), 455 peer, 456 receiver, 457 shutdown_receiver, 458 )); 459 } 460 } 461 462 pub(super) async fn forward_outbox(&self) { 463 let outbox = self.inner.node.lock().await.drain_outbox(); 464 if let Err(error) = self.broadcast(outbox).await { 465 if debug_logging_enabled() { 466 eprintln!("p2p rebroadcast failed: {error:#}"); 467 } 468 } 469 } 470 } 471 472 fn byte_bounded_gossip_batches( 473 envelopes: Vec<GossipEnvelope>, 474 ) -> Result<Vec<(Arc<[GossipEnvelope]>, u32)>> { 475 let mut batches = Vec::new(); 476 let mut batch = Vec::new(); 477 let mut batch_bytes = 0_usize; 478 for envelope in envelopes { 479 let encoded_bytes = serde_json::to_vec(&envelope)?.len().saturating_add(1); 480 if encoded_bytes > MAX_GOSSIP_LINE_BYTES.saturating_add(1) { 481 anyhow::bail!( 482 "p2p message is {} bytes, exceeding {} byte limit", 483 encoded_bytes.saturating_sub(1), 484 MAX_GOSSIP_LINE_BYTES 485 ); 486 } 487 if !batch.is_empty() && batch_bytes.saturating_add(encoded_bytes) > MAX_OUTBOUND_BATCH_BYTES 488 { 489 batches.push((Arc::from(std::mem::take(&mut batch)), batch_bytes as u32)); 490 batch_bytes = 0; 491 } 492 batch_bytes = batch_bytes.saturating_add(encoded_bytes); 493 batch.push(envelope); 494 } 495 if !batch.is_empty() { 496 batches.push((Arc::from(batch), batch_bytes as u32)); 497 } 498 Ok(batches) 499 } 500 501 #[cfg(test)] 502 mod tests { 503 use std::{collections::BTreeMap, sync::Arc}; 504 505 use crate::{ 506 app::{GossipEnvelope, NodeCore, PeerBook}, 507 domain::{Ledger, Wallet}, 508 }; 509 510 use super::super::test_support::{allocations, gossip_network, node}; 511 512 #[test] 513 fn outbound_batches_are_bounded_by_encoded_bytes() { 514 let large_tip = "a".repeat(super::super::MAX_GOSSIP_LINE_BYTES / 2); 515 let envelopes = vec![ 516 GossipEnvelope::PeerStatus { 517 height: 1, 518 tip_hash: large_tip.clone(), 519 time_ms: 1, 520 }, 521 GossipEnvelope::PeerStatus { 522 height: 2, 523 tip_hash: large_tip, 524 time_ms: 2, 525 }, 526 ]; 527 528 let batches = super::byte_bounded_gossip_batches(envelopes).unwrap(); 529 530 assert_eq!(batches.len(), 2); 531 assert!(batches.iter().all(|(_, bytes)| { 532 usize::try_from(*bytes).unwrap() <= super::super::MAX_OUTBOUND_BATCH_BYTES 533 })); 534 } 535 536 #[tokio::test] 537 async fn sync_progress_is_incremental_and_scoped_to_the_active_validation() { 538 let alice = Wallet::from_seed("sync-progress-alice"); 539 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 540 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 541 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 542 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 543 544 let first = network.begin_sync_progress(0, 90); 545 network.update_sync_progress(first.id(), 23); 546 assert_eq!( 547 network.sync_progress(), 548 Some(super::super::SyncProgress { 549 start_height: 0, 550 validated_height: 23, 551 target_height: 90, 552 }) 553 ); 554 555 let second = network.begin_sync_progress(23, 76); 556 network.update_sync_progress(second.id(), 41); 557 assert_eq!(network.sync_progress().unwrap().target_height, 90); 558 drop(first); 559 assert_eq!(network.sync_progress().unwrap().validated_height, 41); 560 561 drop(second); 562 assert_eq!(network.sync_progress(), None); 563 } 564 565 #[tokio::test] 566 async fn invalidating_sync_progress_hides_and_rejects_stale_validation() { 567 let alice = Wallet::from_seed("invalidated-sync-progress-alice"); 568 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 569 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 570 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 571 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 572 573 let stale = network.begin_sync_progress(20, 90); 574 network.update_sync_progress(stale.id(), 41); 575 assert!(stale.is_current()); 576 577 network.invalidate_sync_progress(); 578 579 assert_eq!(network.sync_progress(), None); 580 assert!(!stale.is_current()); 581 network.update_sync_progress(stale.id(), 42); 582 assert_eq!(network.sync_progress(), None); 583 584 let fresh = network.begin_sync_progress(0, 90); 585 assert!(fresh.is_current()); 586 assert_eq!(network.sync_progress().unwrap().validated_height, 0); 587 588 drop(stale); 589 assert_eq!(network.sync_progress().unwrap().validated_height, 0); 590 } 591 592 #[tokio::test] 593 async fn identical_chain_validations_cannot_run_concurrently() { 594 let alice = Wallet::from_seed("validation-coordinator-alice"); 595 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 596 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 597 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 598 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 599 600 let first = network 601 .claim_chain_validation("same-page".to_string()) 602 .await 603 .unwrap(); 604 let waiting_network = network.clone(); 605 let waiting = tokio::spawn(async move { 606 waiting_network 607 .claim_chain_validation("same-page".to_string()) 608 .await 609 .unwrap() 610 }); 611 tokio::time::sleep(std::time::Duration::from_millis(25)).await; 612 assert!(!waiting.is_finished()); 613 614 drop(first); 615 let second = tokio::time::timeout(std::time::Duration::from_secs(1), waiting) 616 .await 617 .unwrap() 618 .unwrap(); 619 drop(second); 620 assert!( 621 network 622 .inner 623 .chain_validation 624 .active 625 .lock() 626 .unwrap() 627 .is_empty() 628 ); 629 } 630 631 #[tokio::test] 632 async fn peer_exchange_does_not_advertise_self_when_outbound_only() { 633 let alice = Wallet::from_seed("px-private-alice"); 634 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 635 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 636 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 637 "127.0.0.1:9545".to_string(), 638 ]))); 639 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 640 match network.peer_exchange().await { 641 GossipEnvelope::PeerList { peers } => { 642 assert!(!peers.contains(&"127.0.0.1:9544".to_string())); 643 assert!(peers.contains(&"127.0.0.1:9545".to_string())); 644 } 645 other => panic!("expected peer list, got {other:?}"), 646 } 647 } 648 649 #[tokio::test] 650 async fn setup_placeholder_only_dials_explicit_outbound_peers() { 651 let alice = Wallet::from_seed("setup-placeholder-outbound-alice"); 652 let node = Arc::new(tokio::sync::Mutex::new(NodeCore::from_ledger( 653 alice, 654 Ledger::new(BTreeMap::new(), 1), 655 0, 656 ))); 657 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 658 "127.0.0.1:9545".to_string(), 659 ]))); 660 peers 661 .lock() 662 .await 663 .add_discovered_peer("127.0.0.1:9546".to_string()); 664 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 665 666 network.ensure_outbound_sessions().await; 667 668 let sessions = network.inner.sessions.lock().await; 669 assert!(sessions.contains_key("127.0.0.1:9545")); 670 assert!(!sessions.contains_key("127.0.0.1:9546")); 671 } 672 673 #[tokio::test] 674 async fn peer_exchange_omits_hostname_bootstrap_peers() { 675 let alice = Wallet::from_seed("px-hostname-alice"); 676 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 677 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 678 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 679 "iuna.jhx.app:9444".to_string(), 680 "127.0.0.1:9545".to_string(), 681 ]))); 682 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 683 match network.peer_exchange().await { 684 GossipEnvelope::PeerList { peers } => { 685 assert!(!peers.contains(&"iuna.jhx.app:9444".to_string())); 686 assert!(peers.contains(&"127.0.0.1:9545".to_string())); 687 } 688 other => panic!("expected peer list, got {other:?}"), 689 } 690 } 691 692 #[tokio::test] 693 async fn peer_exchange_does_not_advertise_discovered_listening_peers() { 694 let alice = Wallet::from_seed("px-discovered-alice"); 695 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 696 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 697 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 698 peers 699 .lock() 700 .await 701 .add_discovered_peer("127.0.0.1:9546".to_string()); 702 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 703 704 match network.peer_exchange().await { 705 GossipEnvelope::PeerList { peers } => { 706 assert!(!peers.contains(&"127.0.0.1:9546".to_string())); 707 } 708 other => panic!("expected peer list, got {other:?}"), 709 } 710 } 711 712 #[tokio::test] 713 async fn peer_exchange_advertises_stable_listen_and_known_peers() { 714 let alice = Wallet::from_seed("px-alice"); 715 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 716 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 717 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 718 "127.0.0.1:9545".to_string(), 719 ]))); 720 let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); 721 network.set_accept_inbound(true).await.unwrap(); 722 723 match network.peer_exchange().await { 724 GossipEnvelope::PeerList { peers } => { 725 assert!(peers.contains(&"127.0.0.1:9544".to_string())); 726 assert!(peers.contains(&"127.0.0.1:9545".to_string())); 727 } 728 other => panic!("expected peer list, got {other:?}"), 729 } 730 network.set_accept_inbound(false).await.unwrap(); 731 } 732 733 #[tokio::test] 734 async fn peer_exchange_filters_announced_self_from_known_peers() { 735 let alice = Wallet::from_seed("px-announced-self-alice"); 736 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 737 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 738 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 739 "8.8.8.8:9444".to_string(), 740 "8.8.4.4:9444".to_string(), 741 ]))); 742 let network = gossip_network( 743 node, 744 peers, 745 "127.0.0.1:0".parse().unwrap(), 746 Some("8.8.8.8:9444".parse().unwrap()), 747 ); 748 network.set_accept_inbound(true).await.unwrap(); 749 750 match network.peer_exchange().await { 751 GossipEnvelope::PeerList { peers } => { 752 assert_eq!( 753 peers.iter().filter(|peer| *peer == "8.8.8.8:9444").count(), 754 1 755 ); 756 assert!(peers.contains(&"8.8.4.4:9444".to_string())); 757 } 758 other => panic!("expected peer list, got {other:?}"), 759 } 760 network.set_accept_inbound(false).await.unwrap(); 761 } 762 }