From 9f82c4726efac185364bfc086185f29a0f44c516 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Tue, 11 Aug 2026 15:26:34 +0000 Subject: [PATCH] Bound the first inbound frame, and let the reaper close the socket it forgets An accepted TCP socket took an inbound slot before a single byte was read. The cap is tested at accept, the pool insert and the counter bump follow with no read in between, and the frame reader's two read_exact calls carry no deadline. So an unauthenticated remote held a slot by connecting and sending nothing, and since pool keys are ip:port, N sockets from one address are N slots rather than one. At the 256 default that locks out inbound peering for as long as the attacker keeps the sockets open. The first frame on an inbound connection now has a deadline. It is a module constant rather than a config key: this branch takes no new operator-facing surface, and a knob is not needed to fix a missing bound. The onion listener has the identical accept-then-count ordering and gets the same treatment; it was not in the original report. The second half is that nothing reclaimed a slot once taken. The node-layer handshake reaper tore down session state but never closed the transport connection, so a peer that sent a real msg1 and then stalled was forgotten by the node while its socket, its pool entry and its slot lived on. The reaper now closes the transport connection too. Closing twice is safe: every close_connection implementation guards on removing the entry from its pool, and the connectionless ones are no-ops, so the handshake paths that already close and then drop a link are undisturbed. What this does not close, and it should be said plainly rather than discovered later: the deadline covers the first frame only. A peer that sends one well-formed frame and then goes silent still holds its slot, and so does one that completes msg1 and stalls beyond the reaper's reach. Closing those needs a rolling idle deadline, which interacts with heartbeats being per peer rather than per link and is a larger decision than this change. --- src/node/handlers/rx_loop.rs | 2 +- src/node/handlers/timeout.rs | 34 ++- src/node/lifecycle.rs | 2 +- src/node/tests/handshake.rs | 4 +- src/transport/tcp/mod.rs | 488 +++++++++++++++++++++++++++++++++-- src/transport/tcp/stream.rs | 58 +++-- src/transport/tor/mod.rs | 216 ++++++++++++++-- 7 files changed, 722 insertions(+), 82 deletions(-) diff --git a/src/node/handlers/rx_loop.rs b/src/node/handlers/rx_loop.rs index 525550fb..518d8522 100644 --- a/src/node/handlers/rx_loop.rs +++ b/src/node/handlers/rx_loop.rs @@ -250,7 +250,7 @@ impl Node { let _ = response_tx.send(response); } _ = tick.tick() => { - self.check_timeouts(); + self.check_timeouts().await; let now_ms = Self::now_ms(); self.reload_peer_acl().await; // The host map hot-reloads on the same tick as the ACL. It diff --git a/src/node/handlers/timeout.rs b/src/node/handlers/timeout.rs index e59c4e84..27437cf2 100644 --- a/src/node/handlers/timeout.rs +++ b/src/node/handlers/timeout.rs @@ -11,7 +11,7 @@ impl Node { /// /// Called periodically by the RX event loop. Removes connections that have /// been idle longer than the configured handshake timeout or are in Failed state. - pub(in crate::node) fn check_timeouts(&mut self) { + pub(in crate::node) async fn check_timeouts(&mut self) { if self.connections.is_empty() { return; } @@ -53,16 +53,21 @@ impl Node { self.schedule_retry(*identity.node_addr(), now_ms); } } - self.cleanup_stale_connection(link_id, now_ms); + self.cleanup_stale_connection(link_id, now_ms).await; } } /// Remove a handshake connection and all associated state. /// - /// Frees the session index, removes pending_outbound entry, and cleans up - /// the link and address mapping. Does not log — callers provide context-appropriate - /// log messages. - pub(in crate::node) fn cleanup_stale_connection(&mut self, link_id: LinkId, _now_ms: u64) { + /// Frees the session index, removes pending_outbound entry, closes the + /// underlying transport connection, and cleans up the link and address + /// mapping. Does not log — callers provide context-appropriate log + /// messages. + pub(in crate::node) async fn cleanup_stale_connection( + &mut self, + link_id: LinkId, + _now_ms: u64, + ) { let conn = match self.connections.remove(&link_id) { Some(c) => c, None => return, @@ -77,6 +82,23 @@ impl Node { let _ = self.index_allocator.free(idx); } + // Tear down the transport connection, not just the node-side state. + // A connection-oriented transport otherwise keeps the socket, its + // pool entry and its inbound-slot accounting alive after the node + // has forgotten the handshake that socket belonged to, so a peer + // that sends msg1 and then stalls holds an inbound slot forever. + // Closing twice is harmless: every `close_connection` implementation + // is `if let Some(conn) = pool.remove(addr)` and the connectionless + // default is a no-op, so the handshake paths that already close and + // then drop a link cannot be disturbed by this. + if let Some(link) = self.links.get(&link_id) { + let tid = link.transport_id(); + let addr = link.remote_addr().clone(); + if let Some(transport) = self.transports.get(&tid) { + transport.close_connection(&addr).await; + } + } + // Remove link and addr_to_link self.remove_link(&link_id); if let Some(transport_id) = transport_id { diff --git a/src/node/lifecycle.rs b/src/node/lifecycle.rs index a4d44b2e..e2523976 100644 --- a/src/node/lifecycle.rs +++ b/src/node/lifecycle.rs @@ -754,7 +754,7 @@ impl Node { .map(|(link_id, _)| *link_id) .collect(); for link_id in stale { - self.cleanup_stale_connection(link_id, now_ms); + self.cleanup_stale_connection(link_id, now_ms).await; } } } diff --git a/src/node/tests/handshake.rs b/src/node/tests/handshake.rs index 822eb3db..996bdc2d 100644 --- a/src/node/tests/handshake.rs +++ b/src/node/tests/handshake.rs @@ -700,7 +700,7 @@ async fn test_stale_connection_cleanup() { // Connection was created at time 1000ms. check_timeouts uses SystemTime::now(), // which is far beyond the 30s timeout. The connection should be cleaned up. - node.check_timeouts(); + node.check_timeouts().await; // Verify everything was cleaned up assert_eq!( @@ -770,7 +770,7 @@ async fn test_failed_connection_cleanup() { assert_eq!(node.connection_count(), 1); // Failed connections should be cleaned up immediately regardless of age - node.check_timeouts(); + node.check_timeouts().await; assert_eq!( node.connection_count(), diff --git a/src/transport/tcp/mod.rs b/src/transport/tcp/mod.rs index 350c8799..beaadb8d 100644 --- a/src/transport/tcp/mod.rs +++ b/src/transport/tcp/mod.rs @@ -126,6 +126,9 @@ pub struct TcpTransport { /// fallback when this transport has no explicit `max_inbound_connections`. /// `None` means "not provided" — fall through to the built-in default. node_max_connections: Option, + /// Deadline from accept to the first complete inbound frame. Defaults to + /// `INBOUND_FIRST_FRAME_TIMEOUT`; overridable only from tests. + first_frame_timeout: Duration, /// Transport statistics. stats: Arc, } @@ -149,10 +152,22 @@ impl TcpTransport { accept_task: None, local_addr: None, node_max_connections: None, + first_frame_timeout: INBOUND_FIRST_FRAME_TIMEOUT, stats: Arc::new(TcpStats::new()), } } + /// Override the accept-to-first-frame deadline. + /// + /// Test-only: the accept loop is reachable from the test module only + /// through `start_async()`, which reads this field when it builds the + /// `AcceptConfig`, so there is no other way to drive the deadline at a + /// duration a unit test can wait for. + #[cfg(test)] + pub(crate) fn set_first_frame_timeout(&mut self, d: Duration) { + self.first_frame_timeout = d; + } + /// Set the node-wide `node.limits.max_connections` value. /// /// Used as the inbound-cap fallback when this transport instance has no @@ -241,6 +256,7 @@ 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, }; let accept_task = tokio::spawn(async move { @@ -459,6 +475,10 @@ impl TcpTransport { mtu, recv_stats, Direction::Outbound, + // Outbound connections hold no inbound slot and are not + // gated on an accept-loop insert. + None, + None, ) .await; }); @@ -708,6 +728,10 @@ impl TcpTransport { mss_mtu, recv_stats, Direction::Outbound, + // Outbound connections hold no inbound slot and are not + // gated on an accept-loop insert. + None, + None, ) .await; }); @@ -801,6 +825,20 @@ impl Transport for TcpTransport { // Accept Loop // ============================================================================ +/// Deadline from accept to the first complete inbound FMP frame. +/// +/// An accepted socket takes an inbound pool slot before any byte is read, +/// so without a deadline a remote that connects and stays silent holds +/// that slot for as long as it keeps the socket open. The node-layer +/// reaper cannot see such a socket: no frame means no link and no node +/// state to time out. The value matches the node-layer handshake reaper +/// (`handshake_timeout_secs`, `src/config/node.rs:101`), so a peer that +/// misses this deadline would have been reaped node-side anyway. +/// +/// Deliberately not a config key: `maint` takes no new operator-facing +/// TOML surface. +pub(crate) const INBOUND_FIRST_FRAME_TIMEOUT: Duration = Duration::from_secs(30); + /// Socket configuration parameters passed to the accept loop. struct AcceptConfig { mtu: u16, @@ -809,6 +847,7 @@ struct AcceptConfig { keepalive_secs: u64, recv_buf: usize, send_buf: usize, + first_frame_timeout: Duration, } /// TCP accept loop — runs as a spawned task when bind_addr is configured. @@ -828,6 +867,7 @@ async fn accept_loop( keepalive_secs, recv_buf, send_buf, + first_frame_timeout, } = cfg; debug!(transport_id = %transport_id, "TCP accept loop starting"); @@ -904,6 +944,12 @@ async fn accept_loop( let recv_stats = stats.clone(); let recv_addr = remote_addr.clone(); + // Readiness barrier: the receive task must not reach its + // cleanup path before the pool insert and counter bump below, + // or it would remove nothing and leave an orphaned entry with + // a permanently incremented inbound counter. + let (ready_tx, ready_rx) = tokio::sync::oneshot::channel(); + let recv_task = tokio::spawn(async move { tcp_receive_loop( read_half, @@ -914,6 +960,8 @@ async fn accept_loop( conn_mtu, recv_stats, Direction::Inbound, + Some(first_frame_timeout), + Some(ready_rx), ) .await; }); @@ -928,10 +976,15 @@ async fn accept_loop( let mut pool_guard = pool.lock().await; pool_guard.insert(remote_addr.clone(), conn); + drop(pool_guard); stats.record_connection_accepted(); stats.record_pool_inbound_added(); + // Release the receive task now that both the pool entry and + // the inbound counter are in place. + let _ = ready_tx.send(()); + debug!( transport_id = %transport_id, remote_addr = %remote_addr, @@ -962,6 +1015,12 @@ async fn accept_loop( /// the cleanup path can decrement the correct `pool_inbound` / /// `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 +/// 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)] async fn tcp_receive_loop( mut reader: tokio::net::tcp::OwnedReadHalf, @@ -972,6 +1031,8 @@ async fn tcp_receive_loop( mtu: u16, stats: Arc, direction: Direction, + first_frame_timeout: Option, + ready_rx: Option>, ) { debug!( transport_id = %transport_id, @@ -979,39 +1040,74 @@ async fn tcp_receive_loop( "TCP receive loop starting" ); - loop { - match read_fmp_packet(&mut reader, mtu).await { - Ok(data) => { - stats.record_recv(data.len()); + // An `Err` here means the accept loop went away between the insert and + // the signal. Fall through to the cleanup below rather than returning, + // so a pooled entry cannot be stranded with the counter incremented. + let admitted = match ready_rx { + Some(rx) => rx.await.is_ok(), + None => true, + }; - trace!( - transport_id = %transport_id, - remote_addr = %remote_addr, - bytes = data.len(), - "TCP packet received" - ); + 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 { + Ok(result) => result, + Err(_) => { + // Not a recv error: `record_recv_error` means + // framing or I/O failure, and folding deadline + // expiries into it corrupts that counter. + 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" + ); + break; + } + } + } + _ => read_fmp_packet(&mut reader, mtu).await, + }; + first = false; - let packet = ReceivedPacket::new(transport_id, remote_addr.clone(), data); + match read { + Ok(data) => { + stats.record_recv(data.len()); - if packet_tx.send(packet).await.is_err() { + trace!( + transport_id = %transport_id, + remote_addr = %remote_addr, + bytes = data.len(), + "TCP packet received" + ); + + let packet = ReceivedPacket::new(transport_id, remote_addr.clone(), data); + + if packet_tx.send(packet).await.is_err() { + debug!( + transport_id = %transport_id, + "Packet channel closed, stopping TCP receive loop" + ); + break; + } + } + Err(e) => { + stats.record_recv_error(); + // EOF or protocol error — remove connection from pool debug!( transport_id = %transport_id, - "Packet channel closed, stopping TCP receive loop" + remote_addr = %remote_addr, + error = %e, + "TCP receive error, removing connection" ); break; } } - Err(e) => { - stats.record_recv_error(); - // EOF or protocol error — remove connection from pool - debug!( - transport_id = %transport_id, - remote_addr = %remote_addr, - error = %e, - "TCP receive error, removing connection" - ); - break; - } } } @@ -1146,10 +1242,34 @@ fn read_mss_mtu(stream: &std::net::TcpStream, default_mtu: u16) -> u16 { #[cfg(test)] mod tests { + use super::stream::build_msg1_frame; use super::*; use crate::transport::packet_channel; use tokio::time::{Duration, timeout}; + /// Poll `f` every 10ms until it holds or `limit` elapses. + async fn wait_until bool>(mut f: F, limit: Duration) -> bool { + let deadline = Instant::now() + limit; + loop { + if f() { + return true; + } + if Instant::now() >= deadline { + return false; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + } + + fn capped_config(max_inbound: usize) -> TcpConfig { + TcpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + mtu: Some(1400), + max_inbound_connections: Some(max_inbound), + ..Default::default() + } + } + fn make_config() -> TcpConfig { TcpConfig { bind_addr: Some("127.0.0.1:0".to_string()), @@ -1752,4 +1872,324 @@ mod tests { t1.stop_async().await.unwrap(); t2.stop_async().await.unwrap(); } + + // ======================================================================== + // Inbound first-frame deadline + // ======================================================================== + + /// A socket that connects and sends nothing must have its inbound slot + /// released by the first-frame deadline. + /// + /// Break-check: with the `tokio::time::timeout` wrapper removed from the + /// first read, the socket parks on an unbounded `read_exact` and the + /// count stays at 1 for as long as the peer keeps the socket open, so + /// the second assertion fails. + #[tokio::test] + async fn idle_inbound_socket_releases_its_slot() { + let (tx, _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.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + // Connect and say nothing. Held open for the whole test so that any + // slot release is the deadline's doing and not a client disconnect. + let squatter = TcpStream::connect(listen).await.unwrap(); + + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 1, + Duration::from_secs(2) + ) + .await, + "an accepted socket should take an inbound slot" + ); + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 0, + Duration::from_secs(2) + ) + .await, + "a silent inbound socket should lose its slot at the first-frame 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 silent socket, a genuine peer is refused + /// until the deadline frees the slot, and admitted afterwards. + /// + /// Break-check: without the deadline the squatter never releases, so the + /// genuine peer's frame is never delivered and the final receive times + /// out. + #[tokio::test] + async fn inbound_cap_recovers_after_first_frame_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_millis(300)); + transport.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + let squatter = TcpStream::connect(listen).await.unwrap(); + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 1, + Duration::from_secs(2) + ) + .await, + "the squatter should fill the cap of one" + ); + + // While the cap is full a genuine peer is rejected outright. + let mut early = TcpStream::connect(listen).await.unwrap(); + let _ = early.write_all(&build_msg1_frame()).await; + assert!( + timeout(Duration::from_millis(200), rx.recv()) + .await + .is_err(), + "a peer arriving while the cap is full must not be admitted" + ); + drop(early); + + // The deadline frees the slot without the squatter disconnecting. + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 0, + Duration::from_secs(2) + ) + .await, + "the 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(); + } + + /// Regression guard, not evidence that the fix works. + /// + /// 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. + #[tokio::test] + async fn established_connection_survives_long_idle() { + 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.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + let mut peer = TcpStream::connect(listen).await.unwrap(); + peer.write_all(&build_msg1_frame()).await.unwrap(); + let packet = timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout") + .expect("packet channel closed"); + assert_eq!(packet.data, build_msg1_frame()); + + // Four deadlines' worth of silence after the first frame. + tokio::time::sleep(Duration::from_millis(800)).await; + + assert_eq!( + transport.stats().pool_inbound_count(), + 1, + "an established connection must not be dropped by the first-frame deadline" + ); + assert!(!transport.pool.lock().await.is_empty()); + + 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] + async fn slow_first_frame_within_deadline_is_admitted() { + 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.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + let mut peer = TcpStream::connect(listen).await.unwrap(); + tokio::time::sleep(Duration::from_millis(300)).await; + peer.write_all(&build_msg1_frame()).await.unwrap(); + + let packet = timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout") + .expect("packet channel closed"); + assert_eq!(packet.data, build_msg1_frame()); + assert_eq!(transport.stats().pool_inbound_count(), 1); + + drop(peer); + transport.stop_async().await.unwrap(); + } + + /// The honest-slow-peer case the wrapper actually kills: a first frame + /// that *starts* inside the deadline but completes after it. The + /// deadline covers the whole frame, not its first byte, so the drip is + /// dropped and its slot released. + #[tokio::test] + async fn byte_dripped_first_frame_past_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_millis(300)); + transport.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + let frame = build_msg1_frame(); + let mut peer = TcpStream::connect(listen).await.unwrap(); + // Prefix inside the 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 first frame completing after the 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(); + } + + /// Break-check for the readiness barrier's error path. + /// + /// Stands in for an accept loop aborted between the pool insert and the + /// `ready_tx.send()`: the sender is dropped, so `ready_rx.await` returns + /// `Err`. The receive loop must still fall through to its cleanup, or + /// the pooled entry and its inbound-counter increment are stranded with + /// no task left to undo them. A bare `return` on the error path fails + /// both assertions below. + #[tokio::test] + async fn receive_loop_cleans_up_when_readiness_signal_is_dropped() { + let (tx, _rx) = packet_channel(10); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let listen = listener.local_addr().unwrap(); + let client = TcpStream::connect(listen).await.unwrap(); + let (server, peer_addr) = listener.accept().await.unwrap(); + let remote = TransportAddr::from_string(&peer_addr.to_string()); + let (read_half, write_half) = server.into_split(); + + let pool: ConnectionPool = Arc::new(Mutex::new(HashMap::new())); + let stats = Arc::new(TcpStats::new()); + pool.lock().await.insert( + remote.clone(), + TcpConnection { + writer: Arc::new(Mutex::new(write_half)), + recv_task: tokio::spawn(async {}), + mtu: 1400, + established_at: Instant::now(), + direction: Direction::Inbound, + }, + ); + stats.record_pool_inbound_added(); + assert_eq!(stats.pool_inbound_count(), 1); + + let (ready_tx, ready_rx) = tokio::sync::oneshot::channel::<()>(); + drop(ready_tx); + + tcp_receive_loop( + read_half, + TransportId::new(1), + remote.clone(), + tx, + pool.clone(), + 1400, + stats.clone(), + Direction::Inbound, + Some(Duration::from_millis(50)), + Some(ready_rx), + ) + .await; + + assert!( + pool.lock().await.is_empty(), + "an aborted accept must not strand a pool entry" + ); + assert_eq!( + stats.pool_inbound_count(), + 0, + "an aborted accept must not strand an inbound-counter increment" + ); + drop(client); + } + + /// Invariant guard: a deadline that expires immediately still leaves no + /// orphaned pool entry or counter increment behind. + /// + /// This is not a break-check for the readiness barrier. On the + /// current-thread test runtime the accept loop queues for the pool lock + /// before the spawned receive task can run at all, so the insert wins + /// the race with or without the barrier. The barrier's error path is + /// break-checked in `receive_loop_cleans_up_when_readiness_signal_is_dropped`. + #[tokio::test] + async fn zero_deadline_leaves_no_orphaned_pool_entry() { + let (tx, _rx) = packet_channel(100); + let mut transport = TcpTransport::new(TransportId::new(1), None, make_config(), tx); + transport.set_first_frame_timeout(Duration::ZERO); + transport.start_async().await.unwrap(); + let listen = transport.local_addr().unwrap(); + + // Hold the pool across the accept so the receive task cannot reach + // its cleanup while the accept loop is mid-insert. + let guard = transport.pool.lock().await; + let client = TcpStream::connect(listen).await.unwrap(); + tokio::time::sleep(Duration::from_millis(100)).await; + drop(guard); + + // Sequence the checks off `connections_accepted`, which the accept + // loop bumps only after its insert. Reading the pool counter first + // would otherwise observe the pre-accept zero and prove nothing. + assert!( + wait_until( + || transport.stats().snapshot().connections_accepted == 1, + Duration::from_secs(2) + ) + .await, + "the accept loop should have admitted the connection" + ); + assert!( + wait_until( + || transport.stats().pool_inbound_count() == 0 + && transport + .pool + .try_lock() + .map(|p| p.is_empty()) + .unwrap_or(false), + Duration::from_secs(2) + ) + .await, + "an immediately expired deadline should leave neither a pool entry nor a counter increment" + ); + + drop(client); + transport.stop_async().await.unwrap(); + } } diff --git a/src/transport/tcp/stream.rs b/src/transport/tcp/stream.rs index 4323fdc3..4d5e03c0 100644 --- a/src/transport/tcp/stream.rs +++ b/src/transport/tcp/stream.rs @@ -178,6 +178,39 @@ pub async fn read_fmp_packet( Ok(packet) } +// ============================================================================ +// Test Frame Builders +// ============================================================================ + +/// Build a minimal established frame with the given payload_len. +/// Layout: [ver+phase:1][flags:1][payload_len:2 LE][12 bytes header][payload_len bytes][16 bytes tag] +/// +/// Lives at module scope so the transport modules that share this reader +/// (tcp, tor) can build wire-shaped frames in their own tests. +#[cfg(test)] +pub(crate) fn build_established_frame(payload_len: u16) -> Vec { + let total = PREFIX_SIZE + ESTABLISHED_REMAINING_HEADER + payload_len as usize + AEAD_TAG_SIZE; + let mut frame = vec![0u8; total]; + frame[0] = 0x00; // ver=0, phase=0 (established) + frame[1] = 0x00; // flags + frame[2..4].copy_from_slice(&payload_len.to_le_bytes()); + // Fill remaining with pattern for verification + for (i, byte) in frame[PREFIX_SIZE..total].iter_mut().enumerate() { + *byte = ((PREFIX_SIZE + i) & 0xFF) as u8; + } + frame +} + +/// Build a msg1 frame (114 bytes total). +#[cfg(test)] +pub(crate) fn build_msg1_frame() -> Vec { + let mut frame = vec![0xAA; MSG1_WIRE_SIZE]; + frame[0] = 0x01; // ver=0, phase=1 + frame[1] = 0x00; // flags + frame[2..4].copy_from_slice(&MSG1_PAYLOAD_LEN.to_le_bytes()); + frame +} + // ============================================================================ // Tests // ============================================================================ @@ -187,31 +220,6 @@ mod tests { use super::*; use std::io::Cursor; - /// Build a minimal established frame with the given payload_len. - /// Layout: [ver+phase:1][flags:1][payload_len:2 LE][12 bytes header][payload_len bytes][16 bytes tag] - fn build_established_frame(payload_len: u16) -> Vec { - let total = - PREFIX_SIZE + ESTABLISHED_REMAINING_HEADER + payload_len as usize + AEAD_TAG_SIZE; - let mut frame = vec![0u8; total]; - frame[0] = 0x00; // ver=0, phase=0 (established) - frame[1] = 0x00; // flags - frame[2..4].copy_from_slice(&payload_len.to_le_bytes()); - // Fill remaining with pattern for verification - for (i, byte) in frame[PREFIX_SIZE..total].iter_mut().enumerate() { - *byte = ((PREFIX_SIZE + i) & 0xFF) as u8; - } - frame - } - - /// Build a msg1 frame (114 bytes total). - fn build_msg1_frame() -> Vec { - let mut frame = vec![0xAA; MSG1_WIRE_SIZE]; - frame[0] = 0x01; // ver=0, phase=1 - frame[1] = 0x00; // flags - frame[2..4].copy_from_slice(&MSG1_PAYLOAD_LEN.to_le_bytes()); - frame - } - /// Build a msg2 frame (69 bytes total). fn build_msg2_frame() -> Vec { let mut frame = vec![0xBB; MSG2_WIRE_SIZE]; diff --git a/src/transport/tor/mod.rs b/src/transport/tor/mod.rs index 5c2b8a92..43527a74 100644 --- a/src/transport/tor/mod.rs +++ b/src/transport/tor/mod.rs @@ -31,6 +31,7 @@ use super::{ TransportError, TransportId, TransportState, TransportType, }; use crate::config::TorConfig; +use crate::transport::tcp::INBOUND_FIRST_FRAME_TIMEOUT; use crate::transport::tcp::stream::read_fmp_packet; use control::{ControlAuth, TorControlClient, TorMonitoringInfo}; use stats::TorStats; @@ -408,6 +409,7 @@ impl TorTransport { pool, mtu, max_inbound, + INBOUND_FIRST_FRAME_TIMEOUT, stats, ) .await; @@ -835,6 +837,10 @@ impl TorTransport { mtu, recv_stats, Direction::Outbound, + // Outbound connections hold no inbound slot and are not + // gated on an accept-loop insert. + None, + None, ) .await; }); @@ -1067,6 +1073,10 @@ impl TorTransport { mtu, recv_stats, Direction::Outbound, + // Outbound connections hold no inbound slot and are not + // gated on an accept-loop insert. + None, + None, ) .await; }); @@ -1178,6 +1188,12 @@ impl Transport for TorTransport { /// connection from the pool and exits. `direction` is captured so the /// cleanup path can decrement the correct `pool_inbound` / /// `pool_outbound` counter. +/// +/// `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 +/// 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)] async fn tor_receive_loop( mut reader: tokio::net::tcp::OwnedReadHalf, @@ -1188,6 +1204,8 @@ async fn tor_receive_loop( mtu: u16, stats: Arc, direction: Direction, + first_frame_timeout: Option, + ready_rx: Option>, ) { debug!( transport_id = %transport_id, @@ -1195,38 +1213,73 @@ async fn tor_receive_loop( "Tor receive loop starting" ); - loop { - match read_fmp_packet(&mut reader, mtu).await { - Ok(data) => { - stats.record_recv(data.len()); + // An `Err` here means the accept loop went away between the insert and + // the signal. Fall through to the cleanup below rather than returning, + // so a pooled entry cannot be stranded with the counter incremented. + let admitted = match ready_rx { + Some(rx) => rx.await.is_ok(), + None => true, + }; - trace!( - transport_id = %transport_id, - remote_addr = %remote_addr, - bytes = data.len(), - "Tor packet received" - ); + 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 { + Ok(result) => result, + Err(_) => { + // Not a recv error: `record_recv_error` means + // framing or I/O failure, and folding deadline + // expiries into it corrupts that counter. + 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 onion connection" + ); + break; + } + } + } + _ => read_fmp_packet(&mut reader, mtu).await, + }; + first = false; - let packet = ReceivedPacket::new(transport_id, remote_addr.clone(), data); + match read { + Ok(data) => { + stats.record_recv(data.len()); - if packet_tx.send(packet).await.is_err() { + trace!( + transport_id = %transport_id, + remote_addr = %remote_addr, + bytes = data.len(), + "Tor packet received" + ); + + let packet = ReceivedPacket::new(transport_id, remote_addr.clone(), data); + + if packet_tx.send(packet).await.is_err() { + debug!( + transport_id = %transport_id, + "Packet channel closed, stopping Tor receive loop" + ); + break; + } + } + Err(e) => { + stats.record_recv_error(); debug!( transport_id = %transport_id, - "Packet channel closed, stopping Tor receive loop" + remote_addr = %remote_addr, + error = %e, + "Tor receive error, removing connection" ); break; } } - Err(e) => { - stats.record_recv_error(); - debug!( - transport_id = %transport_id, - remote_addr = %remote_addr, - error = %e, - "Tor receive error, removing connection" - ); - break; - } } } @@ -1292,6 +1345,7 @@ fn configure_socket( /// connections to a local TCP listener; we accept them, configure /// socket options, split the stream, and spawn a per-connection /// receive task. +#[allow(clippy::too_many_arguments)] async fn tor_accept_loop( listener: TcpListener, transport_id: TransportId, @@ -1299,6 +1353,7 @@ async fn tor_accept_loop( pool: ConnectionPool, mtu: u16, max_inbound: usize, + first_frame_timeout: Duration, stats: Arc, ) { debug!( @@ -1376,6 +1431,12 @@ async fn tor_accept_loop( let recv_addr = remote_addr.clone(); let recv_tx = packet_tx.clone(); + // Readiness barrier: the receive task must not reach its cleanup + // path before the pool insert and counter bump below, or it would + // remove nothing and leave an orphaned entry with a permanently + // incremented inbound counter. + let (ready_tx, ready_rx) = tokio::sync::oneshot::channel(); + let recv_task = tokio::spawn(async move { tor_receive_loop( read_half, @@ -1386,6 +1447,8 @@ async fn tor_accept_loop( mtu, recv_stats, Direction::Inbound, + Some(first_frame_timeout), + Some(ready_rx), ) .await; }); @@ -1406,6 +1469,10 @@ async fn tor_accept_loop( stats.record_connection_accepted(); stats.record_pool_inbound_added(); + // Release the receive task now that both the pool entry and the + // inbound counter are in place. + let _ = ready_tx.send(()); + debug!( transport_id = %transport_id, peer_addr = %peer_addr, @@ -2034,4 +2101,107 @@ mod tests { let err = format!("{}", result.unwrap_err()); assert!(err.contains("directory")); } + + // ======================================================================== + // Inbound first-frame deadline (onion listener) + // ======================================================================== + + /// Poll `f` every 10ms until it holds or `limit` elapses. + async fn wait_until bool>(mut f: F, limit: Duration) -> bool { + let deadline = Instant::now() + limit; + loop { + if f() { + return true; + } + if Instant::now() >= deadline { + return false; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + } + + /// Drives `tor_accept_loop` directly: the only production path to it is + /// `start_directory_mode`, which needs a Tor-managed hostname file and a + /// running daemon, so it is not reachable from a unit test. + fn spawn_onion_accept_loop( + listener: TcpListener, + packet_tx: PacketTx, + first_frame_timeout: Duration, + ) -> (ConnectionPool, Arc, JoinHandle<()>) { + let pool: ConnectionPool = Arc::new(Mutex::new(HashMap::new())); + let stats = Arc::new(TorStats::new()); + let handle = tokio::spawn(tor_accept_loop( + listener, + TransportId::new(1), + packet_tx, + pool.clone(), + 1400, + 64, + first_frame_timeout, + stats.clone(), + )); + (pool, stats, handle) + } + + /// Mirror of the TCP case: a silent onion-side socket must lose its + /// inbound slot at the deadline. Break-check: with the wrapper removed + /// the count stays at 1 and the second assertion fails. + #[tokio::test] + async fn idle_inbound_onion_socket_releases_its_slot() { + 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)); + + // Held open for the whole test: any release is the deadline's doing. + let squatter = TcpStream::connect(listen).await.unwrap(); + + assert!( + wait_until(|| stats.pool_inbound_count() == 1, Duration::from_secs(2)).await, + "an accepted onion socket should take an inbound slot" + ); + assert!( + wait_until(|| stats.pool_inbound_count() == 0, Duration::from_secs(2)).await, + "a silent onion socket should lose its slot at the first-frame deadline" + ); + assert!(pool.lock().await.is_empty()); + + drop(squatter); + accept.abort(); + } + + /// Regression guard only, as for TCP: the deadline is scoped to the + /// first iteration, so this passes by construction. It exists so a + /// future general idle deadline cannot start reaping quiet onion links + /// without a test going red. + #[tokio::test] + async fn established_onion_connection_survives_long_idle() { + 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)); + + let mut peer = TcpStream::connect(listen).await.unwrap(); + peer.write_all(&build_msg1_frame()).await.unwrap(); + let packet = tokio::time::timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout") + .expect("packet channel closed"); + assert_eq!(packet.data, build_msg1_frame()); + + // Four deadlines' worth of silence after the first frame. + tokio::time::sleep(Duration::from_millis(800)).await; + + assert_eq!( + stats.pool_inbound_count(), + 1, + "an established onion connection must not be dropped by the first-frame deadline" + ); + assert!(!pool.lock().await.is_empty()); + + drop(peer); + accept.abort(); + } }