From 5f60d2b02e3b91492059f01231a18d248320251b Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Sun, 4 Oct 2026 22:29:04 +0000 Subject: [PATCH] Finish the non-dialing send on master's own paths, and tidy what the merge left The two node-side inline sends on the UDP encrypt-worker fallback paths use send_existing. Both only ever see a UDP handle, for which send_existing forwards to the same send, so behaviour is unchanged; it leaves no node-side caller of the dialing send. The medium-change heartbeat note said the send after a failed write redials the peer. The link send no longer dials: once the stranded connection is evicted, the next send fails at once and starts a background connect, and only toward an address this node dialed. The classification test lists NotConnected among the terminal errors, so flipping it to transient fails a test that names the change. The rx-stall tests drop the blocking read timeouts their polling read never used, the Nym send_existing tests reuse the existing mock-proxy fixture, and the SOCKS5 pool tests use the shared wait_until helper. --- src/node/handlers/netmon.rs | 11 ++++++----- src/node/handlers/session.rs | 2 +- src/node/mod.rs | 2 +- src/node/tests/rx_stall.rs | 9 +-------- src/transport/mod.rs | 2 ++ src/transport/nym/mod.rs | 38 ++---------------------------------- src/transport/socks5/pool.rs | 15 +------------- 7 files changed, 14 insertions(+), 65 deletions(-) diff --git a/src/node/handlers/netmon.rs b/src/node/handlers/netmon.rs index fb77735d..8e658c10 100644 --- a/src/node/handlers/netmon.rs +++ b/src/node/handlers/netmon.rs @@ -167,11 +167,12 @@ impl Node { /// /// So a peer on TCP, Tor, Nym or BLE keeps the periodic heartbeat it had /// before this detector existed, and `link_dead_timeout_secs` remains the - /// backstop. Note that it does *not* recover by redialling: `send_async` - /// only dials when the pool holds no connection for the address, and a - /// connection stranded by a medium change is still in the pool. It is - /// evicted after a write to it fails, so the redial happens on the send - /// after the failure, not on the first one. Doing better for them means + /// backstop. Note that it does *not* recover by redialling. The link send + /// never dials, and a connection stranded by a medium change is still in + /// the pool until its writer's write to it fails. The send after that + /// finds no connection and fails, starting a background connect only if + /// this node dialed the peer at that address, and a later send uses it; + /// a peer that dialed in is never redialled. Doing better for them means /// dropping the stale connection /// rather than writing into it, which is a different change with a real /// cost behind it — a Tor peer pays a fresh circuit — and is not this one. diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index bf2aaec6..422f9b55 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -2930,7 +2930,7 @@ impl Node { debug!(next_hop = %next_hop_addr, "Transport gone before inline send of session data"); return; }; - if let Err(error) = transport.send(remote_addr, &wire).await { + if let Err(error) = transport.send_existing(remote_addr, &wire).await { debug!(next_hop = %next_hop_addr, %error, "Inline send of session data failed"); } } diff --git a/src/node/mod.rs b/src/node/mod.rs index f12b60e6..91c0d21b 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -3941,7 +3941,7 @@ impl Node { reason: format!("encryption failed: {}", e), })?; transport - .send(&remote_addr, &wire) + .send_existing(&remote_addr, &wire) .await .map_err(|e| link_send_error(*node_addr, e))? } diff --git a/src/node/tests/rx_stall.rs b/src/node/tests/rx_stall.rs index 58711df3..08c2236e 100644 --- a/src/node/tests/rx_stall.rs +++ b/src/node/tests/rx_stall.rs @@ -148,11 +148,7 @@ async fn prime_link(node: &Node, bh: &Blackhole) -> std::net::TcpStream { tokio::time::sleep(Duration::from_millis(10)).await; } let (accepted, _) = bh.listener.accept().unwrap(); - let accepted = std::net::TcpStream::from(accepted); - accepted - .set_read_timeout(Some(Duration::from_millis(1000))) - .unwrap(); - accepted + std::net::TcpStream::from(accepted) } /// Close the node's connection to `bh` from the far end, after taking the @@ -581,9 +577,6 @@ async fn msg1_resend_to_dead_outbound_leg_recovers_after_background_connect() { } let (accepted, _) = bh.listener.accept().unwrap(); let mut accepted = std::net::TcpStream::from(accepted); - accepted - .set_read_timeout(Some(Duration::from_millis(1000))) - .unwrap(); println!( "msg1 resend sent after {:?} over the background connect", start.elapsed() diff --git a/src/transport/mod.rs b/src/transport/mod.rs index cd9286e4..a0227300 100644 --- a/src/transport/mod.rs +++ b/src/transport/mod.rs @@ -1996,6 +1996,8 @@ mod tests { // answer, not this node's inability to transmit. TransportError::Timeout, TransportError::ConnectionRefused, + // The connection is gone and nothing will bring it back. + TransportError::NotConnected, ] { assert!( !terminal.is_transient(), diff --git a/src/transport/nym/mod.rs b/src/transport/nym/mod.rs index c7cbf3c2..0c0a34be 100644 --- a/src/transport/nym/mod.rs +++ b/src/transport/nym/mod.rs @@ -1286,45 +1286,11 @@ mod tests { dest.stop_async().await.unwrap(); } - /// A destination TCP transport behind a mock SOCKS5 proxy, and a - /// started Nym transport pointed at the proxy. - async fn nym_behind_mock_proxy() -> ( - NymTransport, - TcpTransport, - crate::transport::PacketRx, - TransportAddr, - ) { - let (dest_tx, 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 (tx, _rx) = packet_channel(32); - let config = NymConfig { - socks5_addr: Some(proxy_addr.to_string()), - startup_timeout_secs: Some(5), - connect_timeout_ms: Some(5000), - ..Default::default() - }; - let mut t = NymTransport::new(TransportId::new(200), None, config, tx); - t.start_async().await.unwrap(); - let target = TransportAddr::from_string(&dest_addr.to_string()); - (t, dest, dest_rx, target) - } - /// With no pooled connection and no connect under way, `send_existing` /// fails with `NotConnected` and opens nothing. #[tokio::test] async fn send_existing_without_connection_fails_fast_and_dials_nothing() { - let (mut t, mut dest, _dest_rx, target) = nym_behind_mock_proxy().await; + let (mut dest, _dest_rx, mut t, target) = nym_via_mock_proxy().await; let result = t.send_existing(&target, &build_msg1_frame()).await; @@ -1352,7 +1318,7 @@ mod tests { /// carries the send, with no second connection opened. #[tokio::test] async fn send_existing_promotes_a_finished_background_connect_and_sends_on_it() { - let (mut t, mut dest, mut dest_rx, target) = nym_behind_mock_proxy().await; + let (mut dest, mut dest_rx, mut t, target) = nym_via_mock_proxy().await; t.connect_async(&target).await.unwrap(); let mut waited = 0; diff --git a/src/transport/socks5/pool.rs b/src/transport/socks5/pool.rs index d348184f..6163337d 100644 --- a/src/transport/socks5/pool.rs +++ b/src/transport/socks5/pool.rs @@ -417,6 +417,7 @@ pub(crate) async fn proxied_receive_loop( #[cfg(test)] mod tests { use super::*; + use crate::testutil::wait_until; use crate::transport::packet_channel; use crate::transport::stream::{next_conn_id, park_writer}; use portable_atomic::{AtomicU64, Ordering}; @@ -463,20 +464,6 @@ mod tests { } } - /// 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; - } - } - /// A writer whose write fails must not remove a newer connection that has /// taken its address in the pool, nor run `on_remove` for it. #[tokio::test]