Files
fips/src/transport/socks5/pool.rs
T
Johnathan Corgan 514b2ac699 Drop inbound TCP and onion connections that go silent after a frame
An accepted connection takes an inbound slot before any byte is read,
and the receive loop bounded only the wait for the first complete
frame. A remote could send one well-formed 32-byte established frame,
which the node drops without closing the transport because it names no
session, and then hold the slot for as long as it kept the socket open.
Enough of them fill max_inbound_connections and lock genuine peers out.

Every read on an inbound connection is now bounded: the first-frame
deadline until a frame arrives, then an idle deadline that re-arms on
each complete frame. The deadline covers the whole frame, so a remote
cannot keep the connection by dripping bytes.

The idle deadline is the node's own link-silence bound (handshake
resend ladder, three ticks, and the longer of the heartbeat interval
and the link-dead timeout), 64 s at default settings. A connection
carrying a live link receives at least a heartbeat per interval, and
the node would reap a link silent for longer, so the deadline cannot
drop a link the node would keep. It follows the existing liveness keys
and adds no configuration. The node sets it on every TCP and Tor
instance; a transport built outside the node falls back to 64 s, and a
test asserts that the value derived from stock settings equals that
fallback, so a change to one default cannot leave the other behind.

The Tor onion listener shares the receive loop and gets the same
deadline. Nym is outbound-only and keeps no counted inbound slots, so
it has neither the first-frame nor the idle deadline, and its receive
loop's doc comment now says so. Outbound connections on every transport
stay without a deadline.

The two tests that pinned the old first-read-only scoping are replaced
by tests for the new contract: frames every 250 ms against a 1 s idle
deadline keep the connection across three deadlines, and quiet longer
than the first-frame deadline but shorter than the idle one is still
kept. A node-level test runs two real nodes over loopback TCP with
heartbeat-only traffic for three derived deadlines and checks the link
and the inbound connection survive.
2026-10-01 22:40:40 +00:00

267 lines
10 KiB
Rust

