diff --git a/src/node/handlers/rekey.rs b/src/node/handlers/rekey.rs index 317945b8..f641f772 100644 --- a/src/node/handlers/rekey.rs +++ b/src/node/handlers/rekey.rs @@ -96,6 +96,19 @@ pub(in crate::node) fn link_silence_ms(node: &crate::config::NodeConfig) -> u64 .saturating_add(after_ms) } +/// The inbound idle deadline for stream transports: how long an accepted +/// connection may go without delivering a complete frame once it has +/// delivered one. +/// +/// Set to `link_silence_ms`, so a connection carrying a link the node would +/// keep is never dropped by it, whatever the heartbeat, link-dead, tick and +/// resend settings are; a connection carrying no live link is reclaimed. +pub(in crate::node) fn inbound_idle_timeout( + node: &crate::config::NodeConfig, +) -> std::time::Duration { + std::time::Duration::from_millis(link_silence_ms(node)) +} + /// How long a rekey responder holds a pending session its initiator has /// not adopted before retiring it. /// diff --git a/src/node/mod.rs b/src/node/mod.rs index 482eb57c..ee21d2a4 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -1125,10 +1125,15 @@ impl Node { // raising `node.limits.max_connections` actually raises the inbound // ceiling rather than being silently capped at the transport default. let node_max_connections = self.config().node.limits.max_connections; + // Inbound stream connections are dropped after this long without a + // complete frame. Derived from the node's own liveness timers so it + // cannot drop a connection whose link the node would keep. + let idle_timeout = handlers::rekey::inbound_idle_timeout(&self.config().node); for (name, tcp_config) in tcp_instances { let transport_id = self.allocate_transport_id(); let mut tcp = TcpTransport::new(transport_id, name, tcp_config, packet_tx.clone()); tcp.set_node_max_connections(node_max_connections); + tcp.set_inbound_idle_timeout(idle_timeout); transports.push(TransportHandle::Tcp(tcp)); } @@ -1143,7 +1148,8 @@ impl Node { for (name, tor_config) in tor_instances { let transport_id = self.allocate_transport_id(); - let tor = TorTransport::new(transport_id, name, tor_config, packet_tx.clone()); + let mut tor = TorTransport::new(transport_id, name, tor_config, packet_tx.clone()); + tor.set_inbound_idle_timeout(idle_timeout); transports.push(TransportHandle::Tor(tor)); } diff --git a/src/node/tests/tcp.rs b/src/node/tests/tcp.rs index 669106ad..25a37a3d 100644 --- a/src/node/tests/tcp.rs +++ b/src/node/tests/tcp.rs @@ -12,7 +12,8 @@ use crate::transport::{ ConnectionState, TransportAddr, TransportHandle, TransportId, packet_channel, }; use spanning_tree::{ - TestNode, cleanup_nodes, drain_all_packets, initiate_handshake, verify_tree_convergence, + TestNode, cleanup_nodes, drain_all_packets, initiate_handshake, process_available_packets, + verify_tree_convergence, }; use std::time::Duration; @@ -29,6 +30,13 @@ async fn make_test_node_tcp() -> TestNode { /// so immutable fields (e.g. heartbeat/link-dead timeouts) are set before the /// `NodeContext` is built rather than poked afterward. async fn make_test_node_tcp_with(config: Config) -> TestNode { + make_test_node_tcp_idle(config, None).await +} + +/// Like `make_test_node_tcp_with`, and also sets the transport's inbound +/// idle deadline when `idle` is given, as `create_transports` does for a +/// node built from config. +async fn make_test_node_tcp_idle(config: Config, idle: Option) -> TestNode { let mut node = make_node_with(config); let transport_id = TransportId::new(1); @@ -40,6 +48,9 @@ async fn make_test_node_tcp_with(config: Config) -> TestNode { let (packet_tx, packet_rx) = packet_channel(256); let mut transport = TcpTransport::new(transport_id, None, config, packet_tx); + if let Some(d) = idle { + transport.set_inbound_idle_timeout(d); + } transport.start_async().await.unwrap(); let local_addr = transport @@ -218,6 +229,82 @@ async fn test_tcp_connection_loss_detection() { cleanup_nodes(&mut nodes).await; } +/// A TCP link whose only traffic is heartbeats outlives several inbound idle +/// deadlines when the deadline is the one the node derives from its own +/// liveness timers. +/// +/// This is the healthy-path check for the idle deadline: it uses the real +/// derivation (`inbound_idle_timeout`), not a hand-picked value, and +/// short timers so three deadlines fit in a test. Break-check: give the +/// listening node's transport an idle deadline below the heartbeat interval +/// and its inbound connection is dropped, so the inbound count falls to 0. +#[tokio::test] +async fn tcp_link_kept_alive_only_by_heartbeats_survives_several_idle_deadlines() { + use crate::node::handlers::rekey::inbound_idle_timeout; + + let mut config = Config::new(); + config.node.heartbeat_interval_secs = 1; + config.node.link_dead_timeout_secs = 3; + config.node.rate_limit.handshake_resend_interval_ms = 100; + config.node.rate_limit.handshake_max_resends = 1; + let idle = inbound_idle_timeout(&config.node); + // 100 ms of ladder, three 1 s ticks, 3 s link-dead. + assert_eq!(idle, Duration::from_millis(6100)); + let mut nodes = vec![ + make_test_node_tcp_idle(config.clone(), Some(idle)).await, + make_test_node_tcp_idle(config, Some(idle)).await, + ]; + + // Node 0 dials node 1, so node 1 holds the inbound connection. + initiate_handshake(&mut nodes, 0, 1).await; + assert!(drain_all_packets(&mut nodes, false).await > 0); + + let addr_0 = *nodes[0].node.node_addr(); + let addr_1 = *nodes[1].node.node_addr(); + let inbound = |nodes: &[TestNode]| match nodes[1].node.transports.get(&nodes[1].transport_id) { + Some(TransportHandle::Tcp(t)) => t.stats().pool_inbound_count(), + _ => panic!("node 1 should have its TCP transport"), + }; + assert_eq!( + inbound(&nodes), + 1, + "node 1 should hold one inbound connection" + ); + + // Heartbeats only, for three idle deadlines and a little more. + let start = tokio::time::Instant::now(); + let mut processed = 0; + while start.elapsed() < idle * 3 + Duration::from_millis(500) { + nodes[0].node.check_link_heartbeats().await; + nodes[1].node.check_link_heartbeats().await; + processed += process_available_packets(&mut nodes).await; + assert_eq!( + inbound(&nodes), + 1, + "the inbound connection was dropped after {:?} of heartbeat-only traffic", + start.elapsed() + ); + tokio::time::sleep(Duration::from_millis(250)).await; + } + + // Sibling control: heartbeats were actually exchanged, so the kept + // connection is not an artefact of nothing having run. + assert!( + processed >= 18, + "expected at least one heartbeat per second each way, processed {processed}" + ); + assert!( + nodes[0].node.get_peer(&addr_1).is_some(), + "node 0 lost node 1" + ); + assert!( + nodes[1].node.get_peer(&addr_0).is_some(), + "node 1 lost node 0" + ); + + cleanup_nodes(&mut nodes).await; +} + /// TCP reconnection after link death: connect-on-send re-establishes the link. /// /// After both peers detect a dead link and remove each other, a fresh diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index b20cc4e5..eb2d5e42 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -3698,6 +3698,110 @@ async fn a_node_without_ble_config_draws_no_ble_warning() { ); } +/// The inbound idle deadline is the node's link-silence bound: 64 s at stock +/// settings, moving one for one with the link-dead timeout, and never below +/// the link-dead timeout plus a tick, the shortest silence after which the +/// node itself would reap the link. +#[test] +fn the_inbound_idle_deadline_is_the_link_silence_bound_and_tracks_the_link_dead_timeout() { + use crate::node::handlers::rekey::inbound_idle_timeout; + use crate::transport::tcp::INBOUND_IDLE_TIMEOUT; + + // 31 s of msg1 ladder (1+2+4+8+16), three 1 s ticks, 30 s link-dead. + let stock = crate::config::NodeConfig::default(); + assert_eq!(inbound_idle_timeout(&stock), Duration::from_secs(64)); + // The transport default a TCP or Tor instance holds before the node sets + // it must be the same stock value, so the two cannot drift apart. + assert_eq!(inbound_idle_timeout(&stock), INBOUND_IDLE_TIMEOUT); + + let raised = crate::config::NodeConfig { + link_dead_timeout_secs: 300, + ..Default::default() + }; + assert_eq!( + inbound_idle_timeout(&raised), + inbound_idle_timeout(&stock) + Duration::from_secs(270) + ); + + for node in [ + stock, + raised, + crate::config::NodeConfig { + heartbeat_interval_secs: 90, + ..Default::default() + }, + crate::config::NodeConfig { + tick_interval_secs: 5, + link_dead_timeout_secs: 10, + ..Default::default() + }, + ] { + let floor = Duration::from_secs( + node.link_dead_timeout_secs + .max(node.heartbeat_interval_secs) + + node.tick_interval_secs, + ); + assert!( + inbound_idle_timeout(&node) > floor, + "idle deadline {:?} must exceed {:?} for heartbeat {} s, link-dead {} s, tick {} s", + inbound_idle_timeout(&node), + floor, + node.heartbeat_interval_secs, + node.link_dead_timeout_secs, + node.tick_interval_secs, + ); + } +} + +/// `create_transports` hands every TCP and Tor instance the idle deadline +/// derived from the node's own liveness timers, not the transport default. +/// +/// Break-check: remove either `set_inbound_idle_timeout` call and that +/// transport keeps `INBOUND_IDLE_TIMEOUT`, which the non-default timers here +/// are chosen to differ from. +#[tokio::test] +async fn create_transports_sets_the_inbound_idle_deadline_from_the_node_liveness_timers() { + use crate::node::handlers::rekey::inbound_idle_timeout; + use crate::transport::tcp::INBOUND_IDLE_TIMEOUT; + + let mut config = crate::Config::new(); + config.node.control.enabled = false; + config.node.heartbeat_interval_secs = 5; + config.node.link_dead_timeout_secs = 120; + config.transports.tcp = + crate::config::TransportInstances::Single(crate::config::TcpConfig::default()); + config.transports.tor = + crate::config::TransportInstances::Single(crate::config::TorConfig::default()); + let expected = inbound_idle_timeout(&config.node); + // 31 s of ladder, three 1 s ticks, 120 s link-dead. + assert_eq!(expected, Duration::from_secs(154)); + assert_ne!(expected, INBOUND_IDLE_TIMEOUT); + + let mut node = make_node_with(config); + let (tx, _rx) = packet_channel(8); + let transports = node.create_transports(&tx).await; + + let mut seen = (0, 0); + for handle in &transports { + match handle { + crate::transport::TransportHandle::Tcp(t) => { + assert_eq!(t.inbound_idle_timeout(), expected, "TCP idle deadline"); + seen.0 += 1; + } + crate::transport::TransportHandle::Tor(t) => { + assert_eq!(t.inbound_idle_timeout(), expected, "Tor idle deadline"); + seen.1 += 1; + } + _ => {} + } + } + assert_eq!( + seen, + (1, 1), + "one TCP and one Tor transport should be built" + ); +} + #[cfg(all(ble_available, any(target_os = "android", test)))] mod test_radio { use crate::transport::ble::addr::BleAddr; diff --git a/src/transport/nym/mod.rs b/src/transport/nym/mod.rs index 243d09a9..11381df6 100644 --- a/src/transport/nym/mod.rs +++ b/src/transport/nym/mod.rs @@ -643,14 +643,14 @@ fn parse_target_addr(addr: &TransportAddr) -> Result( mtu: u16, stats: Arc, label: &'static str, - first_frame_timeout: Option, + deadline: Option, ready_rx: Option>, on_remove: impl Fn(&S, &M), ) { @@ -186,11 +187,13 @@ pub(crate) async fn proxied_receive_loop( if admitted { let mut first = true; loop { - let read = match first_frame_timeout { - // Bound the first read only. A silent remote otherwise holds its - // inbound slot for as long as it keeps the socket open. - Some(d) if first => { - match tokio::time::timeout(d, read_fmp_packet(&mut reader, mtu)).await { + let read = match deadline { + // Bound every read. A remote that goes silent, before or after + // its first frame, otherwise holds its inbound slot for as long + // as it keeps the socket open. + Some(d) => { + let limit = d.for_read(first); + match tokio::time::timeout(limit, read_fmp_packet(&mut reader, mtu)).await { Ok(result) => result, Err(_) => { // Not a recv error: `record_recv_error` means framing @@ -199,15 +202,16 @@ pub(crate) async fn proxied_receive_loop( debug!( transport_id = %transport_id, remote_addr = %remote_addr, - timeout_secs = d.as_secs_f64(), - "No complete frame within the first-frame deadline, dropping inbound {} connection", + deadline = InboundDeadline::phase(first), + timeout_secs = limit.as_secs_f64(), + "No complete frame within the inbound deadline, dropping inbound {} connection", label ); break; } } } - _ => read_fmp_packet(&mut reader, mtu).await, + None => read_fmp_packet(&mut reader, mtu).await, }; first = false; diff --git a/src/transport/tcp/mod.rs b/src/transport/tcp/mod.rs index e85940ff..268ca729 100644 --- a/src/transport/tcp/mod.rs +++ b/src/transport/tcp/mod.rs @@ -87,6 +87,9 @@ pub struct TcpTransport { /// Deadline from accept to the first complete inbound frame. Defaults to /// `INBOUND_FIRST_FRAME_TIMEOUT`; overridable only from tests. first_frame_timeout: Duration, + /// Longest wait for each later complete inbound frame. Defaults to + /// `INBOUND_IDLE_TIMEOUT`; the node sets it from its liveness timers. + idle_timeout: Duration, /// Transport statistics. stats: Arc, } @@ -111,6 +114,7 @@ impl TcpTransport { local_addr: None, node_max_connections: None, first_frame_timeout: INBOUND_FIRST_FRAME_TIMEOUT, + idle_timeout: INBOUND_IDLE_TIMEOUT, stats: Arc::new(TcpStats::new()), } } @@ -126,6 +130,22 @@ impl TcpTransport { self.first_frame_timeout = d; } + /// Set the inbound idle deadline: the longest an accepted connection may + /// go, after its first frame, without delivering another complete frame. + /// + /// The node derives it from its own link-liveness timers, so it never + /// drops a connection carrying a link the node would keep. Takes effect + /// at the next `start_async()`. + pub fn set_inbound_idle_timeout(&mut self, d: Duration) { + self.idle_timeout = d; + } + + /// The inbound idle deadline the accept loop will use. + #[cfg(test)] + pub(crate) fn inbound_idle_timeout(&self) -> Duration { + self.idle_timeout + } + /// Set the node-wide `node.limits.max_connections` value. /// /// Used as the inbound-cap fallback when this transport instance has no @@ -214,7 +234,10 @@ impl TcpTransport { keepalive_secs: self.config.keepalive_secs(), recv_buf: self.config.recv_buf_size(), send_buf: self.config.send_buf_size(), - first_frame_timeout: self.first_frame_timeout, + deadline: InboundDeadline { + first_frame: self.first_frame_timeout, + idle: self.idle_timeout, + }, }; let accept_task = tokio::spawn(async move { @@ -767,6 +790,45 @@ impl Transport for TcpTransport { /// TOML surface. pub(crate) const INBOUND_FIRST_FRAME_TIMEOUT: Duration = Duration::from_secs(30); +/// Default deadline for each complete inbound frame after the first. +/// +/// Without it, a remote that sends one well-formed frame and then goes +/// silent holds its inbound slot for as long as it keeps the socket open: +/// a frame that names no session is dropped by the node without closing +/// the transport. The node replaces this with the bound derived from its +/// own liveness timers (`link_silence_ms`); 64 s is that bound at stock +/// settings, kept here so a transport built outside the node still has a +/// deadline. +pub(crate) const INBOUND_IDLE_TIMEOUT: Duration = Duration::from_secs(64); + +/// Read deadlines for an inbound connection, which holds a capped pool +/// slot from accept. +/// +/// Each deadline covers one complete frame, not a byte: a remote that +/// drips a frame slower than the deadline is dropped. The idle deadline +/// re-arms on every frame, so a connection carrying a live link, which +/// receives at least a heartbeat per interval, is never dropped by it. +#[derive(Clone, Copy, Debug)] +pub(crate) struct InboundDeadline { + /// Deadline from accept to the first complete frame. + pub(crate) first_frame: Duration, + /// Deadline for each later complete frame, from the end of the last. + pub(crate) idle: Duration, +} + +impl InboundDeadline { + /// The deadline for the next read: `first` is true until a frame has + /// been received. + pub(crate) fn for_read(&self, first: bool) -> Duration { + if first { self.first_frame } else { self.idle } + } + + /// The log word for an expiry of the deadline `for_read(first)` gave. + pub(crate) fn phase(first: bool) -> &'static str { + if first { "first-frame" } else { "idle" } + } +} + /// Socket configuration parameters passed to the accept loop. struct AcceptConfig { mtu: u16, @@ -775,7 +837,7 @@ struct AcceptConfig { keepalive_secs: u64, recv_buf: usize, send_buf: usize, - first_frame_timeout: Duration, + deadline: InboundDeadline, } /// TCP accept loop — runs as a spawned task when bind_addr is configured. @@ -795,7 +857,7 @@ async fn accept_loop( keepalive_secs, recv_buf, send_buf, - first_frame_timeout, + deadline, } = cfg; debug!(transport_id = %transport_id, "TCP accept loop starting"); @@ -907,7 +969,7 @@ async fn accept_loop( conn_mtu, recv_stats, Direction::Inbound, - Some(first_frame_timeout), + Some(deadline), Some(ready_rx), ) .await; @@ -967,9 +1029,10 @@ async fn accept_loop( /// `pool_outbound` counter regardless of whether the matching pool /// entry survived to be removed. /// -/// `first_frame_timeout` bounds the wait for the *first* complete frame -/// only, and is `Some` for inbound connections (which hold a capped pool -/// slot from accept) and `None` for outbound ones. `ready_rx`, when +/// `deadline` bounds the wait for every complete frame: the first-frame +/// deadline until one arrives, the idle deadline for each one after. It is +/// `Some` for inbound connections (which hold a capped pool slot from +/// accept) and `None` for outbound ones. `ready_rx`, when /// present, is the accept loop's readiness barrier: the loop must not run /// its cleanup before the accept loop has inserted the pool entry. #[allow(clippy::too_many_arguments)] @@ -982,7 +1045,7 @@ async fn tcp_receive_loop( mtu: u16, stats: Arc, direction: Direction, - first_frame_timeout: Option, + deadline: Option, ready_rx: Option>, ) { let remote_addr = &key.remote; @@ -1003,11 +1066,13 @@ async fn tcp_receive_loop( if admitted { let mut first = true; loop { - let read = match first_frame_timeout { - // Bound the first read only. A silent remote otherwise holds - // its inbound slot for as long as it keeps the socket open. - Some(d) if first => { - match tokio::time::timeout(d, read_fmp_packet(&mut reader, mtu)).await { + let read = match deadline { + // Bound every read. A remote that goes silent, before or after + // its first frame, otherwise holds its inbound slot for as + // long as it keeps the socket open. + Some(d) => { + let limit = d.for_read(first); + match tokio::time::timeout(limit, read_fmp_packet(&mut reader, mtu)).await { Ok(result) => result, Err(_) => { // Not a recv error: `record_recv_error` means @@ -1016,14 +1081,15 @@ async fn tcp_receive_loop( debug!( transport_id = %transport_id, remote_addr = %remote_addr, - timeout_secs = d.as_secs_f64(), - "No complete frame within the first-frame deadline, dropping inbound connection" + deadline = InboundDeadline::phase(first), + timeout_secs = limit.as_secs_f64(), + "No complete frame within the inbound deadline, dropping inbound connection" ); break; } } } - _ => read_fmp_packet(&mut reader, mtu).await, + None => read_fmp_packet(&mut reader, mtu).await, }; first = false; @@ -1232,7 +1298,7 @@ fn read_mss_mtu(stream: &std::net::TcpStream, default_mtu: u16) -> u16 { #[cfg(test)] mod tests { use super::*; - use crate::transport::framing::build_msg1_frame; + use crate::transport::framing::{build_established_frame, build_msg1_frame}; use crate::transport::packet_channel; use tokio::time::{Duration, timeout}; @@ -2102,18 +2168,170 @@ mod tests { transport.stop_async().await.unwrap(); } - /// Regression guard, not evidence that the fix works. + /// The smallest frame the reader accepts as established: a 16-byte + /// header, no payload, a 16-byte tag. The node drops one naming no + /// session without closing the transport, so it is what a squatter + /// sends to get past the first-frame deadline. + fn squatter_frame() -> Vec { + let frame = build_established_frame(0); + assert_eq!(frame.len(), crate::proto::fmp::wire::ENCRYPTED_MIN_SIZE); + frame + } + + /// A remote that sends one well-formed frame and then goes silent must + /// lose its inbound slot at the idle deadline, without disconnecting. /// - /// The deadline is scoped to the first iteration, so an established - /// connection that then goes quiet cannot be dropped by it: this test - /// passes by construction under the current design. It is kept so that a - /// future general (every-read) idle deadline cannot silently start - /// reaping quiet links without a test going red. + /// The first-frame deadline is set far above the idle one, so the + /// release can only be the idle deadline's doing. Break-check: scope the + /// deadline back to the first read only and the count stays at 1. #[tokio::test] - async fn established_connection_survives_long_idle() { + async fn inbound_connection_that_goes_silent_after_one_established_frame_releases_its_slot() { + let (tx, mut rx) = packet_channel(100); + let mut transport = TcpTransport::new(TransportId::new(1), None, make_config(), tx); + transport.set_first_frame_timeout(Duration::from_secs(5)); + transport.set_inbound_idle_timeout(Duration::from_millis(300)); + transport.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + let mut squatter = TcpStream::connect(listen).await.unwrap(); + squatter.write_all(&squatter_frame()).await.unwrap(); + let packet = timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout waiting for the squatter's frame") + .expect("packet channel closed"); + assert_eq!(packet.data, squatter_frame()); + assert_eq!( + transport.stats().pool_inbound_count(), + 1, + "the connection should hold its slot once its first frame is in" + ); + + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 0, + Duration::from_secs(2) + ) + .await, + "a connection silent after its first frame should lose its slot at the idle deadline" + ); + assert!( + transport.pool.lock().await.is_empty(), + "the pool entry should go with the slot" + ); + + drop(squatter); + transport.stop_async().await.unwrap(); + } + + /// With the cap filled by a one-frame squatter, a genuine peer is refused + /// until the idle deadline frees the slot, and admitted afterwards. + /// + /// Break-check: without the idle deadline the squatter never releases, + /// so the genuine peer's frame is never delivered. + #[tokio::test] + async fn inbound_cap_filled_by_one_frame_squatters_admits_a_genuine_peer_after_the_idle_deadline() + { + let (tx, mut rx) = packet_channel(100); + let mut transport = TcpTransport::new(TransportId::new(1), None, capped_config(1), tx); + transport.set_first_frame_timeout(Duration::from_secs(5)); + transport.set_inbound_idle_timeout(Duration::from_secs(1)); + transport.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + let mut squatter = TcpStream::connect(listen).await.unwrap(); + squatter.write_all(&squatter_frame()).await.unwrap(); + timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout waiting for the squatter's frame") + .expect("packet channel closed"); + assert_eq!(transport.stats().pool_inbound_count(), 1); + + // While the cap is full a genuine peer is rejected outright. The + // 250 ms window ends well inside the squatter's 1 s idle deadline. + let mut early = TcpStream::connect(listen).await.unwrap(); + let _ = early.write_all(&build_msg1_frame()).await; + assert!( + timeout(Duration::from_millis(250), rx.recv()) + .await + .is_err(), + "a peer arriving while the cap is full must not be admitted" + ); + drop(early); + + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 0, + Duration::from_secs(3) + ) + .await, + "the idle deadline should free the slot the squatter took" + ); + + let mut genuine = TcpStream::connect(listen).await.unwrap(); + genuine.write_all(&build_msg1_frame()).await.unwrap(); + let packet = timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout waiting for the genuine peer's frame") + .expect("packet channel closed"); + assert_eq!(packet.data, build_msg1_frame()); + + drop(squatter); + drop(genuine); + transport.stop_async().await.unwrap(); + } + + /// The healthy path: a connection that delivers a frame more often than + /// the idle deadline keeps its slot across many deadlines, and every + /// frame is delivered. This is the shape of a link kept alive only by + /// heartbeats. + /// + /// Break-check: a deadline that does not re-arm on each frame (a single + /// deadline from accept) drops the connection after the first second. + #[tokio::test] + async fn inbound_connection_sending_a_frame_every_interval_below_the_idle_deadline_is_kept() { + let (tx, mut rx) = packet_channel(100); + let mut transport = TcpTransport::new(TransportId::new(1), None, make_config(), tx); + transport.set_first_frame_timeout(Duration::from_secs(1)); + transport.set_inbound_idle_timeout(Duration::from_secs(1)); + transport.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + let mut peer = TcpStream::connect(listen).await.unwrap(); + // Twelve frames 250 ms apart: 3 s, three idle deadlines, with 750 ms + // of slack between each frame and the deadline it re-arms. + for i in 0..12 { + peer.write_all(&squatter_frame()).await.unwrap(); + let packet = timeout(Duration::from_secs(2), rx.recv()) + .await + .unwrap_or_else(|_| panic!("timeout waiting for frame {i}")) + .expect("packet channel closed"); + assert_eq!(packet.data, squatter_frame()); + assert_eq!( + transport.stats().pool_inbound_count(), + 1, + "a connection delivering frames inside the idle deadline must keep its slot (frame {i})" + ); + tokio::time::sleep(Duration::from_millis(250)).await; + } + assert_eq!(transport.stats().pool_inbound_count(), 1); + assert!(!transport.pool.lock().await.is_empty()); + + drop(peer); + transport.stop_async().await.unwrap(); + } + + /// An established connection quiet for longer than the first-frame + /// deadline, but not the idle deadline, keeps its slot: once a frame is + /// in, the first-frame deadline no longer applies. + /// + /// Break-check: apply the first-frame deadline to every read and the + /// connection is dropped after 200 ms of quiet. + #[tokio::test] + async fn established_connection_quiet_past_the_first_frame_deadline_is_kept() { let (tx, mut rx) = packet_channel(100); let mut transport = TcpTransport::new(TransportId::new(1), None, make_config(), tx); transport.set_first_frame_timeout(Duration::from_millis(200)); + transport.set_inbound_idle_timeout(Duration::from_secs(5)); transport.start_async().await.unwrap(); let listen = transport.local_addr().unwrap(); @@ -2125,7 +2343,7 @@ mod tests { .expect("packet channel closed"); assert_eq!(packet.data, build_msg1_frame()); - // Four deadlines' worth of silence after the first frame. + // Four first-frame deadlines' worth of silence after the first frame. tokio::time::sleep(Duration::from_millis(800)).await; assert_eq!( @@ -2139,6 +2357,54 @@ mod tests { transport.stop_async().await.unwrap(); } + /// The idle deadline covers a complete frame, not its first bytes: a + /// frame after the first whose prefix arrives inside the deadline and + /// whose remainder arrives after it is not delivered, and the slot is + /// released. + /// + /// Break-check: bound only the prefix read and the body read waits + /// forever, so the late frame is delivered. + #[tokio::test] + async fn inbound_connection_dripping_a_frame_slower_than_the_idle_deadline_is_dropped() { + let (tx, mut rx) = packet_channel(100); + let mut transport = TcpTransport::new(TransportId::new(1), None, make_config(), tx); + transport.set_first_frame_timeout(Duration::from_secs(5)); + transport.set_inbound_idle_timeout(Duration::from_millis(300)); + transport.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + let mut peer = TcpStream::connect(listen).await.unwrap(); + peer.write_all(&squatter_frame()).await.unwrap(); + timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout waiting for the first frame") + .expect("packet channel closed"); + + let frame = squatter_frame(); + // Prefix inside the idle deadline, remainder well past it. + peer.write_all(&frame[..4]).await.unwrap(); + tokio::time::sleep(Duration::from_millis(600)).await; + let _ = peer.write_all(&frame[4..]).await; + + assert!( + timeout(Duration::from_millis(500), rx.recv()) + .await + .is_err(), + "a frame completing after the idle deadline must not be delivered" + ); + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 0, + Duration::from_secs(2) + ) + .await, + "the dripped connection should have released its slot" + ); + + drop(peer); + transport.stop_async().await.unwrap(); + } + /// A genuine peer that is slow to start, but finishes its first frame /// inside the deadline, is admitted. #[tokio::test] @@ -2250,7 +2516,10 @@ mod tests { 1400, stats.clone(), Direction::Inbound, - Some(Duration::from_millis(50)), + Some(InboundDeadline { + first_frame: Duration::from_millis(50), + idle: Duration::from_millis(50), + }), Some(ready_rx), ) .await; diff --git a/src/transport/tor/mod.rs b/src/transport/tor/mod.rs index 53fa65ae..6730f6b6 100644 --- a/src/transport/tor/mod.rs +++ b/src/transport/tor/mod.rs @@ -33,7 +33,7 @@ use crate::transport::socks5::{ ConnectingEntry, ConnectingPool, DialError, ProxiedConnection, ProxiedPool, Socks5Auth, Socks5Dialer, SocksTarget, poll_connecting, proxied_receive_loop, }; -use crate::transport::tcp::INBOUND_FIRST_FRAME_TIMEOUT; +use crate::transport::tcp::{INBOUND_FIRST_FRAME_TIMEOUT, INBOUND_IDLE_TIMEOUT, InboundDeadline}; use control::{ControlAuth, TorControlClient, TorMonitoringInfo}; use stats::TorStats; @@ -151,6 +151,10 @@ pub struct TorTransport { cached_monitoring: Arc>>, /// Background monitoring task handle. monitoring_task: Option>, + /// Longest wait for each complete inbound onion frame after the first. + /// Defaults to `INBOUND_IDLE_TIMEOUT`; the node sets it from its + /// liveness timers. + idle_timeout: Duration, } impl TorTransport { @@ -175,9 +179,27 @@ impl TorTransport { control_client: None, cached_monitoring: Arc::new(std::sync::RwLock::new(None)), monitoring_task: None, + idle_timeout: INBOUND_IDLE_TIMEOUT, } } + /// Set the inbound idle deadline for the onion listener: the longest an + /// accepted connection may go, after its first frame, without delivering + /// another complete frame. + /// + /// The node derives it from its own link-liveness timers, so it never + /// drops a connection carrying a link the node would keep. Takes effect + /// at the next `start_async()`. + pub fn set_inbound_idle_timeout(&mut self, d: Duration) { + self.idle_timeout = d; + } + + /// The inbound idle deadline the onion accept loop will use. + #[cfg(test)] + pub(crate) fn inbound_idle_timeout(&self) -> Duration { + self.idle_timeout + } + /// Get the instance name (if configured as a named instance). pub fn name(&self) -> Option<&str> { self.name.as_deref() @@ -359,6 +381,10 @@ impl TorTransport { let mtu = self.config.mtu(); let max_inbound = self.config.max_inbound_connections(); let stats = self.stats.clone(); + let deadline = InboundDeadline { + first_frame: INBOUND_FIRST_FRAME_TIMEOUT, + idle: self.idle_timeout, + }; let accept_handle = tokio::spawn(async move { tor_accept_loop( @@ -368,7 +394,7 @@ impl TorTransport { pool, mtu, max_inbound, - INBOUND_FIRST_FRAME_TIMEOUT, + deadline, stats, ) .await; @@ -1057,9 +1083,9 @@ impl Transport for TorTransport { /// below zero). `direction` is retained for the terminal "receive loop /// stopped" debug field the shared loop deliberately leaves to each transport. /// -/// `first_frame_timeout` is `Some` for an inbound connection, which holds a -/// capped pool slot from the moment it is accepted, and `None` for an -/// outbound one, which holds no such slot. `ready_rx`, when present, is the +/// `deadline` is `Some` for an inbound connection, which holds a capped pool +/// slot from the moment it is accepted, and `None` for an outbound one, +/// which holds no such slot. `ready_rx`, when present, is the /// accept loop's readiness barrier: the loop must not run its cleanup before /// the accept loop has inserted the pool entry and bumped its counter. #[allow(clippy::too_many_arguments)] @@ -1072,7 +1098,7 @@ async fn tor_receive_loop( mtu: u16, stats: Arc, direction: Direction, - first_frame_timeout: Option, + deadline: Option, ready_rx: Option>, ) { proxied_receive_loop( @@ -1084,7 +1110,7 @@ async fn tor_receive_loop( mtu, stats, "Tor", - first_frame_timeout, + deadline, ready_rx, |stats, meta| match meta { Direction::Inbound => stats.record_pool_inbound_removed(), @@ -1112,11 +1138,11 @@ async fn tor_receive_loop( /// socket options, split the stream, and spawn a per-connection /// receive task. /// -/// `first_frame_timeout` is the deadline from accept to the first complete -/// inbound frame, handed to each spawned receive loop. An accepted socket -/// takes an inbound slot against `max_inbound` before any byte is read, so -/// without it a remote that connects and stays silent holds that slot for as -/// long as it keeps the socket open. +/// `deadline` holds the first-frame and idle deadlines handed to each spawned +/// receive loop. An accepted socket takes an inbound slot against +/// `max_inbound` before any byte is read, so without them a remote that +/// connects and stays silent, or sends one frame and then goes silent, holds +/// that slot for as long as it keeps the socket open. #[allow(clippy::too_many_arguments)] async fn tor_accept_loop( listener: TcpListener, @@ -1125,7 +1151,7 @@ async fn tor_accept_loop( pool: ProxiedPool, mtu: u16, max_inbound: usize, - first_frame_timeout: Duration, + deadline: InboundDeadline, stats: Arc, ) { debug!( @@ -1219,7 +1245,7 @@ async fn tor_accept_loop( mtu, recv_stats, Direction::Inbound, - Some(first_frame_timeout), + Some(deadline), Some(ready_rx), ) .await; @@ -1971,7 +1997,8 @@ mod tests { fn spawn_onion_accept_loop( listener: TcpListener, packet_tx: PacketTx, - first_frame_timeout: Duration, + first_frame: Duration, + idle: Duration, ) -> (ProxiedPool, Arc, JoinHandle<()>) { let pool: ProxiedPool = Arc::new(Mutex::new(HashMap::new())); let stats = Arc::new(TorStats::new()); @@ -1982,7 +2009,7 @@ mod tests { pool.clone(), 1400, 64, - first_frame_timeout, + InboundDeadline { first_frame, idle }, stats.clone(), )); (pool, stats, handle) @@ -1997,8 +2024,12 @@ mod tests { let (tx, _rx) = packet_channel(32); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let listen = listener.local_addr().unwrap(); - let (pool, stats, accept) = - spawn_onion_accept_loop(listener, tx, Duration::from_millis(200)); + let (pool, stats, accept) = spawn_onion_accept_loop( + listener, + tx, + Duration::from_millis(200), + Duration::from_secs(5), + ); // Held open for the whole test: any release is the deadline's doing. let squatter = TcpStream::connect(listen).await.unwrap(); @@ -2026,8 +2057,12 @@ mod tests { let (tx, mut rx) = packet_channel(32); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let listen = listener.local_addr().unwrap(); - let (_pool, stats, accept) = - spawn_onion_accept_loop(listener, tx, Duration::from_millis(300)); + let (_pool, stats, accept) = spawn_onion_accept_loop( + listener, + tx, + Duration::from_millis(300), + Duration::from_secs(5), + ); let frame = build_msg1_frame(); let mut peer = TcpStream::connect(listen).await.unwrap(); @@ -2051,18 +2086,99 @@ mod tests { accept.abort(); } - /// The healthy path, and a regression guard as for TCP: the deadline is - /// scoped to the first iteration, so an established onion connection that - /// then goes quiet keeps its slot. It exists so a future general idle - /// deadline cannot start reaping quiet onion links without a test going - /// red. + /// The smallest frame the reader accepts as established, which the node + /// drops without closing the transport when it names no session. + fn squatter_frame() -> Vec { + let frame = crate::transport::framing::build_established_frame(0); + assert_eq!(frame.len(), crate::proto::fmp::wire::ENCRYPTED_MIN_SIZE); + frame + } + + /// Mirror of the TCP case: an onion-side remote that sends one frame and + /// then goes silent must lose its slot at the idle deadline. The + /// first-frame deadline is far above the idle one, so the release is the + /// idle deadline's doing. Break-check: scope the deadline in the shared + /// loop back to the first read and the count stays at 1. #[tokio::test] - async fn established_onion_connection_survives_long_idle() { + async fn onion_connection_that_goes_silent_after_one_frame_releases_its_slot() { + let (tx, mut rx) = packet_channel(32); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let listen = listener.local_addr().unwrap(); + let (pool, stats, accept) = spawn_onion_accept_loop( + listener, + tx, + Duration::from_secs(5), + Duration::from_millis(300), + ); + + let mut squatter = TcpStream::connect(listen).await.unwrap(); + squatter.write_all(&squatter_frame()).await.unwrap(); + let packet = tokio::time::timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout waiting for the squatter's frame") + .expect("packet channel closed"); + assert_eq!(packet.data, squatter_frame()); + assert_eq!(stats.pool_inbound_count(), 1); + + assert!( + wait_until(|| stats.pool_inbound_count() == 0, Duration::from_secs(2)).await, + "an onion connection silent after its first frame should lose its slot at the idle deadline" + ); + assert!(pool.lock().await.is_empty()); + + drop(squatter); + accept.abort(); + } + + /// The healthy path: an onion connection delivering a frame more often + /// than the idle deadline keeps its slot across many deadlines. + /// + /// Break-check: a deadline that does not re-arm on each frame (a single + /// deadline from accept) drops the connection after the first second. + #[tokio::test] + async fn onion_connection_sending_a_frame_every_interval_below_the_idle_deadline_is_kept() { let (tx, mut rx) = packet_channel(32); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let listen = listener.local_addr().unwrap(); let (pool, stats, accept) = - spawn_onion_accept_loop(listener, tx, Duration::from_millis(200)); + spawn_onion_accept_loop(listener, tx, Duration::from_secs(1), Duration::from_secs(1)); + + let mut peer = TcpStream::connect(listen).await.unwrap(); + // Twelve frames 250 ms apart: 3 s, three idle deadlines, with 750 ms + // of slack between each frame and the deadline it re-arms. + for i in 0..12 { + peer.write_all(&squatter_frame()).await.unwrap(); + let packet = tokio::time::timeout(Duration::from_secs(2), rx.recv()) + .await + .unwrap_or_else(|_| panic!("timeout waiting for frame {i}")) + .expect("packet channel closed"); + assert_eq!(packet.data, squatter_frame()); + assert_eq!( + stats.pool_inbound_count(), + 1, + "an onion connection delivering frames inside the idle deadline must keep its slot (frame {i})" + ); + tokio::time::sleep(Duration::from_millis(250)).await; + } + assert!(!pool.lock().await.is_empty()); + + drop(peer); + accept.abort(); + } + + /// An established onion connection quiet for longer than the first-frame + /// deadline, but not the idle deadline, keeps its slot. + #[tokio::test] + async fn established_onion_connection_quiet_past_the_first_frame_deadline_is_kept() { + let (tx, mut rx) = packet_channel(32); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let listen = listener.local_addr().unwrap(); + let (pool, stats, accept) = spawn_onion_accept_loop( + listener, + tx, + Duration::from_millis(200), + Duration::from_secs(5), + ); let mut peer = TcpStream::connect(listen).await.unwrap(); peer.write_all(&build_msg1_frame()).await.unwrap(); @@ -2072,7 +2188,7 @@ mod tests { .expect("packet channel closed"); assert_eq!(packet.data, build_msg1_frame()); - // Four deadlines' worth of silence after the first frame. + // Four first-frame deadlines' worth of silence after the first frame. tokio::time::sleep(Duration::from_millis(800)).await; assert_eq!( @@ -2143,7 +2259,10 @@ mod tests { 1400, stats.clone(), Direction::Inbound, - Some(Duration::from_millis(50)), + Some(InboundDeadline { + first_frame: Duration::from_millis(50), + idle: Duration::from_millis(50), + }), Some(ready_rx), ) .await; @@ -2197,7 +2316,10 @@ mod tests { 1400, recv_stats, Direction::Inbound, - Some(Duration::from_secs(5)), + Some(InboundDeadline { + first_frame: Duration::from_secs(5), + idle: Duration::from_secs(5), + }), Some(ready_rx), ) .await; @@ -2286,7 +2408,10 @@ mod tests { pool.clone(), 1400, 64, - Duration::from_secs(5), + InboundDeadline { + first_frame: Duration::from_secs(5), + idle: Duration::from_secs(5), + }, stats.clone(), ));