diff --git a/CHANGELOG.md b/CHANGELOG.md index 706e136f..22fcfd3b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/src/node/tests/tcp.rs b/src/node/tests/tcp.rs index 669106ad..6c14bbf5 100644 --- a/src/node/tests/tcp.rs +++ b/src/node/tests/tcp.rs @@ -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; +} diff --git a/src/transport/nym/mod.rs b/src/transport/nym/mod.rs index 821a1b03..f1a5dea0 100644 --- a/src/transport/nym/mod.rs +++ b/src/transport/nym/mod.rs @@ -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(); + } } diff --git a/src/transport/stream.rs b/src/transport/stream.rs index 3bb62864..103ed68d 100644 --- a/src/transport/stream.rs +++ b/src/transport/stream.rs @@ -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( } 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(); + } +} diff --git a/src/transport/tcp/mod.rs b/src/transport/tcp/mod.rs index ec7b48e3..1c63d545 100644 --- a/src/transport/tcp/mod.rs +++ b/src/transport/tcp/mod.rs @@ -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(); + } } diff --git a/src/transport/tor/mod.rs b/src/transport/tor/mod.rs index b194f600..ed93dcb6 100644 --- a/src/transport/tor/mod.rs +++ b/src/transport/tor/mod.rs @@ -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(); + } }