handshake.rs (19626B)
1 use std::{collections::BTreeMap, net::SocketAddr}; 2 3 use anyhow::Result; 4 use tokio::{ 5 net::{TcpStream, tcp::OwnedWriteHalf}, 6 time::timeout, 7 }; 8 9 use super::identity::{ 10 new_verification_nonce, peer_verification_response, peer_verification_response_is_valid, 11 }; 12 use super::line_codec::{LimitedLineReader, parse_envelope, read_session_envelope}; 13 use super::metrics::P2pMetricsCounters; 14 use super::peer_addr::{advertised_peer_is_discoverable, normalize_advertised_peer}; 15 use super::{ 16 CONNECT_TIMEOUT, GossipNetwork, HANDSHAKE_TIMEOUT, MAX_PEER_VERIFICATION_ENVELOPES, PeerStatus, 17 write_envelope, 18 }; 19 use crate::{ 20 app::{ 21 GossipEnvelope, NETWORK_ID, PROTOCOL_VERSION, PeerDirection, ProtocolHello, 22 debug_logging_enabled, now_ms, validate_protocol_capabilities, 23 validate_transaction_v2_peer_capability, 24 }, 25 domain::Ledger, 26 }; 27 28 pub(super) struct PeerVerificationSession<'a> { 29 pub(super) writer: &'a mut OwnedWriteHalf, 30 pub(super) reader: &'a mut LimitedLineReader<tokio::net::tcp::OwnedReadHalf>, 31 pub(super) connection_label: &'a str, 32 } 33 34 pub(super) async fn record_peer_status( 35 network: &GossipNetwork, 36 known_peer: &Option<String>, 37 remote_addr: SocketAddr, 38 peer_status: &PeerStatus, 39 ) { 40 let local_receive_time_ms = now_ms(); 41 if let Some(peer) = known_peer { 42 let mut peers = network.inner.peers.lock().await; 43 peers.record_status(peer, peer_status.height, peer_status.tip_hash.clone()); 44 peers.record_clock_observation( 45 peer, 46 PeerDirection::Outbound, 47 peer_status.time_ms, 48 local_receive_time_ms, 49 ); 50 } else { 51 let peer = remote_addr.to_string(); 52 let mut peers = network.inner.peers.lock().await; 53 peers.record_clock_observation( 54 &peer, 55 PeerDirection::Inbound, 56 peer_status.time_ms, 57 local_receive_time_ms, 58 ); 59 peers.record_received(&peer, 1); 60 } 61 } 62 63 async fn record_peer_hello( 64 network: &GossipNetwork, 65 known_peer: &Option<String>, 66 remote_addr: SocketAddr, 67 hello: ProtocolHello, 68 ) { 69 let (peer, direction) = match known_peer { 70 Some(peer) => (peer.clone(), PeerDirection::Outbound), 71 None => (remote_addr.to_string(), PeerDirection::Inbound), 72 }; 73 network 74 .inner 75 .peers 76 .lock() 77 .await 78 .record_hello(&peer, direction, hello); 79 } 80 81 pub(super) async fn process_hello( 82 network: &GossipNetwork, 83 remote_addr: SocketAddr, 84 known_peer: &mut Option<String>, 85 hello: ProtocolHello, 86 ) -> Result<PeerStatus> { 87 process_hello_inner(network, None, remote_addr, known_peer, hello).await 88 } 89 90 pub(super) async fn process_hello_with_verification( 91 network: &GossipNetwork, 92 writer: &mut OwnedWriteHalf, 93 reader: &mut LimitedLineReader<tokio::net::tcp::OwnedReadHalf>, 94 connection_label: &str, 95 remote_addr: SocketAddr, 96 known_peer: &mut Option<String>, 97 hello: ProtocolHello, 98 ) -> Result<PeerStatus> { 99 let mut verification_session = PeerVerificationSession { 100 writer, 101 reader, 102 connection_label, 103 }; 104 process_hello_inner( 105 network, 106 Some(&mut verification_session), 107 remote_addr, 108 known_peer, 109 hello, 110 ) 111 .await 112 } 113 114 async fn process_hello_inner( 115 network: &GossipNetwork, 116 mut verification_session: Option<&mut PeerVerificationSession<'_>>, 117 remote_addr: SocketAddr, 118 known_peer: &mut Option<String>, 119 hello: ProtocolHello, 120 ) -> Result<PeerStatus> { 121 if hello.protocol_version != PROTOCOL_VERSION { 122 anyhow::bail!( 123 "unsupported protocol version {}; expected {}", 124 hello.protocol_version, 125 PROTOCOL_VERSION 126 ); 127 } 128 validate_protocol_capabilities(&hello.capabilities)?; 129 if hello.network_id != NETWORK_ID { 130 anyhow::bail!( 131 "wrong network {}; expected {}", 132 hello.network_id, 133 NETWORK_ID 134 ); 135 } 136 if hello 137 .node_id 138 .as_deref() 139 .is_some_and(|node_id| node_id == network.inner.node_id) 140 { 141 P2pMetricsCounters::inc(&network.inner.metrics.self_peer_rejections); 142 forget_stale_self_peer(network, known_peer).await; 143 return Ok(PeerStatus::rejected( 144 hello.height, 145 hello.tip_hash, 146 hello.time_ms, 147 )); 148 } 149 let (local_genesis, local_accepts_remote_genesis, local_height) = { 150 let node = network.inner.node.lock().await; 151 ( 152 node.ledger().genesis_hash().to_string(), 153 node.ledger().is_setup_placeholder(), 154 node.ledger().height(), 155 ) 156 }; 157 validate_transaction_v2_peer_capability(&hello.capabilities, local_height, hello.height)?; 158 let peer_capabilities = hello.capabilities.clone(); 159 let genesis_mismatch = hello.genesis_hash != local_genesis; 160 let remote_is_setup_placeholder = 161 hello.height == 0 && hello.genesis_hash == setup_placeholder_genesis_hash(); 162 let request_bootstrap = genesis_mismatch && local_accepts_remote_genesis; 163 if genesis_mismatch && !local_accepts_remote_genesis && !remote_is_setup_placeholder { 164 anyhow::bail!( 165 "wrong genesis {}; expected {local_genesis}", 166 hello.genesis_hash 167 ); 168 } 169 170 let remote_node_id = hello.node_id.clone(); 171 let mut reject_session = false; 172 if let Some(listen_addr) = &hello.listen_addr { 173 let peer = normalize_advertised_peer(listen_addr, remote_addr)?; 174 if network.is_self_peer(&peer).await { 175 P2pMetricsCounters::inc(&network.inner.metrics.self_peer_rejections); 176 forget_stale_self_peer(network, known_peer).await; 177 reject_session = true; 178 } else { 179 let verified = match verification_session.as_mut() { 180 Some(session) => { 181 remember_verified_advertised_peer( 182 network, 183 session, 184 remote_addr, 185 known_peer, 186 peer.clone(), 187 remote_node_id.as_deref(), 188 ) 189 .await? 190 } 191 None => false, 192 }; 193 if !verified && debug_logging_enabled() { 194 eprintln!( 195 "p2p advertised address {peer} ignored because ownership was not verified" 196 ); 197 } 198 } 199 } 200 record_peer_status( 201 network, 202 known_peer, 203 remote_addr, 204 &PeerStatus::with_time(hello.height, hello.tip_hash.clone(), hello.time_ms), 205 ) 206 .await; 207 record_peer_hello(network, known_peer, remote_addr, hello.clone()).await; 208 let mut status = if request_bootstrap { 209 PeerStatus::with_bootstrap_request(hello.height, hello.tip_hash, hello.time_ms) 210 .with_capabilities(peer_capabilities) 211 } else { 212 PeerStatus::with_time(hello.height, hello.tip_hash, hello.time_ms) 213 .with_capabilities(peer_capabilities) 214 }; 215 status.reject_session = reject_session; 216 Ok(status) 217 } 218 219 fn setup_placeholder_genesis_hash() -> String { 220 Ledger::new(BTreeMap::new(), 1).genesis_hash().to_string() 221 } 222 223 async fn remember_verified_advertised_peer( 224 network: &GossipNetwork, 225 session: &mut PeerVerificationSession<'_>, 226 remote_addr: SocketAddr, 227 known_peer: &mut Option<String>, 228 peer: String, 229 expected_node_id: Option<&str>, 230 ) -> Result<bool> { 231 if !advertised_peer_is_discoverable(&peer, remote_addr)? { 232 return Ok(false); 233 } 234 if known_peer.as_deref() != Some(peer.as_str()) { 235 let Some(expected_node_id) = expected_node_id else { 236 return Ok(false); 237 }; 238 if !verify_connected_peer_node_id(network, session, &peer, expected_node_id).await? { 239 return Ok(false); 240 } 241 if !verify_advertised_peer_node_id(network, &peer, expected_node_id).await { 242 return Ok(false); 243 } 244 } 245 remember_discoverable_advertised_peer(network, remote_addr, known_peer, peer).await 246 } 247 248 async fn verify_connected_peer_node_id( 249 network: &GossipNetwork, 250 session: &mut PeerVerificationSession<'_>, 251 peer: &str, 252 expected_node_id: &str, 253 ) -> Result<bool> { 254 let nonce = new_verification_nonce(); 255 write_envelope( 256 session.writer, 257 &GossipEnvelope::PeerVerificationChallenge { 258 address: peer.to_string(), 259 nonce: nonce.clone(), 260 }, 261 ) 262 .await?; 263 264 for _ in 0..MAX_PEER_VERIFICATION_ENVELOPES { 265 let envelope = match timeout( 266 HANDSHAKE_TIMEOUT, 267 read_session_envelope(network, session.connection_label, session.reader), 268 ) 269 .await 270 { 271 Ok(Ok(Some(envelope))) => envelope, 272 Ok(Ok(None)) | Err(_) => return Ok(false), 273 Ok(Err(error)) => return Err(error), 274 }; 275 match envelope { 276 GossipEnvelope::PeerVerificationResponse { 277 address, 278 nonce: response_nonce, 279 node_id, 280 signature, 281 } => { 282 return Ok(peer_verification_response_is_valid( 283 &address, 284 &response_nonce, 285 &node_id, 286 &signature, 287 peer, 288 &nonce, 289 expected_node_id, 290 )); 291 } 292 GossipEnvelope::PeerVerificationChallenge { address, nonce } => { 293 if let Some(response) = peer_verification_response(network, &address, &nonce) { 294 write_envelope(session.writer, &response).await?; 295 } 296 } 297 _ => {} 298 } 299 } 300 Ok(false) 301 } 302 303 pub(super) async fn verify_advertised_peer_node_id( 304 network: &GossipNetwork, 305 peer: &str, 306 expected_node_id: &str, 307 ) -> bool { 308 let stream = match timeout(CONNECT_TIMEOUT, TcpStream::connect(peer)).await { 309 Ok(Ok(stream)) => stream, 310 Ok(Err(error)) => { 311 if debug_logging_enabled() { 312 eprintln!("p2p announced address {peer} failed verification: {error}"); 313 } 314 return false; 315 } 316 Err(_) => { 317 if debug_logging_enabled() { 318 eprintln!("p2p announced address {peer} failed verification: timeout"); 319 } 320 return false; 321 } 322 }; 323 let (reader, mut writer) = stream.into_split(); 324 let mut reader = LimitedLineReader::new(reader); 325 let line = match timeout(HANDSHAKE_TIMEOUT, reader.read_line()).await { 326 Ok(Ok(Some(line))) => line, 327 Ok(Ok(None)) => return false, 328 Ok(Err(error)) => { 329 if debug_logging_enabled() { 330 eprintln!( 331 "p2p announced address {peer} sent invalid verification hello: {error:#}" 332 ); 333 } 334 return false; 335 } 336 Err(_) => return false, 337 }; 338 let hello = match parse_envelope(&line) { 339 Ok(GossipEnvelope::Hello(hello)) => hello, 340 Ok(_) | Err(_) => return false, 341 }; 342 343 if !advertised_peer_hello_is_compatible(network, &hello).await 344 || hello.node_id.as_deref() != Some(expected_node_id) 345 { 346 return false; 347 } 348 349 let nonce = new_verification_nonce(); 350 if write_envelope( 351 &mut writer, 352 &GossipEnvelope::PeerVerificationChallenge { 353 address: peer.to_string(), 354 nonce: nonce.clone(), 355 }, 356 ) 357 .await 358 .is_err() 359 { 360 return false; 361 } 362 for _ in 0..MAX_PEER_VERIFICATION_ENVELOPES { 363 let line = match timeout(HANDSHAKE_TIMEOUT, reader.read_line()).await { 364 Ok(Ok(Some(line))) => line, 365 Ok(Ok(None)) | Ok(Err(_)) | Err(_) => return false, 366 }; 367 let envelope = match parse_envelope(&line) { 368 Ok(envelope) => envelope, 369 Err(_) => return false, 370 }; 371 if let GossipEnvelope::PeerVerificationResponse { 372 address, 373 nonce: response_nonce, 374 node_id, 375 signature, 376 } = envelope 377 { 378 return peer_verification_response_is_valid( 379 &address, 380 &response_nonce, 381 &node_id, 382 &signature, 383 peer, 384 &nonce, 385 expected_node_id, 386 ); 387 } 388 } 389 false 390 } 391 392 async fn advertised_peer_hello_is_compatible( 393 network: &GossipNetwork, 394 hello: &ProtocolHello, 395 ) -> bool { 396 if hello.protocol_version != PROTOCOL_VERSION 397 || hello.network_id != NETWORK_ID 398 || validate_protocol_capabilities(&hello.capabilities).is_err() 399 { 400 return false; 401 } 402 let (local_genesis, local_accepts_remote_genesis, local_height) = { 403 let node = network.inner.node.lock().await; 404 ( 405 node.ledger().genesis_hash().to_string(), 406 node.ledger().is_setup_placeholder(), 407 node.ledger().height(), 408 ) 409 }; 410 if validate_transaction_v2_peer_capability(&hello.capabilities, local_height, hello.height) 411 .is_err() 412 { 413 return false; 414 } 415 let remote_is_setup_placeholder = 416 hello.height == 0 && hello.genesis_hash == setup_placeholder_genesis_hash(); 417 hello.genesis_hash == local_genesis 418 || local_accepts_remote_genesis 419 || remote_is_setup_placeholder 420 } 421 422 pub(super) async fn remember_discoverable_advertised_peer( 423 network: &GossipNetwork, 424 remote_addr: SocketAddr, 425 known_peer: &mut Option<String>, 426 peer: String, 427 ) -> Result<bool> { 428 if !advertised_peer_is_discoverable(&peer, remote_addr)? { 429 return Ok(false); 430 } 431 if let Some(previous_peer) = known_peer.as_deref() { 432 network 433 .inner 434 .peers 435 .lock() 436 .await 437 .replace_peer_address(previous_peer, peer.clone()); 438 } else { 439 network 440 .inner 441 .peers 442 .lock() 443 .await 444 .add_discovered_peer(peer.clone()); 445 } 446 *known_peer = Some(peer); 447 Ok(true) 448 } 449 450 pub(super) async fn forget_stale_self_peer( 451 network: &GossipNetwork, 452 known_peer: &mut Option<String>, 453 ) { 454 if let Some(previous_peer) = known_peer.take() { 455 network.inner.peers.lock().await.remove_peer(&previous_peer); 456 } 457 } 458 459 #[cfg(test)] 460 mod tests { 461 use std::sync::Arc; 462 463 use crate::{ 464 app::{PeerBook, PeerDirection}, 465 domain::Wallet, 466 }; 467 468 use super::super::{ 469 PeerStatus, 470 test_support::{allocations, gossip_network, node}, 471 }; 472 use super::{ 473 forget_stale_self_peer, record_peer_status, remember_discoverable_advertised_peer, 474 }; 475 476 #[tokio::test] 477 async fn inbound_announced_address_replaces_gateway_address_for_ui() { 478 let alice = Wallet::from_seed("hello-public-announced-inbound-alice"); 479 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 480 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 481 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 482 let network = gossip_network( 483 node, 484 Arc::clone(&peers), 485 "0.0.0.0:9444".parse().unwrap(), 486 None, 487 ); 488 let mut known_peer = None; 489 490 let remembered = remember_discoverable_advertised_peer( 491 &network, 492 "10.42.0.1:51234".parse().unwrap(), 493 &mut known_peer, 494 "142.132.164.59:9444".to_string(), 495 ) 496 .await 497 .unwrap(); 498 record_peer_status( 499 &network, 500 &known_peer, 501 "10.42.0.1:51234".parse().unwrap(), 502 &PeerStatus::with_time(7, "tip".to_string(), 1_000), 503 ) 504 .await; 505 506 assert!(remembered); 507 assert_eq!(known_peer.as_deref(), Some("142.132.164.59:9444")); 508 let listed = peers.lock().await.list(); 509 assert_eq!(listed.len(), 1); 510 let peer = &listed[0]; 511 assert_eq!(peer.address, "142.132.164.59:9444"); 512 assert_eq!(peer.direction, PeerDirection::Outbound); 513 assert_eq!(peer.last_known_height, Some(7)); 514 assert_eq!(peer.messages_received, 0); 515 516 let repeated = remember_discoverable_advertised_peer( 517 &network, 518 "10.42.0.1:51234".parse().unwrap(), 519 &mut known_peer, 520 "142.132.164.59:9444".to_string(), 521 ) 522 .await 523 .unwrap(); 524 525 assert!(repeated); 526 assert_eq!( 527 peers.lock().await.list()[0].direction, 528 PeerDirection::Outbound 529 ); 530 } 531 532 #[tokio::test] 533 async fn inbound_status_does_not_create_outbound_ephemeral_peer() { 534 let alice = Wallet::from_seed("inbound-status-alice"); 535 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 536 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 537 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 538 let network = gossip_network( 539 node, 540 Arc::clone(&peers), 541 "127.0.0.1:9544".parse().unwrap(), 542 None, 543 ); 544 545 record_peer_status( 546 &network, 547 &None, 548 "127.0.0.1:51729".parse().unwrap(), 549 &PeerStatus::new(4, "tip".to_string()), 550 ) 551 .await; 552 553 let peers = peers.lock().await; 554 assert!(peers.addresses().is_empty()); 555 let listed = peers.list(); 556 assert_eq!(listed.len(), 1); 557 assert_eq!(listed[0].direction, PeerDirection::Inbound); 558 } 559 560 #[tokio::test] 561 async fn peer_announcement_ignores_private_ephemeral_address() { 562 let alice = Wallet::from_seed("px-private-announcement-alice"); 563 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 564 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 565 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); 566 let network = gossip_network( 567 node, 568 Arc::clone(&peers), 569 "0.0.0.0:9444".parse().unwrap(), 570 None, 571 ); 572 let mut known_peer = None; 573 574 let remembered = remember_discoverable_advertised_peer( 575 &network, 576 "142.132.164.59:51234".parse().unwrap(), 577 &mut known_peer, 578 "10.42.1.1:10091".to_string(), 579 ) 580 .await 581 .unwrap(); 582 583 assert!(!remembered); 584 assert!(known_peer.is_none()); 585 assert!(peers.lock().await.addresses().is_empty()); 586 } 587 588 #[tokio::test] 589 async fn peer_announcement_removes_outbound_peer_that_announces_self_address() { 590 let alice = Wallet::from_seed("px-self-announcement-alice"); 591 let allocations = allocations(std::slice::from_ref(&alice), 1_000); 592 let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); 593 let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ 594 "10.42.1.1:30508".to_string(), 595 ]))); 596 let network = gossip_network( 597 node, 598 Arc::clone(&peers), 599 "0.0.0.0:9444".parse().unwrap(), 600 None, 601 ); 602 let mut known_peer = Some("10.42.1.1:30508".to_string()); 603 604 forget_stale_self_peer(&network, &mut known_peer).await; 605 606 assert_eq!(known_peer, None); 607 assert!(peers.lock().await.addresses().is_empty()); 608 } 609 }