iuna

iuna

iuna - experimental mainnet-candidate protocol
git clone https://getiuna.org/git/iuna.git
Log | Files | Refs | README | LICENSE

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 }