fix(transport): let a deliberate close finish writing queued frames

Sending on a TCP, Tor or Nym connection only queues the frame for the
connection's writer task, and closing the connection aborted that task. A
frame sent just before a close was never written if the writer had not run in
between, and on a single-threaded runtime it had not. A control-API disconnect
queues a Disconnect and then closes the connection, so the Disconnect was
lost and the peer kept this node until its own link-dead detection removed it.

A deliberate close now drops the connection's queue and aborts only its
receive task. The writer writes what was queued and exits, which closes the
stream. A detached timer aborts the writer if it is still writing after five
seconds, so the close never waits on the peer. Stopping the transport, and
every teardown after a connection has failed, still abort the writer at once.
A writer that outlives its pool entry this way removes only an entry carrying
its own connection id, so it cannot tear down a newer connection at the same
address.

The new tests were run on the unfixed tree first. Red there: a frame queued
just before close did not arrive within two seconds over TCP, Tor or Nym, and
after an api_disconnect over TCP the peer still had this node after two
seconds of packet processing. Also added, and not red before this change
because they guard what it adds: a close whose writer is parked on a peer that
has stopped reading returns within 200 ms, and the drain timer both aborts a
writer that outlives its bound and ends as soon as the writer exits.

Break-checks, each reverting one element and running the tests named for it:
aborting the writer again in the TCP close reds the TCP test and the
api_disconnect test, and in the Tor or Nym close reds that transport's test;
a drain timer that never aborts reds its abort test; a close that awaits the
drain reds the 200 ms test; and giving the TCP accept loop's receive loop a
different id from its entry reds the TCP test's check that the peer releases
its inbound slot after the close.
This commit is contained in:
Johnathan Corgan
2026-09-14 16:52:00 +00:00
parent 3e77aa41ef
commit 0c932635ae
6 changed files with 350 additions and 15 deletions
+5 -1
View File
@@ -399,7 +399,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
receive loop act only on their own connection: one that outlives its
connection can no longer tear down a newer connection that has taken the
same address, and a receive loop that ends stops its writer rather than
leaving it writing to a peer that has gone.
leaving it writing to a peer that has gone. Closing one of their connections
on purpose, as a control-API disconnect does, now lets the writer finish the
frames already queued, within five seconds, instead of discarding them, so a
Disconnect sent just before the close reaches the peer. Stopping the
transport, or a connection that has failed, still discards them.
- A per-peer `connect()`-ed UDP socket is no longer left pinned to an interface
the host has moved off. Established UDP peers get their own socket for the
+40
View File
@@ -459,3 +459,43 @@ async fn test_api_disconnect_closes_the_tcp_connection() {
cleanup_nodes(&mut nodes).await;
}
/// The Disconnect `api_disconnect` sends must reach the peer before the TCP
/// connection closes, so the peer forgets this node at once.
///
/// The send only queues the Disconnect for the connection's writer, and the
/// close follows in the same call. Nothing else removes the peer on node 1
/// within the wait: its TCP EOF removes only the pool entry, and link-dead
/// detection takes far longer.
#[tokio::test]
async fn api_disconnect_delivers_the_disconnect_before_closing() {
let mut nodes = vec![make_test_node_tcp().await, make_test_node_tcp().await];
initiate_handshake(&mut nodes, 0, 1).await;
drain_all_packets(&mut nodes, false).await;
let addr_0 = *nodes[0].node.node_addr();
let node1_npub = nodes[1].node.npub();
assert!(
nodes[1].node.get_peer(&addr_0).is_some(),
"node 1 should have node 0 as peer"
);
nodes[0]
.node
.api_disconnect(&node1_npub)
.await
.expect("api_disconnect should succeed");
let deadline = std::time::Instant::now() + Duration::from_secs(2);
while nodes[1].node.get_peer(&addr_0).is_some() && std::time::Instant::now() < deadline {
spanning_tree::process_available_packets(&mut nodes).await;
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert!(
nodes[1].node.get_peer(&addr_0).is_none(),
"node 1 never received the Disconnect sent before the close"
);
cleanup_nodes(&mut nodes).await;
}
+54 -3
View File
@@ -24,7 +24,7 @@ use crate::transport::socks5::{
Socks5Auth, Socks5Dialer, SocksTarget, poll_connecting, proxied_receive_loop,
proxied_send_loop,
};
use crate::transport::stream::{ConnId, next_conn_id};
use crate::transport::stream::{ConnId, WRITER_DRAIN_TIMEOUT, drain_writer, next_conn_id};
use stats::NymStats;
use std::collections::HashMap;
@@ -585,11 +585,22 @@ impl NymTransport {
}
/// Close a specific connection asynchronously.
///
/// Aborts the receive task and lets the writer finish the frames already
/// queued, within [`WRITER_DRAIN_TIMEOUT`], without waiting for it. This
/// mirrors `TcpTransport::close_connection_async`.
pub async fn close_connection_async(&self, addr: &TransportAddr) {
let mut pool = self.pool.lock().await;
if let Some(conn) = pool.remove(addr) {
conn.recv_task.abort();
conn.send_task.abort();
let ProxiedConnection {
send_tx,
send_task,
recv_task,
..
} = conn;
drop(send_tx);
recv_task.abort();
drain_writer(send_task, WRITER_DRAIN_TIMEOUT);
debug!(
transport_id = %self.transport_id,
remote_addr = %addr,
@@ -1188,4 +1199,44 @@ mod tests {
nym.stop_async().await.unwrap();
}
// ========================================================================
// Deliberate close finishes the frames already queued
// ========================================================================
/// A frame queued immediately before a deliberate close must still reach
/// the peer through the proxy, and the close must still end the
/// connection at the far side.
#[tokio::test]
async fn nym_frame_queued_just_before_close_still_reaches_the_peer() {
let (mut dest, mut dest_rx, mut nym, target) = nym_via_mock_proxy().await;
let frame = build_msg1_frame();
nym.send_async(&target, &frame).await.unwrap();
let first = tokio::time::timeout(Duration::from_secs(2), dest_rx.recv())
.await
.expect("timeout waiting for the first frame")
.expect("channel closed");
assert_eq!(first.data, frame);
nym.send_async(&target, &frame).await.unwrap();
nym.close_connection_async(&target).await;
let second = tokio::time::timeout(Duration::from_secs(2), dest_rx.recv())
.await
.expect("a frame queued just before close was never written")
.expect("channel closed");
assert_eq!(second.data, frame);
assert!(
wait_until(
|| dest.stats().snapshot().pool_inbound == 0,
Duration::from_secs(5)
)
.await,
"the close must still end the connection once the queue is written"
);
nym.stop_async().await.unwrap();
dest.stop_async().await.unwrap();
}
}
+68
View File
@@ -6,8 +6,10 @@
//! may touch the pool are written once here.
use std::collections::HashMap;
use std::time::Duration;
use portable_atomic::{AtomicU64, Ordering};
use tokio::task::JoinHandle;
use crate::transport::TransportAddr;
@@ -49,3 +51,69 @@ pub(crate) fn remove_own<C: PooledConn>(
}
pool.remove(addr)
}
/// How long a deliberately closed connection's writer may keep writing the
/// frames already queued before it is stopped.
///
/// It bounds only how long the socket and the writer task outlive the close;
/// no caller waits on it. A peer that is still reading drains a full queue in
/// far less, and one that has not drained it by then has stopped reading.
pub(crate) const WRITER_DRAIN_TIMEOUT: Duration = Duration::from_secs(5);
/// Let a closed connection's writer finish the frames already queued, and stop
/// it if it is still running after `bound`.
///
/// The caller must already have dropped the connection's queue, so the writer
/// exits once it has written what was queued. The wait runs on its own task,
/// which ends as soon as the writer does; the returned handle is that task's.
pub(crate) fn drain_writer(send_task: JoinHandle<()>, bound: Duration) -> JoinHandle<()> {
tokio::spawn(async move {
let mut send_task = send_task;
if tokio::time::timeout(bound, &mut send_task).await.is_err() {
send_task.abort();
}
})
}
#[cfg(test)]
mod tests {
use super::*;
/// A writer still running when the bound expires is stopped, and the
/// timer ends with it.
///
/// The writer holds a oneshot sender and never finishes, so the sender is
/// dropped only if the task is aborted.
#[tokio::test]
async fn drain_writer_aborts_a_writer_that_outlives_the_bound() {
let (guard_tx, guard_rx) = tokio::sync::oneshot::channel::<()>();
let writer = tokio::spawn(async move {
let _guard = guard_tx;
std::future::pending::<()>().await
});
let timer = drain_writer(writer, Duration::from_millis(50));
assert!(
matches!(
tokio::time::timeout(Duration::from_secs(1), guard_rx).await,
Ok(Err(_))
),
"a writer still running at the bound was left running"
);
tokio::time::timeout(Duration::from_secs(1), timer)
.await
.expect("the drain timer outlived the writer it stopped")
.unwrap();
}
/// The wait ends when the writer does, not when the bound expires.
#[tokio::test]
async fn drain_writer_ends_as_soon_as_the_writer_does() {
let writer = tokio::spawn(async {});
let timer = drain_writer(writer, Duration::from_secs(30));
tokio::time::timeout(Duration::from_millis(100), timer)
.await
.expect("the drain timer kept running after the writer exited")
.unwrap();
}
}
+107 -7
View File
@@ -32,7 +32,9 @@ use super::{
};
use crate::config::TcpConfig;
use crate::transport::framing::read_fmp_packet;
use crate::transport::stream::{ConnId, next_conn_id, remove_own};
use crate::transport::stream::{
ConnId, WRITER_DRAIN_TIMEOUT, drain_writer, next_conn_id, remove_own,
};
use pool::{ConnectingEntry, ConnectingPool, ConnectionPool, Direction, TcpConnection};
use stats::TcpStats;
@@ -486,21 +488,35 @@ impl TcpTransport {
/// Close a specific connection asynchronously.
///
/// Removes the connection from the pool, aborts its receive task,
/// and drops the write half (sends FIN to remote).
/// Removes the connection from the pool and aborts its receive task. The
/// writer is not aborted: dropping the queue lets it finish writing the
/// frames already queued, such as a Disconnect sent just before this close,
/// and then exit, which drops the write half and sends FIN. A detached
/// timer aborts it if it is still writing after [`WRITER_DRAIN_TIMEOUT`],
/// so this call never waits on the peer. Stopping the transport, and every
/// teardown after a connection has failed, abort the writer instead and
/// discard what it had queued.
pub async fn close_connection_async(&self, addr: &TransportAddr) {
let mut pool = self.pool.lock().await;
if let Some(conn) = pool.remove(addr) {
conn.recv_task.abort();
conn.send_task.abort();
match conn.direction {
let TcpConnection {
send_tx,
send_task,
recv_task,
direction,
..
} = conn;
drop(send_tx);
recv_task.abort();
drain_writer(send_task, WRITER_DRAIN_TIMEOUT);
match direction {
Direction::Inbound => self.stats.record_pool_inbound_removed(),
Direction::Outbound => self.stats.record_pool_outbound_removed(),
}
debug!(
transport_id = %self.transport_id,
remote_addr = %addr,
direction = ?conn.direction,
direction = ?direction,
"TCP connection closed (close_connection)"
);
}
@@ -2831,4 +2847,88 @@ mod tests {
transport.stop_async().await.unwrap();
}
// ========================================================================
// Deliberate close finishes the frames already queued
// ========================================================================
/// A frame queued immediately before a deliberate close must still be
/// written.
///
/// Sending only queues the frame for the connection's writer task. A close
/// that aborts that task before it has run discards the frame, which is
/// how a Disconnect sent just before a close, or a handshake message sent
/// just before the losing side of a crossed connection is closed, never
/// reaches the peer. The second half checks the close still closes: the
/// peer sees FIN and releases its inbound slot, so a writer that never
/// exits cannot pass.
#[tokio::test]
async fn a_frame_queued_just_before_close_still_reaches_the_peer() {
let (tx1, _rx1) = packet_channel(100);
let (tx2, mut rx2) = packet_channel(100);
let mut t1 = TcpTransport::new(TransportId::new(1), None, make_outbound_config(), tx1);
let mut t2 = TcpTransport::new(TransportId::new(2), None, make_config(), tx2);
t1.start_async().await.unwrap();
t2.start_async().await.unwrap();
let remote = TransportAddr::from_string(&t2.local_addr().unwrap().to_string());
let frame = build_msg1_frame();
// Pool the connection and let its writer go idle.
t1.send_async(&remote, &frame).await.unwrap();
let first = timeout(Duration::from_secs(2), rx2.recv())
.await
.expect("timeout waiting for the first frame")
.expect("packet channel closed");
assert_eq!(first.data, frame);
// Queue, then close with nothing in between.
t1.send_async(&remote, &frame).await.unwrap();
t1.close_connection_async(&remote).await;
let second = timeout(Duration::from_secs(2), rx2.recv())
.await
.expect("a frame queued just before close was never written")
.expect("packet channel closed");
assert_eq!(second.data, frame);
assert!(
wait_until(
|| t2.stats().snapshot().pool_inbound == 0,
Duration::from_secs(5)
)
.await,
"the close must still end the connection once the queue is written"
);
t1.stop_async().await.unwrap();
t2.stop_async().await.unwrap();
}
/// A deliberate close must return at once even when the writer cannot
/// finish, because the peer has stopped reading. Draining happens after
/// the close returns, never inside it.
#[tokio::test]
async fn close_does_not_wait_for_a_writer_parked_on_a_deaf_peer() {
let (tx1, _rx1) = packet_channel(100);
let mut t1 = TcpTransport::new(TransportId::new(1), None, make_outbound_config(), tx1);
t1.start_async().await.unwrap();
let listener = capped_deaf_listener();
let remote = TransportAddr::from_string(&listener.local_addr().unwrap().to_string());
let frame = vec![0xAB; 1400];
park_tcp_writer(&t1, &remote, &frame).await;
let (_peer, _) = listener.accept().await.unwrap();
assert!(
timeout(
Duration::from_millis(200),
t1.close_connection_async(&remote)
)
.await
.is_ok(),
"close waited on a writer that cannot finish"
);
t1.stop_async().await.unwrap();
}
}
+76 -4
View File
@@ -34,7 +34,7 @@ use crate::transport::socks5::{
Socks5Auth, Socks5Dialer, SocksTarget, poll_connecting, proxied_receive_loop,
proxied_send_loop,
};
use crate::transport::stream::{ConnId, next_conn_id};
use crate::transport::stream::{ConnId, WRITER_DRAIN_TIMEOUT, drain_writer, next_conn_id};
use crate::transport::tcp::INBOUND_FIRST_FRAME_TIMEOUT;
use control::{ControlAuth, TorControlClient, TorMonitoringInfo};
use stats::TorStats;
@@ -1016,12 +1016,24 @@ impl TorTransport {
}
/// Close a specific connection asynchronously.
///
/// Aborts the receive task and lets the writer finish the frames already
/// queued, within [`WRITER_DRAIN_TIMEOUT`], without waiting for it. This
/// mirrors `TcpTransport::close_connection_async`.
pub async fn close_connection_async(&self, addr: &TransportAddr) {
let mut pool = self.pool.lock().await;
if let Some(conn) = pool.remove(addr) {
conn.recv_task.abort();
conn.send_task.abort();
match conn.meta {
let ProxiedConnection {
send_tx,
send_task,
recv_task,
meta,
..
} = conn;
drop(send_tx);
recv_task.abort();
drain_writer(send_task, WRITER_DRAIN_TIMEOUT);
match meta {
Direction::Inbound => self.stats.record_pool_inbound_removed(),
Direction::Outbound => self.stats.record_pool_outbound_removed(),
}
@@ -2556,4 +2568,64 @@ mod tests {
accept.abort();
drop(sock);
}
// ========================================================================
// Deliberate close finishes the frames already queued
// ========================================================================
/// A frame queued immediately before a deliberate close must still reach
/// the peer through the proxy, and the close must still end the
/// connection at the far side.
#[tokio::test]
async fn tor_frame_queued_just_before_close_still_reaches_the_peer() {
let (dest_tx, mut dest_rx) = packet_channel(32);
let dest_config = TcpConfig {
bind_addr: Some("127.0.0.1:0".to_string()),
..Default::default()
};
let mut dest = TcpTransport::new(TransportId::new(100), None, dest_config, dest_tx);
dest.start_async().await.unwrap();
let dest_addr = dest.local_addr().unwrap();
let mock = MockSocks5Server::new(dest_addr).await.unwrap();
let proxy_addr = mock.addr();
let _proxy_handle = mock.spawn();
let (tor_tx, _tor_rx) = packet_channel(32);
let tor_config = TorConfig {
socks5_addr: Some(proxy_addr.to_string()),
..Default::default()
};
let mut tor = TorTransport::new(TransportId::new(200), None, tor_config, tor_tx);
tor.start_async().await.unwrap();
let target = TransportAddr::from_string(&dest_addr.to_string());
let frame = build_msg1_frame();
tor.send_async(&target, &frame).await.unwrap();
let first = tokio::time::timeout(Duration::from_secs(2), dest_rx.recv())
.await
.expect("timeout waiting for the first frame")
.expect("channel closed");
assert_eq!(first.data, frame);
tor.send_async(&target, &frame).await.unwrap();
tor.close_connection_async(&target).await;
let second = tokio::time::timeout(Duration::from_secs(2), dest_rx.recv())
.await
.expect("a frame queued just before close was never written")
.expect("channel closed");
assert_eq!(second.data, frame);
assert!(
wait_until(
|| dest.stats().snapshot().pool_inbound == 0,
Duration::from_secs(5)
)
.await,
"the close must still end the connection once the queue is written"
);
tor.stop_async().await.unwrap();
dest.stop_async().await.unwrap();
}
}