diff --git a/src/node/dataplane/peer_actions.rs b/src/node/dataplane/peer_actions.rs index 64208b0e..c2e4730f 100644 --- a/src/node/dataplane/peer_actions.rs +++ b/src/node/dataplane/peer_actions.rs @@ -184,10 +184,12 @@ impl Node { let frame = build_msg2(sender_idx, receiver_idx, &bytes); // Surface the send Result. A missing transport skips // the send and continues (mirrors `handle_msg1`'s - // `if let Some(transport)` guard); a send *error* runs the - // pre-refactor msg2-send-failure cleanup (`handle_msg1` - // L494-503) and ABORTS the remaining queue so the queued - // `PromoteToActive` never runs. + // `if let Some(transport)` guard); a send *error* is + // settled by the rule `handle_msg1` uses + // (`Node::msg2_failed`: a transient one keeps the leg, + // a terminal one disposes it) and, either way, ABORTS + // the remaining queue so the queued `PromoteToActive` + // never runs. let send_err = match self.transports.get(&ambient.transport_id) { Some(transport) => { transport.send(&ambient.remote_addr, &frame).await.err() @@ -195,42 +197,7 @@ impl Node { None => None, }; if let Some(e) = send_err { - // A transient refusal is not a failed handshake. The - // interface under this transport is absent or - // mid-rebind, and the binder is already working to - // bring it back — so the half-built link is left - // exactly as it is for the initiator's msg1 resend to - // land on, rather than being torn down and rebuilt. - // - // Tearing down here charged a *local* interface flap - // to the remote: the reject counter it recorded means - // "the peer sent something invalid", which is a - // different thing entirely and one an operator reads - // as the peer's fault. - // - // Nothing leaks by staying. An initiator that never - // resends leaves a stale connection, which - // `check_timeouts` reaps at `handshake_timeout_secs` - // exactly as it reaps every other abandoned handshake. - if e.is_transient() { - debug!( - link_id = %link, - error = %e, - "Deferred msg2: the transport is between interfaces" - ); - return; - } - // Restored pre-refactor msg2-send-failure warn! - // (`handle_msg1` L665): the send error text is surfaced - // at the executor point where the failure is now handled. - warn!(link_id = %link, error = %e, "Failed to send msg2"); - self.remove_link(&link); - if let Some(idx) = ambient.our_index { - let _ = self.index_allocator.free(idx); - } - self.remove_peer_machine(link); - self.stats_mut() - .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); + self.msg2_failed(link, sender_idx, &e); return; } } else { diff --git a/src/node/handlers/handshake.rs b/src/node/handlers/handshake.rs index 38a8f982..3b8b080b 100644 --- a/src/node/handlers/handshake.rs +++ b/src/node/handlers/handshake.rs @@ -23,7 +23,7 @@ use crate::proto::fmp::{ NegotiationPayload, OutboundSnapshot, PromotionResult, RekeyClaim, RekeyMsg2Reject, RekeyMsg2Snapshot, WireOutcome, cross_connection_winner, decide_fmp_negotiation, }; -use crate::transport::{Link, LinkDirection, LinkId, ReceivedPacket}; +use crate::transport::{Link, LinkDirection, LinkId, ReceivedPacket, TransportError}; use crate::utils::index::SessionIndex; use std::time::{Duration, Instant}; use tracing::{debug, info, warn}; @@ -481,29 +481,60 @@ impl Node { ); } Err(e) => { - warn!( - link_id = %link_id, - error = %e, - "Failed to send msg2" - ); - // Clean up on failure (disposing the machine drops the - // Noise handles it carries with it) - self.remove_link(&link_id); - let _ = self.index_allocator.free(our_index); - self.remove_peer_machine(link_id); - self.stats_mut() - .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); - return; + if !self.msg2_failed(link_id, our_index, &e) { + return; + } } } } - // XX: handshake NOT complete yet — need msg3. + // XX: handshake NOT complete yet — need msg3. A leg whose msg2 send + // was deferred waits here too, for the initiator's msg1 resend. // Store in pending_inbound for msg3 dispatch. self.pending_inbound .insert((packet.transport_id, our_index.as_u32()), link_id); } + /// Apply the msg2-send-failure rule to the inbound leg on `link`, which + /// owns `our_index`. Returns whether the leg was kept. + /// + /// A transient refusal is not a failed handshake. The interface under the + /// transport is absent or mid-rebind and the binder is already working to + /// bring it back, so the leg is left exactly as it is for the initiator's + /// msg1 resend to land on, rather than being torn down and rebuilt. + /// Tearing down charged a local interface flap to the remote: the reject + /// counter it recorded means "the peer sent something invalid", which an + /// operator reads as the peer's fault. Nothing leaks by staying: a leg + /// nobody resends to is reaped at `handshake_timeout_secs` like every + /// other abandoned handshake. + /// + /// A terminal error disposes the leg (handing its address back to any + /// link it displaced), frees the index, removes the machine and records + /// the reject. + pub(in crate::node) fn msg2_failed( + &mut self, + link: LinkId, + our_index: SessionIndex, + error: &TransportError, + ) -> bool { + if error.is_transient() { + debug!( + link_id = %link, + error = %error, + "Deferred msg2: the transport is between interfaces" + ); + return true; + } + warn!(link_id = %link, error = %error, "Failed to send msg2"); + // Disposing the machine drops the Noise handles it carries. + self.remove_link(&link); + let _ = self.index_allocator.free(our_index); + self.remove_peer_machine(link); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); + false + } + /// Whether `msg1` is a different attempt from the one the pending inbound /// leg on `link` answered. A leg with no machine or no recorded msg1 is /// treated as answering it, which keeps the duplicate path's behaviour for diff --git a/src/node/tests/handshake.rs b/src/node/tests/handshake.rs index a7f8104f..1a783640 100644 --- a/src/node/tests/handshake.rs +++ b/src/node/tests/handshake.rs @@ -5405,3 +5405,210 @@ async fn a_stale_leg_at_a_peer_address_does_not_block_the_peer_rekey() { stop_hs(&mut initiator).await; stop_hs(&mut responder).await; } + +// =========================================================================== +// A msg2 send refused because the interface under the transport is absent or +// mid-rebind is transient: the responder keeps the inbound leg for the +// initiator's msg1 resend and does not charge the remote with a reject. +// =========================================================================== + +/// A running node built from `config` whose only transport is Ethernet on an +/// interface no host has, so every send off it reports +/// `InterfaceUnavailable`: the real error from the real code path, rather +/// than a stub that merely returns something transient. +#[cfg(any(target_os = "linux", target_os = "macos"))] +async fn absent_node(transport_id: TransportId, config: Config) -> Node { + use crate::config::EthernetConfig; + use crate::transport::ethernet::EthernetTransport; + + let mut node = make_node_with(config); + node.supervisor.state = NodeState::Running; + let eth_config = EthernetConfig { + interface: "fips-absent-x0".to_string(), + ethertype: None, + mtu: None, + recv_buf_size: None, + send_buf_size: None, + listen: Some(true), + announce: Some(false), + auto_connect: None, + accept_connections: Some(true), + beacon_interval_secs: None, + optional: Some(true), + }; + let (tx, _rx) = packet_channel(8); + let mut eth = EthernetTransport::new(transport_id, Some("lab".into()), eth_config, tx); + eth.start_async() + .await + .expect("an absent interface is not a start failure"); + node.transports + .insert(transport_id, TransportHandle::Ethernet(eth)); + node +} + +/// Install an established peer on `node` whose inbound link is at `addr` on +/// `transport_id`, with the address key naming it. Returns the peer's +/// `NodeAddr`, its link and its session index. +#[cfg(any(target_os = "linux", target_os = "macos"))] +fn seed_peer( + node: &mut Node, + transport_id: TransportId, + addr: &TransportAddr, +) -> (NodeAddr, LinkId, SessionIndex) { + let link = node.allocate_link_id(); + let index = node.index_allocator.allocate().unwrap(); + let identity = seed_completed_connection_with(node, link, 1_000, |seed| { + seed.with_our_index(index) + .with_their_index(SessionIndex::new(42)) + .with_transport_id(transport_id) + .with_source_addr(addr.clone()) + }); + node.links.insert( + link, + Link::connectionless( + link, + transport_id, + addr.clone(), + LinkDirection::Inbound, + std::time::Duration::from_millis(100), + ), + ); + node.addr_to_link.insert((transport_id, addr.clone()), link); + node.promote_connection(link, identity, 2_000).unwrap(); + (*identity.node_addr(), link, index) +} + +/// Assert the leg the msg1 from `addr` built is pending as after a sent msg2: +/// its machine parked at `SentMsg2` holding its index, its `pending_inbound` +/// entry, the address key naming it and its stored msg2. Returns the leg. +#[cfg(any(target_os = "linux", target_os = "macos"))] +fn deferred_leg(node: &Node, transport_id: TransportId, addr: &TransportAddr) -> LinkId { + use crate::peer::machine::{HandshakePhase, PeerState}; + + let leg = node + .addr_to_link + .get(&(transport_id, addr.clone())) + .copied() + .expect("the address key must name the kept leg"); + assert!(node.links.contains_key(&leg), "the leg's link stays"); + let machine = node + .peer_machines + .get(&leg) + .expect("a deferred msg2 must leave the leg's machine registered"); + assert!( + matches!( + machine.state(), + PeerState::Handshaking { + phase: HandshakePhase::SentMsg2, + .. + } + ), + "the leg waits for msg3 exactly as after a sent msg2" + ); + let index = machine.our_index().expect("the leg holds its index"); + assert!(node.index_allocator.is_allocated(index)); + assert_eq!( + node.pending_inbound.get(&(transport_id, index.as_u32())), + Some(&leg), + "a later msg3 must find the leg" + ); + assert!( + machine.conn_handshake_msg2().is_some(), + "the stored msg2 answers the initiator's resend" + ); + leg +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +#[tokio::test] +async fn a_transient_msg2_failure_keeps_the_inbound_leg_for_the_retry() { + let transport_id = TransportId::new(1); + let mut node = absent_node(transport_id, Config::new()).await; + let initiator = make_node(); + let source_addr = TransportAddr::from_string("aa:bb:cc:dd:ee:ff"); + let rejects_before = node.stats().handshake.bad_state; + + let (_hs, msg1) = bare_msg1( + initiator.identity().keypair(), + initiator.startup_epoch(), + SessionIndex::new(7), + ); + let deliver = |data: Vec| { + ReceivedPacket::with_timestamp(transport_id, source_addr.clone(), data, Node::now_ms()) + }; + node.handle_msg1(deliver(msg1.clone())).await; + + assert_eq!( + node.link_count(), + 1, + "the leg must survive a transport that is merely between interfaces" + ); + let leg = deferred_leg(&node, transport_id, &source_addr); + assert_eq!( + node.stats().handshake.bad_state, + rejects_before, + "a local interface flap must not be recorded as the peer's misbehaviour" + ); + + // The initiator's resend of the same msg1 lands on the kept leg. + node.handle_msg1(deliver(msg1)).await; + assert_eq!( + node.link_count(), + 1, + "the resend must not build a second leg" + ); + assert_eq!(deferred_leg(&node, transport_id, &source_addr), leg); + assert_eq!(node.stats().handshake.bad_state, rejects_before); + #[cfg(debug_assertions)] + node.debug_assert_peer_maps_coherent(); +} + +/// On XX a msg1 carries no epoch or identity, so at msg1 a restart, a rekey +/// and a fresh dial from an established peer's address look alike, and a +/// restart is decided at msg3. What the deferral means here is that the fresh +/// leg waits at `SentMsg2` and the established peering is untouched. +#[cfg(any(target_os = "linux", target_os = "macos"))] +#[tokio::test] +async fn a_transient_msg2_failure_from_an_established_peer_address_leaves_the_fresh_leg_pending() { + let transport_id = TransportId::new(1); + let mut node = absent_node(transport_id, Config::new()).await; + let source_addr = TransportAddr::from_string("aa:bb:cc:dd:ee:ff"); + + let (peer_addr, peer_link, peer_index) = seed_peer(&mut node, transport_id, &source_addr); + let rejects_before = node.stats().handshake.bad_state; + + let initiator = make_node(); + let (_hs, msg1) = bare_msg1( + initiator.identity().keypair(), + initiator.startup_epoch(), + SessionIndex::new(7), + ); + node.handle_msg1(ReceivedPacket::with_timestamp( + transport_id, + source_addr.clone(), + msg1, + Node::now_ms(), + )) + .await; + + let leg = deferred_leg(&node, transport_id, &source_addr); + assert_ne!(leg, peer_link, "the msg1 built a fresh leg"); + assert_eq!( + node.displaced_links.get(&leg).map(|(_, prior)| *prior), + Some(peer_link), + "the fresh leg records the peer's link it displaced" + ); + let peer = node + .get_peer(&peer_addr) + .expect("the established peering is untouched"); + assert_eq!(peer.link_id(), peer_link); + assert_eq!(peer.our_index(), Some(peer_index)); + assert!(node.links.contains_key(&peer_link), "the peer's link stays"); + assert_eq!( + node.stats().handshake.bad_state, + rejects_before, + "a local interface flap must not be recorded as the peer's misbehaviour" + ); + #[cfg(debug_assertions)] + node.debug_assert_peer_maps_coherent(); +}