//! Shared connection pool for the proxied (Tor / Nym) transports.
//!
//! Both transports keep the same two maps — an established-connection pool and
//! a pending-connection ("connecting") pool — and poll a completed background
//! connect the same way. The only per-transport difference is the metadata
//! carried on each pooled connection (`Direction` for tor's inbound/outbound
//! pool accounting, `()` for nym), captured by the generic `M` type parameter.
use std::collections::HashMap;
use std::sync::Arc;
use futures::FutureExt;
use tokio::net::TcpStream;
use tokio::net::tcp::{OwnedReadHalf, OwnedWriteHalf};
use tokio::sync::Mutex;
use tokio::task::JoinHandle;
use tokio::time::Instant;
use tracing::{debug, trace};
use crate::transport::framing::read_fmp_packet;
use crate::transport::tcp::InboundDeadline;
use crate::transport::{
ConnectionState, PacketTx, ReceivedPacket, TransportAddr, TransportError, TransportId,
};
/// State for a single pooled connection to a peer.
///
/// `M` is per-transport metadata: `Direction` for tor (drives
/// inbound/outbound pool accounting), `()` for nym.
pub(crate) struct ProxiedConnection<M> {
/// Write half of the split stream.
pub writer: Arc<Mutex<OwnedWriteHalf>>,
/// Receive task for this connection.
pub recv_task: JoinHandle<()>,
/// MTU for this connection.
#[allow(dead_code)]
pub mtu: u16,
/// When the connection was established.
#[allow(dead_code)]
pub established_at: Instant,
/// Per-transport metadata (tor: `Direction`; nym: `()`).
pub meta: M,
}
/// Shared connection pool: addr -> per-connection state.
pub(crate) type ProxiedPool<M> = Arc<Mutex<HashMap<TransportAddr, ProxiedConnection<M>>>>;
/// A pending background connection attempt.
///
/// Holds the JoinHandle for a spawned SOCKS5 connect task. The task
/// produces a configured `TcpStream` and MTU on success.
pub(crate) struct ConnectingEntry {
/// Background task performing SOCKS5 connect + socket configuration.
pub task: JoinHandle<Result<(TcpStream, u16), TransportError>>,
}
/// Map of addresses with background connection attempts in progress.
pub(crate) type ConnectingPool = Arc<Mutex<HashMap<TransportAddr, ConnectingEntry>>>;
/// Poll the state of a connection to a remote address.
///
/// Checks both established and connecting pools. If a background connect task
/// has completed successfully, invokes `promote` (which spawns a receive loop
/// and inserts into the established pool) and reports `Connected`; on failure
/// reports it. Synchronous — uses `try_lock` internally and returns
/// `ConnectionState::Connecting` if a lock can't be acquired.
///
/// This is the byte-for-byte former `connection_state_sync` body, with the
/// per-transport `promote_connection` call abstracted behind `promote`.
pub(crate) fn poll_connecting<M>(
pool: &ProxiedPool<M>,
connecting: &ConnectingPool,
addr: &TransportAddr,
promote: impl FnOnce(TcpStream, u16),
) -> ConnectionState {
// Check established pool first
if let Ok(pool) = pool.try_lock() {
if pool.contains_key(addr) {
return ConnectionState::Connected;
}
} else {
return ConnectionState::Connecting; // can't tell, assume still going
}
// Check connecting pool
let mut connecting = match connecting.try_lock() {
Ok(c) => c,
Err(_) => return ConnectionState::Connecting,
};
let entry = match connecting.get_mut(addr) {
Some(e) => e,
None => return ConnectionState::None,
};
// Check if the background task has completed
if !entry.task.is_finished() {
return ConnectionState::Connecting;
}
// Task is done — take the result and remove from connecting pool.
let addr_clone = addr.clone();
let task = connecting.remove(&addr_clone).unwrap().task;
// Since the task is finished, we can safely poll it with now_or_never.
match task.now_or_never() {
Some(Ok(Ok((stream, mtu)))) => {
promote(stream, mtu);
ConnectionState::Connected
}
Some(Ok(Err(e))) => ConnectionState::Failed(format!("{}", e)),
Some(Err(e)) => ConnectionState::Failed(format!("task failed: {}", e)),
None => ConnectionState::Connecting,
}
}
/// Minimal stats surface the shared receive loop needs.
///
/// The per-transport stats structs implement this by delegating to their
/// shared counter base; the loop records received bytes and receive errors
/// without knowing the concrete transport.
pub(crate) trait ProxiedStats: Send + Sync + 'static {
/// Record a successful receive of `bytes` bytes.
fn record_recv(&self, bytes: usize);
/// Record a receive error.
fn record_recv_error(&self);
}
/// Shared per-connection receive loop for the proxied transports.
///
/// Reads complete FMP packets, delivers them to the node, and on error/EOF
/// removes the connection from the pool and runs `on_remove` for any
/// per-transport teardown accounting. The `label` is the in-loop log word
/// ("Nym" / "Tor").
///
/// Teardown/cleanup contract (reproduced exactly to stay behavior-neutral):
/// the pool entry is removed, and `on_remove` fires **only** when the removal
/// returned `Some`, taking the metadata from the removed entry, and after the
/// pool guard is dropped. Firing on `Some` only means a concurrent
/// `close`/`stop` teardown of the same address can never double-count.
///
/// The terminal "receive loop stopped" log is **not** emitted here — it is
/// hoisted into each per-transport wrapper (tor carries a `direction` field
/// nym lacks), so this loop is silent on exit.
///
/// `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 a connection that takes a capped inbound slot from accept
/// — today only tor's onion listener — and `None` everywhere else, which
/// covers every outbound connection and the whole of the nym transport (nym
/// is outbound-only and keeps no counted slots). A deadline expiry is not a
/// receive error and is deliberately not recorded as one.
///
/// `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, or the removal finds nothing, `on_remove` never
/// fires, and the increment is stranded for the life of the process.
#[allow(clippy::too_many_arguments)]
pub(crate) async fn proxied_receive_loop<S: ProxiedStats, M>(
mut reader: OwnedReadHalf,
transport_id: TransportId,
remote_addr: TransportAddr,
packet_tx: PacketTx,
pool: ProxiedPool<M>,
mtu: u16,
stats: Arc<S>,
label: &'static str,
deadline: Option<InboundDeadline>,
ready_rx: Option<tokio::sync::oneshot::Receiver<()>>,
on_remove: impl Fn(&S, &M),
) {
debug!(
transport_id = %transport_id,
remote_addr = %remote_addr,
"{} receive loop starting",
label
);
// 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,
};
if admitted {
let mut first = true;
loop {
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
// or I/O failure, and folding deadline expiries into
// it corrupts that counter.
debug!(
transport_id = %transport_id,
remote_addr = %remote_addr,
deadline = InboundDeadline::phase(first),
timeout_secs = limit.as_secs_f64(),
"No complete frame within the inbound deadline, dropping inbound {} connection",
label
);
break;
}
}
}
None => read_fmp_packet(&mut reader, mtu).await,
};
first = false;
match read {
Ok(data) => {
stats.record_recv(data.len());
trace!(
transport_id = %transport_id,
remote_addr = %remote_addr,
bytes = data.len(),
"{} packet received",
label
);
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 {} receive loop",
label
);
break;
}
}
Err(e) => {
stats.record_recv_error();
debug!(
transport_id = %transport_id,
remote_addr = %remote_addr,
error = %e,
"{} receive error, removing connection",
label
);
break;
}
}
}
}
// Clean up: remove ourselves from the pool, then run per-transport
// teardown accounting. The teardown fires only when this loop actually
// removed the entry, using the metadata from the removed entry, so a
// concurrent close/stop teardown of the same address can never
// double-count.
let mut pool_guard = pool.lock().await;
if let Some(removed) = pool_guard.remove(&remote_addr) {
drop(pool_guard);
on_remove(&*stats, &removed.meta);
}
}