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"); +}