iuna

iuna - experimental devnet protocol
git clone https://iuna.jhx.app/git/iuna.git
Log | Files | Refs | README | LICENSE

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 }