session.rs (16010B)
1 use std::{net::SocketAddr, time::Duration}; 2 3 use anyhow::Result; 4 use tokio::{ 5 net::{TcpListener, TcpStream}, 6 sync::mpsc, 7 time::{Instant, interval, interval_at, sleep, timeout}, 8 }; 9 10 use crate::app::{GossipEnvelope, debug_logging_enabled}; 11 12 use super::{ 13 CONNECT_TIMEOUT, GossipNetwork, HANDSHAKE_TIMEOUT, INITIAL_RECONNECT_DELAY, 14 MAX_RECONNECT_DELAY, PEER_EXCHANGE_INTERVAL, PEER_QUEUE_SIZE, PeerStatus, 15 SESSION_SYNC_INTERVAL, is_self_peer_address_for, next_reconnect_delay_with_max, 16 process_envelope, process_hello_with_verification, push_catchup_to_peer, read_session_envelope, 17 record_peer_status, respond_to_peer_verification_challenge, write_envelope, write_payload, 18 write_peer_exchange, 19 }; 20 21 type OutboundBatch = Vec<GossipEnvelope>; 22 23 pub(super) async fn accept_loop(network: GossipNetwork, listener: TcpListener) { 24 loop { 25 match listener.accept().await { 26 Ok((stream, remote_addr)) => { 27 let network = network.clone(); 28 let permit = match network.try_acquire_inbound_session(remote_addr.ip()) { 29 Ok(permit) => permit, 30 Err(rejection) => { 31 super::P2pMetricsCounters::inc( 32 &network.inner.metrics.inbound_sessions_rejected, 33 ); 34 super::P2pMetricsCounters::set_last( 35 &network.inner.metrics.last_session_failure, 36 format!("{remote_addr}: {}", rejection.label()), 37 ); 38 if debug_logging_enabled() { 39 eprintln!( 40 "p2p inbound connection from {remote_addr} rejected: {}", 41 rejection.label() 42 ); 43 } 44 drop(stream); 45 continue; 46 } 47 }; 48 super::P2pMetricsCounters::inc(&network.inner.metrics.inbound_sessions_started); 49 tokio::spawn(async move { 50 let _permit = permit; 51 let result = session_loop( 52 network.clone(), 53 stream, 54 remote_addr, 55 None, 56 mpsc::channel(1).1, 57 ) 58 .await; 59 match result { 60 Ok(()) => { 61 super::P2pMetricsCounters::inc(&network.inner.metrics.sessions_closed); 62 } 63 Err(error) if super::is_quiet_disconnect(&error) => { 64 super::P2pMetricsCounters::inc( 65 &network.inner.metrics.quiet_disconnects, 66 ); 67 } 68 Err(error) => { 69 super::P2pMetricsCounters::inc(&network.inner.metrics.session_failures); 70 super::P2pMetricsCounters::set_last( 71 &network.inner.metrics.last_session_failure, 72 format!("{remote_addr}: {error:#}"), 73 ); 74 if debug_logging_enabled() { 75 eprintln!( 76 "p2p inbound connection from {remote_addr} failed: {error:#}" 77 ); 78 } 79 } 80 } 81 }); 82 } 83 Err(error) if debug_logging_enabled() => eprintln!("p2p accept failed: {error:#}"), 84 Err(_) => {} 85 } 86 } 87 } 88 89 pub(super) async fn outbound_supervisor(network: GossipNetwork) { 90 let mut tick = interval(Duration::from_secs(2)); 91 loop { 92 tick.tick().await; 93 network.ensure_outbound_sessions().await; 94 } 95 } 96 97 pub(super) async fn outbound_session( 98 network: GossipNetwork, 99 peer: String, 100 mut receiver: mpsc::Receiver<OutboundBatch>, 101 ) { 102 let mut reconnect_delay = INITIAL_RECONNECT_DELAY; 103 loop { 104 let self_filter_addr = network.self_filter_addr().await; 105 if !peer_is_connectable(&network, &peer).await 106 || is_self_peer_address_for(&peer, network.inner.listen_addr, self_filter_addr) 107 { 108 network.inner.sessions.lock().await.remove(&peer); 109 return; 110 } 111 if network.inner.peers.lock().await.is_banned(&peer) { 112 sleep(MAX_RECONNECT_DELAY).await; 113 continue; 114 } 115 super::P2pMetricsCounters::inc(&network.inner.metrics.outbound_connect_attempts); 116 let stream = match timeout(CONNECT_TIMEOUT, TcpStream::connect(&peer)).await { 117 Ok(Ok(stream)) => { 118 super::P2pMetricsCounters::inc(&network.inner.metrics.outbound_connect_successes); 119 stream 120 } 121 Ok(Err(error)) => { 122 super::P2pMetricsCounters::inc(&network.inner.metrics.outbound_connect_failures); 123 network 124 .inner 125 .peers 126 .lock() 127 .await 128 .record_error(&peer, format!("connecting to peer {peer}: {error}")); 129 sleep(reconnect_delay).await; 130 reconnect_delay = next_reconnect_delay(reconnect_delay); 131 continue; 132 } 133 Err(_) => { 134 super::P2pMetricsCounters::inc(&network.inner.metrics.outbound_connect_failures); 135 network 136 .inner 137 .peers 138 .lock() 139 .await 140 .record_error(&peer, format!("connecting to peer {peer}: timeout")); 141 sleep(reconnect_delay).await; 142 reconnect_delay = next_reconnect_delay(reconnect_delay); 143 continue; 144 } 145 }; 146 147 reconnect_delay = INITIAL_RECONNECT_DELAY; 148 let remote_addr = stream.peer_addr().unwrap_or_else(|_| { 149 peer.parse() 150 .unwrap_or_else(|_| SocketAddr::from(([0, 0, 0, 0], 0))) 151 }); 152 super::P2pMetricsCounters::inc(&network.inner.metrics.outbound_sessions_started); 153 let result = session_loop( 154 network.clone(), 155 stream, 156 remote_addr, 157 Some(peer.clone()), 158 receiver, 159 ) 160 .await; 161 match result { 162 Ok(()) => { 163 super::P2pMetricsCounters::inc(&network.inner.metrics.sessions_closed); 164 } 165 Err(error) if super::is_quiet_disconnect(&error) => { 166 super::P2pMetricsCounters::inc(&network.inner.metrics.quiet_disconnects); 167 } 168 Err(error) => { 169 super::P2pMetricsCounters::inc(&network.inner.metrics.session_failures); 170 let message = format!("{error:#}"); 171 super::P2pMetricsCounters::set_last( 172 &network.inner.metrics.last_session_failure, 173 format!("{peer}: {message}"), 174 ); 175 network 176 .inner 177 .peers 178 .lock() 179 .await 180 .record_error(&peer, message.clone()); 181 if debug_logging_enabled() { 182 eprintln!("p2p session with {peer} failed: {message}"); 183 } 184 } 185 } 186 187 let (sender, next_receiver) = mpsc::channel(PEER_QUEUE_SIZE); 188 receiver = next_receiver; 189 if !peer_is_connectable(&network, &peer).await { 190 network.inner.sessions.lock().await.remove(&peer); 191 return; 192 } 193 network 194 .inner 195 .sessions 196 .lock() 197 .await 198 .insert(peer.clone(), sender); 199 sleep(reconnect_delay).await; 200 reconnect_delay = next_reconnect_delay(reconnect_delay); 201 } 202 } 203 204 async fn session_loop( 205 network: GossipNetwork, 206 stream: TcpStream, 207 remote_addr: SocketAddr, 208 stable_peer: Option<String>, 209 mut outbound: mpsc::Receiver<OutboundBatch>, 210 ) -> Result<()> { 211 let (reader, mut writer) = stream.into_split(); 212 let connection_label = stable_peer 213 .as_ref() 214 .map(|peer| format!("outbound {peer}")) 215 .unwrap_or_else(|| format!("inbound {remote_addr}")); 216 let advertised_addr = network.advertised_addr().await; 217 let hello = network.inner.node.lock().await.hello( 218 advertised_addr.map(|addr| addr.to_string()), 219 Some(network.inner.node_id.clone()), 220 ); 221 write_envelope(&mut writer, &hello).await?; 222 let mut reader = super::LimitedLineReader::new(reader); 223 let mut sync_tick = interval_at( 224 Instant::now() + SESSION_SYNC_INTERVAL, 225 SESSION_SYNC_INTERVAL, 226 ); 227 let mut peer_exchange_tick = interval_at( 228 Instant::now() + PEER_EXCHANGE_INTERVAL, 229 PEER_EXCHANGE_INTERVAL, 230 ); 231 let mut outbound_closed = false; 232 let mut peer_status: Option<PeerStatus> = None; 233 let is_outbound_session = stable_peer.is_some(); 234 let mut known_peer = stable_peer; 235 236 if known_peer.is_some() { 237 if let Ok(Ok(Some(envelope))) = timeout( 238 HANDSHAKE_TIMEOUT, 239 read_session_envelope(&network, &connection_label, &mut reader), 240 ) 241 .await 242 { 243 if let GossipEnvelope::Hello(hello) = envelope { 244 peer_status = Some( 245 process_hello_with_verification( 246 &network, 247 &mut writer, 248 &mut reader, 249 &connection_label, 250 remote_addr, 251 &mut known_peer, 252 hello, 253 ) 254 .await?, 255 ); 256 if is_outbound_session && known_peer.is_none() { 257 return Ok(()); 258 } 259 super::maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()) 260 .await?; 261 write_peer_exchange(&network, &mut writer, &known_peer).await?; 262 } else if let GossipEnvelope::PeerStatus { 263 height, 264 tip_hash, 265 time_ms, 266 } = envelope 267 { 268 let status = PeerStatus::from_envelope(height, tip_hash, time_ms); 269 record_peer_status(&network, &known_peer, remote_addr, &status).await; 270 peer_status = Some(status); 271 super::maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()) 272 .await?; 273 write_peer_exchange(&network, &mut writer, &known_peer).await?; 274 } else if respond_to_peer_verification_challenge(&network, &mut writer, &envelope) 275 .await? 276 { 277 if known_peer.is_none() { 278 return Ok(()); 279 } 280 } else { 281 process_envelope( 282 &network, 283 &mut writer, 284 remote_addr, 285 &mut known_peer, 286 envelope, 287 ) 288 .await?; 289 if is_outbound_session && known_peer.is_none() { 290 return Ok(()); 291 } 292 } 293 } 294 } 295 296 loop { 297 tokio::select! { 298 maybe_batch = outbound.recv(), if !outbound_closed => { 299 match maybe_batch { 300 Some(batch) => { 301 let payload = super::envelopes_for_peer( 302 Some(&network.inner.node), 303 peer_status.clone(), 304 &batch, 305 ).await; 306 write_payload(&mut writer, &payload).await?; 307 if let Some(peer) = &known_peer { 308 network.inner.peers.lock().await.record_sent(peer, payload.len() as u64); 309 } 310 } 311 None => outbound_closed = true, 312 } 313 } 314 _ = sync_tick.tick() => { 315 let status = network.inner.node.lock().await.peer_status(); 316 write_envelope(&mut writer, &status).await?; 317 if let Some(status) = peer_status.as_mut() { 318 if let Some(updated_status) = push_catchup_to_peer(&network, &mut writer, status).await? { 319 *status = updated_status; 320 } 321 } 322 } 323 _ = peer_exchange_tick.tick() => { 324 write_peer_exchange(&network, &mut writer, &known_peer).await?; 325 } 326 envelope = read_session_envelope(&network, &connection_label, &mut reader) => { 327 let Some(envelope) = envelope? else { 328 return Ok(()); 329 }; 330 if let GossipEnvelope::Hello(hello) = envelope { 331 peer_status = Some( 332 process_hello_with_verification( 333 &network, 334 &mut writer, 335 &mut reader, 336 &connection_label, 337 remote_addr, 338 &mut known_peer, 339 hello, 340 ) 341 .await?, 342 ); 343 if is_outbound_session && known_peer.is_none() { 344 return Ok(()); 345 } 346 super::maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()).await?; 347 write_peer_exchange(&network, &mut writer, &known_peer).await?; 348 continue; 349 } 350 if let GossipEnvelope::PeerStatus { 351 height, 352 tip_hash, 353 time_ms, 354 } = &envelope 355 { 356 let status = PeerStatus::from_envelope(*height, tip_hash.clone(), *time_ms); 357 record_peer_status(&network, &known_peer, remote_addr, &status).await; 358 peer_status = Some(status); 359 super::maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()).await?; 360 write_peer_exchange(&network, &mut writer, &known_peer).await?; 361 continue; 362 } 363 364 if respond_to_peer_verification_challenge(&network, &mut writer, &envelope).await? { 365 if known_peer.is_none() { 366 return Ok(()); 367 } 368 continue; 369 } 370 process_envelope( 371 &network, 372 &mut writer, 373 remote_addr, 374 &mut known_peer, 375 envelope, 376 ).await?; 377 if is_outbound_session && known_peer.is_none() { 378 return Ok(()); 379 } 380 } 381 } 382 } 383 } 384 385 async fn peer_is_connectable(network: &GossipNetwork, peer: &str) -> bool { 386 network.inner.peers.lock().await.is_connectable_peer(peer) 387 } 388 389 pub(super) fn next_reconnect_delay(current: Duration) -> Duration { 390 next_reconnect_delay_with_max(current, MAX_RECONNECT_DELAY) 391 } 392 393 #[cfg(test)] 394 mod tests { 395 use std::time::Duration; 396 397 use super::super::{INITIAL_RECONNECT_DELAY, MAX_RECONNECT_DELAY}; 398 use super::next_reconnect_delay; 399 400 #[test] 401 fn reconnect_backoff_is_capped() { 402 assert_eq!( 403 next_reconnect_delay(INITIAL_RECONNECT_DELAY), 404 Duration::from_secs(2) 405 ); 406 assert_eq!( 407 next_reconnect_delay(MAX_RECONNECT_DELAY), 408 MAX_RECONNECT_DELAY 409 ); 410 } 411 }