From e65950bc5e1b5425b20f807c14a1bde8f3506ea3 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Sun, 4 Oct 2026 22:29:04 +0000 Subject: [PATCH] Send the XX handshake replies and remaining tick sends without dialing The merge from master brought the non-dialing send and switched the call sites whose code next shares with master: the rekey msg1 and its resend, the msg1 resend on an outbound handshake, the executor's stored msg1 send, the encrypted link send and the executor's msg2 arm (unreachable on next, since handle_msg1 sends its msg2 inline). The XX handshake's own sites still called the dialing transport send, so a reply to a connection that had closed held the rx loop for the whole connect timeout, as in #176. The replies now use send_existing and never dial: the msg2 to a msg1, the msg2 resent for a duplicate msg1, the msg3 to a dial's msg2 and to a rekey msg2, the msg2 resent for a duplicate handshake at msg3, and the Disconnect sent to a peer the access list refuses at msg3. The rekey msg3 resend and the anonymous dial's inline msg1 go through send_nowait, which starts a background connect only toward an address this node dialed, and the tests cover that for a dialed peer and an inbound one. process_packet and start_handshake are visible to the node tests, which call them directly. The tests in src/node/tests/rx_stall.rs are the XX counterparts of master's. --- src/node/dataplane/rx_loop.rs | 2 +- src/node/handlers/handshake.rs | 38 +- src/node/handlers/rekey.rs | 10 +- src/node/lifecycle/mod.rs | 12 +- src/node/mod.rs | 5 +- src/node/tests/mod.rs | 1 + src/node/tests/rx_stall.rs | 1408 ++++++++++++++++++++++++++++++++ 7 files changed, 1459 insertions(+), 17 deletions(-) create mode 100644 src/node/tests/rx_stall.rs diff --git a/src/node/dataplane/rx_loop.rs b/src/node/dataplane/rx_loop.rs index 26948e4b..de3f6a81 100644 --- a/src/node/dataplane/rx_loop.rs +++ b/src/node/dataplane/rx_loop.rs @@ -619,7 +619,7 @@ impl Node { /// Process a single received packet. /// /// Dispatches based on the phase field in the 4-byte common prefix. - async fn process_packet(&mut self, packet: ReceivedPacket) { + pub(in crate::node) async fn process_packet(&mut self, packet: ReceivedPacket) { if packet.data.len() < COMMON_PREFIX_SIZE { return; // Drop packets too short for common prefix } diff --git a/src/node/handlers/handshake.rs b/src/node/handlers/handshake.rs index 3b8b080b..fe84fab8 100644 --- a/src/node/handlers/handshake.rs +++ b/src/node/handlers/handshake.rs @@ -329,11 +329,15 @@ impl Node { "Msg1 differs from the one the pending handshake at this address answered; starting a new handshake" ); } else { - // Genuinely pending handshake — resend msg2 + // Genuinely pending handshake — resend msg2. Like every + // reply on the rx loop it never dials: with the msg1's + // connection gone, a dial to its address (an inbound + // connection's ephemeral port) would hold the loop for up + // to the connect timeout. let msg2_bytes = self.find_stored_msg2(existing_link_id); if let Some(msg2) = msg2_bytes { if let Some(transport) = self.transports.get(&packet.transport_id) { - match transport.send(&packet.remote_addr, &msg2).await { + match transport.send_existing(&packet.remote_addr, &msg2).await { Ok(_) => debug!( remote_addr = %packet.remote_addr, "Resent msg2 for duplicate msg1" @@ -469,8 +473,15 @@ impl Node { machine.set_conn_handshake_msg1(packet.data.clone(), 0); self.peer_machines.insert(link_id, machine); + // The msg2 goes back on the msg1's connection and never dials: if that + // connection has closed, a dial to its address (for an inbound + // connection, the initiator's ephemeral port) would hold the rx loop + // for up to the connect timeout. The failure tears the leg down below. if let Some(transport) = self.transports.get(&packet.transport_id) { - match transport.send(&packet.remote_addr, &wire_msg2).await { + match transport + .send_existing(&packet.remote_addr, &wire_msg2) + .await + { Ok(bytes) => { debug!( link_id = %link_id, @@ -686,14 +697,18 @@ impl Node { peer.set_remote_epoch(remote_epoch); } - // Send msg3 before setting pending session + // Send msg3 before setting pending session. + // A reply on the rx loop, so it never dials; + // with the link's connection gone the send + // fails at once and the rekey is abandoned + // below, as on any send failure. let wire_msg3 = build_msg3(our_index, header.sender_idx, &msg3_bytes); let msg3_sent = if let (Some(tid), Some(addr)) = (transport_id, &remote_addr) && let Some(transport) = self.transports.get(&tid) { - match transport.send(addr, &wire_msg3).await { + match transport.send_existing(addr, &wire_msg3).await { Ok(_) => { debug!( peer = %display_name, @@ -1117,12 +1132,17 @@ impl Node { return; } - // Build and send msg3 + // Build and send msg3, on the msg2's connection. A reply on the rx + // loop, so it never dials: with that connection gone the send fails + // at once and the sweep reclaims the handshake, as on any send failure. let our_index = our_index.unwrap_or(header.receiver_idx); let wire_msg3 = build_msg3(our_index, header.sender_idx, &msg3_bytes); if let Some(transport) = self.transports.get(&packet.transport_id) { - match transport.send(&packet.remote_addr, &wire_msg3).await { + match transport + .send_existing(&packet.remote_addr, &wire_msg3) + .await + { Ok(bytes) => { debug!( peer = %self.peer_display_name(&peer_node_addr), @@ -1886,11 +1906,11 @@ impl Node { // Not a rekey — duplicate handshake from same epoch. Resend the // stored msg2 bytes as-is (a driver mechanism: replaying the // stored frame, not rebuilding it), leaving the active peer - // untouched. + // untouched. On the msg3's connection, never dialing. if let Some(msg2) = msg2 && let Some(transport) = self.transports.get(&packet.transport_id) { - match transport.send(&packet.remote_addr, &msg2).await { + match transport.send_existing(&packet.remote_addr, &msg2).await { Ok(_) => debug!( peer = %self.peer_display_name(&peer_node_addr), "Resent msg2 for duplicate handshake (same epoch)" diff --git a/src/node/handlers/rekey.rs b/src/node/handlers/rekey.rs index b8cba908..fb4a3b1d 100644 --- a/src/node/handlers/rekey.rs +++ b/src/node/handlers/rekey.rs @@ -650,16 +650,20 @@ impl Node { } for (node_addr, payload) in to_resend { - let (transport_id, remote_addr) = match self.peers.get(&node_addr) { + let (link_id, transport_id, remote_addr) = match self.peers.get(&node_addr) { Some(p) => match (p.transport_id(), p.current_addr()) { - (Some(tid), Some(addr)) => (tid, addr.clone()), + (Some(tid), Some(addr)) => (p.link_id(), tid, addr.clone()), _ => continue, }, None => continue, }; + // A failed send records no resend, so the msg3 stays due and is + // retried next tick, over any connection the failed send started. let sent = if let Some(transport) = self.transports.get(&transport_id) { - transport.send(&remote_addr, &payload).await.is_ok() + self.send_nowait(transport, link_id, &remote_addr, &payload) + .await + .is_ok() } else { false }; diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index ab248ecd..9f9eb687 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -782,7 +782,7 @@ impl Node { /// connect-resolution). Anonymous discovery (no `peer_identity`) leaves /// identity to be learned from the XX msg2, which crystallizes it onto the /// leg-born machine. - pub(super) async fn start_handshake( + pub(in crate::node) async fn start_handshake( &mut self, link_id: LinkId, transport_id: TransportId, @@ -877,9 +877,15 @@ impl Node { // already carry the msg1-prep provenance. self.peer_machines.insert(link_id, machine); - // Send the wire format handshake message + // Send the wire format handshake message. It never dials: if the + // connection the dial resolved to has gone, the send fails at once + // into the failure path below, after starting a background connect + // to the dial address. if let Some(transport) = self.transports.get(&transport_id) { - match transport.send(&remote_addr, &wire_msg1).await { + match self + .send_nowait(transport, link_id, &remote_addr, &wire_msg1) + .await + { Ok(bytes) => { debug!( link_id = %link_id, diff --git a/src/node/mod.rs b/src/node/mod.rs index 1d277561..5f63cdfc 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -4182,8 +4182,11 @@ impl Node { .get(&transport_id) .ok_or(NodeError::TransportNotFound(transport_id))?; + // The one caller answers a msg3 on the rx loop, so this never dials: + // a dial to the msg3's address, with its connection gone, would hold + // the loop for up to the connect timeout. transport - .send(remote_addr, &wire_packet) + .send_existing(remote_addr, &wire_packet) .await .map(|_| ()) .map_err(|e| match e { diff --git a/src/node/tests/mod.rs b/src/node/tests/mod.rs index eacb0894..075306b7 100644 --- a/src/node/tests/mod.rs +++ b/src/node/tests/mod.rs @@ -24,6 +24,7 @@ mod mmp_chartests; mod netmon; mod probe; mod routing; +mod rx_stall; mod session; mod spanning_tree; mod tcp; diff --git a/src/node/tests/rx_stall.rs b/src/node/tests/rx_stall.rs new file mode 100644 index 00000000..b0833b87 --- /dev/null +++ b/src/node/tests/rx_stall.rs @@ -0,0 +1,1408 @@ +//! A handshake send to a TCP connection that has gone away must not dial. +//! +//! When a handshake message arrives on an inbound TCP connection that has +//! since closed, the reply has nowhere to go. The handlers that answer it run +//! inline on the rx loop, so a reply that fell through to TCP connect-on-send +//! held every other frame the loop owns, including frames that arrived on +//! UDP, for the whole connect timeout. On the XX handshake the replies are the +//! msg2 to a msg1 (fresh or resent), the msg3 to a msg2 (a dial's or a +//! rekey's), the msg2 resent for a duplicate handshake at msg3, and the +//! Disconnect sent to a peer refused at msg3. The tick's handshake sends (the +//! rekey msg1 and its resends, the rekey msg3 resend, the msg1 resend on an +//! outbound handshake), the dial path's msg1 sends and the encrypted link send +//! are awaited by the same loop and had the same exposure. These tests assert +//! each send now fails at once instead: no connect attempt is counted and the +//! call returns well inside a bound far below the timeout. Each one also runs +//! a healthy control, so a run that skips the send path entirely fails too. +//! +//! The tick and dial-path sends may start a background connect, but only to +//! an address this node dialed; the tests check that it starts there, that +//! it does not start toward an inbound peer's address, and that a later send +//! uses the connection once it is up. +//! +//! The far end of each connection is played by the test: it reads the node's +//! frames off the socket and answers them with the XX messages a peer would +//! send. +//! +//! The unanswered SYN is constructed locally: a listener with a backlog of +//! zero whose single accept slot is already taken. Linux drops further SYNs to +//! a listener whose accept queue is full, so a connect to it times out rather +//! than being refused. `Blackhole::silent()` checks that before any test +//! relies on it, which is what lets a regression show up at its real size: +//! one connect timeout per reply, counted in `connect_timeouts`. +//! +//! The tests print their measurements; run with `--nocapture` to see them. + +use super::*; +use crate::config::{TcpConfig, UdpConfig}; +use crate::node::acl::PeerAclReloader; +use crate::peer::machine::{PeerEvent, PeerMachine, TimerKind}; +use crate::proto::fmp::NegotiationPayload; +use crate::proto::fmp::wire::{ + CommonPrefix, Msg1Header, Msg2Header, PHASE_ESTABLISHED, PHASE_MSG1, PHASE_MSG2, PHASE_MSG3, + build_msg1, build_msg2, build_msg3, +}; +use crate::proto::link::LinkMessageType; +use crate::testutil::Blackhole; +use crate::transport::framing::read_fmp_packet; +use crate::transport::tcp::TcpTransport; +use crate::transport::tcp::stats::TcpStatsSnapshot; +use crate::transport::udp::UdpTransport; +use crate::transport::{ConnectionState, PacketTx, TransportHandle, TransportId}; +use std::net::SocketAddr; +use std::time::Instant; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpStream; +use tokio::time::timeout; + +const UDP_ID: u32 = 1; +const TCP_ID: u32 = 2; +const EPOCH: [u8; 8] = [7u8; 8]; +const TCP_MTU: u16 = 1400; + +/// The connect timeout every dead-link test runs at: the shipped default, so +/// a send that dials shows up at the size it has in the field. +const CONNECT_TIMEOUT_MS: u64 = 5000; + +/// How long a handler answering a dead link may take. Far below any connect +/// timeout, far above the microseconds a failed pool lookup costs. +const BOUND: Duration = Duration::from_millis(250); + +/// The negotiation payload a peer appends to its msg2 or msg3, optionally +/// declaring the handshake a rekey of the session the node indexes as +/// `rekey_of`. +fn negotiation(rekey_of: Option) -> Vec { + let payload = NegotiationPayload::fmp(1, 1, crate::proto::fmp::NodeProfile::Full); + match rekey_of { + Some(idx) => payload.with_rekey_of(idx).encode(), + None => payload.encode(), + } +} + +/// An XX initiator for `sender` dialing `node`: the machine that carries its +/// handshake, and the framed msg1 it sends under `sender_index`. +fn initiator(node: &Node, sender: &Identity, sender_index: u32) -> (PeerMachine, Vec) { + let target = PeerIdentity::from_pubkey_full(node.identity().pubkey_full()); + let mut leg = outbound_leg(LinkId::new(0x5EED), target, 1000); + let noise_msg1 = leg + .start_handshake(sender.keypair(), EPOCH, 1000) + .expect("start_handshake produces noise msg1"); + ( + leg, + build_msg1(SessionIndex::new(sender_index), &noise_msg1), + ) +} + +/// A genuine wire msg1 from a fresh identity to `node`. +fn craft_msg1(node: &Node, sender_index: u32) -> Vec { + initiator(node, &Identity::generate(), sender_index).1 +} + +/// The framed msg3 an initiator answers the framed `msg2` with, sent under +/// `sender_index`. +fn msg3_for( + leg: &mut PeerMachine, + msg2: &[u8], + sender_index: u32, + rekey_of: Option, +) -> Vec { + let header = Msg2Header::parse(msg2).expect("msg2 header"); + let neg = negotiation(rekey_of); + let (noise_msg3, _) = leg + .complete_handshake(header.noise_msg2(msg2), Some(&neg), 1100) + .expect("the initiator reads the node's msg2"); + build_msg3( + SessionIndex::new(sender_index), + header.sender_idx, + &noise_msg3, + ) +} + +/// The framed msg2 `responder` answers the framed `msg1` with, sent under +/// `sender_index`. +fn msg2_for(responder: &Identity, msg1: &[u8], sender_index: u32) -> Vec { + let header = Msg1Header::parse(msg1).expect("msg1 header"); + let mut leg = inbound_leg(LinkId::new(0x77), 1000); + let neg = negotiation(None); + let noise_msg2 = leg + .receive_handshake_init( + responder.keypair(), + EPOCH, + &msg1[header.noise_msg1_offset..], + Some(&neg), + 1000, + ) + .expect("the responder reads the node's msg1"); + build_msg2( + SessionIndex::new(sender_index), + header.sender_idx, + &noise_msg2, + ) +} + +/// A node with a UDP and a TCP transport feeding one packet channel, as a +/// node built from config has. Returns the node, a sender into that channel +/// (to inject frames as if a transport had delivered them), and the UDP +/// transport's local address. +async fn node_with_udp_and_tcp(connect_timeout_ms: u64) -> (Node, PacketTx, SocketAddr) { + node_from(make_node(), connect_timeout_ms).await +} + +/// [`node_with_udp_and_tcp`] on a node the caller built. +async fn node_from(mut node: Node, connect_timeout_ms: u64) -> (Node, PacketTx, SocketAddr) { + let (tx, rx) = packet_channel(1024); + + let udp_cfg = UdpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + mtu: Some(1280), + ..Default::default() + }; + let mut udp = UdpTransport::new(TransportId::new(UDP_ID), None, udp_cfg, tx.clone()); + udp.start_async().await.unwrap(); + let udp_addr = udp.local_addr().unwrap(); + + let tcp_cfg = TcpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + mtu: Some(TCP_MTU), + connect_timeout_ms: Some(connect_timeout_ms), + ..Default::default() + }; + let mut tcp = TcpTransport::new(TransportId::new(TCP_ID), None, tcp_cfg, tx.clone()); + tcp.start_async().await.unwrap(); + + node.transports + .insert(TransportId::new(UDP_ID), TransportHandle::Udp(udp)); + node.transports + .insert(TransportId::new(TCP_ID), TransportHandle::Tcp(tcp)); + node.packet_rx = Some(rx); + node.supervisor.state = NodeState::Running; + (node, tx, udp_addr) +} + +/// The node's TCP transport. +fn tcp(node: &Node) -> &TransportHandle { + node.transports + .get(&TransportId::new(TCP_ID)) + .expect("no TCP transport") +} + +/// The TCP transport's live counters. +fn tcp_stats(node: &Node) -> TcpStatsSnapshot { + match tcp(node) { + TransportHandle::Tcp(t) => t.stats().snapshot(), + _ => panic!("transport {TCP_ID} is not TCP"), + } +} + +/// Stop every transport the node holds. +async fn stop_all(node: &mut Node) { + for (_, t) in node.transports.iter_mut() { + t.stop().await.ok(); + } +} + +/// Wait until the TCP transport holds no connection to `addr`. +async fn wait_pool_gone(node: &Node, addr: &TransportAddr) { + let start = Instant::now(); + while tcp(node).connection_state(addr) != ConnectionState::None { + assert!( + start.elapsed() < Duration::from_secs(3), + "pool entry for {addr} never dropped" + ); + tokio::time::sleep(Duration::from_millis(10)).await; + } +} + +/// Take the connection waiting in `bh`'s accept queue as a runtime socket. +fn accept_end(bh: &Blackhole) -> TcpStream { + let (accepted, _) = bh.listener.accept().unwrap(); + let accepted = std::net::TcpStream::from(accepted); + accepted.set_nonblocking(true).unwrap(); + TcpStream::from_std(accepted).unwrap() +} + +/// Open a connection from the node to `bh`'s free accept slot, the way a +/// node ends up holding a live connection to a peer's address, and return +/// the far end of it. +async fn prime_link(node: &Node, bh: &Blackhole) -> TcpStream { + let addr = bh.transport_addr(); + tcp(node).connect(&addr).await.unwrap(); + let start = Instant::now(); + while tcp(node).connection_state(&addr) != ConnectionState::Connected { + assert!( + start.elapsed() < Duration::from_secs(3), + "connection to {addr} never came up" + ); + tokio::time::sleep(Duration::from_millis(10)).await; + } + accept_end(bh) +} + +/// Close the node's connection to `bh` from the far end, after taking the +/// listener's accept slot so that any later dial to the address hangs. +async fn kill_link(node: &Node, bh: &mut Blackhole, far_end: TcpStream) { + bh.fill(); + drop(far_end); + wait_pool_gone(node, &bh.transport_addr()).await; +} + +/// The next frame the node wrote to `far_end`, if one arrives within a +/// second. +async fn next_frame(far_end: &mut TcpStream) -> Option> { + timeout( + Duration::from_millis(1000), + read_fmp_packet(far_end, TCP_MTU), + ) + .await + .ok()? + .ok() +} + +/// The phase of a frame. +fn phase(frame: &[u8]) -> Option { + CommonPrefix::parse(frame).map(|p| p.phase) +} + +/// The next frame of `wanted` phase the node wrote to `far_end`, skipping +/// any others, if one arrives within a second of the last. +async fn frame_of(far_end: &mut TcpStream, wanted: u8) -> Option> { + loop { + let frame = next_frame(far_end).await?; + if phase(&frame) == Some(wanted) { + return Some(frame); + } + } +} + +/// [`frame_of`] on `far_end` when there is one. +async fn frame_maybe(far_end: Option<&mut TcpStream>, wanted: u8) -> Option> { + match far_end { + Some(stream) => frame_of(stream, wanted).await, + None => None, + } +} + +/// Read and discard what the node has written to `far_end` until it goes +/// quiet, so a later read sees only what the call under test sends. +async fn drain(far_end: &mut TcpStream) { + while let Ok(Ok(_)) = timeout( + Duration::from_millis(150), + read_fmp_packet(far_end, TCP_MTU), + ) + .await + {} +} + +/// What one handler call against a TCP reply address did. +#[derive(Debug)] +struct Reply { + elapsed: Duration, + connect_timeouts: u64, + connect_refused: u64, + connections_established: u64, +} + +/// One call into a node that may send on TCP. +enum Trigger { + /// A frame for the packet handler. + Packet(ReceivedPacket), + /// The tick's rekey msg1 check. + RekeyCheck, + /// The tick's rekey msg1 resend, with every resend due. + RekeyResend, + /// The tick's rekey msg3 resend, with every resend due. + RekeyMsg3Resend, + /// The tick's handshake timers, at this time. + PeerTimers(u64), + /// The executor's send of the msg1 armed on this outbound link. + StoredMsg1(LinkId, TransportAddr), + /// An anonymous dial's inline handshake start on this outbound link. + StartHandshake(LinkId, TransportAddr), + /// An encrypted link message (a heartbeat) to this peer. + LinkMessage(NodeAddr), +} + +/// Fire `trigger` once and measure it against the TCP counters. +async fn timed_fire(node: &mut Node, trigger: Trigger) -> Reply { + let before = tcp_stats(node); + let t0 = Instant::now(); + match trigger { + Trigger::Packet(packet) => node.process_packet(packet).await, + Trigger::RekeyCheck => node.check_rekey().await, + Trigger::RekeyResend => node.resend_pending_rekeys(Node::now_ms() + 60_000).await, + Trigger::RekeyMsg3Resend => { + node.resend_pending_fmp_rekey_msg3(Node::now_ms() + 60_000) + .await + } + Trigger::PeerTimers(now_ms) => node.drive_peer_timers(now_ms).await, + Trigger::StoredMsg1(link, addr) => { + node.send_stored_msg1(link, TransportId::new(TCP_ID), &addr, Node::now_ms()) + .await + } + Trigger::StartHandshake(link, addr) => { + let _ = node + .start_handshake(link, TransportId::new(TCP_ID), addr, None) + .await; + } + Trigger::LinkMessage(peer) => { + let heartbeat = [LinkMessageType::Heartbeat.to_byte()]; + let _ = node.send_encrypted_link_message(&peer, &heartbeat).await; + } + } + let elapsed = t0.elapsed(); + let after = tcp_stats(node); + Reply { + elapsed, + connect_timeouts: after.connect_timeouts - before.connect_timeouts, + connect_refused: after.connect_refused - before.connect_refused, + connections_established: after.connections_established - before.connections_established, + } +} + +/// Run `process_packet` once and measure it against the TCP counters. +async fn timed_process(node: &mut Node, packet: ReceivedPacket) -> Reply { + timed_fire(node, Trigger::Packet(packet)).await +} + +/// Every way a reply to a dead link failed to be fast and dial-free. +fn dial_findings(what: &str, r: &Reply) -> Vec { + let mut found = Vec::new(); + if r.elapsed >= BOUND { + found.push(format!( + "{what}: handler took {:?} (bound {BOUND:?}); a connect timeout is {CONNECT_TIMEOUT_MS} ms", + r.elapsed + )); + } + if r.connect_timeouts != 0 { + found.push(format!("{what}: {} connect timeouts", r.connect_timeouts)); + } + if r.connect_refused != 0 { + found.push(format!("{what}: {} connects refused", r.connect_refused)); + } + if r.connections_established != 0 { + found.push(format!( + "{what}: {} connections dialed", + r.connections_established + )); + } + found +} + +/// Assert a reply to a dead link failed at once and attempted no connect. +fn assert_no_dial(what: &str, r: &Reply) { + let found = dial_findings(what, r); + assert!(found.is_empty(), "{}", found.join("; ")); +} + +/// A real TCP client connected to the node's listener that has sent one +/// msg1. Returns the client and the frame as the node's receive task +/// delivered it, which carries the client's address as the reply address. +async fn msg1_over_real_tcp(node: &mut Node) -> (TcpStream, ReceivedPacket) { + let listen = tcp(node).local_addr().expect("TCP listener bound"); + let mut client = TcpStream::connect(listen).await.unwrap(); + let data = craft_msg1(node, 0x51); + client.write_all(&data).await.unwrap(); + let rx = node.packet_rx.as_mut().expect("packet channel"); + let packet = timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("msg1 never reached the packet channel") + .expect("packet channel closed"); + assert_eq!(packet.transport_id, TransportId::new(TCP_ID)); + (client, packet) +} + +/// Whether `client` receives a msg2 within a second. +async fn client_gets_msg2(client: &mut TcpStream) -> bool { + let mut buf = [0u8; 2048]; + match timeout(Duration::from_secs(1), client.read(&mut buf)).await { + Ok(Ok(n)) if n > 0 => phase(&buf[..n]) == Some(PHASE_MSG2), + _ => false, + } +} + +/// A msg2 reply whose TCP connection is gone fails at once, attempts no +/// connect, and tears down the half-built link. A msg1 on a live inbound +/// connection is still answered on that connection. +#[tokio::test] +async fn msg2_reply_to_dead_tcp_link_returns_without_dialing() { + // Control: a live inbound connection gets its msg2. + let (mut node, _tx, _) = node_with_udp_and_tcp(CONNECT_TIMEOUT_MS).await; + let (mut client, packet) = msg1_over_real_tcp(&mut node).await; + let r = timed_process(&mut node, packet).await; + let answered = client_gets_msg2(&mut client).await; + println!("msg2 control live inbound connection: {r:?}, msg2 received {answered}"); + assert_no_dial("msg2 control", &r); + assert!(answered, "control: the live connection got no msg2"); + assert_eq!( + node.pending_inbound.len(), + 1, + "control: the leg waits for its msg3" + ); + stop_all(&mut node).await; + + // The reply address has no connection and does not answer SYNs. + let bh = Blackhole::silent(); + let (mut node, _tx, _) = node_with_udp_and_tcp(CONNECT_TIMEOUT_MS).await; + let data = craft_msg1(&node, 0x11); + let packet = ReceivedPacket::new(TransportId::new(TCP_ID), bh.transport_addr(), data); + let r = timed_process(&mut node, packet).await; + println!("msg2 dead blackholed reply address: {r:?}"); + assert_no_dial("msg2 to a dead link", &r); + assert!(node.links.is_empty(), "the half-built link is torn down"); + assert!( + node.addr_to_link.is_empty(), + "the half-built link is unindexed" + ); + assert!( + node.pending_inbound.is_empty(), + "no leg waits for a msg3 that cannot come" + ); + assert_eq!(node.index_allocator.count(), 0, "the index is freed"); + stop_all(&mut node).await; +} + +/// Establish a peer on TCP at `bh`'s address over a live connection the node +/// holds there, with the node as the XX responder. Returns the node, the +/// peer's identity and address, and the far end of the connection, with +/// everything the node sent during setup read off it. +async fn peer_on_tcp(bh: &Blackhole) -> (Node, Identity, NodeAddr, TcpStream) { + peer_from(make_node(), bh).await +} + +/// [`peer_on_tcp`] on a node the caller built. +async fn peer_from(node: Node, bh: &Blackhole) -> (Node, Identity, NodeAddr, TcpStream) { + let (mut node, _tx, _) = node_from(node, CONNECT_TIMEOUT_MS).await; + let mut far_end = prime_link(&node, bh).await; + let sender = Identity::generate(); + let sender_addr = *PeerIdentity::from_pubkey_full(sender.pubkey_full()).node_addr(); + let link = bh.transport_addr(); + let tcp_id = TransportId::new(TCP_ID); + + let (mut leg, msg1) = initiator(&node, &sender, 0x01); + node.process_packet(ReceivedPacket::new(tcp_id, link.clone(), msg1)) + .await; + let msg2 = frame_of(&mut far_end, PHASE_MSG2) + .await + .expect("the msg2 went out on the connection"); + let msg3 = msg3_for(&mut leg, &msg2, 0x01, None); + node.process_packet(ReceivedPacket::new(tcp_id, link.clone(), msg3)) + .await; + assert_eq!(node.peer_count(), 1, "peer established over TCP"); + let p = node.get_peer(&sender_addr).unwrap(); + assert_eq!(p.transport_id(), Some(tcp_id)); + assert_eq!(p.current_addr(), Some(&link)); + drain(&mut far_end).await; + (node, sender, sender_addr, far_end) +} + +/// Record `peer`'s link as one this node dialed at the address it holds, +/// as if the node had been the initiator. +fn make_link_outbound(node: &mut Node, peer: &NodeAddr) { + let link_id = node.get_peer(peer).unwrap().link_id(); + let old = node.links.get(&link_id).expect("the peer's link"); + let link = Link::new( + link_id, + old.transport_id(), + old.remote_addr().clone(), + LinkDirection::Outbound, + old.base_rtt(), + ); + node.links.insert(link_id, link); +} + +/// Age `peer`'s session past the rekey trigger of a default config. +fn age_past_rekey(node: &mut Node, peer: &NodeAddr) { + let after = node.config().node.rekey.after_secs + crate::node::REKEY_JITTER_SECS as u64 + 1; + node.get_peer_mut(peer) + .unwrap() + .test_backdate_session_established(Duration::from_secs(after)); +} + +/// An outbound handshake to `addr` on TCP as a dial arms it. +struct DialLeg { + link: LinkId, + /// The framed msg1 the dial stored for sending. + wire: Vec, +} + +/// Arm an outbound handshake to `target` at `addr` on TCP as a dial does: a +/// link this node dialed, a machine that has sent its msg1 and holds the +/// wire, and a retransmit timer due at `due_ms`. +fn dial_leg( + node: &mut Node, + addr: &TransportAddr, + target: &Identity, + now_ms: u64, + due_ms: u64, +) -> DialLeg { + let tcp_id = TransportId::new(TCP_ID); + let target = PeerIdentity::from_pubkey_full(target.pubkey_full()); + let link_id = node.allocate_link_id(); + let mut leg = outbound_leg(link_id, target, now_ms); + let our_index = node.index_allocator.allocate().unwrap(); + let noise_msg1 = leg + .start_handshake(node.identity().keypair(), node.startup_epoch(), now_ms) + .unwrap(); + let wire = build_msg1(our_index, &noise_msg1); + node.links.insert( + link_id, + Link::new( + link_id, + tcp_id, + addr.clone(), + LinkDirection::Outbound, + Duration::from_millis(100), + ), + ); + node.addr_to_link.insert((tcp_id, addr.clone()), link_id); + node.pending_outbound + .insert((tcp_id, our_index.as_u32()), link_id); + let mut machine = PeerMachine::new_outbound(link_id, Some(target), now_ms); + let _ = machine.step( + PeerEvent::Dial { + transport_id: tcp_id, + remote_addr: addr.clone(), + peer_identity: target, + connection_oriented: false, + }, + now_ms, + &mut node.index_allocator, + ); + machine.set_conn_handshake_msg1(wire.clone(), due_ms); + machine.set_conn_our_index(our_index); + machine.set_conn_transport_id(tcp_id); + machine.set_conn_source_addr(addr.clone()); + machine.set_leg(leg.take_leg().unwrap()); + assert!(machine.is_handshaking_sent_msg1()); + node.peer_machines.insert(link_id, machine); + node.peer_timers + .entry(link_id) + .or_default() + .insert(TimerKind::HandshakeRetransmit, due_ms); + DialLeg { + link: link_id, + wire, + } +} + +/// What a tick's rekey check did to a TCP peer. +#[derive(Debug)] +struct RekeyStart { + reply: Reply, + /// The rekey msg1 went out and the cycle is in flight. + started: bool, + /// The far end received the msg1. + delivered: bool, + /// The transport's connection state for the peer's address afterwards. + state: ConnectionState, +} + +/// Run the tick's rekey check for a TCP peer due to rekey, whose link is +/// alive or (with `dead`) closed, and which this node dialed (`outbound`) +/// or accepted. +async fn rekey_start(dead: bool, outbound: bool) -> RekeyStart { + let mut bh = Blackhole::open(false); + let (mut node, _sender, sender_addr, far_end) = peer_on_tcp(&bh).await; + if outbound { + make_link_outbound(&mut node, &sender_addr); + } + let mut far_end = if dead { + kill_link(&node, &mut bh, far_end).await; + None + } else { + Some(far_end) + }; + age_past_rekey(&mut node, &sender_addr); + let reply = timed_fire(&mut node, Trigger::RekeyCheck).await; + let started = node + .get_peer(&sender_addr) + .is_some_and(|p| p.rekey_in_progress()); + let delivered = frame_maybe(far_end.as_mut(), PHASE_MSG1).await.is_some(); + let state = tcp(&node).connection_state(&bh.transport_addr()); + stop_all(&mut node).await; + RekeyStart { + reply, + started, + delivered, + state, + } +} + +/// The tick's rekey msg1 to a TCP peer whose connection has closed fails at +/// once without dialing. A background connect is started toward a peer this +/// node dialed, and never toward an inbound peer's address. With the link +/// alive the msg1 goes out and the cycle starts. +#[tokio::test] +async fn rekey_msg1_to_dead_tcp_link_does_not_hold_tick() { + let r = rekey_start(false, false).await; + println!("rekey msg1 control link alive: {r:?}"); + assert_no_dial("rekey msg1 control", &r.reply); + assert!( + r.started && r.delivered, + "control: the rekey msg1 did not go out" + ); + + let r = rekey_start(true, false).await; + println!("rekey msg1 dead inbound peer: {r:?}"); + assert_no_dial("rekey msg1 to a closed inbound link", &r.reply); + assert!(!r.started, "a failed rekey msg1 starts no cycle"); + assert_eq!( + r.state, + ConnectionState::None, + "a connect was started toward an inbound peer's address" + ); + + let r = rekey_start(true, true).await; + println!("rekey msg1 dead outbound peer: {r:?}"); + assert_no_dial("rekey msg1 to a closed outbound link", &r.reply); + assert!(!r.started, "a failed rekey msg1 starts no cycle"); + assert_eq!( + r.state, + ConnectionState::Connecting, + "no background connect toward the address this node dialed" + ); +} + +/// The tick's msg1 resend on an outbound handshake whose address does not +/// answer fails at once without dialing and starts a background connect. +/// Once the address answers and that connect finishes, a later tick sends +/// the msg1 on it: the connection is not left stranded unused. +#[tokio::test] +async fn msg1_resend_to_dead_outbound_leg_recovers_after_background_connect() { + let mut bh = Blackhole::silent(); + let (mut node, _tx, _) = node_with_udp_and_tcp(CONNECT_TIMEOUT_MS).await; + let addr = bh.transport_addr(); + let now_ms = Node::now_ms(); + let leg = dial_leg( + &mut node, + &addr, + &Identity::generate(), + now_ms, + now_ms + 1000, + ); + + let r = timed_fire(&mut node, Trigger::PeerTimers(now_ms + 1000)).await; + println!("msg1 resend blackholed dial address: {r:?}"); + assert_no_dial("msg1 resend to a dead outbound leg", &r); + assert_eq!( + node.connection_resend_count(leg.link), + 0, + "nothing was sent" + ); + assert_eq!( + tcp(&node).connection_state(&addr), + ConnectionState::Connecting, + "no background connect toward the dial address" + ); + + // The address starts answering: empty the accept queue, and the + // background connect's retransmitted SYN completes. + let _filler_ends = bh.drain(); + let start = Instant::now(); + let mut tick = 1; + while node.connection_resend_count(leg.link) == 0 { + assert!( + start.elapsed() < Duration::from_secs(4), + "the msg1 resend never went out over the background connect" + ); + tokio::time::sleep(Duration::from_millis(100)).await; + let r = timed_fire(&mut node, Trigger::PeerTimers(now_ms + 1000 + tick * 100)).await; + tick += 1; + assert!(r.elapsed < BOUND, "a tick took {:?}", r.elapsed); + assert_eq!(r.connect_timeouts, 0, "a tick dialed and timed out"); + } + let mut accepted = accept_end(&bh); + println!( + "msg1 resend sent after {:?} over the background connect", + start.elapsed() + ); + assert_eq!( + next_frame(&mut accepted).await, + Some(leg.wire), + "the msg1 did not arrive on the background connection" + ); + stop_all(&mut node).await; +} + +/// Run the node's real rx loop, inject `poisoned` TCP msg1 frames whose reply +/// address is blackholed, then send one genuine msg1 over UDP and return how +/// long the UDP initiator waits for its msg2. +async fn udp_msg2_latency(connect_timeout_ms: u64, poisoned: usize) -> Duration { + let (mut node, tx, udp_addr) = node_with_udp_and_tcp(connect_timeout_ms).await; + let holes: Vec = (0..poisoned).map(|_| Blackhole::silent()).collect(); + let poison: Vec = holes + .iter() + .enumerate() + .map(|(i, bh)| { + let data = craft_msg1(&node, 0x100 + i as u32); + ReceivedPacket::new(TransportId::new(TCP_ID), bh.transport_addr(), data) + }) + .collect(); + let udp_msg1 = craft_msg1(&node, 0x200); + let peer = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap(); + + // Long enough to measure a regression, one timeout per poisoned frame, + // rather than only reporting that it was slow. + let budget = Duration::from_millis(connect_timeout_ms * poisoned as u64 + 3000); + let measure = async { + for p in poison { + tx.send(p).await.unwrap(); + } + // Let the loop pick up the TCP frames before the UDP one arrives. + tokio::time::sleep(Duration::from_millis(20)).await; + let t0 = Instant::now(); + peer.send_to(&udp_msg1, udp_addr).await.unwrap(); + let mut buf = [0u8; 2048]; + loop { + let (n, _) = timeout(budget, peer.recv_from(&mut buf)) + .await + .expect("no msg2 over UDP within budget") + .unwrap(); + if phase(&buf[..n]) == Some(PHASE_MSG2) { + return t0.elapsed(); + } + } + }; + + let latency = tokio::select! { + r = node.run_rx_loop() => panic!("rx loop exited: {r:?}"), + l = measure => l, + }; + stop_all(&mut node).await; + drop(holes); + latency +} + +/// Inside the real rx loop, a UDP initiator's msg2 does not wait +/// behind TCP msg1s whose replies have nowhere to go. +#[tokio::test] +async fn udp_handshake_is_not_delayed_by_dead_tcp_replies_in_rx_loop() { + for poisoned in [0usize, 3] { + let latency = udp_msg2_latency(CONNECT_TIMEOUT_MS, poisoned).await; + println!( + "rx loop connect_timeout_ms={CONNECT_TIMEOUT_MS} blackholed TCP msg1 ahead={poisoned}: UDP msg2 after {latency:?}" + ); + assert!( + latency < BOUND, + "UDP msg2 took {latency:?} behind {poisoned} dead TCP replies (bound {BOUND:?})" + ); + } +} + +/// A counter from the off-loop `show_transports` view. +fn snapshot_stat( + handle: &crate::control::read_handle::ControlReadHandle, + id: u32, + key: &str, +) -> u64 { + let v = crate::control::queries::show_transports_from_handle(handle); + v["transports"] + .as_array() + .unwrap() + .iter() + .find(|t| t["transport_id"] == id) + .and_then(|t| t["stats"][key].as_u64()) + .unwrap_or_else(|| panic!("no stats.{key} for transport {id}: {v}")) +} + +/// The longest the tick may go without publishing during the burst. It +/// runs every second, so a longer gap means a tick was held. +const SNAPSHOT_LAG: Duration = Duration::from_millis(1500); + +/// What the off-loop view and the tick did during a burst of UDP traffic. +#[derive(Debug)] +struct TickProgress { + /// UDP frames sent during the burst. + sent: u64, + /// Off-loop reads of UDP `packets_recv` taken once traffic had arrived. + checks: usize, + /// The reads that fell outside the live count sampled just before and + /// just after, as (ms into the burst, live before, view, live after). + off: Vec<(u128, u64, u64, u64)>, + /// Entity snapshot publishes seen during the burst. + publishes: usize, + /// The longest stretch of the burst with no publish, counting from its + /// start and to its end. + max_gap: Duration, + /// Off-loop TCP `connect_timeouts` at the end of the burst. + snapshot_timeouts: u64, + /// Live TCP `connect_timeouts` at the end of the burst. + live_timeouts: u64, +} + +/// Queue `poisoned` TCP msg1s whose replies are blackholed, then send junk +/// UDP for a few seconds while the real rx loop runs. Throughout the burst, +/// read the off-loop `show_transports` view between two live samples, and +/// watch for the tick's entity snapshot publishes. +async fn tick_progress(poisoned: usize) -> TickProgress { + let ms = 300u64; + let (mut node, tx, udp_addr) = node_with_udp_and_tcp(ms).await; + let handle = node.control_read_handle(); + let live_tcp = match tcp(&node) { + TransportHandle::Tcp(t) => t.stats().clone(), + _ => unreachable!(), + }; + let live_udp = match node.transports.get(&TransportId::new(UDP_ID)) { + Some(TransportHandle::Udp(t)) => t.stats().clone(), + _ => unreachable!(), + }; + let bh = Blackhole::silent(); + let poison: Vec = (0..poisoned) + .map(|i| { + let data = craft_msg1(&node, 0x300 + i as u32); + ReceivedPacket::new(TransportId::new(TCP_ID), bh.transport_addr(), data) + }) + .collect(); + let peer = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap(); + + let measure = async { + // The interval's first tick fires at once and publishes a snapshot. + tokio::time::sleep(Duration::from_millis(100)).await; + for p in poison { + tx.send(p).await.unwrap(); + } + // UDP keeps arriving. Junk frames are enough: the transport counts + // them before the rx loop ever sees them. Fourteen replies that each + // dialed for 300 ms would hold the loop longer than the whole burst. + let t0 = Instant::now(); + let mut sent = 0u64; + let mut checks = 0; + let mut off = Vec::new(); + // Holding the last publish seen keeps its allocation alive, so a + // later publish cannot reuse the address and pass for the same one. + let mut last = std::sync::Arc::clone(&*handle.entities()); + let mut last_at = Duration::ZERO; + let mut publishes = 0; + let mut max_gap = Duration::ZERO; + while t0.elapsed() < Duration::from_millis(3600) { + peer.send_to(b"junk-frame", udp_addr).await.unwrap(); + sent += 1; + let before = live_udp.snapshot().packets_recv; + let view = snapshot_stat(&handle, UDP_ID, "packets_recv"); + let after = live_udp.snapshot().packets_recv; + let now = t0.elapsed(); + // A zero count says nothing about whether the view is live. + if before > 0 { + checks += 1; + if view < before || view > after { + off.push((now.as_millis(), before, view, after)); + } + } + let current = std::sync::Arc::clone(&*handle.entities()); + if !std::sync::Arc::ptr_eq(¤t, &last) { + publishes += 1; + max_gap = max_gap.max(now - last_at); + last = current; + last_at = now; + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + max_gap = max_gap.max(t0.elapsed() - last_at); + TickProgress { + sent, + checks, + off, + publishes, + max_gap, + snapshot_timeouts: snapshot_stat(&handle, TCP_ID, "connect_timeouts"), + live_timeouts: live_tcp.snapshot().connect_timeouts, + } + }; + + let progress = tokio::select! { + r = node.run_rx_loop() => panic!("rx loop exited: {r:?}"), + m = measure => m, + }; + stop_all(&mut node).await; + progress +} + +/// While dead TCP replies are queued, the rx loop's tick keeps running and +/// publishing, and the off-loop `show_transports` view tracks the live +/// counters throughout: it reads them at request time rather than from the +/// tick's copy, so it would stay current even if the tick were held. +#[tokio::test] +async fn tick_runs_and_snapshot_tracks_live_counters_under_dead_tcp_replies() { + for poisoned in [0usize, 14] { + let p = tick_progress(poisoned).await; + println!("tick {poisoned} dead TCP replies queued: {p:?}"); + assert!( + p.checks >= 20, + "only {} off-loop reads made from {} UDP frames; the burst did not exercise the view", + p.checks, + p.sent + ); + assert!( + p.off.is_empty(), + "with {poisoned} dead replies queued the off-loop view did not track the live UDP \ + packets_recv in {} of {} reads, first and last (ms, live before, view, live after) \ + {:?} {:?}", + p.off.len(), + p.checks, + p.off.first(), + p.off.last() + ); + assert!( + p.max_gap < SNAPSHOT_LAG, + "with {poisoned} dead replies queued the tick was held: {} publishes, longest gap \ + {:?} (bound {SNAPSHOT_LAG:?})", + p.publishes, + p.max_gap + ); + assert_eq!(p.live_timeouts, 0, "a reply dialed and timed out"); + assert_eq!(p.snapshot_timeouts, p.live_timeouts); + } +} + +/// Through the real accept and receive tasks: a client that sends a +/// msg1 and closes before the node answers draws no connect attempt. With +/// the client still connected, the msg2 goes back on its connection. +/// +/// Over loopback a SYN to the closed client port is answered with a reset, +/// so a dialing reply fails fast with a refusal here rather than stalling; +/// this checks the trigger, and the counters are what show a dial. +#[tokio::test] +async fn msg1_then_close_over_real_tcp_makes_no_connect_attempt() { + // Control: the client stays connected. + let (mut node, _tx, _) = node_with_udp_and_tcp(CONNECT_TIMEOUT_MS).await; + let (mut client, packet) = msg1_over_real_tcp(&mut node).await; + let r = timed_process(&mut node, packet).await; + let answered = client_gets_msg2(&mut client).await; + println!("open-close control client connected: {r:?}, msg2 received {answered}"); + assert_no_dial("open-close control", &r); + assert!(answered, "control: the connected client got no msg2"); + stop_all(&mut node).await; + + let (mut node, _tx, _) = node_with_udp_and_tcp(CONNECT_TIMEOUT_MS).await; + let (client, packet) = msg1_over_real_tcp(&mut node).await; + drop(client); + wait_pool_gone(&node, &packet.remote_addr).await; + let pool_outbound = tcp_stats(&node).pool_outbound; + let r = timed_process(&mut node, packet).await; + println!("open-close client closed first: {r:?}"); + assert_no_dial("msg1 then close", &r); + assert_eq!( + tcp_stats(&node).pool_outbound, + pool_outbound, + "a new outbound pool entry appeared" + ); + assert!(node.links.is_empty(), "the half-built link is torn down"); + stop_all(&mut node).await; +} + +/// The handshake and link sends the rx loop awaits, other than those +/// covered above, that may reach a TCP link which has gone away. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum ReplySite { + /// A second copy of the msg1 a pending inbound handshake answered: the + /// stored msg2 is resent before any crypto. + DuplicateMsg1, + /// A msg3 from an established peer declaring a rekey of a session the + /// node does not hold: the peer's stored msg2 is resent. + ResendMsg2, + /// A msg3 from a peer the node's access list refuses: an encrypted + /// Disconnect is sent before the leg is torn down. + AclDisconnect, + /// The msg2 answering this node's dial: the msg3 is sent. + Msg3, + /// The msg2 answering this node's rekey msg1: the rekey msg3 is sent. + RekeyMsg3, + /// The tick's resend of a rekey msg1 to an inbound peer. + RekeyMsg1Resend, + /// The tick's resend of a rekey msg3 to an inbound peer. + RekeyMsg3Resend, + /// The tick's resend of a rekey msg1 to a peer this node dialed. + RekeyMsg1ResendOutbound, + /// The tick's resend of a rekey msg3 to a peer this node dialed. + RekeyMsg3ResendOutbound, + /// The executor's send of the msg1 armed by an outbound dial. + StoredMsg1, + /// An anonymous dial's inline msg1 send once its connection is up. + StartHandshake, + /// An encrypted link message to a peer this node dialed. + LinkMessage, + /// An encrypted link message to a peer that dialed this node. + LinkMessageInbound, +} + +impl ReplySite { + /// Whether the site, finding no connection, starts a background + /// connect: only toward an address this node dialed. + fn connects(self) -> bool { + matches!( + self, + ReplySite::StoredMsg1 + | ReplySite::StartHandshake + | ReplySite::LinkMessage + | ReplySite::RekeyMsg1ResendOutbound + | ReplySite::RekeyMsg3ResendOutbound + ) + } + + /// The phase of the frame the site sends. + fn sends(self) -> u8 { + match self { + ReplySite::DuplicateMsg1 | ReplySite::ResendMsg2 => PHASE_MSG2, + ReplySite::Msg3 + | ReplySite::RekeyMsg3 + | ReplySite::RekeyMsg3Resend + | ReplySite::RekeyMsg3ResendOutbound => PHASE_MSG3, + ReplySite::RekeyMsg1Resend + | ReplySite::RekeyMsg1ResendOutbound + | ReplySite::StoredMsg1 + | ReplySite::StartHandshake => PHASE_MSG1, + ReplySite::AclDisconnect | ReplySite::LinkMessage | ReplySite::LinkMessageInbound => { + PHASE_ESTABLISHED + } + } + } +} + +/// A site brought within one call of firing. +struct Armed { + node: Node, + far_end: TcpStream, + trigger: Trigger, + /// The access-list files, kept for as long as the node may read them. + _acl: Option, +} + +/// Start a rekey from the node to the peer at the far end of `far_end`, and +/// answer its msg1 with the peer's msg2. Returns the msg2, not yet delivered. +async fn rekey_msg2( + node: &mut Node, + sender: &Identity, + sender_addr: &NodeAddr, + far_end: &mut TcpStream, + link: &TransportAddr, +) -> ReceivedPacket { + age_past_rekey(node, sender_addr); + node.check_rekey().await; + let msg1 = frame_of(far_end, PHASE_MSG1) + .await + .expect("the rekey msg1 went out on the connection"); + let msg2 = msg2_for(sender, &msg1, 0x33); + ReceivedPacket::new(TransportId::new(TCP_ID), link.clone(), msg2) +} + +/// Bring `row`'s site within one call of firing against a live connection +/// at `bh`. +async fn arm_site(row: ReplySite, bh: &Blackhole) -> Armed { + let tcp_id = TransportId::new(TCP_ID); + let link = bh.transport_addr(); + let mut acl = None; + let (node, far_end, trigger) = match row { + ReplySite::DuplicateMsg1 => { + let (mut node, _tx, _) = node_with_udp_and_tcp(CONNECT_TIMEOUT_MS).await; + let mut far_end = prime_link(&node, bh).await; + let (_leg, msg1) = initiator(&node, &Identity::generate(), 0x21); + node.process_packet(ReceivedPacket::new(tcp_id, link.clone(), msg1.clone())) + .await; + assert!( + frame_of(&mut far_end, PHASE_MSG2).await.is_some(), + "the first msg2 went out on the connection" + ); + let packet = ReceivedPacket::new(tcp_id, link, msg1); + (node, far_end, Trigger::Packet(packet)) + } + ReplySite::ResendMsg2 => { + let (mut node, sender, _, mut far_end) = peer_on_tcp(bh).await; + let (mut leg, msg1) = initiator(&node, &sender, 0x22); + node.process_packet(ReceivedPacket::new(tcp_id, link.clone(), msg1)) + .await; + let msg2 = frame_of(&mut far_end, PHASE_MSG2) + .await + .expect("the second handshake's msg2 went out on the connection"); + let msg3 = msg3_for(&mut leg, &msg2, 0x22, Some(SessionIndex::new(0xDEAD))); + let packet = ReceivedPacket::new(tcp_id, link, msg3); + (node, far_end, Trigger::Packet(packet)) + } + ReplySite::AclDisconnect => { + let dir = tempfile::tempdir().unwrap(); + let sender = Identity::generate(); + let mut base = make_node(); + base.peer_acl = PeerAclReloader::with_paths( + dir.path().join("peers.allow"), + dir.path().join("peers.deny"), + ); + std::fs::write( + dir.path().join("peers.deny"), + format!("{}\n", sender.npub()), + ) + .unwrap(); + assert!(base.reload_peer_acl().await, "the access list loaded"); + acl = Some(dir); + let (mut node, _tx, _) = node_from(base, CONNECT_TIMEOUT_MS).await; + let mut far_end = prime_link(&node, bh).await; + let (mut leg, msg1) = initiator(&node, &sender, 0x24); + node.process_packet(ReceivedPacket::new(tcp_id, link.clone(), msg1)) + .await; + let msg2 = frame_of(&mut far_end, PHASE_MSG2) + .await + .expect("the msg2 went out on the connection"); + let msg3 = msg3_for(&mut leg, &msg2, 0x24, None); + let packet = ReceivedPacket::new(tcp_id, link, msg3); + (node, far_end, Trigger::Packet(packet)) + } + ReplySite::Msg3 => { + let (mut node, _tx, _) = node_with_udp_and_tcp(CONNECT_TIMEOUT_MS).await; + let far_end = prime_link(&node, bh).await; + let responder = Identity::generate(); + let now_ms = Node::now_ms(); + let leg = dial_leg(&mut node, &link, &responder, now_ms, now_ms + 60_000); + let msg2 = msg2_for(&responder, &leg.wire, 0x25); + let packet = ReceivedPacket::new(tcp_id, link, msg2); + (node, far_end, Trigger::Packet(packet)) + } + ReplySite::RekeyMsg3 => { + let (mut node, sender, sender_addr, mut far_end) = peer_on_tcp(bh).await; + let packet = rekey_msg2(&mut node, &sender, &sender_addr, &mut far_end, &link).await; + (node, far_end, Trigger::Packet(packet)) + } + ReplySite::RekeyMsg1Resend | ReplySite::RekeyMsg1ResendOutbound => { + let (mut node, _sender, sender_addr, mut far_end) = peer_on_tcp(bh).await; + if row == ReplySite::RekeyMsg1ResendOutbound { + make_link_outbound(&mut node, &sender_addr); + } + age_past_rekey(&mut node, &sender_addr); + node.check_rekey().await; + assert!( + node.get_peer(&sender_addr) + .is_some_and(|p| p.rekey_in_progress()), + "the rekey cycle started" + ); + assert!( + frame_of(&mut far_end, PHASE_MSG1).await.is_some(), + "the rekey msg1 went out on the connection" + ); + (node, far_end, Trigger::RekeyResend) + } + ReplySite::RekeyMsg3Resend | ReplySite::RekeyMsg3ResendOutbound => { + let (mut node, sender, sender_addr, mut far_end) = peer_on_tcp(bh).await; + if row == ReplySite::RekeyMsg3ResendOutbound { + make_link_outbound(&mut node, &sender_addr); + } + let packet = rekey_msg2(&mut node, &sender, &sender_addr, &mut far_end, &link).await; + node.process_packet(packet).await; + assert!( + frame_of(&mut far_end, PHASE_MSG3).await.is_some(), + "the rekey msg3 went out on the connection" + ); + assert!( + node.get_peer(&sender_addr) + .is_some_and(|p| p.rekey_msg3_payload().is_some()), + "the rekey msg3 is retained for resending" + ); + (node, far_end, Trigger::RekeyMsg3Resend) + } + ReplySite::StoredMsg1 => { + let (mut node, _tx, _) = node_with_udp_and_tcp(CONNECT_TIMEOUT_MS).await; + let far_end = prime_link(&node, bh).await; + let now_ms = Node::now_ms(); + let leg = dial_leg( + &mut node, + &link, + &Identity::generate(), + now_ms, + now_ms + 1000, + ); + (node, far_end, Trigger::StoredMsg1(leg.link, link)) + } + ReplySite::StartHandshake => { + let (mut node, _tx, _) = node_with_udp_and_tcp(CONNECT_TIMEOUT_MS).await; + let far_end = prime_link(&node, bh).await; + let link_id = node.allocate_link_id(); + node.links.insert( + link_id, + Link::new( + link_id, + tcp_id, + link.clone(), + LinkDirection::Outbound, + Duration::from_millis(100), + ), + ); + node.addr_to_link.insert((tcp_id, link.clone()), link_id); + (node, far_end, Trigger::StartHandshake(link_id, link)) + } + ReplySite::LinkMessage => { + let (mut node, _sender, sender_addr, far_end) = peer_on_tcp(bh).await; + make_link_outbound(&mut node, &sender_addr); + (node, far_end, Trigger::LinkMessage(sender_addr)) + } + ReplySite::LinkMessageInbound => { + let (node, _sender, sender_addr, far_end) = peer_on_tcp(bh).await; + (node, far_end, Trigger::LinkMessage(sender_addr)) + } + }; + Armed { + node, + far_end, + trigger, + _acl: acl, + } +} + +/// Fire `row`'s site with its established connection closed (or, with +/// `dead` false, still open). Returns the measurement, whether the far end +/// received the frame the site sends, and the transport's connection state +/// for the far end's address afterwards. +async fn fire_site(row: ReplySite, dead: bool) -> (Reply, bool, ConnectionState) { + let mut bh = Blackhole::open(false); + let Armed { + mut node, + mut far_end, + trigger, + _acl, + } = arm_site(row, &bh).await; + drain(&mut far_end).await; + let mut far_end = if dead { + kill_link(&node, &mut bh, far_end).await; + None + } else { + Some(far_end) + }; + let r = timed_fire(&mut node, trigger).await; + let delivered = frame_maybe(far_end.as_mut(), row.sends()).await.is_some(); + let state = tcp(&node).connection_state(&bh.transport_addr()); + stop_all(&mut node).await; + (r, delivered, state) +} + +/// Every send the rx loop awaits on a TCP link that has gone away returns +/// at once without a connect attempt, and with the link alive the send is +/// delivered. Only a send toward an address this node dialed leaves a +/// background connect behind. +#[tokio::test] +async fn every_rx_loop_handshake_send_to_dead_tcp_link_is_bounded() { + // Every row runs before any assertion, so one red names all the sites + // that dial rather than only the first. + let mut found = Vec::new(); + for row in [ + ReplySite::DuplicateMsg1, + ReplySite::ResendMsg2, + ReplySite::AclDisconnect, + ReplySite::Msg3, + ReplySite::RekeyMsg3, + ReplySite::RekeyMsg1Resend, + ReplySite::RekeyMsg3Resend, + ReplySite::RekeyMsg1ResendOutbound, + ReplySite::RekeyMsg3ResendOutbound, + ReplySite::StoredMsg1, + ReplySite::StartHandshake, + ReplySite::LinkMessage, + ReplySite::LinkMessageInbound, + ] { + let (r, delivered, _) = fire_site(row, false).await; + println!("{row:?} control link alive: {r:?}, delivered {delivered}"); + found.extend(dial_findings(&format!("{row:?} control"), &r)); + if !delivered { + found.push(format!("{row:?} control: the send was not delivered")); + } + + let (r, _, state) = fire_site(row, true).await; + println!("{row:?} dead link closed: {r:?}, afterwards {state:?}"); + found.extend(dial_findings(&format!("{row:?} to a closed link"), &r)); + let expected = if row.connects() { + ConnectionState::Connecting + } else { + ConnectionState::None + }; + if state != expected { + found.push(format!( + "{row:?} to a closed link: connection state {state:?}, expected {expected:?}" + )); + } + } + assert!(found.is_empty(), "{}", found.join("\n")); +} + +/// Run one send toward an outbound TCP peer whose current address an +/// authenticated frame has moved away from the address the link was dialed +/// at: the tick's rekey msg1 (`rekey`) or a link message. The connection at +/// the moved address is open, or (with `dead`) closed with the address no +/// longer answering. Returns the measurement, whether the send arrived at +/// the moved address, and the connection state there afterwards. +async fn moved_send(rekey: bool, dead: bool) -> (Reply, bool, ConnectionState) { + let bh = Blackhole::open(false); + let (mut node, _sender, sender_addr, _dialed_end) = peer_on_tcp(&bh).await; + make_link_outbound(&mut node, &sender_addr); + let mut moved = Blackhole::open(false); + let moved_end = prime_link(&node, &moved).await; + node.get_peer_mut(&sender_addr) + .unwrap() + .set_current_addr(TransportId::new(TCP_ID), moved.transport_addr()); + let mut moved_end = if dead { + kill_link(&node, &mut moved, moved_end).await; + None + } else { + Some(moved_end) + }; + let (trigger, sends) = if rekey { + age_past_rekey(&mut node, &sender_addr); + (Trigger::RekeyCheck, PHASE_MSG1) + } else { + (Trigger::LinkMessage(sender_addr), PHASE_ESTABLISHED) + }; + let r = timed_fire(&mut node, trigger).await; + let delivered = frame_maybe(moved_end.as_mut(), sends).await.is_some(); + let state = tcp(&node).connection_state(&moved.transport_addr()); + stop_all(&mut node).await; + (r, delivered, state) +} + +/// A peer this node dialed, whose current address has moved, is sent to at +/// the moved address. With the connection there gone the send fails at once +/// and starts no connect toward it: only the address the link was dialed at +/// is known to have a listener, and the moved one may be an ephemeral port. +#[tokio::test] +async fn send_to_moved_outbound_peer_does_not_connect_to_its_moved_address() { + for rekey in [false, true] { + let what = if rekey { "rekey msg1" } else { "link message" }; + let (r, delivered, _) = moved_send(rekey, false).await; + println!("moved peer {what} control connection open: {r:?}, delivered {delivered}"); + assert_no_dial(&format!("moved peer {what} control"), &r); + assert!( + delivered, + "control: the {what} did not go to the moved address" + ); + + let (r, _, state) = moved_send(rekey, true).await; + println!("moved peer {what} connection closed: {r:?}, afterwards {state:?}"); + assert_no_dial(&format!("moved peer {what} to a closed connection"), &r); + assert_eq!( + state, + ConnectionState::None, + "the {what} started a connect toward the peer's moved address" + ); + } +} + +/// `may_dial` allows a connect only toward the address an outbound link was +/// dialed at, on that link's transport: never for an inbound link, an +/// address the link has since moved to, another transport, or an unknown +/// link. A plain test, with no runtime or sockets. +#[test] +fn may_dial_allows_only_an_outbound_link_s_dial_address_on_its_transport() { + let mut node = make_node(); + let tcp = TransportId::new(1); + let dialed = TransportAddr::from_string("192.0.2.1:443"); + let moved = TransportAddr::from_string("192.0.2.1:50123"); + let (out, inb) = (LinkId::new(1), LinkId::new(2)); + for (id, dir) in [ + (out, LinkDirection::Outbound), + (inb, LinkDirection::Inbound), + ] { + let link = Link::new(id, tcp, dialed.clone(), dir, Duration::from_millis(100)); + node.links.insert(id, link); + } + + assert!(node.may_dial(out, tcp, &dialed), "outbound, dial address"); + assert!(!node.may_dial(out, tcp, &moved), "outbound, moved address"); + assert!( + !node.may_dial(out, TransportId::new(2), &dialed), + "outbound, other transport" + ); + assert!(!node.may_dial(inb, tcp, &dialed), "inbound link"); + assert!(!node.may_dial(LinkId::new(3), tcp, &dialed), "unknown link"); +}