tests.rs (68133B)
1 use std::{ 2 collections::BTreeMap, 3 net::SocketAddr, 4 sync::{Arc, Mutex as StdMutex}, 5 }; 6 7 use crate::{ 8 app::{ 9 GossipEnvelope, MAINNET_CANDIDATE_NETWORK_ID, NETWORK_ID, NodeCore, PROTOCOL_VERSION, 10 PeerBook, PeerDirection, ProtocolHello, 11 }, 12 domain::{Ledger, Wallet, run_vdf}, 13 }; 14 use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; 15 16 use super::test_support::{allocations, gossip_network, node, queue_plaintext_burn}; 17 18 #[tokio::test] 19 async fn block_batch_validation_reports_each_validated_height() { 20 let alice = Wallet::from_seed("batch-validation-progress-alice"); 21 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 22 let local = node("batch-progress-local", alice.clone(), allocations.clone()); 23 let mut remote = node("batch-progress-remote", alice.clone(), allocations); 24 for timestamp_ms in [1, 2] { 25 queue_plaintext_burn(&mut remote, &alice, 1); 26 remote.drain_outbox(); 27 remote.mine_one_at(timestamp_ms).unwrap(); 28 remote.drain_outbox(); 29 } 30 let blocks = remote.ledger().blocks_from(1, 10); 31 let progress = Arc::new(StdMutex::new(Vec::new())); 32 let reported = Arc::clone(&progress); 33 34 let adopted = super::validate_blocks_extension( 35 local.clone_ledger(), 36 blocks, 37 crate::app::now_ms(), 38 move |height| reported.lock().unwrap().push(height), 39 ) 40 .await 41 .unwrap(); 42 43 assert_eq!(*progress.lock().unwrap(), vec![1, 2]); 44 assert_eq!(adopted.height(), 2); 45 } 46 47 #[tokio::test] 48 async fn session_serializes_inventory_and_paginated_catchup_requests() { 49 let alice = Wallet::from_seed("immediate-page-sync-alice"); 50 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 51 let local = Arc::new(tokio::sync::Mutex::new(node( 52 "immediate-page-local", 53 alice.clone(), 54 allocations.clone(), 55 ))); 56 let mut remote = node("immediate-page-remote", alice.clone(), allocations); 57 for timestamp_ms in [1, 2] { 58 queue_plaintext_burn(&mut remote, &alice, 1); 59 remote.drain_outbox(); 60 remote.mine_one_at(timestamp_ms).unwrap(); 61 remote.drain_outbox(); 62 } 63 let hello = remote.hello(None, None); 64 let blocks = remote.ledger().blocks_from(1, 2); 65 66 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 67 let peer_addr = listener.local_addr().unwrap(); 68 let server = tokio::spawn(async move { 69 let (stream, _) = listener.accept().await.unwrap(); 70 let (reader, mut writer) = stream.into_split(); 71 let mut lines = BufReader::new(reader).lines(); 72 let first = lines.next_line().await.unwrap().unwrap(); 73 assert!(matches!( 74 super::parse_envelope(&first).unwrap(), 75 GossipEnvelope::Hello(_) 76 )); 77 super::write_envelope(&mut writer, &hello).await.unwrap(); 78 79 let mut sent_first_page = false; 80 loop { 81 let line = tokio::time::timeout(std::time::Duration::from_secs(30), lines.next_line()) 82 .await 83 .expect("next block request did not follow the imported page") 84 .unwrap() 85 .unwrap(); 86 match super::parse_envelope(&line).unwrap() { 87 GossipEnvelope::BlockRangeRequest { from_height: 1, .. } => { 88 assert!(!sent_first_page); 89 sent_first_page = true; 90 super::write_envelope( 91 &mut writer, 92 &GossipEnvelope::Inventory { 93 blocks: vec![crate::app::BlockInventory { 94 height: blocks[1].height, 95 hash: blocks[1].hash.clone(), 96 }], 97 }, 98 ) 99 .await 100 .unwrap(); 101 let duplicate_request = 102 tokio::time::timeout(std::time::Duration::from_millis(200), async { 103 loop { 104 let line = lines.next_line().await.unwrap().unwrap(); 105 let envelope = super::parse_envelope(&line).unwrap(); 106 if matches!( 107 envelope, 108 GossipEnvelope::BlockRequest { .. } 109 | GossipEnvelope::BlockRangeRequest { .. } 110 | GossipEnvelope::BlockLocatorRequest { .. } 111 ) { 112 break envelope; 113 } 114 } 115 }) 116 .await; 117 assert!( 118 duplicate_request.is_err(), 119 "Inventory bypassed the outstanding paginated catchup request: {duplicate_request:?}" 120 ); 121 super::write_envelope( 122 &mut writer, 123 &GossipEnvelope::Blocks { 124 blocks: vec![blocks[0].clone()], 125 }, 126 ) 127 .await 128 .unwrap(); 129 } 130 GossipEnvelope::BlockRangeRequest { from_height: 2, .. } => { 131 assert!(sent_first_page); 132 super::write_envelope( 133 &mut writer, 134 &GossipEnvelope::Blocks { 135 blocks: vec![blocks[1].clone()], 136 }, 137 ) 138 .await 139 .unwrap(); 140 break; 141 } 142 _ => {} 143 } 144 } 145 }); 146 147 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 148 peer_addr.to_string(), 149 ]))); 150 let _network = super::GossipNetwork::start( 151 Arc::clone(&local), 152 peers, 153 "127.0.0.1:0".parse().unwrap(), 154 None, 155 false, 156 ) 157 .await 158 .unwrap(); 159 160 tokio::time::timeout(std::time::Duration::from_secs(60), server) 161 .await 162 .unwrap() 163 .unwrap(); 164 tokio::time::timeout(std::time::Duration::from_secs(30), async { 165 loop { 166 if local.lock().await.ledger().height() == 2 { 167 break; 168 } 169 tokio::task::yield_now().await; 170 } 171 }) 172 .await 173 .unwrap(); 174 } 175 176 #[tokio::test] 177 async fn invalid_bootstrap_response_is_not_retried_without_backoff() { 178 let wallet = Wallet::from_seed("invalid-bootstrap-backoff"); 179 let remote = node( 180 "invalid-bootstrap-remote", 181 wallet.clone(), 182 allocations(std::slice::from_ref(&wallet), 1_000), 183 ); 184 let hello = remote.hello(None, None); 185 let mut invalid_bootstrap = remote.chain_bootstrap(); 186 invalid_bootstrap.genesis_block.hash = "0".repeat(64); 187 188 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 189 let peer_addr = listener.local_addr().unwrap(); 190 let server = tokio::spawn(async move { 191 let (stream, _) = listener.accept().await.unwrap(); 192 let (reader, mut writer) = stream.into_split(); 193 let mut lines = BufReader::new(reader).lines(); 194 let first = lines.next_line().await.unwrap().unwrap(); 195 assert!(matches!( 196 super::parse_envelope(&first).unwrap(), 197 GossipEnvelope::Hello(_) 198 )); 199 super::write_envelope(&mut writer, &hello).await.unwrap(); 200 201 loop { 202 let line = lines.next_line().await.unwrap().unwrap(); 203 if matches!( 204 super::parse_envelope(&line).unwrap(), 205 GossipEnvelope::ChainBootstrapRequest 206 ) { 207 break; 208 } 209 } 210 super::write_envelope( 211 &mut writer, 212 &GossipEnvelope::ChainBootstrap(invalid_bootstrap), 213 ) 214 .await 215 .unwrap(); 216 217 tokio::time::timeout(std::time::Duration::from_secs(3), async { 218 loop { 219 let Some(line) = lines.next_line().await.unwrap() else { 220 break; 221 }; 222 assert!( 223 !matches!( 224 super::parse_envelope(&line).unwrap(), 225 GossipEnvelope::ChainBootstrapRequest 226 ), 227 "invalid bootstrap was retried immediately" 228 ); 229 } 230 }) 231 .await 232 .ok(); 233 }); 234 235 let local = Arc::new(tokio::sync::Mutex::new(NodeCore::from_ledger( 236 wallet, 237 Ledger::new(BTreeMap::new(), 1), 238 0, 239 ))); 240 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 241 peer_addr.to_string(), 242 ]))); 243 let _network = 244 super::GossipNetwork::start(local, peers, "127.0.0.1:0".parse().unwrap(), None, false) 245 .await 246 .unwrap(); 247 248 tokio::time::timeout(std::time::Duration::from_secs(5), server) 249 .await 250 .unwrap() 251 .unwrap(); 252 } 253 254 #[tokio::test] 255 async fn inbound_session_receives_new_block_inventory_without_waiting_for_status_tick() { 256 let alice = Wallet::from_seed("inbound-fast-block-relay-alice"); 257 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 258 let mut source = node( 259 "inbound-fast-block-relay-source", 260 alice.clone(), 261 allocations.clone(), 262 ); 263 let remote = node( 264 "inbound-fast-block-relay-remote", 265 alice.clone(), 266 allocations, 267 ); 268 queue_plaintext_burn(&mut source, &alice, 1); 269 source.drain_outbox(); 270 let block = source.mine_one_at(1).unwrap(); 271 source.drain_outbox(); 272 273 let reserved = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 274 let listen_addr = reserved.local_addr().unwrap(); 275 drop(reserved); 276 let network = super::GossipNetwork::start( 277 Arc::new(tokio::sync::Mutex::new(source)), 278 Arc::new(tokio::sync::Mutex::new(PeerBook::default())), 279 listen_addr, 280 None, 281 true, 282 ) 283 .await 284 .unwrap(); 285 286 let stream = tokio::net::TcpStream::connect(listen_addr).await.unwrap(); 287 let (reader, mut writer) = stream.into_split(); 288 let mut lines = BufReader::new(reader).lines(); 289 let hello = lines.next_line().await.unwrap().unwrap(); 290 assert!(matches!( 291 super::parse_envelope(&hello).unwrap(), 292 GossipEnvelope::Hello(_) 293 )); 294 295 network 296 .broadcast(vec![GossipEnvelope::Block(block.clone())]) 297 .await 298 .unwrap(); 299 assert!( 300 tokio::time::timeout(std::time::Duration::from_millis(150), lines.next_line()) 301 .await 302 .is_err(), 303 "inbound peer received gossip before completing Hello" 304 ); 305 306 super::write_envelope(&mut writer, &remote.hello(None, None)) 307 .await 308 .unwrap(); 309 tokio::time::timeout(std::time::Duration::from_secs(1), async { 310 loop { 311 if !network.inner.sessions.lock().await.is_empty() { 312 break; 313 } 314 tokio::task::yield_now().await; 315 } 316 }) 317 .await 318 .expect("inbound session was not registered after Hello"); 319 320 network 321 .broadcast(vec![GossipEnvelope::Block(block.clone())]) 322 .await 323 .unwrap(); 324 let relayed = tokio::time::timeout(std::time::Duration::from_secs(1), async { 325 loop { 326 let line = lines.next_line().await.unwrap().unwrap(); 327 if let GossipEnvelope::Inventory { blocks } = super::parse_envelope(&line).unwrap() { 328 break blocks; 329 } 330 } 331 }) 332 .await 333 .expect("inbound block relay waited for periodic anti-entropy"); 334 335 assert!(relayed.iter().any(|item| item.hash == block.hash)); 336 network.set_accept_inbound(false).await.unwrap(); 337 } 338 339 #[tokio::test] 340 async fn inbound_self_connection_is_closed_before_relay_registration() { 341 let wallet = Wallet::from_seed("inbound-self-connection"); 342 let reserved = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 343 let listen_addr = reserved.local_addr().unwrap(); 344 drop(reserved); 345 let network = super::GossipNetwork::start( 346 Arc::new(tokio::sync::Mutex::new(node( 347 "inbound-self-connection", 348 wallet.clone(), 349 allocations(&[wallet], 1_000), 350 ))), 351 Arc::new(tokio::sync::Mutex::new(PeerBook::default())), 352 listen_addr, 353 None, 354 true, 355 ) 356 .await 357 .unwrap(); 358 359 let stream = tokio::net::TcpStream::connect(listen_addr).await.unwrap(); 360 let (reader, mut writer) = stream.into_split(); 361 let mut lines = BufReader::new(reader).lines(); 362 let local_hello = 363 match super::parse_envelope(&lines.next_line().await.unwrap().unwrap()).unwrap() { 364 GossipEnvelope::Hello(hello) => hello, 365 other => panic!("expected Hello, got {other:?}"), 366 }; 367 super::write_envelope(&mut writer, &GossipEnvelope::Hello(local_hello)) 368 .await 369 .unwrap(); 370 371 let closed = tokio::time::timeout(std::time::Duration::from_secs(1), lines.next_line()) 372 .await 373 .expect("self connection remained open") 374 .unwrap(); 375 assert!(closed.is_none()); 376 assert!(network.inner.sessions.lock().await.is_empty()); 377 assert_eq!(network.metrics().self_peer_rejections, 1); 378 network.set_accept_inbound(false).await.unwrap(); 379 } 380 381 #[tokio::test] 382 async fn full_outbound_queue_is_metric_not_peer_error() { 383 let wallet = Wallet::from_seed("full-outbound-queue"); 384 let node = Arc::new(tokio::sync::Mutex::new(node( 385 "full-outbound-queue", 386 wallet.clone(), 387 allocations(&[wallet], 1_000), 388 ))); 389 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 390 "127.0.0.1:9444".to_string(), 391 ]))); 392 let network = super::GossipNetwork { 393 inner: Arc::new(super::GossipNetworkInner { 394 node, 395 peers: Arc::clone(&peers), 396 listen_addr: "127.0.0.1:9544".parse().unwrap(), 397 p2p_announce_addr: tokio::sync::Mutex::new(None), 398 node_id: super::new_node_id(), 399 accept_task: tokio::sync::Mutex::new(None), 400 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 401 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 402 metrics: super::P2pMetricsCounters::default(), 403 sync_progress: StdMutex::new(super::SyncProgressState::default()), 404 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 405 }), 406 }; 407 let (sender, _receiver) = tokio::sync::mpsc::channel(1); 408 let queue_bytes = Arc::new(tokio::sync::Semaphore::new(super::PEER_QUEUE_BYTES)); 409 let queued_bytes = Arc::clone(&queue_bytes).try_acquire_many_owned(1).unwrap(); 410 sender 411 .try_send(super::OutboundBatch { 412 envelopes: Arc::from(vec![GossipEnvelope::PeerStatus { 413 height: 1, 414 tip_hash: "queued".to_string(), 415 time_ms: 1_000, 416 }]), 417 _queued_bytes: queued_bytes, 418 }) 419 .unwrap(); 420 network.inner.sessions.lock().await.insert( 421 "127.0.0.1:9444".to_string(), 422 super::GossipSession { 423 peer: "127.0.0.1:9444".to_string(), 424 sender, 425 shutdown: tokio::sync::watch::channel(false).0, 426 queue_bytes, 427 }, 428 ); 429 430 network 431 .broadcast(vec![GossipEnvelope::PeerStatus { 432 height: 2, 433 tip_hash: "new".to_string(), 434 time_ms: 2_000, 435 }]) 436 .await 437 .unwrap(); 438 439 assert_eq!(network.metrics().outbound_queue_full, 1); 440 let peer = peers.lock().await.list().pop().unwrap(); 441 assert_eq!(peer.last_error, None); 442 assert_eq!(peer.last_error_ms, None); 443 } 444 445 #[tokio::test] 446 async fn full_inbound_queue_disconnects_the_session() { 447 let wallet = Wallet::from_seed("full-inbound-queue"); 448 let network = gossip_network( 449 Arc::new(tokio::sync::Mutex::new(node( 450 "full-inbound-queue", 451 wallet.clone(), 452 allocations(&[wallet], 1_000), 453 ))), 454 Arc::new(tokio::sync::Mutex::new(PeerBook::default())), 455 "127.0.0.1:9544".parse().unwrap(), 456 None, 457 ); 458 let session_id = format!("{}127.0.0.1:51234", super::INBOUND_SESSION_PREFIX); 459 let (sender, _receiver) = tokio::sync::mpsc::channel(1); 460 let (shutdown, mut shutdown_receiver) = tokio::sync::watch::channel(false); 461 let queue_bytes = Arc::new(tokio::sync::Semaphore::new(super::INBOUND_PEER_QUEUE_BYTES)); 462 let queued_bytes = Arc::clone(&queue_bytes).try_acquire_many_owned(1).unwrap(); 463 sender 464 .try_send(super::OutboundBatch { 465 envelopes: Arc::from(vec![GossipEnvelope::PeerStatus { 466 height: 1, 467 tip_hash: "queued".to_string(), 468 time_ms: 1_000, 469 }]), 470 _queued_bytes: queued_bytes, 471 }) 472 .unwrap(); 473 network.inner.sessions.lock().await.insert( 474 session_id.clone(), 475 super::GossipSession { 476 peer: "127.0.0.1:51234".to_string(), 477 sender, 478 shutdown, 479 queue_bytes, 480 }, 481 ); 482 483 network 484 .broadcast(vec![GossipEnvelope::PeerStatus { 485 height: 2, 486 tip_hash: "new".to_string(), 487 time_ms: 2_000, 488 }]) 489 .await 490 .unwrap(); 491 492 assert_eq!(network.metrics().outbound_queue_full, 1); 493 assert!( 494 !network 495 .inner 496 .sessions 497 .lock() 498 .await 499 .contains_key(&session_id) 500 ); 501 shutdown_receiver.changed().await.unwrap(); 502 assert!(*shutdown_receiver.borrow()); 503 } 504 505 #[tokio::test] 506 async fn exhausted_inbound_byte_budget_disconnects_the_session() { 507 let wallet = Wallet::from_seed("inbound-byte-budget"); 508 let network = gossip_network( 509 Arc::new(tokio::sync::Mutex::new(node( 510 "inbound-byte-budget", 511 wallet.clone(), 512 allocations(&[wallet], 1_000), 513 ))), 514 Arc::new(tokio::sync::Mutex::new(PeerBook::default())), 515 "127.0.0.1:9544".parse().unwrap(), 516 None, 517 ); 518 let session_id = format!("{}127.0.0.1:51235", super::INBOUND_SESSION_PREFIX); 519 let (sender, mut receiver) = tokio::sync::mpsc::channel(4); 520 let (shutdown, mut shutdown_receiver) = tokio::sync::watch::channel(false); 521 network.inner.sessions.lock().await.insert( 522 session_id.clone(), 523 super::GossipSession { 524 peer: "127.0.0.1:51235".to_string(), 525 sender, 526 shutdown, 527 queue_bytes: Arc::new(tokio::sync::Semaphore::new(1)), 528 }, 529 ); 530 531 network 532 .broadcast(vec![GossipEnvelope::PeerStatus { 533 height: 2, 534 tip_hash: "larger-than-one-byte".to_string(), 535 time_ms: 2_000, 536 }]) 537 .await 538 .unwrap(); 539 540 assert_eq!(network.metrics().outbound_queue_full, 1); 541 assert!( 542 !network 543 .inner 544 .sessions 545 .lock() 546 .await 547 .contains_key(&session_id) 548 ); 549 shutdown_receiver.changed().await.unwrap(); 550 assert!(*shutdown_receiver.borrow()); 551 assert!(matches!( 552 receiver.try_recv(), 553 Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) 554 )); 555 } 556 557 #[tokio::test] 558 async fn broadcast_skips_banned_inbound_session_identity() { 559 let wallet = Wallet::from_seed("banned-inbound-relay"); 560 let node = Arc::new(tokio::sync::Mutex::new(node( 561 "banned-inbound-relay", 562 wallet.clone(), 563 allocations(&[wallet], 1_000), 564 ))); 565 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 566 let network = gossip_network( 567 node, 568 Arc::clone(&peers), 569 "127.0.0.1:9544".parse().unwrap(), 570 None, 571 ); 572 let peer = "127.0.0.1:51234"; 573 for _ in 0..crate::app::PEER_MISBEHAVIOR_BAN_SCORE { 574 peers 575 .lock() 576 .await 577 .record_inbound_misbehavior(peer, "invalid block"); 578 } 579 let (sender, mut receiver) = tokio::sync::mpsc::channel(1); 580 let queue_bytes = Arc::new(tokio::sync::Semaphore::new(super::INBOUND_PEER_QUEUE_BYTES)); 581 network.inner.sessions.lock().await.insert( 582 format!("{}{peer}", super::INBOUND_SESSION_PREFIX), 583 super::GossipSession { 584 peer: peer.to_string(), 585 sender, 586 shutdown: tokio::sync::watch::channel(false).0, 587 queue_bytes, 588 }, 589 ); 590 591 network 592 .broadcast(vec![GossipEnvelope::PeerStatus { 593 height: 1, 594 tip_hash: "tip".to_string(), 595 time_ms: 1_000, 596 }]) 597 .await 598 .unwrap(); 599 600 assert!(matches!( 601 receiver.try_recv(), 602 Err(tokio::sync::mpsc::error::TryRecvError::Empty) 603 )); 604 } 605 606 #[tokio::test] 607 async fn single_block_fork_error_requests_blocks_by_locator() { 608 let alice = Wallet::from_seed("single-block-fork-alice"); 609 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 610 let mut local_node = node( 611 "local-single-block-fork", 612 alice.clone(), 613 allocations.clone(), 614 ); 615 let mut remote_node = node("remote-single-block-fork", alice.clone(), allocations); 616 617 queue_plaintext_burn(&mut local_node, &alice, 1); 618 local_node.drain_outbox(); 619 local_node.mine_one_at(1).unwrap(); 620 local_node.drain_outbox(); 621 622 queue_plaintext_burn(&mut remote_node, &alice, 1); 623 remote_node.drain_outbox(); 624 remote_node.mine_one_at(2).unwrap(); 625 remote_node.drain_outbox(); 626 queue_plaintext_burn(&mut remote_node, &alice, 1); 627 remote_node.drain_outbox(); 628 let remote_block = remote_node.mine_one_at(3).unwrap(); 629 assert_eq!(remote_block.height, 2); 630 assert_ne!( 631 remote_block.prev_hash, 632 local_node.ledger().tip_hash().to_string() 633 ); 634 635 let network = gossip_network( 636 Arc::new(tokio::sync::Mutex::new(local_node)), 637 Arc::new(tokio::sync::Mutex::new(PeerBook::default())), 638 "127.0.0.1:9544".parse().unwrap(), 639 None, 640 ); 641 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 642 let client = tokio::net::TcpStream::connect(listener.local_addr().unwrap()) 643 .await 644 .unwrap(); 645 let (server, remote_addr) = listener.accept().await.unwrap(); 646 let (_server_reader, mut server_writer) = server.into_split(); 647 let (client_reader, _client_writer) = client.into_split(); 648 let mut client_reader = super::LimitedLineReader::new(client_reader); 649 let mut known_peer = None; 650 651 let requested_chain_data = super::process_envelope( 652 &network, 653 &mut server_writer, 654 remote_addr, 655 &mut known_peer, 656 GossipEnvelope::Block(remote_block), 657 ) 658 .await 659 .unwrap(); 660 assert!(requested_chain_data); 661 662 let line = tokio::time::timeout(std::time::Duration::from_secs(1), client_reader.read_line()) 663 .await 664 .unwrap() 665 .unwrap() 666 .unwrap(); 667 let GossipEnvelope::BlockLocatorRequest { locator, limit } = 668 super::parse_envelope(&line).unwrap() 669 else { 670 panic!("expected block locator request"); 671 }; 672 let fork_blocks = remote_node.blocks_after_locator(&locator, limit); 673 assert_eq!(fork_blocks.len(), 2); 674 let local_ledger = network.inner.node.lock().await.clone_ledger(); 675 let adopted = 676 super::validate_blocks_extension(local_ledger, fork_blocks, crate::app::now_ms(), |_| {}) 677 .await 678 .unwrap(); 679 assert_eq!(adopted.tip_hash(), remote_node.ledger().tip_hash()); 680 } 681 682 #[tokio::test] 683 async fn block_page_without_local_ancestor_requests_blocks_by_locator() { 684 let alice = Wallet::from_seed("block-page-fork-alice"); 685 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 686 let mut local_node = node("local-block-page-fork", alice.clone(), allocations.clone()); 687 let mut remote_node = node("remote-block-page-fork", alice.clone(), allocations); 688 689 queue_plaintext_burn(&mut local_node, &alice, 1); 690 local_node.drain_outbox(); 691 local_node.mine_one_at(1).unwrap(); 692 local_node.drain_outbox(); 693 694 for timestamp_ms in [2, 3] { 695 queue_plaintext_burn(&mut remote_node, &alice, 1); 696 remote_node.drain_outbox(); 697 remote_node.mine_one_at(timestamp_ms).unwrap(); 698 remote_node.drain_outbox(); 699 } 700 let remote_page = remote_node.blocks_from(2, 10); 701 assert_eq!(remote_page.len(), 1); 702 assert_ne!( 703 remote_page[0].prev_hash, 704 local_node.ledger().tip_hash().to_string() 705 ); 706 707 let network = gossip_network( 708 Arc::new(tokio::sync::Mutex::new(local_node)), 709 Arc::new(tokio::sync::Mutex::new(PeerBook::default())), 710 "127.0.0.1:9545".parse().unwrap(), 711 None, 712 ); 713 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 714 let client = tokio::net::TcpStream::connect(listener.local_addr().unwrap()) 715 .await 716 .unwrap(); 717 let (server, remote_addr) = listener.accept().await.unwrap(); 718 let (_server_reader, mut server_writer) = server.into_split(); 719 let (client_reader, _client_writer) = client.into_split(); 720 let mut client_reader = super::LimitedLineReader::new(client_reader); 721 let mut known_peer = None; 722 723 let requested_chain_data = super::process_envelope( 724 &network, 725 &mut server_writer, 726 remote_addr, 727 &mut known_peer, 728 GossipEnvelope::Blocks { 729 blocks: remote_page, 730 }, 731 ) 732 .await 733 .unwrap(); 734 assert!(requested_chain_data); 735 736 let line = tokio::time::timeout(std::time::Duration::from_secs(1), client_reader.read_line()) 737 .await 738 .unwrap() 739 .unwrap() 740 .unwrap(); 741 let GossipEnvelope::BlockLocatorRequest { locator, limit } = 742 super::parse_envelope(&line).unwrap() 743 else { 744 panic!("expected block locator request"); 745 }; 746 let fork_blocks = remote_node.blocks_after_locator(&locator, limit); 747 assert_eq!(fork_blocks.len(), 2); 748 let local_ledger = network.inner.node.lock().await.clone_ledger(); 749 let adopted = 750 super::validate_blocks_extension(local_ledger, fork_blocks, crate::app::now_ms(), |_| {}) 751 .await 752 .unwrap(); 753 assert_eq!(adopted.tip_hash(), remote_node.ledger().tip_hash()); 754 } 755 756 #[tokio::test] 757 async fn future_block_rejection_does_not_poison_peer_or_later_acceptance() { 758 let alice = Wallet::from_seed("future-block-p2p-alice"); 759 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 760 let mut producer = node("future-block-producer", alice.clone(), allocations.clone()); 761 queue_plaintext_burn(&mut producer, &alice, 1); 762 producer.drain_outbox(); 763 let future_timestamp = crate::app::now_ms().saturating_add(10 * 60 * 1_000); 764 let prepared = producer 765 .ledger() 766 .prepare_next_block(alice.address(), future_timestamp) 767 .unwrap(); 768 let vdf_output = run_vdf(prepared.vdf_seed(), prepared.vdf_rounds()); 769 let future_block = prepared.finish(&alice, vdf_output); 770 771 let local_node = node("future-block-local", alice.clone(), allocations); 772 let original_tip = local_node.ledger().tip_hash().to_string(); 773 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 774 let network = gossip_network( 775 Arc::new(tokio::sync::Mutex::new(local_node)), 776 Arc::clone(&peers), 777 "127.0.0.1:9544".parse().unwrap(), 778 None, 779 ); 780 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 781 let client = tokio::net::TcpStream::connect(listener.local_addr().unwrap()) 782 .await 783 .unwrap(); 784 let (server, remote_addr) = listener.accept().await.unwrap(); 785 let (_server_reader, mut server_writer) = server.into_split(); 786 let (_client_reader, _client_writer) = client.into_split(); 787 let peer = remote_addr.to_string(); 788 let mut known_peer = Some(peer.clone()); 789 790 super::process_envelope( 791 &network, 792 &mut server_writer, 793 remote_addr, 794 &mut known_peer, 795 GossipEnvelope::Block(future_block.clone()), 796 ) 797 .await 798 .unwrap(); 799 800 assert_eq!( 801 network.inner.node.lock().await.ledger().tip_hash(), 802 original_tip 803 ); 804 assert_eq!(network.metrics().rejected_blocks, 1); 805 let recorded_peer = peers 806 .lock() 807 .await 808 .list() 809 .into_iter() 810 .find(|entry| entry.address == peer) 811 .expect("future block sender should be recorded"); 812 assert_eq!(recorded_peer.misbehavior_score, 0); 813 assert!(recorded_peer.banned_until_ms.is_none()); 814 assert!( 815 recorded_peer 816 .last_error 817 .as_deref() 818 .is_some_and(|error| error.contains("too far in the future")) 819 ); 820 821 let verified_block = super::verify_block_vdf(future_block.clone()).await.unwrap(); 822 network 823 .inner 824 .node 825 .lock() 826 .await 827 .receive_preverified_block_at(verified_block, future_block.timestamp_ms) 828 .unwrap(); 829 assert_eq!( 830 network.inner.node.lock().await.ledger().tip_hash(), 831 future_block.hash 832 ); 833 } 834 835 #[tokio::test] 836 async fn invalid_block_batch_is_rejected_atomically_without_partial_import() { 837 let alice = Wallet::from_seed("invalid-batch-p2p-alice"); 838 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 839 let mut remote_node = node("invalid-batch-remote", alice.clone(), allocations.clone()); 840 queue_plaintext_burn(&mut remote_node, &alice, 1); 841 remote_node.drain_outbox(); 842 let first_block = remote_node.mine_one_at(1).unwrap(); 843 remote_node.drain_outbox(); 844 queue_plaintext_burn(&mut remote_node, &alice, 1); 845 remote_node.drain_outbox(); 846 let second_block = remote_node.mine_one_at(2).unwrap(); 847 848 let mut invalid_second_block = second_block.clone(); 849 invalid_second_block.reward = invalid_second_block.reward.saturating_add(1); 850 invalid_second_block.hash = invalid_second_block.compute_hash(); 851 852 let local_node = node("invalid-batch-local", alice, allocations); 853 let original_tip = local_node.ledger().tip_hash().to_string(); 854 let original_height = local_node.ledger().height(); 855 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 856 let network = gossip_network( 857 Arc::new(tokio::sync::Mutex::new(local_node)), 858 Arc::clone(&peers), 859 "127.0.0.1:9544".parse().unwrap(), 860 None, 861 ); 862 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 863 let client = tokio::net::TcpStream::connect(listener.local_addr().unwrap()) 864 .await 865 .unwrap(); 866 let (server, remote_addr) = listener.accept().await.unwrap(); 867 let (_server_reader, mut server_writer) = server.into_split(); 868 let (_client_reader, _client_writer) = client.into_split(); 869 let peer = remote_addr.to_string(); 870 let mut known_peer = Some(peer.clone()); 871 872 super::process_envelope( 873 &network, 874 &mut server_writer, 875 remote_addr, 876 &mut known_peer, 877 GossipEnvelope::Blocks { 878 blocks: vec![first_block.clone(), invalid_second_block], 879 }, 880 ) 881 .await 882 .unwrap(); 883 884 let node = network.inner.node.lock().await; 885 assert_eq!(node.ledger().height(), original_height); 886 assert_eq!(node.ledger().tip_hash(), original_tip); 887 assert!(!node.ledger().has_block(&first_block.hash)); 888 drop(node); 889 890 assert_eq!(network.metrics().rejected_block_batches, 1); 891 assert!( 892 network 893 .metrics() 894 .last_chain_payload_error 895 .as_deref() 896 .is_some_and(|error| error.contains("block batch: block reward is invalid")) 897 ); 898 let recorded_peer = peers 899 .lock() 900 .await 901 .list() 902 .into_iter() 903 .find(|entry| entry.address == peer) 904 .expect("invalid batch sender should be recorded"); 905 assert_eq!(recorded_peer.misbehavior_score, 1); 906 assert!( 907 recorded_peer 908 .last_error 909 .as_deref() 910 .is_some_and(|error| error.contains("block reward is invalid")) 911 ); 912 } 913 914 #[tokio::test] 915 async fn chain_bootstrap_request_only_writes_bootstrap_without_mutating_local_state() { 916 let alice = Wallet::from_seed("snapshot-request-spam-alice"); 917 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 918 let mut local_node = node("snapshot-request-spam", alice.clone(), allocations); 919 queue_plaintext_burn(&mut local_node, &alice, 1); 920 local_node.drain_outbox(); 921 local_node.mine_one_at(1).unwrap(); 922 local_node.drain_outbox(); 923 let before = local_node.chain_snapshot(); 924 925 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 926 let network = gossip_network( 927 Arc::new(tokio::sync::Mutex::new(local_node)), 928 Arc::clone(&peers), 929 "127.0.0.1:9544".parse().unwrap(), 930 None, 931 ); 932 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 933 let client = tokio::net::TcpStream::connect(listener.local_addr().unwrap()) 934 .await 935 .unwrap(); 936 let (server, remote_addr) = listener.accept().await.unwrap(); 937 let (_server_reader, mut server_writer) = server.into_split(); 938 let (client_reader, _client_writer) = client.into_split(); 939 let mut client_reader = super::LimitedLineReader::new(client_reader); 940 let mut known_peer = None; 941 942 super::process_envelope( 943 &network, 944 &mut server_writer, 945 remote_addr, 946 &mut known_peer, 947 GossipEnvelope::ChainBootstrapRequest, 948 ) 949 .await 950 .unwrap(); 951 952 let line = tokio::time::timeout(std::time::Duration::from_secs(1), client_reader.read_line()) 953 .await 954 .unwrap() 955 .unwrap() 956 .unwrap(); 957 let GossipEnvelope::ChainBootstrap(bootstrap) = super::parse_envelope(&line).unwrap() else { 958 panic!("expected chain bootstrap"); 959 }; 960 assert_eq!(bootstrap.genesis_block, before.blocks[0]); 961 assert_eq!(bootstrap.height, before.blocks.last().unwrap().height); 962 assert_eq!(bootstrap.tip_hash, before.blocks.last().unwrap().hash); 963 assert_eq!(network.inner.node.lock().await.chain_snapshot(), before); 964 assert!(peers.lock().await.list().is_empty()); 965 assert!(known_peer.is_none()); 966 } 967 968 #[tokio::test] 969 async fn hello_rejects_wrong_network_or_genesis_without_banning() { 970 let alice = Wallet::from_seed("hello-alice"); 971 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 972 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 973 let network = super::GossipNetwork { 974 inner: Arc::new(super::GossipNetworkInner { 975 node, 976 peers: Arc::new(tokio::sync::Mutex::new(PeerBook::default())), 977 listen_addr: "127.0.0.1:9544".parse().unwrap(), 978 p2p_announce_addr: tokio::sync::Mutex::new(None), 979 node_id: super::new_node_id(), 980 accept_task: tokio::sync::Mutex::new(None), 981 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 982 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 983 metrics: super::P2pMetricsCounters::default(), 984 sync_progress: StdMutex::new(super::SyncProgressState::default()), 985 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 986 }), 987 }; 988 989 let wrong_network = ProtocolHello { 990 protocol_version: PROTOCOL_VERSION, 991 capabilities: Vec::new(), 992 network_id: "other-network".to_string(), 993 genesis_hash: network 994 .inner 995 .node 996 .lock() 997 .await 998 .ledger() 999 .genesis_hash() 1000 .to_string(), 1001 listen_addr: None, 1002 node_id: None, 1003 height: 0, 1004 tip_hash: "tip".to_string(), 1005 time_ms: 1_000, 1006 }; 1007 assert!( 1008 super::process_hello( 1009 &network, 1010 "127.0.0.1:9545".parse().unwrap(), 1011 &mut None, 1012 wrong_network, 1013 ) 1014 .await 1015 .unwrap_err() 1016 .to_string() 1017 .contains("wrong network") 1018 ); 1019 1020 let wrong_genesis = ProtocolHello { 1021 protocol_version: PROTOCOL_VERSION, 1022 capabilities: Vec::new(), 1023 network_id: NETWORK_ID.to_string(), 1024 genesis_hash: "not-local-genesis".to_string(), 1025 listen_addr: Some("127.0.0.1:9545".to_string()), 1026 node_id: None, 1027 height: 0, 1028 tip_hash: "tip".to_string(), 1029 time_ms: 1_000, 1030 }; 1031 assert!( 1032 super::process_hello( 1033 &network, 1034 "127.0.0.1:9545".parse().unwrap(), 1035 &mut None, 1036 wrong_genesis, 1037 ) 1038 .await 1039 .unwrap_err() 1040 .to_string() 1041 .contains("wrong genesis") 1042 ); 1043 1044 let wrong_protocol = ProtocolHello { 1045 protocol_version: PROTOCOL_VERSION + 1, 1046 capabilities: Vec::new(), 1047 network_id: NETWORK_ID.to_string(), 1048 genesis_hash: network 1049 .inner 1050 .node 1051 .lock() 1052 .await 1053 .ledger() 1054 .genesis_hash() 1055 .to_string(), 1056 listen_addr: Some("127.0.0.1:9545".to_string()), 1057 node_id: None, 1058 height: 0, 1059 tip_hash: "tip".to_string(), 1060 time_ms: 1_000, 1061 }; 1062 assert!( 1063 super::process_hello( 1064 &network, 1065 "127.0.0.1:9545".parse().unwrap(), 1066 &mut None, 1067 wrong_protocol, 1068 ) 1069 .await 1070 .unwrap_err() 1071 .to_string() 1072 .contains("unsupported protocol version") 1073 ); 1074 1075 assert!(network.inner.peers.lock().await.list().is_empty()); 1076 } 1077 1078 #[tokio::test] 1079 async fn hello_records_remote_clock_observation() { 1080 let alice = Wallet::from_seed("hello-clock-alice"); 1081 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 1082 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 1083 let network = super::GossipNetwork { 1084 inner: Arc::new(super::GossipNetworkInner { 1085 node, 1086 peers: Arc::new(tokio::sync::Mutex::new(PeerBook::default())), 1087 listen_addr: "127.0.0.1:9544".parse().unwrap(), 1088 p2p_announce_addr: tokio::sync::Mutex::new(None), 1089 node_id: super::new_node_id(), 1090 accept_task: tokio::sync::Mutex::new(None), 1091 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 1092 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 1093 metrics: super::P2pMetricsCounters::default(), 1094 sync_progress: StdMutex::new(super::SyncProgressState::default()), 1095 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 1096 }), 1097 }; 1098 let remote_time_ms = crate::app::now_ms().saturating_add(60_000); 1099 let hello = ProtocolHello { 1100 protocol_version: PROTOCOL_VERSION, 1101 capabilities: Vec::new(), 1102 network_id: NETWORK_ID.to_string(), 1103 genesis_hash: network 1104 .inner 1105 .node 1106 .lock() 1107 .await 1108 .ledger() 1109 .genesis_hash() 1110 .to_string(), 1111 listen_addr: Some("127.0.0.1:9545".to_string()), 1112 node_id: None, 1113 height: 0, 1114 tip_hash: "tip".to_string(), 1115 time_ms: remote_time_ms, 1116 }; 1117 let expected_hello = hello.clone(); 1118 1119 let mut known_peer = None; 1120 super::process_hello( 1121 &network, 1122 "127.0.0.1:9545".parse().unwrap(), 1123 &mut known_peer, 1124 hello, 1125 ) 1126 .await 1127 .unwrap(); 1128 1129 let peers = network.inner.peers.lock().await.list(); 1130 let peer = peers 1131 .iter() 1132 .find(|peer| peer.address == "127.0.0.1:9545") 1133 .unwrap(); 1134 assert!(peer.last_clock_offset_ms.unwrap() > 30_000); 1135 assert_eq!(peer.last_clock_offset_accepted, Some(true)); 1136 assert_eq!(peer.last_hello.as_ref(), Some(&expected_hello)); 1137 } 1138 1139 #[tokio::test] 1140 async fn hello_remembers_advertised_address_after_signed_session_and_dialback() { 1141 let alice = Wallet::from_seed("hello-dialback-alice"); 1142 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 1143 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 1144 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 1145 let network = super::GossipNetwork { 1146 inner: Arc::new(super::GossipNetworkInner { 1147 node: Arc::clone(&node), 1148 peers: Arc::clone(&peers), 1149 listen_addr: "127.0.0.1:9544".parse().unwrap(), 1150 p2p_announce_addr: tokio::sync::Mutex::new(None), 1151 node_id: super::new_node_id(), 1152 accept_task: tokio::sync::Mutex::new(None), 1153 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 1154 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 1155 metrics: super::P2pMetricsCounters::default(), 1156 sync_progress: StdMutex::new(super::SyncProgressState::default()), 1157 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 1158 }), 1159 }; 1160 let remote_node_id = super::new_node_id(); 1161 let remote_addr = spawn_hello_server(ProtocolHello { 1162 protocol_version: PROTOCOL_VERSION, 1163 capabilities: Vec::new(), 1164 network_id: NETWORK_ID.to_string(), 1165 genesis_hash: node.lock().await.ledger().genesis_hash().to_string(), 1166 listen_addr: None, 1167 node_id: Some(remote_node_id.clone()), 1168 height: 0, 1169 tip_hash: "tip".to_string(), 1170 time_ms: 1_000, 1171 }) 1172 .await; 1173 let original_addr = spawn_verification_responder(remote_node_id.clone()).await; 1174 let stream = tokio::net::TcpStream::connect(original_addr).await.unwrap(); 1175 let remote_socket = stream.peer_addr().unwrap(); 1176 let (reader, mut writer) = stream.into_split(); 1177 let mut reader = super::LimitedLineReader::new(reader); 1178 let hello = ProtocolHello { 1179 protocol_version: PROTOCOL_VERSION, 1180 capabilities: Vec::new(), 1181 network_id: NETWORK_ID.to_string(), 1182 genesis_hash: node.lock().await.ledger().genesis_hash().to_string(), 1183 listen_addr: Some(remote_addr.to_string()), 1184 node_id: Some(remote_node_id), 1185 height: 0, 1186 tip_hash: "tip".to_string(), 1187 time_ms: 1_000, 1188 }; 1189 let mut known_peer = None; 1190 1191 super::process_hello_with_verification( 1192 &network, 1193 &mut writer, 1194 &mut reader, 1195 "test-original-peer", 1196 remote_socket, 1197 &mut known_peer, 1198 hello, 1199 ) 1200 .await 1201 .unwrap(); 1202 1203 assert_eq!(known_peer, Some(remote_addr.to_string())); 1204 let listed = peers.lock().await.list(); 1205 let peer = listed 1206 .iter() 1207 .find(|peer| peer.address == remote_addr.to_string()) 1208 .unwrap(); 1209 assert_eq!(peer.direction, PeerDirection::Outbound); 1210 assert!( 1211 peers 1212 .lock() 1213 .await 1214 .addresses() 1215 .contains(&remote_addr.to_string()) 1216 ); 1217 } 1218 1219 #[tokio::test] 1220 async fn hello_ignores_advertised_address_when_connected_peer_cannot_sign_claimed_node_id() { 1221 let alice = Wallet::from_seed("hello-dialback-spoof-alice"); 1222 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 1223 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 1224 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 1225 let network = super::GossipNetwork { 1226 inner: Arc::new(super::GossipNetworkInner { 1227 node: Arc::clone(&node), 1228 peers: Arc::clone(&peers), 1229 listen_addr: "127.0.0.1:9544".parse().unwrap(), 1230 p2p_announce_addr: tokio::sync::Mutex::new(None), 1231 node_id: super::new_node_id(), 1232 accept_task: tokio::sync::Mutex::new(None), 1233 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 1234 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 1235 metrics: super::P2pMetricsCounters::default(), 1236 sync_progress: StdMutex::new(super::SyncProgressState::default()), 1237 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 1238 }), 1239 }; 1240 let victim_node_id = super::new_node_id(); 1241 let attacker_node_id = super::new_node_id(); 1242 let remote_addr = spawn_hello_server(ProtocolHello { 1243 protocol_version: PROTOCOL_VERSION, 1244 capabilities: Vec::new(), 1245 network_id: NETWORK_ID.to_string(), 1246 genesis_hash: node.lock().await.ledger().genesis_hash().to_string(), 1247 listen_addr: None, 1248 node_id: Some(victim_node_id.clone()), 1249 height: 0, 1250 tip_hash: "tip".to_string(), 1251 time_ms: 1_000, 1252 }) 1253 .await; 1254 let original_addr = spawn_verification_responder(attacker_node_id).await; 1255 let stream = tokio::net::TcpStream::connect(original_addr).await.unwrap(); 1256 let remote_socket = stream.peer_addr().unwrap(); 1257 let (reader, mut writer) = stream.into_split(); 1258 let mut reader = super::LimitedLineReader::new(reader); 1259 let hello = ProtocolHello { 1260 protocol_version: PROTOCOL_VERSION, 1261 capabilities: Vec::new(), 1262 network_id: NETWORK_ID.to_string(), 1263 genesis_hash: node.lock().await.ledger().genesis_hash().to_string(), 1264 listen_addr: Some(remote_addr.to_string()), 1265 node_id: Some(victim_node_id), 1266 height: 0, 1267 tip_hash: "tip".to_string(), 1268 time_ms: 1_000, 1269 }; 1270 let mut known_peer = None; 1271 1272 super::process_hello_with_verification( 1273 &network, 1274 &mut writer, 1275 &mut reader, 1276 "test-attacker-peer", 1277 remote_socket, 1278 &mut known_peer, 1279 hello, 1280 ) 1281 .await 1282 .unwrap(); 1283 1284 assert!(known_peer.is_none()); 1285 assert!( 1286 !peers 1287 .lock() 1288 .await 1289 .addresses() 1290 .contains(&remote_addr.to_string()) 1291 ); 1292 } 1293 1294 #[tokio::test] 1295 async fn dialback_rejects_address_that_signs_with_different_node_id() { 1296 let alice = Wallet::from_seed("hello-dialback-mismatch-alice"); 1297 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 1298 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 1299 let network = super::GossipNetwork { 1300 inner: Arc::new(super::GossipNetworkInner { 1301 node: Arc::clone(&node), 1302 peers: Arc::new(tokio::sync::Mutex::new(PeerBook::default())), 1303 listen_addr: "127.0.0.1:9544".parse().unwrap(), 1304 p2p_announce_addr: tokio::sync::Mutex::new(None), 1305 node_id: super::new_node_id(), 1306 accept_task: tokio::sync::Mutex::new(None), 1307 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 1308 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 1309 metrics: super::P2pMetricsCounters::default(), 1310 sync_progress: StdMutex::new(super::SyncProgressState::default()), 1311 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 1312 }), 1313 }; 1314 let honest_node_id = super::new_node_id(); 1315 let claimed_node_id = super::new_node_id(); 1316 let remote_addr = spawn_hello_server(ProtocolHello { 1317 protocol_version: PROTOCOL_VERSION, 1318 capabilities: Vec::new(), 1319 network_id: NETWORK_ID.to_string(), 1320 genesis_hash: node.lock().await.ledger().genesis_hash().to_string(), 1321 listen_addr: None, 1322 node_id: Some(honest_node_id), 1323 height: 0, 1324 tip_hash: "tip".to_string(), 1325 time_ms: 1_000, 1326 }) 1327 .await; 1328 1329 assert!( 1330 !super::verify_advertised_peer_node_id( 1331 &network, 1332 &remote_addr.to_string(), 1333 &claimed_node_id 1334 ) 1335 .await 1336 ); 1337 } 1338 1339 #[tokio::test] 1340 async fn inbound_verification_only_session_closes_after_response() { 1341 let alice = Wallet::from_seed("verification-only-close-alice"); 1342 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 1343 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 1344 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 1345 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 1346 let listen_addr = listener.local_addr().unwrap(); 1347 drop(listener); 1348 let network = super::GossipNetwork::start(node, peers, listen_addr, None, true) 1349 .await 1350 .unwrap(); 1351 1352 let stream = tokio::net::TcpStream::connect(listen_addr).await.unwrap(); 1353 let (reader, mut writer) = stream.into_split(); 1354 let mut reader = super::LimitedLineReader::new(reader); 1355 let hello_line = reader.read_line().await.unwrap().unwrap(); 1356 let node_id = match super::parse_envelope(&hello_line).unwrap() { 1357 GossipEnvelope::Hello(hello) => hello.node_id.unwrap(), 1358 other => panic!("expected hello, got {other:?}"), 1359 }; 1360 let nonce = super::new_verification_nonce(); 1361 super::write_envelope( 1362 &mut writer, 1363 &GossipEnvelope::PeerVerificationChallenge { 1364 address: listen_addr.to_string(), 1365 nonce: nonce.clone(), 1366 }, 1367 ) 1368 .await 1369 .unwrap(); 1370 1371 let response_line = tokio::time::timeout(std::time::Duration::from_secs(1), reader.read_line()) 1372 .await 1373 .unwrap() 1374 .unwrap() 1375 .unwrap(); 1376 match super::parse_envelope(&response_line).unwrap() { 1377 GossipEnvelope::PeerVerificationResponse { 1378 address, 1379 nonce: response_nonce, 1380 node_id: response_node_id, 1381 signature, 1382 } => assert!(super::peer_verification_response_is_valid( 1383 &address, 1384 &response_nonce, 1385 &response_node_id, 1386 &signature, 1387 &listen_addr.to_string(), 1388 &nonce, 1389 &node_id, 1390 )), 1391 other => panic!("expected verification response, got {other:?}"), 1392 } 1393 1394 let closed = tokio::time::timeout(std::time::Duration::from_secs(1), reader.read_line()) 1395 .await 1396 .unwrap() 1397 .unwrap(); 1398 assert!(closed.is_none()); 1399 network.set_accept_inbound(false).await.unwrap(); 1400 } 1401 1402 #[tokio::test] 1403 async fn setup_placeholder_rejects_bootstrap_with_unpinned_candidate_genesis() { 1404 let local_wallet = Wallet::from_seed("setup-placeholder-local"); 1405 let local_ledger = Ledger::new(BTreeMap::new(), 1); 1406 let local_node = Arc::new(tokio::sync::Mutex::new(NodeCore::from_ledger( 1407 local_wallet, 1408 local_ledger.clone(), 1409 0, 1410 ))); 1411 let network = super::GossipNetwork { 1412 inner: Arc::new(super::GossipNetworkInner { 1413 node: local_node, 1414 peers: Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 1415 "iuna.jhx.app:9444".to_string(), 1416 ]))), 1417 listen_addr: "127.0.0.1:9544".parse().unwrap(), 1418 p2p_announce_addr: tokio::sync::Mutex::new(None), 1419 node_id: super::new_node_id(), 1420 accept_task: tokio::sync::Mutex::new(None), 1421 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 1422 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 1423 metrics: super::P2pMetricsCounters::default(), 1424 sync_progress: StdMutex::new(super::SyncProgressState::default()), 1425 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 1426 }), 1427 }; 1428 1429 let remote_wallet = Wallet::from_seed("setup-placeholder-remote"); 1430 let remote_node = node( 1431 "remote", 1432 remote_wallet.clone(), 1433 allocations(std::slice::from_ref(&remote_wallet), 1_000), 1434 ); 1435 let remote_genesis = remote_node.ledger().genesis_hash().to_string(); 1436 let remote_bootstrap = remote_node.chain_bootstrap(); 1437 let hello = ProtocolHello { 1438 protocol_version: PROTOCOL_VERSION, 1439 capabilities: Vec::new(), 1440 network_id: NETWORK_ID.to_string(), 1441 genesis_hash: remote_genesis.clone(), 1442 listen_addr: Some("142.132.164.59:9444".to_string()), 1443 node_id: None, 1444 height: 5, 1445 tip_hash: "remote-tip".to_string(), 1446 time_ms: 1_000, 1447 }; 1448 let mut known_peer = Some("iuna.jhx.app:9444".to_string()); 1449 let peer_status = super::process_hello( 1450 &network, 1451 "142.132.164.59:51234".parse().unwrap(), 1452 &mut known_peer, 1453 hello, 1454 ) 1455 .await 1456 .unwrap(); 1457 1458 assert!(peer_status.request_bootstrap); 1459 assert_eq!(known_peer.as_deref(), Some("iuna.jhx.app:9444")); 1460 let listed = network.inner.peers.lock().await.list(); 1461 assert_eq!(listed.len(), 1); 1462 let peer = listed 1463 .into_iter() 1464 .find(|peer| peer.address == "iuna.jhx.app:9444") 1465 .unwrap(); 1466 assert_eq!(peer.misbehavior_score, 0); 1467 assert!(!peer.is_banned_at(crate::app::now_ms())); 1468 1469 let error = super::validate_chain_bootstrap( 1470 MAINNET_CANDIDATE_NETWORK_ID, 1471 remote_bootstrap, 1472 crate::app::now_ms(), 1473 ) 1474 .await 1475 .unwrap_err(); 1476 assert!(error.to_string().contains("does not match pinned genesis")); 1477 assert_eq!( 1478 network.inner.node.lock().await.ledger().genesis_hash(), 1479 local_ledger.genesis_hash() 1480 ); 1481 } 1482 1483 #[tokio::test] 1484 async fn real_node_accepts_setup_placeholder_peer_without_requesting_its_chain() { 1485 let wallet = Wallet::from_seed("setup-placeholder-peer-real-node"); 1486 let node = Arc::new(tokio::sync::Mutex::new(node( 1487 "real", 1488 wallet.clone(), 1489 allocations(std::slice::from_ref(&wallet), 1_000), 1490 ))); 1491 let network = super::GossipNetwork { 1492 inner: Arc::new(super::GossipNetworkInner { 1493 node: Arc::clone(&node), 1494 peers: Arc::new(tokio::sync::Mutex::new(PeerBook::default())), 1495 listen_addr: "127.0.0.1:9544".parse().unwrap(), 1496 p2p_announce_addr: tokio::sync::Mutex::new(None), 1497 node_id: super::new_node_id(), 1498 accept_task: tokio::sync::Mutex::new(None), 1499 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 1500 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 1501 metrics: super::P2pMetricsCounters::default(), 1502 sync_progress: StdMutex::new(super::SyncProgressState::default()), 1503 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 1504 }), 1505 }; 1506 let setup_ledger = Ledger::new(BTreeMap::new(), 1); 1507 let hello = ProtocolHello { 1508 protocol_version: PROTOCOL_VERSION, 1509 capabilities: Vec::new(), 1510 network_id: NETWORK_ID.to_string(), 1511 genesis_hash: setup_ledger.genesis_hash().to_string(), 1512 listen_addr: Some("127.0.0.1:9545".to_string()), 1513 node_id: None, 1514 height: 0, 1515 tip_hash: setup_ledger.status().tip_hash, 1516 time_ms: 1_000, 1517 }; 1518 1519 let peer_status = super::process_hello( 1520 &network, 1521 "127.0.0.1:51234".parse().unwrap(), 1522 &mut None, 1523 hello, 1524 ) 1525 .await 1526 .unwrap(); 1527 1528 assert!(!peer_status.request_bootstrap); 1529 } 1530 1531 #[tokio::test] 1532 async fn periodic_broadcast_does_not_push_duplicate_block_pages_to_lagging_peer() { 1533 let wallet = Wallet::from_seed("lagging-peer-pull-only"); 1534 let mut local = node( 1535 "lagging-peer-pull-only-source", 1536 wallet.clone(), 1537 allocations(std::slice::from_ref(&wallet), 1_000), 1538 ); 1539 let genesis_hash = local.ledger().genesis_hash().to_string(); 1540 queue_plaintext_burn(&mut local, &wallet, 1); 1541 local.drain_outbox(); 1542 local.mine_one_at(1).unwrap(); 1543 local.drain_outbox(); 1544 let local = Arc::new(tokio::sync::Mutex::new(local)); 1545 1546 let payload = super::envelopes_for_peer( 1547 Some(&local), 1548 Some(super::PeerStatus::new(0, genesis_hash)), 1549 &[GossipEnvelope::PeerStatus { 1550 height: 1, 1551 tip_hash: "tip".to_string(), 1552 time_ms: 1, 1553 }], 1554 ) 1555 .await; 1556 1557 assert!(matches!( 1558 payload.as_slice(), 1559 [GossipEnvelope::PeerStatus { .. }] 1560 )); 1561 } 1562 1563 #[tokio::test] 1564 async fn transaction_v2_mempool_gossip_requires_explicit_peer_capability() { 1565 let wallet = Wallet::from_seed("v2-mempool-capability-filter"); 1566 let node = Arc::new(tokio::sync::Mutex::new(node( 1567 "v2-mempool-capability-filter", 1568 wallet.clone(), 1569 allocations(std::slice::from_ref(&wallet), 1_000), 1570 ))); 1571 let envelopes = vec![ 1572 GossipEnvelope::TransactionV2 { 1573 envelope: "00".to_string(), 1574 }, 1575 GossipEnvelope::PeerStatus { 1576 height: 1, 1577 tip_hash: "tip".to_string(), 1578 time_ms: 1, 1579 }, 1580 ]; 1581 1582 let legacy = super::PeerStatus::new(1, "tip".to_string()).with_capabilities(vec![ 1583 crate::app::CAPABILITY_ADDRESS_V1_READ.to_string(), 1584 crate::app::CAPABILITY_SIGNATURE_SCHEMES_V1.to_string(), 1585 crate::app::CAPABILITY_TRANSACTION_V2_BLOCKS.to_string(), 1586 ]); 1587 let legacy_payload = super::envelopes_for_peer(Some(&node), Some(legacy), &envelopes).await; 1588 assert!(matches!( 1589 legacy_payload.as_slice(), 1590 [GossipEnvelope::PeerStatus { .. }] 1591 )); 1592 1593 let upgraded = super::PeerStatus::new(1, "tip".to_string()).with_capabilities(vec![ 1594 crate::app::CAPABILITY_TRANSACTION_V2_MEMPOOL.to_string(), 1595 ]); 1596 let upgraded_payload = super::envelopes_for_peer(Some(&node), Some(upgraded), &envelopes).await; 1597 assert_eq!(upgraded_payload, envelopes); 1598 } 1599 1600 #[tokio::test] 1601 async fn hello_ignores_private_advertised_listen_address() { 1602 let alice = Wallet::from_seed("hello-private-listen-alice"); 1603 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 1604 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 1605 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 1606 let network = super::GossipNetwork { 1607 inner: Arc::new(super::GossipNetworkInner { 1608 node: Arc::clone(&node), 1609 peers: Arc::clone(&peers), 1610 listen_addr: "0.0.0.0:9444".parse().unwrap(), 1611 p2p_announce_addr: tokio::sync::Mutex::new(None), 1612 node_id: super::new_node_id(), 1613 accept_task: tokio::sync::Mutex::new(None), 1614 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 1615 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 1616 metrics: super::P2pMetricsCounters::default(), 1617 sync_progress: StdMutex::new(super::SyncProgressState::default()), 1618 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 1619 }), 1620 }; 1621 let status = node.lock().await.ledger().status(); 1622 let hello = ProtocolHello { 1623 protocol_version: PROTOCOL_VERSION, 1624 capabilities: Vec::new(), 1625 network_id: NETWORK_ID.to_string(), 1626 genesis_hash: node.lock().await.ledger().genesis_hash().to_string(), 1627 listen_addr: Some("10.42.1.1:12138".to_string()), 1628 node_id: None, 1629 height: status.height, 1630 tip_hash: status.tip_hash, 1631 time_ms: 1_000, 1632 }; 1633 1634 let mut known_peer = None; 1635 super::process_hello( 1636 &network, 1637 "142.132.164.59:51234".parse().unwrap(), 1638 &mut known_peer, 1639 hello, 1640 ) 1641 .await 1642 .unwrap(); 1643 1644 assert!(known_peer.is_none()); 1645 assert!(peers.lock().await.addresses().is_empty()); 1646 } 1647 1648 #[tokio::test] 1649 async fn hello_ignores_loopback_alias_for_unspecified_self() { 1650 let alice = Wallet::from_seed("hello-self-alias-alice"); 1651 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 1652 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 1653 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 1654 let network = super::GossipNetwork { 1655 inner: Arc::new(super::GossipNetworkInner { 1656 node, 1657 peers: Arc::clone(&peers), 1658 listen_addr: "0.0.0.0:9545".parse().unwrap(), 1659 p2p_announce_addr: tokio::sync::Mutex::new(None), 1660 node_id: super::new_node_id(), 1661 accept_task: tokio::sync::Mutex::new(None), 1662 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 1663 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 1664 metrics: super::P2pMetricsCounters::default(), 1665 sync_progress: StdMutex::new(super::SyncProgressState::default()), 1666 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 1667 }), 1668 }; 1669 let hello = ProtocolHello { 1670 protocol_version: PROTOCOL_VERSION, 1671 capabilities: Vec::new(), 1672 network_id: NETWORK_ID.to_string(), 1673 genesis_hash: network 1674 .inner 1675 .node 1676 .lock() 1677 .await 1678 .ledger() 1679 .genesis_hash() 1680 .to_string(), 1681 listen_addr: Some("127.0.0.1:9545".to_string()), 1682 node_id: None, 1683 height: 0, 1684 tip_hash: "tip".to_string(), 1685 time_ms: 1_000, 1686 }; 1687 1688 super::process_hello( 1689 &network, 1690 "127.0.0.1:52144".parse().unwrap(), 1691 &mut None, 1692 hello, 1693 ) 1694 .await 1695 .unwrap(); 1696 1697 assert_eq!(network.metrics().self_peer_rejections, 1); 1698 assert!(peers.lock().await.addresses().is_empty()); 1699 let listed = peers.lock().await.list(); 1700 assert_eq!(listed.len(), 1); 1701 assert_eq!(listed[0].direction, PeerDirection::Inbound); 1702 } 1703 1704 #[tokio::test] 1705 async fn hello_removes_outbound_peer_that_announces_self_address() { 1706 let alice = Wallet::from_seed("hello-self-outbound-alice"); 1707 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 1708 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 1709 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 1710 "10.42.1.1:16987".to_string(), 1711 ]))); 1712 let network = super::GossipNetwork { 1713 inner: Arc::new(super::GossipNetworkInner { 1714 node, 1715 peers: Arc::clone(&peers), 1716 listen_addr: "0.0.0.0:9444".parse().unwrap(), 1717 p2p_announce_addr: tokio::sync::Mutex::new(None), 1718 node_id: super::new_node_id(), 1719 accept_task: tokio::sync::Mutex::new(None), 1720 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 1721 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 1722 metrics: super::P2pMetricsCounters::default(), 1723 sync_progress: StdMutex::new(super::SyncProgressState::default()), 1724 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 1725 }), 1726 }; 1727 let hello = ProtocolHello { 1728 protocol_version: PROTOCOL_VERSION, 1729 capabilities: Vec::new(), 1730 network_id: NETWORK_ID.to_string(), 1731 genesis_hash: network 1732 .inner 1733 .node 1734 .lock() 1735 .await 1736 .ledger() 1737 .genesis_hash() 1738 .to_string(), 1739 listen_addr: Some("127.0.0.1:9444".to_string()), 1740 node_id: None, 1741 height: 0, 1742 tip_hash: "tip".to_string(), 1743 time_ms: 1_000, 1744 }; 1745 let mut known_peer = Some("10.42.1.1:16987".to_string()); 1746 1747 super::process_hello( 1748 &network, 1749 "10.42.1.1:16987".parse().unwrap(), 1750 &mut known_peer, 1751 hello, 1752 ) 1753 .await 1754 .unwrap(); 1755 1756 assert_eq!(network.metrics().self_peer_rejections, 1); 1757 assert!(known_peer.is_none()); 1758 assert!(peers.lock().await.addresses().is_empty()); 1759 } 1760 1761 #[tokio::test] 1762 async fn hello_removes_outbound_peer_with_same_node_id() { 1763 let alice = Wallet::from_seed("hello-self-node-id-alice"); 1764 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 1765 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 1766 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 1767 "142.132.164.59:9444".to_string(), 1768 ]))); 1769 let network = super::GossipNetwork { 1770 inner: Arc::new(super::GossipNetworkInner { 1771 node, 1772 peers: Arc::clone(&peers), 1773 listen_addr: "0.0.0.0:9444".parse().unwrap(), 1774 p2p_announce_addr: tokio::sync::Mutex::new(None), 1775 node_id: super::new_node_id(), 1776 accept_task: tokio::sync::Mutex::new(None), 1777 sessions: tokio::sync::Mutex::new(BTreeMap::new()), 1778 inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), 1779 metrics: super::P2pMetricsCounters::default(), 1780 sync_progress: StdMutex::new(super::SyncProgressState::default()), 1781 chain_validation: Arc::new(super::ChainValidationCoordinator::default()), 1782 }), 1783 }; 1784 let hello = ProtocolHello { 1785 protocol_version: PROTOCOL_VERSION, 1786 capabilities: Vec::new(), 1787 network_id: NETWORK_ID.to_string(), 1788 genesis_hash: network 1789 .inner 1790 .node 1791 .lock() 1792 .await 1793 .ledger() 1794 .genesis_hash() 1795 .to_string(), 1796 listen_addr: Some("0.0.0.0:9444".to_string()), 1797 node_id: Some(network.inner.node_id.clone()), 1798 height: 0, 1799 tip_hash: "tip".to_string(), 1800 time_ms: 1_000, 1801 }; 1802 let mut known_peer = Some("142.132.164.59:9444".to_string()); 1803 1804 super::process_hello( 1805 &network, 1806 "142.132.164.59:52144".parse().unwrap(), 1807 &mut known_peer, 1808 hello, 1809 ) 1810 .await 1811 .unwrap(); 1812 1813 assert_eq!(network.metrics().self_peer_rejections, 1); 1814 assert!(known_peer.is_none()); 1815 assert!(peers.lock().await.addresses().is_empty()); 1816 } 1817 1818 async fn spawn_hello_server(hello: ProtocolHello) -> SocketAddr { 1819 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 1820 let addr = listener.local_addr().unwrap(); 1821 tokio::spawn(async move { 1822 let Ok((stream, _)) = listener.accept().await else { 1823 return; 1824 }; 1825 let node_id = hello.node_id.clone(); 1826 let (reader, mut writer) = stream.into_split(); 1827 let line = serde_json::to_string(&GossipEnvelope::Hello(hello)).unwrap(); 1828 let _ = writer.write_all(line.as_bytes()).await; 1829 let _ = writer.write_all(b"\n").await; 1830 let Some(node_id) = node_id else { 1831 return; 1832 }; 1833 let mut reader = super::LimitedLineReader::new(reader); 1834 let Ok(Some(line)) = reader.read_line().await else { 1835 return; 1836 }; 1837 let Ok(GossipEnvelope::PeerVerificationChallenge { address, nonce }) = 1838 super::parse_envelope(&line) 1839 else { 1840 return; 1841 }; 1842 let Some(response) = 1843 super::peer_verification_response_for_node_id(&node_id, &address, &nonce) 1844 else { 1845 return; 1846 }; 1847 let line = serde_json::to_string(&response).unwrap(); 1848 let _ = writer.write_all(line.as_bytes()).await; 1849 let _ = writer.write_all(b"\n").await; 1850 }); 1851 addr 1852 } 1853 1854 async fn spawn_verification_responder(node_id: String) -> SocketAddr { 1855 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 1856 let addr = listener.local_addr().unwrap(); 1857 tokio::spawn(async move { 1858 let Ok((stream, _)) = listener.accept().await else { 1859 return; 1860 }; 1861 let (reader, mut writer) = stream.into_split(); 1862 let mut reader = super::LimitedLineReader::new(reader); 1863 let Ok(Some(line)) = reader.read_line().await else { 1864 return; 1865 }; 1866 let Ok(GossipEnvelope::PeerVerificationChallenge { address, nonce }) = 1867 super::parse_envelope(&line) 1868 else { 1869 return; 1870 }; 1871 let Some(response) = 1872 super::peer_verification_response_for_node_id(&node_id, &address, &nonce) 1873 else { 1874 return; 1875 }; 1876 let line = serde_json::to_string(&response).unwrap(); 1877 let _ = writer.write_all(line.as_bytes()).await; 1878 let _ = writer.write_all(b"\n").await; 1879 }); 1880 addr 1881 }