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