session.rs (26853B)
1 use std::{net::SocketAddr, time::Duration}; 2 3 use anyhow::{Context, Result, bail}; 4 use tokio::{ 5 net::{TcpListener, TcpStream}, 6 sync::{mpsc, watch}, 7 time::{Instant, interval, interval_at, sleep, timeout}, 8 }; 9 10 use crate::app::{GossipEnvelope, debug_logging_enabled}; 11 12 use super::{ 13 CATCHUP_REQUEST_TIMEOUT, CONNECT_TIMEOUT, GossipNetwork, GossipSession, HANDSHAKE_TIMEOUT, 14 INBOUND_PEER_QUEUE_BYTES, INBOUND_PEER_QUEUE_SIZE, INBOUND_SESSION_PREFIX, 15 INITIAL_RECONNECT_DELAY, MAX_RECONNECT_DELAY, OutboundBatch, PEER_EXCHANGE_INTERVAL, 16 PEER_QUEUE_BYTES, PEER_QUEUE_SIZE, PeerStatus, SESSION_SYNC_INTERVAL, is_self_peer_address_for, 17 next_reconnect_delay_with_max, process_envelope, process_hello_with_verification, 18 read_session_envelope, record_peer_status, respond_to_peer_verification_challenge, 19 write_envelope, write_payload, write_peer_exchange, 20 }; 21 22 struct InboundRegistration { 23 key: String, 24 sender: mpsc::Sender<OutboundBatch>, 25 shutdown: watch::Sender<bool>, 26 queue_bytes: std::sync::Arc<tokio::sync::Semaphore>, 27 } 28 29 pub(super) async fn accept_loop(network: GossipNetwork, listener: TcpListener) { 30 loop { 31 match listener.accept().await { 32 Ok((stream, remote_addr)) => { 33 let network = network.clone(); 34 let permit = match network.try_acquire_inbound_session(remote_addr.ip()) { 35 Ok(permit) => permit, 36 Err(rejection) => { 37 super::P2pMetricsCounters::inc( 38 &network.inner.metrics.inbound_sessions_rejected, 39 ); 40 super::P2pMetricsCounters::set_last( 41 &network.inner.metrics.last_session_failure, 42 format!("{remote_addr}: {}", rejection.label()), 43 ); 44 if debug_logging_enabled() { 45 eprintln!( 46 "p2p inbound connection from {remote_addr} rejected: {}", 47 rejection.label() 48 ); 49 } 50 drop(stream); 51 continue; 52 } 53 }; 54 super::P2pMetricsCounters::inc(&network.inner.metrics.inbound_sessions_started); 55 let session_key = format!("{INBOUND_SESSION_PREFIX}{remote_addr}"); 56 let (sender, receiver) = mpsc::channel(INBOUND_PEER_QUEUE_SIZE); 57 let (shutdown, mut shutdown_receiver) = watch::channel(false); 58 let queue_bytes = 59 std::sync::Arc::new(tokio::sync::Semaphore::new(INBOUND_PEER_QUEUE_BYTES)); 60 let registration = InboundRegistration { 61 key: session_key.clone(), 62 sender, 63 shutdown, 64 queue_bytes, 65 }; 66 tokio::spawn(async move { 67 let _permit = permit; 68 let result = session_loop( 69 network.clone(), 70 stream, 71 remote_addr, 72 None, 73 receiver, 74 &mut shutdown_receiver, 75 Some(registration), 76 ) 77 .await; 78 network.inner.sessions.lock().await.remove(&session_key); 79 match result { 80 Ok(()) => { 81 super::P2pMetricsCounters::inc(&network.inner.metrics.sessions_closed); 82 } 83 Err(error) if super::is_quiet_disconnect(&error) => { 84 super::P2pMetricsCounters::inc( 85 &network.inner.metrics.quiet_disconnects, 86 ); 87 } 88 Err(error) => { 89 super::P2pMetricsCounters::inc(&network.inner.metrics.session_failures); 90 super::P2pMetricsCounters::set_last( 91 &network.inner.metrics.last_session_failure, 92 format!("{remote_addr}: {error:#}"), 93 ); 94 if debug_logging_enabled() { 95 eprintln!( 96 "p2p inbound connection from {remote_addr} failed: {error:#}" 97 ); 98 } 99 } 100 } 101 }); 102 } 103 Err(error) if debug_logging_enabled() => eprintln!("p2p accept failed: {error:#}"), 104 Err(_) => {} 105 } 106 } 107 } 108 109 pub(super) async fn outbound_supervisor(network: GossipNetwork) { 110 let mut tick = interval(Duration::from_secs(2)); 111 loop { 112 tick.tick().await; 113 network.ensure_outbound_sessions().await; 114 } 115 } 116 117 pub(super) async fn outbound_session( 118 network: GossipNetwork, 119 peer: String, 120 mut receiver: mpsc::Receiver<OutboundBatch>, 121 mut shutdown: watch::Receiver<bool>, 122 ) { 123 let mut reconnect_delay = INITIAL_RECONNECT_DELAY; 124 loop { 125 let self_filter_addr = network.self_filter_addr().await; 126 if !peer_is_connectable(&network, &peer).await 127 || is_self_peer_address_for(&peer, network.inner.listen_addr, self_filter_addr) 128 { 129 network.inner.sessions.lock().await.remove(&peer); 130 return; 131 } 132 if network.inner.peers.lock().await.is_banned(&peer) { 133 sleep(MAX_RECONNECT_DELAY).await; 134 continue; 135 } 136 super::P2pMetricsCounters::inc(&network.inner.metrics.outbound_connect_attempts); 137 let stream = match timeout(CONNECT_TIMEOUT, TcpStream::connect(&peer)).await { 138 Ok(Ok(stream)) => { 139 super::P2pMetricsCounters::inc(&network.inner.metrics.outbound_connect_successes); 140 stream 141 } 142 Ok(Err(error)) => { 143 super::P2pMetricsCounters::inc(&network.inner.metrics.outbound_connect_failures); 144 network 145 .inner 146 .peers 147 .lock() 148 .await 149 .record_error(&peer, format!("connecting to peer {peer}: {error}")); 150 sleep(reconnect_delay).await; 151 reconnect_delay = next_reconnect_delay(reconnect_delay); 152 continue; 153 } 154 Err(_) => { 155 super::P2pMetricsCounters::inc(&network.inner.metrics.outbound_connect_failures); 156 network 157 .inner 158 .peers 159 .lock() 160 .await 161 .record_error(&peer, format!("connecting to peer {peer}: timeout")); 162 sleep(reconnect_delay).await; 163 reconnect_delay = next_reconnect_delay(reconnect_delay); 164 continue; 165 } 166 }; 167 168 reconnect_delay = INITIAL_RECONNECT_DELAY; 169 let remote_addr = stream.peer_addr().unwrap_or_else(|_| { 170 peer.parse() 171 .unwrap_or_else(|_| SocketAddr::from(([0, 0, 0, 0], 0))) 172 }); 173 super::P2pMetricsCounters::inc(&network.inner.metrics.outbound_sessions_started); 174 let result = session_loop( 175 network.clone(), 176 stream, 177 remote_addr, 178 Some(peer.clone()), 179 receiver, 180 &mut shutdown, 181 None, 182 ) 183 .await; 184 if *shutdown.borrow() { 185 network.inner.sessions.lock().await.remove(&peer); 186 return; 187 } 188 match result { 189 Ok(()) => { 190 super::P2pMetricsCounters::inc(&network.inner.metrics.sessions_closed); 191 } 192 Err(error) if super::is_quiet_disconnect(&error) => { 193 super::P2pMetricsCounters::inc(&network.inner.metrics.quiet_disconnects); 194 } 195 Err(error) => { 196 super::P2pMetricsCounters::inc(&network.inner.metrics.session_failures); 197 let message = format!("{error:#}"); 198 super::P2pMetricsCounters::set_last( 199 &network.inner.metrics.last_session_failure, 200 format!("{peer}: {message}"), 201 ); 202 network 203 .inner 204 .peers 205 .lock() 206 .await 207 .record_error(&peer, message.clone()); 208 if debug_logging_enabled() { 209 eprintln!("p2p session with {peer} failed: {message}"); 210 } 211 } 212 } 213 214 let (sender, next_receiver) = mpsc::channel(PEER_QUEUE_SIZE); 215 let (next_shutdown, next_shutdown_receiver) = watch::channel(false); 216 let queue_bytes = std::sync::Arc::new(tokio::sync::Semaphore::new(PEER_QUEUE_BYTES)); 217 receiver = next_receiver; 218 shutdown = next_shutdown_receiver; 219 if !peer_is_connectable(&network, &peer).await { 220 network.inner.sessions.lock().await.remove(&peer); 221 return; 222 } 223 network.inner.sessions.lock().await.insert( 224 peer.clone(), 225 GossipSession { 226 peer: peer.clone(), 227 sender, 228 shutdown: next_shutdown, 229 queue_bytes, 230 }, 231 ); 232 sleep(reconnect_delay).await; 233 reconnect_delay = next_reconnect_delay(reconnect_delay); 234 } 235 } 236 237 async fn session_loop( 238 network: GossipNetwork, 239 stream: TcpStream, 240 remote_addr: SocketAddr, 241 stable_peer: Option<String>, 242 mut outbound: mpsc::Receiver<OutboundBatch>, 243 shutdown: &mut watch::Receiver<bool>, 244 inbound_registration: Option<InboundRegistration>, 245 ) -> Result<()> { 246 let (reader, mut writer) = stream.into_split(); 247 let connection_label = stable_peer 248 .as_ref() 249 .map(|peer| format!("outbound {peer}")) 250 .unwrap_or_else(|| format!("inbound {remote_addr}")); 251 let advertised_addr = network.advertised_addr().await; 252 let hello = network.inner.node.lock().await.hello( 253 advertised_addr.map(|addr| addr.to_string()), 254 Some(network.inner.node_id.clone()), 255 ); 256 write_envelope(&mut writer, &hello).await?; 257 let mut reader = super::LimitedLineReader::new(reader); 258 let mut sync_tick = interval_at( 259 Instant::now() + SESSION_SYNC_INTERVAL, 260 SESSION_SYNC_INTERVAL, 261 ); 262 let mut peer_exchange_tick = interval_at( 263 Instant::now() + PEER_EXCHANGE_INTERVAL, 264 PEER_EXCHANGE_INTERVAL, 265 ); 266 let mut outbound_closed = false; 267 let mut peer_status: Option<PeerStatus> = None; 268 let mut catchup_requested_at: Option<Instant> = None; 269 let is_outbound_session = stable_peer.is_some(); 270 let mut known_peer = stable_peer; 271 let mut handshake_complete = false; 272 let mut shutdown_closed = false; 273 274 if !is_outbound_session { 275 let envelope = timeout( 276 HANDSHAKE_TIMEOUT, 277 read_session_envelope(&network, &connection_label, &mut reader), 278 ) 279 .await 280 .context("inbound p2p handshake timed out")?? 281 .context("inbound peer closed before sending Hello")?; 282 let hello = match envelope { 283 GossipEnvelope::Hello(hello) => hello, 284 challenge @ GossipEnvelope::PeerVerificationChallenge { .. } => { 285 respond_to_peer_verification_challenge(&network, &mut writer, &challenge).await?; 286 return Ok(()); 287 } 288 _ => bail!("inbound peer sent data before Hello"), 289 }; 290 let status = process_hello_with_verification( 291 &network, 292 &mut writer, 293 &mut reader, 294 &connection_label, 295 remote_addr, 296 &mut known_peer, 297 hello, 298 ) 299 .await?; 300 if status.reject_session { 301 return Ok(()); 302 } 303 peer_status = Some(status); 304 handshake_complete = true; 305 if let Some(registration) = inbound_registration { 306 let peer = known_peer 307 .clone() 308 .unwrap_or_else(|| remote_addr.to_string()); 309 network.inner.sessions.lock().await.insert( 310 registration.key, 311 GossipSession { 312 peer, 313 sender: registration.sender, 314 shutdown: registration.shutdown, 315 queue_bytes: registration.queue_bytes, 316 }, 317 ); 318 } 319 maybe_start_catchup( 320 &network, 321 &mut writer, 322 peer_status.as_ref().unwrap(), 323 &mut catchup_requested_at, 324 ) 325 .await?; 326 write_peer_exchange(&network, &mut writer, &known_peer).await?; 327 } 328 329 if is_outbound_session { 330 if let Ok(Ok(Some(envelope))) = timeout( 331 HANDSHAKE_TIMEOUT, 332 read_session_envelope(&network, &connection_label, &mut reader), 333 ) 334 .await 335 { 336 if let GossipEnvelope::Hello(hello) = envelope { 337 let status = process_hello_with_verification( 338 &network, 339 &mut writer, 340 &mut reader, 341 &connection_label, 342 remote_addr, 343 &mut known_peer, 344 hello, 345 ) 346 .await?; 347 if status.reject_session { 348 return Ok(()); 349 } 350 peer_status = Some(status); 351 handshake_complete = true; 352 if is_outbound_session && known_peer.is_none() { 353 return Ok(()); 354 } 355 maybe_start_catchup( 356 &network, 357 &mut writer, 358 peer_status.as_ref().unwrap(), 359 &mut catchup_requested_at, 360 ) 361 .await?; 362 write_peer_exchange(&network, &mut writer, &known_peer).await?; 363 } else if let GossipEnvelope::PeerStatus { 364 height, 365 tip_hash, 366 time_ms, 367 } = envelope 368 { 369 let status = PeerStatus::from_envelope(height, tip_hash, time_ms) 370 .with_capabilities( 371 peer_status 372 .as_ref() 373 .map(|status| status.capabilities.clone()) 374 .unwrap_or_default(), 375 ); 376 record_peer_status(&network, &known_peer, remote_addr, &status).await; 377 peer_status = Some(status); 378 handshake_complete = true; 379 maybe_start_catchup( 380 &network, 381 &mut writer, 382 peer_status.as_ref().unwrap(), 383 &mut catchup_requested_at, 384 ) 385 .await?; 386 write_peer_exchange(&network, &mut writer, &known_peer).await?; 387 } else if respond_to_peer_verification_challenge(&network, &mut writer, &envelope) 388 .await? 389 { 390 if known_peer.is_none() { 391 return Ok(()); 392 } 393 } else if !maybe_request_inventory( 394 &network, 395 &mut writer, 396 &envelope, 397 &mut catchup_requested_at, 398 ) 399 .await? 400 { 401 let requested_chain_data = process_envelope( 402 &network, 403 &mut writer, 404 remote_addr, 405 &mut known_peer, 406 envelope, 407 ) 408 .await?; 409 if requested_chain_data { 410 catchup_requested_at = Some(Instant::now()); 411 } 412 if is_outbound_session && known_peer.is_none() { 413 return Ok(()); 414 } 415 } 416 } 417 } 418 419 loop { 420 tokio::select! { 421 result = shutdown.changed(), if !shutdown_closed => { 422 match result { 423 Ok(()) if *shutdown.borrow() => return Ok(()), 424 Ok(()) => {} 425 Err(_) if !is_outbound_session => return Ok(()), 426 Err(_) => shutdown_closed = true, 427 } 428 } 429 maybe_batch = outbound.recv(), if !outbound_closed => { 430 match maybe_batch { 431 Some(batch) => { 432 let payload = super::envelopes_for_peer( 433 Some(&network.inner.node), 434 peer_status.clone(), 435 &batch.envelopes, 436 ).await; 437 write_payload(&mut writer, &payload).await?; 438 if let Some(peer) = &known_peer { 439 network.inner.peers.lock().await.record_sent(peer, payload.len() as u64); 440 } 441 } 442 None => outbound_closed = true, 443 } 444 } 445 _ = sync_tick.tick() => { 446 let status = network.inner.node.lock().await.peer_status(); 447 write_envelope(&mut writer, &status).await?; 448 if catchup_requested_at.is_some_and(|started| started.elapsed() >= CATCHUP_REQUEST_TIMEOUT) { 449 catchup_requested_at = None; 450 } 451 if let Some(status) = peer_status.as_ref() { 452 maybe_start_catchup(&network, &mut writer, status, &mut catchup_requested_at).await?; 453 } 454 } 455 _ = peer_exchange_tick.tick() => { 456 write_peer_exchange(&network, &mut writer, &known_peer).await?; 457 } 458 envelope = read_session_envelope(&network, &connection_label, &mut reader) => { 459 let Some(envelope) = envelope? else { 460 return Ok(()); 461 }; 462 if let GossipEnvelope::Hello(hello) = envelope { 463 if handshake_complete { 464 bail!("peer sent duplicate Hello"); 465 } 466 let status = process_hello_with_verification( 467 &network, 468 &mut writer, 469 &mut reader, 470 &connection_label, 471 remote_addr, 472 &mut known_peer, 473 hello, 474 ) 475 .await?; 476 if status.reject_session { 477 return Ok(()); 478 } 479 peer_status = Some(status); 480 handshake_complete = true; 481 if is_outbound_session && known_peer.is_none() { 482 return Ok(()); 483 } 484 maybe_start_catchup( 485 &network, 486 &mut writer, 487 peer_status.as_ref().unwrap(), 488 &mut catchup_requested_at, 489 ).await?; 490 write_peer_exchange(&network, &mut writer, &known_peer).await?; 491 continue; 492 } 493 if let GossipEnvelope::PeerStatus { 494 height, 495 tip_hash, 496 time_ms, 497 } = &envelope 498 { 499 let status = PeerStatus::from_envelope(*height, tip_hash.clone(), *time_ms) 500 .with_capabilities( 501 peer_status 502 .as_ref() 503 .map(|status| status.capabilities.clone()) 504 .unwrap_or_default(), 505 ); 506 record_peer_status(&network, &known_peer, remote_addr, &status).await; 507 peer_status = Some(status); 508 maybe_start_catchup( 509 &network, 510 &mut writer, 511 peer_status.as_ref().unwrap(), 512 &mut catchup_requested_at, 513 ).await?; 514 write_peer_exchange(&network, &mut writer, &known_peer).await?; 515 continue; 516 } 517 518 if respond_to_peer_verification_challenge(&network, &mut writer, &envelope).await? { 519 if known_peer.is_none() { 520 return Ok(()); 521 } 522 continue; 523 } 524 if maybe_request_inventory( 525 &network, 526 &mut writer, 527 &envelope, 528 &mut catchup_requested_at, 529 ) 530 .await? 531 { 532 continue; 533 } 534 let continue_catchup = matches!( 535 &envelope, 536 GossipEnvelope::Block(_) | GossipEnvelope::Blocks { .. } | GossipEnvelope::ChainBootstrap(_) 537 ); 538 let completes_catchup_request = matches!( 539 &envelope, 540 GossipEnvelope::Blocks { .. } | GossipEnvelope::ChainBootstrap(_) 541 ); 542 let (height_before, response_already_applied) = if continue_catchup { 543 let node = network.inner.node.lock().await; 544 let status = node.ledger().status(); 545 let already_applied = match &envelope { 546 GossipEnvelope::Blocks { blocks } => { 547 !blocks.is_empty() && node.ledger().contains_block_sequence(blocks) 548 } 549 _ => false, 550 }; 551 (Some((status.height, status.tip_hash)), already_applied) 552 } else { 553 (None, false) 554 }; 555 let requested_chain_data = process_envelope( 556 &network, 557 &mut writer, 558 remote_addr, 559 &mut known_peer, 560 envelope, 561 ).await?; 562 let chain_changed = if let Some(height_before) = height_before { 563 let status = network.inner.node.lock().await.ledger().status(); 564 status.height != height_before.0 || status.tip_hash != height_before.1 565 } else { 566 false 567 }; 568 let response_satisfied = chain_changed || response_already_applied; 569 update_catchup_request_state( 570 &mut catchup_requested_at, 571 completes_catchup_request, 572 response_satisfied, 573 requested_chain_data, 574 ); 575 if continue_catchup && response_satisfied && !requested_chain_data { 576 if let Some(status) = peer_status.as_ref() { 577 maybe_start_catchup( 578 &network, 579 &mut writer, 580 status, 581 &mut catchup_requested_at, 582 ).await?; 583 } 584 } 585 if session_peer_is_banned(&network, &known_peer, remote_addr).await { 586 return Ok(()); 587 } 588 if is_outbound_session && known_peer.is_none() { 589 return Ok(()); 590 } 591 } 592 } 593 } 594 } 595 596 async fn maybe_start_catchup( 597 network: &GossipNetwork, 598 writer: &mut tokio::net::tcp::OwnedWriteHalf, 599 peer_status: &PeerStatus, 600 requested_at: &mut Option<Instant>, 601 ) -> Result<()> { 602 if requested_at.is_none() && super::maybe_request_catchup(network, writer, peer_status).await? { 603 *requested_at = Some(Instant::now()); 604 } 605 Ok(()) 606 } 607 608 fn update_catchup_request_state( 609 requested_at: &mut Option<Instant>, 610 completes_request: bool, 611 response_satisfied: bool, 612 requested_chain_data: bool, 613 ) { 614 if requested_chain_data || (completes_request && !response_satisfied) { 615 *requested_at = Some(Instant::now()); 616 } else if completes_request { 617 *requested_at = None; 618 } 619 } 620 621 async fn maybe_request_inventory( 622 network: &GossipNetwork, 623 writer: &mut tokio::net::tcp::OwnedWriteHalf, 624 envelope: &GossipEnvelope, 625 requested_at: &mut Option<Instant>, 626 ) -> Result<bool> { 627 let GossipEnvelope::Inventory { blocks } = envelope else { 628 return Ok(false); 629 }; 630 if requested_at.is_some() { 631 return Ok(true); 632 } 633 let request = network 634 .inner 635 .node 636 .lock() 637 .await 638 .missing_inventory_request(blocks); 639 if let Some(request) = request { 640 write_envelope(writer, &request).await?; 641 *requested_at = Some(Instant::now()); 642 } 643 Ok(true) 644 } 645 646 async fn peer_is_connectable(network: &GossipNetwork, peer: &str) -> bool { 647 network.inner.peers.lock().await.is_connectable_peer(peer) 648 } 649 650 async fn session_peer_is_banned( 651 network: &GossipNetwork, 652 known_peer: &Option<String>, 653 remote_addr: SocketAddr, 654 ) -> bool { 655 let peer = known_peer 656 .as_deref() 657 .map(str::to_owned) 658 .unwrap_or_else(|| remote_addr.to_string()); 659 network.inner.peers.lock().await.is_banned(&peer) 660 } 661 662 pub(super) fn next_reconnect_delay(current: Duration) -> Duration { 663 next_reconnect_delay_with_max(current, MAX_RECONNECT_DELAY) 664 } 665 666 #[cfg(test)] 667 mod tests { 668 use std::time::Duration; 669 670 use tokio::time::Instant; 671 672 use super::super::{INITIAL_RECONNECT_DELAY, MAX_RECONNECT_DELAY}; 673 use super::{next_reconnect_delay, update_catchup_request_state}; 674 675 #[test] 676 fn reconnect_backoff_is_capped() { 677 assert_eq!( 678 next_reconnect_delay(INITIAL_RECONNECT_DELAY), 679 Duration::from_secs(2) 680 ); 681 assert_eq!( 682 next_reconnect_delay(MAX_RECONNECT_DELAY), 683 MAX_RECONNECT_DELAY 684 ); 685 } 686 687 #[test] 688 fn already_applied_response_clears_in_flight_request() { 689 let mut requested_at = Some(Instant::now()); 690 691 update_catchup_request_state(&mut requested_at, true, true, false); 692 693 assert!(requested_at.is_none()); 694 } 695 696 #[test] 697 fn fork_recovery_request_is_marked_in_flight() { 698 let mut requested_at = None; 699 700 update_catchup_request_state(&mut requested_at, false, false, true); 701 702 assert!(requested_at.is_some()); 703 } 704 }