diff --git a/src/node/handlers/rekey.rs b/src/node/handlers/rekey.rs index ac0fd7f..9d052dc 100644 --- a/src/node/handlers/rekey.rs +++ b/src/node/handlers/rekey.rs @@ -374,20 +374,20 @@ impl Node { Some(e) => e, None => return, }; - let dest_pubkey = *entry.remote_pubkey(); + let _dest_pubkey = *entry.remote_pubkey(); - // Create Noise XK initiator handshake + // Create Noise XX initiator handshake (rekey: no negotiation payload) let our_keypair = self.identity.keypair(); - let mut handshake = HandshakeState::new_xk_initiator(our_keypair, dest_pubkey); + let mut handshake = HandshakeState::new_xx_initiator(our_keypair); handshake.set_local_epoch(self.startup_epoch); - let msg1 = match handshake.write_xk_message_1() { + let msg1 = match handshake.write_xx_message_1() { Ok(m) => m, Err(e) => { warn!( peer = %self.peer_display_name(dest_addr), error = %e, - "Failed to generate FSP rekey XK msg1" + "Failed to generate FSP rekey XX msg1" ); return; } diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index 9fca0e5..b58d807 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -2,30 +2,29 @@ //! //! Handles locally-delivered session payloads from SessionDatagram envelopes. //! Dispatches based on FSP common prefix phase to specific handlers for -//! SessionSetup (Noise XK msg1), SessionAck (msg2), SessionMsg3 (msg3), +//! SessionSetup (Noise XX msg1), SessionAck (msg2), SessionMsg3 (msg3), //! encrypted data, and error signals (CoordsRequired, PathBroken). -use crate::NodeAddr; -use crate::mmp::report::ReceiverReport; -use crate::mmp::{MAX_SESSION_REPORT_INTERVAL_MS, MIN_SESSION_REPORT_INTERVAL_MS}; use crate::node::session::{EndToEndState, SessionEntry}; use crate::node::session_wire::{ - FSP_COMMON_PREFIX_SIZE, FSP_FLAG_CP, FSP_FLAG_K, FSP_HEADER_SIZE, FSP_PHASE_ESTABLISHED, - FSP_PHASE_MSG1, FSP_PHASE_MSG2, FSP_PHASE_MSG3, FSP_PORT_HEADER_SIZE, FSP_PORT_IPV6_SHIM, - FspCommonPrefix, FspEncryptedHeader, build_fsp_header, fsp_prepend_inner_header, - fsp_strip_inner_header, parse_encrypted_coords, + build_fsp_header, fsp_prepend_inner_header, fsp_strip_inner_header, + parse_encrypted_coords, FspCommonPrefix, FspEncryptedHeader, FSP_COMMON_PREFIX_SIZE, + FSP_FLAG_CP, FSP_FLAG_K, FSP_HEADER_SIZE, FSP_PHASE_ESTABLISHED, FSP_PHASE_MSG1, + FSP_PHASE_MSG2, FSP_PHASE_MSG3, FSP_PORT_HEADER_SIZE, FSP_PORT_IPV6_SHIM, }; +use crate::protocol::{coords_wire_size, encode_coords}; +use crate::upper::icmp::FIPS_OVERHEAD; use crate::node::{Node, NodeError}; -use crate::noise::{ - HandshakeState, XK_HANDSHAKE_MSG1_SIZE, XK_HANDSHAKE_MSG2_SIZE, XK_HANDSHAKE_MSG3_SIZE, -}; +use crate::noise::{HandshakeState, XX_HANDSHAKE_MSG1_SIZE, XX_HANDSHAKE_MSG2_SIZE, XX_HANDSHAKE_MSG3_SIZE}; +use crate::protocol::NegotiationPayload; +use crate::mmp::report::ReceiverReport; +use crate::mmp::{MAX_SESSION_REPORT_INTERVAL_MS, MIN_SESSION_REPORT_INTERVAL_MS}; use crate::protocol::{ CoordsRequired, FspInnerFlags, MtuExceeded, PathBroken, PathMtuNotification, SessionAck, SessionDatagram, SessionMessageType, SessionMsg3, SessionReceiverReport, SessionSenderReport, SessionSetup, }; -use crate::protocol::{coords_wire_size, encode_coords}; -use crate::upper::icmp::FIPS_OVERHEAD; +use crate::NodeAddr; use secp256k1::PublicKey; use tracing::{debug, info, trace}; @@ -37,7 +36,7 @@ impl Node { /// /// - Phase 0x1 → SessionSetup (handshake msg1) /// - Phase 0x2 → SessionAck (handshake msg2) - /// - Phase 0x3 → SessionMsg3 (XK handshake msg3) + /// - Phase 0x3 → SessionMsg3 (XX handshake msg3) /// - Phase 0x0 + U flag → plaintext error signal (CoordsRequired/PathBroken) /// - Phase 0x0 + !U → encrypted session message (data, reports, etc.) pub(in crate::node) async fn handle_session_payload( @@ -50,10 +49,7 @@ impl Node { let prefix = match FspCommonPrefix::parse(payload) { Some(p) => p, None => { - debug!( - len = payload.len(), - "Session payload too short for FSP prefix" - ); + debug!(len = payload.len(), "Session payload too short for FSP prefix"); return; } }; @@ -94,8 +90,7 @@ impl Node { } } FSP_PHASE_ESTABLISHED => { - self.handle_encrypted_session_msg(src_addr, payload, path_mtu, ce_flag) - .await; + self.handle_encrypted_session_msg(src_addr, payload, path_mtu, ce_flag).await; } _ => { debug!(phase = prefix.phase, "Unknown FSP phase"); @@ -112,21 +107,12 @@ impl Node { /// 4. AEAD decrypt with AAD = header_bytes /// 5. Strip FSP inner header → timestamp, msg_type, inner_flags /// 6. Dispatch by msg_type - async fn handle_encrypted_session_msg( - &mut self, - src_addr: &NodeAddr, - payload: &[u8], - path_mtu: u16, - ce_flag: bool, - ) { + async fn handle_encrypted_session_msg(&mut self, src_addr: &NodeAddr, payload: &[u8], path_mtu: u16, ce_flag: bool) { // Parse the 12-byte encrypted header (includes the 4-byte prefix) let header = match FspEncryptedHeader::parse(payload) { Some(h) => h, None => { - debug!( - len = payload.len(), - "Encrypted session message too short for FSP header" - ); + debug!(len = payload.len(), "Encrypted session message too short for FSP header"); return; } }; @@ -181,8 +167,8 @@ impl Node { let received_k_bit = header.flags & FSP_FLAG_K != 0; { let entry = self.sessions.get(src_addr).unwrap(); - let k_bit_flipped = - received_k_bit != entry.current_k_bit() && entry.pending_new_session().is_some(); + let k_bit_flipped = received_k_bit != entry.current_k_bit() + && entry.pending_new_session().is_some(); if k_bit_flipped { let display_name = self.peer_display_name(src_addr); @@ -249,8 +235,7 @@ impl Node { self.sessions.insert(*src_addr, entry); // Strip FSP inner header (6 bytes) - let (timestamp, msg_type, inner_flags_byte, rest) = match fsp_strip_inner_header(&plaintext) - { + let (timestamp, msg_type, inner_flags_byte, rest) = match fsp_strip_inner_header(&plaintext) { Some(parts) => parts, None => { debug!(src = %self.peer_display_name(src_addr), "Decrypted payload too short for FSP inner header"); @@ -263,15 +248,16 @@ impl Node { && let Some(mmp) = entry.mmp_mut() { let now = std::time::Instant::now(); - mmp.receiver - .record_recv(header.counter, timestamp, plaintext.len(), ce_flag, now); + mmp.receiver.record_recv( + header.counter, timestamp, plaintext.len(), ce_flag, now, + ); // Spin bit: advance state machine for correct TX reflection. // RTT samples not fed into SRTT — timestamp-echo provides // accurate RTT; spin bit includes variable inter-frame delays. let inner_flags = FspInnerFlags::from_byte(inner_flags_byte); - let _spin_rtt = mmp - .spin_bit - .rx_observe(inner_flags.spin_bit, header.counter, now); + let _spin_rtt = mmp.spin_bit.rx_observe( + inner_flags.spin_bit, header.counter, now, + ); } // Feed path_mtu from datagram envelope to MMP path MTU tracking. @@ -298,15 +284,9 @@ impl Node { FSP_PORT_IPV6_SHIM => { use crate::FipsAddress; let src_ipv6 = FipsAddress::from_node_addr(src_addr).to_ipv6().octets(); - let dst_ipv6 = FipsAddress::from_node_addr(self.node_addr()) - .to_ipv6() - .octets(); + let dst_ipv6 = FipsAddress::from_node_addr(self.node_addr()).to_ipv6().octets(); - match crate::upper::ipv6_shim::decompress_ipv6( - service_payload, - src_ipv6, - dst_ipv6, - ) { + match crate::upper::ipv6_shim::decompress_ipv6(service_payload, src_ipv6, dst_ipv6) { Some(mut packet) => { if ce_flag { mark_ipv6_ecn_ce(&mut packet); @@ -374,7 +354,7 @@ impl Node { self.flush_pending_packets(src_addr).await; } - /// Handle an incoming SessionSetup (Noise XK msg1). + /// Handle an incoming SessionSetup (Noise XX msg1). /// /// The remote node wants to establish an end-to-end session with us. /// We create an XK responder handshake, process msg1, send SessionAck with msg2, @@ -388,10 +368,10 @@ impl Node { } }; - if setup.handshake_payload.len() != XK_HANDSHAKE_MSG1_SIZE { + if setup.handshake_payload.len() != XX_HANDSHAKE_MSG1_SIZE { debug!( len = setup.handshake_payload.len(), - expected = XK_HANDSHAKE_MSG1_SIZE, + expected = XX_HANDSHAKE_MSG1_SIZE, "Invalid handshake payload size in SessionSetup" ); return; @@ -463,19 +443,19 @@ impl Node { return; } let our_keypair = self.identity.keypair(); - let mut handshake = HandshakeState::new_xk_responder(our_keypair); + let mut handshake = HandshakeState::new_xx_responder(our_keypair); handshake.set_local_epoch(self.startup_epoch); - if let Err(e) = handshake.read_xk_message_1(&setup.handshake_payload) { - debug!(error = %e, "Failed to process rekey XK msg1"); + if let Err(e) = handshake.read_xx_message_1(&setup.handshake_payload) { + debug!(error = %e, "Failed to process rekey XX msg1"); return; } // Generate msg2 - let msg2 = match handshake.write_xk_message_2() { + let msg2 = match handshake.write_xx_message_2() { Ok(m) => m, Err(e) => { - debug!(error = %e, "Failed to generate rekey XK msg2"); + debug!(error = %e, "Failed to generate rekey XX msg2"); return; } }; @@ -513,11 +493,11 @@ impl Node { // Create XK responder handshake and process msg1 let our_keypair = self.identity.keypair(); - let mut handshake = HandshakeState::new_xk_responder(our_keypair); + let mut handshake = HandshakeState::new_xx_responder(our_keypair); handshake.set_local_epoch(self.startup_epoch); - if let Err(e) = handshake.read_xk_message_1(&setup.handshake_payload) { - debug!(error = %e, "Failed to process Noise XK msg1 in SessionSetup"); + if let Err(e) = handshake.read_xx_message_1(&setup.handshake_payload) { + debug!(error = %e, "Failed to process Noise XX msg1 in SessionSetup"); return; } @@ -525,15 +505,25 @@ impl Node { // Use a placeholder pubkey from src_addr for the session entry. // The real pubkey will be registered when msg3 arrives. - // Generate msg2 - let msg2 = match handshake.write_xk_message_2() { + // Generate msg2 with negotiation payload + let mut msg2 = match handshake.write_xx_message_2() { Ok(m) => m, Err(e) => { - debug!(error = %e, "Failed to generate Noise XK msg2 for SessionAck"); + debug!(error = %e, "Failed to generate Noise XX msg2 for SessionAck"); return; } }; + // Encrypt FSP negotiation payload (version [0,0], features=0) + let neg_payload = NegotiationPayload::new(0, 0, 0).encode(); + match handshake.encrypt_payload(&neg_payload) { + Ok(encrypted) => msg2.extend_from_slice(&encrypted), + Err(e) => { + debug!(error = %e, "Failed to encrypt negotiation payload for SessionAck"); + return; + } + } + // Build and send SessionAck (include initiator's coords for return-path warming) let our_coords = self.tree_state.my_coords().clone(); let ack = SessionAck::new(our_coords, setup.src_coords).with_handshake(msg2); @@ -554,20 +544,14 @@ impl Node { let placeholder_pubkey = self.identity.keypair().public_key(); let now_ms = Self::now_ms(); let resend_interval = self.config.node.rate_limit.handshake_resend_interval_ms; - let mut entry = SessionEntry::new( - *src_addr, - placeholder_pubkey, - EndToEndState::AwaitingMsg3(handshake), - now_ms, - false, - ); + let mut entry = SessionEntry::new(*src_addr, placeholder_pubkey, EndToEndState::AwaitingMsg3(handshake), now_ms, false); entry.set_handshake_payload(ack_payload, now_ms + resend_interval); self.sessions.insert(*src_addr, entry); - debug!(src = %self.peer_display_name(src_addr), "SessionSetup processed (XK), SessionAck sent, awaiting msg3"); + debug!(src = %self.peer_display_name(src_addr), "SessionSetup processed (XX), SessionAck sent, awaiting msg3"); } - /// Handle an incoming SessionAck (Noise XK msg2). + /// Handle an incoming SessionAck (Noise XX msg2). /// /// Processes msg2, generates and sends msg3, then transitions to Established. async fn handle_session_ack(&mut self, src_addr: &NodeAddr, inner: &[u8]) { @@ -579,11 +563,11 @@ impl Node { } }; - if ack.handshake_payload.len() != XK_HANDSHAKE_MSG2_SIZE { + if ack.handshake_payload.len() < XX_HANDSHAKE_MSG2_SIZE { debug!( len = ack.handshake_payload.len(), - expected = XK_HANDSHAKE_MSG2_SIZE, - "Invalid handshake payload size in SessionAck" + min = XX_HANDSHAKE_MSG2_SIZE, + "Handshake payload too short in SessionAck" ); return; } @@ -608,18 +592,18 @@ impl Node { }; // Process XK msg2 - if let Err(e) = handshake.read_xk_message_2(&ack.handshake_payload) { - debug!(error = %e, "Failed to process rekey XK msg2"); + if let Err(e) = handshake.read_xx_message_2(&ack.handshake_payload) { + debug!(error = %e, "Failed to process rekey XX msg2"); entry.abandon_rekey(); self.sessions.insert(*src_addr, entry); return; } // Generate XK msg3 - let msg3 = match handshake.write_xk_message_3() { + let msg3 = match handshake.write_xx_message_3() { Ok(m) => m, Err(e) => { - debug!(error = %e, "Failed to generate rekey XK msg3"); + debug!(error = %e, "Failed to generate rekey XX msg3"); entry.abandon_rekey(); self.sessions.insert(*src_addr, entry); return; @@ -644,7 +628,7 @@ impl Node { let session = match handshake.into_session() { Ok(s) => s, Err(e) => { - debug!(error = %e, "Failed to create session from rekey XK"); + debug!(error = %e, "Failed to create session from rekey XX"); entry.abandon_rekey(); self.sessions.insert(*src_addr, entry); return; @@ -673,21 +657,65 @@ impl Node { _ => unreachable!("checked is_initiating above"), }; - // Process XK msg2: read_xk_message_2 (extracts responder's epoch) - if let Err(e) = handshake.read_xk_message_2(&ack.handshake_payload) { - debug!(error = %e, "Failed to process Noise XK msg2 in SessionAck"); - return; // Entry was already removed, don't put back a broken session + // Split msg2 into base XX part and optional negotiation payload + let (base_msg2, neg_bytes) = if ack.handshake_payload.len() > XX_HANDSHAKE_MSG2_SIZE { + (&ack.handshake_payload[..XX_HANDSHAKE_MSG2_SIZE], Some(&ack.handshake_payload[XX_HANDSHAKE_MSG2_SIZE..])) + } else { + (ack.handshake_payload.as_slice(), None) + }; + + // Process XX msg2 (learns responder's identity and epoch) + if let Err(e) = handshake.read_xx_message_2(base_msg2) { + debug!(error = %e, "Failed to process Noise XX msg2 in SessionAck"); + return; } - // Generate XK msg3: write_xk_message_3 (sends encrypted static + epoch) - let msg3 = match handshake.write_xk_message_3() { + // Decrypt negotiation payload from msg2 if present + if let Some(encrypted_neg) = neg_bytes { + match handshake.decrypt_payload(encrypted_neg) { + Ok(_negotiation) => { + // FSP negotiation payload received — currently unused (version [0,0]) + } + Err(e) => { + debug!(error = %e, "Failed to decrypt negotiation payload from SessionAck"); + return; + } + } + } + + // XX: verify responder's identity matches the target we intended to reach. + // Compare x-only keys to avoid parity mismatch: npub-derived keys always + // have even parity (0x02), but the Noise handshake reveals the real parity. + let expected_xonly = entry.remote_pubkey().x_only_public_key().0; + if let Some(remote_pk) = handshake.remote_static() + && remote_pk.x_only_public_key().0 != expected_xonly + { + debug!( + src = %self.peer_display_name(src_addr), + "Responder identity mismatch in SessionAck — disconnecting" + ); + return; + } + + // Generate XX msg3 with negotiation payload + let mut msg3 = match handshake.write_xx_message_3() { Ok(m) => m, Err(e) => { - debug!(error = %e, "Failed to generate Noise XK msg3"); + debug!(error = %e, "Failed to generate Noise XX msg3"); return; } }; + // Encrypt FSP negotiation payload for msg3 + let neg_payload = NegotiationPayload::new(0, 0, 0).encode(); + match handshake.encrypt_payload(&neg_payload) { + Ok(encrypted) => msg3.extend_from_slice(&encrypted), + Err(e) => { + debug!(error = %e, "Failed to encrypt negotiation payload for SessionMsg3"); + return; + } + } + // Send SessionMsg3 (phase 0x3) let msg3_wire = SessionMsg3::new(msg3); let msg3_payload = msg3_wire.encode(); @@ -725,7 +753,7 @@ impl Node { info!(src = %self.peer_display_name(src_addr), "Session established (initiator, XK)"); } - /// Handle an incoming SessionMsg3 (Noise XK msg3). + /// Handle an incoming SessionMsg3 (Noise XX msg3). /// /// The initiator reveals their encrypted static key. The responder /// processes msg3, learns the initiator's identity, and transitions @@ -739,11 +767,11 @@ impl Node { } }; - if msg3.handshake_payload.len() != XK_HANDSHAKE_MSG3_SIZE { + if msg3.handshake_payload.len() < XX_HANDSHAKE_MSG3_SIZE { debug!( len = msg3.handshake_payload.len(), - expected = XK_HANDSHAKE_MSG3_SIZE, - "Invalid handshake payload size in SessionMsg3" + min = XX_HANDSHAKE_MSG3_SIZE, + "Handshake payload too short in SessionMsg3" ); return; } @@ -768,8 +796,8 @@ impl Node { }; // Process XK msg3 - if let Err(e) = handshake.read_xk_message_3(&msg3.handshake_payload) { - debug!(error = %e, "Failed to process rekey XK msg3"); + if let Err(e) = handshake.read_xx_message_3(&msg3.handshake_payload) { + debug!(error = %e, "Failed to process rekey XX msg3"); entry.abandon_rekey(); self.sessions.insert(*src_addr, entry); return; @@ -779,7 +807,7 @@ impl Node { let session = match handshake.into_session() { Ok(s) => s, Err(e) => { - debug!(error = %e, "Failed to create session from rekey XK msg3"); + debug!(error = %e, "Failed to create session from rekey XX msg3"); entry.abandon_rekey(); self.sessions.insert(*src_addr, entry); return; @@ -807,10 +835,30 @@ impl Node { _ => unreachable!("checked is_awaiting_msg3 above"), }; - // Process XK msg3: read_xk_message_3 (extracts initiator's static key and epoch) - if let Err(e) = handshake.read_xk_message_3(&msg3.handshake_payload) { - debug!(error = %e, "Failed to process Noise XK msg3"); - return; // Entry was already removed + // Split msg3 into base XX part and optional negotiation payload + let (base_msg3, neg_bytes) = if msg3.handshake_payload.len() > XX_HANDSHAKE_MSG3_SIZE { + (&msg3.handshake_payload[..XX_HANDSHAKE_MSG3_SIZE], Some(&msg3.handshake_payload[XX_HANDSHAKE_MSG3_SIZE..])) + } else { + (msg3.handshake_payload.as_slice(), None) + }; + + // Process XX msg3 (learns initiator's identity and epoch) + if let Err(e) = handshake.read_xx_message_3(base_msg3) { + debug!(error = %e, "Failed to process Noise XX msg3"); + return; + } + + // Decrypt negotiation payload from msg3 if present + if let Some(encrypted_neg) = neg_bytes { + match handshake.decrypt_payload(encrypted_neg) { + Ok(_negotiation) => { + // FSP negotiation payload received — currently unused (version [0,0]) + } + Err(e) => { + debug!(error = %e, "Failed to decrypt negotiation payload from SessionMsg3"); + return; + } + } } // Extract the initiator's static public key (now available after msg3) @@ -836,13 +884,7 @@ impl Node { let now_ms = Self::now_ms(); // Replace the placeholder pubkey with the real one - let mut new_entry = SessionEntry::new( - *src_addr, - remote_pubkey, - EndToEndState::Established(session), - now_ms, - false, - ); + let mut new_entry = SessionEntry::new(*src_addr, remote_pubkey, EndToEndState::Established(session), now_ms, false); new_entry.set_coords_warmup_remaining(self.config.node.session.coords_warmup_packets); new_entry.mark_established(now_ms); new_entry.init_mmp(&self.config.node.session_mmp); @@ -911,8 +953,7 @@ impl Node { }; let now = std::time::Instant::now(); - mmp.metrics - .process_receiver_report(&rr, our_timestamp_ms, now); + mmp.metrics.process_receiver_report(&rr, our_timestamp_ms, now); // Feed SRTT back to sender/receiver report interval tuning (session-layer bounds) if let Some(srtt_ms) = mmp.metrics.srtt_ms() { @@ -934,8 +975,7 @@ impl Node { // Update reverse delivery ratio from our own receiver state, using per-interval deltas. let our_recv_packets = mmp.receiver.cumulative_packets_recv(); let peer_highest = mmp.receiver.highest_counter(); - mmp.metrics - .update_reverse_delivery(our_recv_packets, peer_highest); + mmp.metrics.update_reverse_delivery(our_recv_packets, peer_highest); trace!( src = %peer_name, @@ -1010,10 +1050,7 @@ impl Node { ); // Send standalone CoordsWarmup immediately (rate-limited) - if self - .coords_response_rate_limiter - .should_send(&msg.dest_addr) - { + if self.coords_response_rate_limiter.should_send(&msg.dest_addr) { if let Some(entry) = self.sessions.get(&msg.dest_addr) && entry.is_established() && let Err(e) = self.send_coords_warmup(&msg.dest_addr).await @@ -1071,10 +1108,7 @@ impl Node { ); // Send standalone CoordsWarmup immediately (rate-limited) - if self - .coords_response_rate_limiter - .should_send(&msg.dest_addr) - { + if self.coords_response_rate_limiter.should_send(&msg.dest_addr) { if let Some(entry) = self.sessions.get(&msg.dest_addr) && entry.is_established() && let Err(e) = self.send_coords_warmup(&msg.dest_addr).await @@ -1161,7 +1195,7 @@ impl Node { /// Initiate an end-to-end session with a remote node. /// - /// Creates a Noise XK handshake as initiator, wraps msg1 in a + /// Creates a Noise XX handshake as initiator, wraps msg1 in a /// SessionSetup, encapsulates in a SessionDatagram, and routes /// toward the destination. pub(in crate::node) async fn initiate_session( @@ -1176,21 +1210,20 @@ impl Node { return Ok(()); } - // Create Noise XK initiator handshake + // Create Noise XX initiator handshake let our_keypair = self.identity.keypair(); - let mut handshake = HandshakeState::new_xk_initiator(our_keypair, dest_pubkey); + let mut handshake = HandshakeState::new_xx_initiator(our_keypair); handshake.set_local_epoch(self.startup_epoch); - let msg1 = handshake - .write_xk_message_1() - .map_err(|e| NodeError::SendFailed { - node_addr: dest_addr, - reason: format!("Noise XK msg1 generation failed: {}", e), - })?; + let msg1 = handshake.write_xx_message_1().map_err(|e| NodeError::SendFailed { + node_addr: dest_addr, + reason: format!("Noise XX msg1 generation failed: {}", e), + })?; // Build SessionSetup with coordinates let our_coords = self.tree_state.my_coords().clone(); let dest_coords = self.get_dest_coords(&dest_addr); - let setup = SessionSetup::new(our_coords, dest_coords).with_handshake(msg1); + let setup = SessionSetup::new(our_coords, dest_coords) + .with_handshake(msg1); let setup_payload = setup.encode(); // Wrap in SessionDatagram @@ -1207,13 +1240,7 @@ impl Node { // Store session entry with handshake payload for potential resend let now_ms = Self::now_ms(); let resend_interval = self.config.node.rate_limit.handshake_resend_interval_ms; - let mut entry = SessionEntry::new( - dest_addr, - dest_pubkey, - EndToEndState::Initiating(handshake), - now_ms, - true, - ); + let mut entry = SessionEntry::new(dest_addr, dest_pubkey, EndToEndState::Initiating(handshake), now_ms, true); entry.set_handshake_payload(setup_payload, now_ms + resend_interval); self.sessions.insert(dest_addr, entry); @@ -1240,13 +1267,10 @@ impl Node { let now_ms = Self::now_ms(); // First borrow: read session metadata (NLL releases before coord decision) - let entry = self - .sessions - .get(dest_addr) - .ok_or_else(|| NodeError::SendFailed { - node_addr: *dest_addr, - reason: "no session".into(), - })?; + let entry = self.sessions.get(dest_addr).ok_or_else(|| NodeError::SendFailed { + node_addr: *dest_addr, + reason: "no session".into(), + })?; let wants_coords = entry.coords_warmup_remaining() > 0; let timestamp = entry.session_timestamp(now_ms); let spin_bit = entry.mmp().is_some_and(|m| m.spin_bit.tx_bit()); @@ -1266,8 +1290,7 @@ impl Node { // Build inner plaintext (doesn't depend on counter) let msg_type = SessionMessageType::DataPacket.to_byte(); // 0x10 let inner_flags = FspInnerFlags { spin_bit }.to_byte(); - let inner_plaintext = - fsp_prepend_inner_header(timestamp, msg_type, inner_flags, &port_payload); + let inner_plaintext = fsp_prepend_inner_header(timestamp, msg_type, inner_flags, &port_payload); // Determine whether coords fit within transport MTU. // If not, send standalone CoordsWarmup before the data packet. @@ -1275,8 +1298,7 @@ impl Node { let src = self.tree_state.my_coords().clone(); let dst = self.get_dest_coords(dest_addr); let coords_size = coords_wire_size(&src) + coords_wire_size(&dst); - let total_wire = - FIPS_OVERHEAD as usize + FSP_PORT_HEADER_SIZE + coords_size + payload.len(); + let total_wire = FIPS_OVERHEAD as usize + FSP_PORT_HEADER_SIZE + coords_size + payload.len(); if total_wire <= self.transport_mtu() as usize { (true, Some(src), Some(dst)) } else { @@ -1292,7 +1314,9 @@ impl Node { }; // Decrement warmup counter if we sent coords (piggybacked or standalone) - if wants_coords && let Some(entry) = self.sessions.get_mut(dest_addr) { + if wants_coords + && let Some(entry) = self.sessions.get_mut(dest_addr) + { entry.set_coords_warmup_remaining(entry.coords_warmup_remaining() - 1); } @@ -1305,13 +1329,10 @@ impl Node { } // Borrow session for counter + encryption (after potential standalone send) - let entry = self - .sessions - .get_mut(dest_addr) - .ok_or_else(|| NodeError::SendFailed { - node_addr: *dest_addr, - reason: "no session".into(), - })?; + let entry = self.sessions.get_mut(dest_addr).ok_or_else(|| NodeError::SendFailed { + node_addr: *dest_addr, + reason: "no session".into(), + })?; let session = match entry.state_mut() { EndToEndState::Established(s) => s, _ => { @@ -1328,12 +1349,12 @@ impl Node { let header = build_fsp_header(counter, flags, payload_len); // Encrypt with AAD binding to the FSP header - let ciphertext = session - .encrypt_with_aad(&inner_plaintext, &header) - .map_err(|e| NodeError::SendFailed { + let ciphertext = session.encrypt_with_aad(&inner_plaintext, &header).map_err(|e| { + NodeError::SendFailed { node_addr: *dest_addr, reason: format!("session encrypt failed: {}", e), - })?; + } + })?; // Assemble: header(12) + [coords] + ciphertext let mut fsp_payload = Vec::with_capacity(FSP_HEADER_SIZE + ciphertext.len() + 200); @@ -1371,19 +1392,13 @@ impl Node { dest_addr: &NodeAddr, ipv6_packet: &[u8], ) -> Result<(), NodeError> { - let compressed = crate::upper::ipv6_shim::compress_ipv6(ipv6_packet).ok_or_else(|| { - NodeError::SendFailed { + let compressed = crate::upper::ipv6_shim::compress_ipv6(ipv6_packet) + .ok_or_else(|| NodeError::SendFailed { node_addr: *dest_addr, reason: "IPv6 header compression failed".into(), - } - })?; - self.send_session_data( - dest_addr, - FSP_PORT_IPV6_SHIM, - FSP_PORT_IPV6_SHIM, - &compressed, - ) - .await + })?; + self.send_session_data(dest_addr, FSP_PORT_IPV6_SHIM, FSP_PORT_IPV6_SHIM, &compressed) + .await } /// Send a non-data session message (reports, notifications) over an established session. @@ -1402,13 +1417,10 @@ impl Node { let now_ms = Self::now_ms(); // Read spin bit and session timestamp from entry - let entry = self - .sessions - .get(dest_addr) - .ok_or_else(|| NodeError::SendFailed { - node_addr: *dest_addr, - reason: "no session".into(), - })?; + let entry = self.sessions.get(dest_addr).ok_or_else(|| NodeError::SendFailed { + node_addr: *dest_addr, + reason: "no session".into(), + })?; let timestamp = entry.session_timestamp(now_ms); let spin_bit = entry.mmp().is_some_and(|m| m.spin_bit.tx_bit()); @@ -1416,13 +1428,10 @@ impl Node { let inner_flags = FspInnerFlags { spin_bit }.to_byte(); // Get mutable access for encryption - let entry = self - .sessions - .get_mut(dest_addr) - .ok_or_else(|| NodeError::SendFailed { - node_addr: *dest_addr, - reason: "no session".into(), - })?; + let entry = self.sessions.get_mut(dest_addr).ok_or_else(|| NodeError::SendFailed { + node_addr: *dest_addr, + reason: "no session".into(), + })?; // Read K-bit before mutable borrow of session state let k_flags = if entry.current_k_bit() { FSP_FLAG_K } else { 0 }; @@ -1447,12 +1456,12 @@ impl Node { let header = build_fsp_header(counter, k_flags, payload_len); // Encrypt with AAD - let ciphertext = session - .encrypt_with_aad(&inner_plaintext, &header) - .map_err(|e| NodeError::SendFailed { + let ciphertext = session.encrypt_with_aad(&inner_plaintext, &header).map_err(|e| { + NodeError::SendFailed { node_addr: *dest_addr, reason: format!("session encrypt failed: {}", e), - })?; + } + })?; // Assemble: header(12) + ciphertext (no coords) let mut fsp_payload = Vec::with_capacity(FSP_HEADER_SIZE + ciphertext.len()); @@ -1482,31 +1491,28 @@ impl Node { /// coordinates via `try_warm_coord_cache()` (same as CP-flagged data /// packets). The encrypted inner payload is the 6-byte inner header /// with no application data. - async fn send_coords_warmup(&mut self, dest_addr: &NodeAddr) -> Result<(), NodeError> { + async fn send_coords_warmup( + &mut self, + dest_addr: &NodeAddr, + ) -> Result<(), NodeError> { let now_ms = Self::now_ms(); let my_coords = self.tree_state.my_coords().clone(); let dest_coords = self.get_dest_coords(dest_addr); // Read session metadata - let entry = self - .sessions - .get(dest_addr) - .ok_or_else(|| NodeError::SendFailed { - node_addr: *dest_addr, - reason: "no session".into(), - })?; + let entry = self.sessions.get(dest_addr).ok_or_else(|| NodeError::SendFailed { + node_addr: *dest_addr, + reason: "no session".into(), + })?; let timestamp = entry.session_timestamp(now_ms); let spin_bit = entry.mmp().is_some_and(|m| m.spin_bit.tx_bit()); // Get mutable access for encryption - let entry = self - .sessions - .get_mut(dest_addr) - .ok_or_else(|| NodeError::SendFailed { - node_addr: *dest_addr, - reason: "no session".into(), - })?; + let entry = self.sessions.get_mut(dest_addr).ok_or_else(|| NodeError::SendFailed { + node_addr: *dest_addr, + reason: "no session".into(), + })?; let session = match entry.state_mut() { EndToEndState::Established(s) => s, _ => { @@ -1529,12 +1535,12 @@ impl Node { let header = build_fsp_header(counter, FSP_FLAG_CP, payload_len); // Encrypt with AAD - let ciphertext = session - .encrypt_with_aad(&inner_plaintext, &header) - .map_err(|e| NodeError::SendFailed { + let ciphertext = session.encrypt_with_aad(&inner_plaintext, &header).map_err(|e| { + NodeError::SendFailed { node_addr: *dest_addr, reason: format!("session encrypt failed: {}", e), - })?; + } + })?; // Assemble: header(12) + coords + ciphertext let coords_size = coords_wire_size(&my_coords) + coords_wire_size(&dest_coords); @@ -1601,8 +1607,7 @@ impl Node { } let encoded = datagram.encode(); - self.send_encrypted_link_message(&next_hop_addr, &encoded) - .await?; + self.send_encrypted_link_message(&next_hop_addr, &encoded).await?; self.stats_mut().forwarding.record_originated(encoded.len()); Ok(()) } @@ -1709,20 +1714,19 @@ impl Node { /// Send ICMPv6 Destination Unreachable back through TUN. pub(in crate::node) fn send_icmpv6_dest_unreachable(&self, original_packet: &[u8]) { + use crate::upper::icmp::{build_dest_unreachable, should_send_icmp_error, DestUnreachableCode}; use crate::FipsAddress; - use crate::upper::icmp::{ - DestUnreachableCode, build_dest_unreachable, should_send_icmp_error, - }; if !should_send_icmp_error(original_packet) { return; } let our_ipv6 = FipsAddress::from_node_addr(self.node_addr()).to_ipv6(); - if let Some(response) = - build_dest_unreachable(original_packet, DestUnreachableCode::NoRoute, our_ipv6) - && let Some(tun_tx) = &self.tun_tx - { + if let Some(response) = build_dest_unreachable( + original_packet, + DestUnreachableCode::NoRoute, + our_ipv6, + ) && let Some(tun_tx) = &self.tun_tx { let _ = tun_tx.send(response); } } @@ -1779,7 +1783,10 @@ impl Node { return; } - let queue = self.pending_tun_packets.entry(dest_addr).or_default(); + let queue = self + .pending_tun_packets + .entry(dest_addr) + .or_default(); if queue.len() >= self.config.node.session.pending_packets_per_dest { queue.pop_front(); // Drop oldest } diff --git a/src/node/handlers/timeout.rs b/src/node/handlers/timeout.rs index 547c494..c8c05c4 100644 --- a/src/node/handlers/timeout.rs +++ b/src/node/handlers/timeout.rs @@ -22,9 +22,7 @@ impl Node { .unwrap_or(0); let timeout_ms = self.config.node.rate_limit.handshake_timeout_secs * 1000; - let stale: Vec = self - .connections - .iter() + let stale: Vec = self.connections.iter() .filter(|(_, conn)| conn.is_timed_out(now_ms, timeout_ms) || conn.is_failed()) .map(|(link_id, _)| *link_id) .collect(); @@ -99,15 +97,19 @@ impl Node { // Collect resend candidates: outbound, in SentMsg1, with stored msg1, // under max resends, and past the scheduled time. - let candidates: Vec<(LinkId, Vec)> = self - .connections - .iter() + // Skip resend if the target peer is already promoted — a cross-connection + // was resolved via the inbound path and resending msg1 would start a new + // handshake on the peer, creating a session mismatch. + let candidates: Vec<(LinkId, Vec)> = self.connections.iter() .filter(|(_, conn)| { conn.is_outbound() && conn.handshake_state() == HandshakeState::SentMsg1 && conn.resend_count() < max_resends && conn.next_resend_at_ms() > 0 && now_ms >= conn.next_resend_at_ms() + && !conn.expected_identity() + .map(|id| self.peers.contains_key(id.node_addr())) + .unwrap_or(false) }) .filter_map(|(link_id, conn)| { conn.handshake_msg1().map(|msg1| (*link_id, msg1.to_vec())) @@ -141,7 +143,9 @@ impl Node { false }; - if sent && let Some(conn) = self.connections.get_mut(&link_id) { + if sent + && let Some(conn) = self.connections.get_mut(&link_id) + { let count = conn.resend_count() + 1; let next = now_ms + (interval_ms as f64 * backoff.powi(count as i32)) as u64; conn.record_resend(next); @@ -172,11 +176,10 @@ impl Node { let ttl = self.config.node.session.default_ttl; // First pass: find timed-out sessions to remove - let timed_out: Vec = self - .sessions - .iter() + let timed_out: Vec = self.sessions.iter() .filter(|(_, entry)| { - !entry.is_established() && now_ms.saturating_sub(entry.last_activity()) > timeout_ms + !entry.is_established() + && now_ms.saturating_sub(entry.last_activity()) > timeout_ms }) .map(|(addr, _)| *addr) .collect(); @@ -190,9 +193,7 @@ impl Node { // Second pass: collect resend candidates let my_addr = *self.node_addr(); - let candidates: Vec<(crate::NodeAddr, Vec)> = self - .sessions - .iter() + let candidates: Vec<(crate::NodeAddr, Vec)> = self.sessions.iter() .filter(|(_, entry)| { !entry.is_established() && entry.handshake_payload().is_some() @@ -206,7 +207,8 @@ impl Node { for (dest_addr, payload) in candidates { use crate::protocol::SessionDatagram; - let mut datagram = SessionDatagram::new(my_addr, dest_addr, payload).with_ttl(ttl); + let mut datagram = SessionDatagram::new(my_addr, dest_addr, payload) + .with_ttl(ttl); let sent = match self.send_session_datagram(&mut datagram).await { Ok(_) => true, Err(e) => { @@ -219,7 +221,9 @@ impl Node { } }; - if sent && let Some(entry) = self.sessions.get_mut(&dest_addr) { + if sent + && let Some(entry) = self.sessions.get_mut(&dest_addr) + { let count = entry.resend_count() + 1; let next = now_ms + (interval_ms as f64 * backoff.powi(count as i32)) as u64; entry.record_resend(next); @@ -242,11 +246,10 @@ impl Node { return; // disabled } - let idle: Vec<_> = self - .sessions - .iter() + let idle: Vec<_> = self.sessions.iter() .filter(|(_, entry)| { - entry.is_established() && now_ms.saturating_sub(entry.last_activity()) > timeout_ms + entry.is_established() + && now_ms.saturating_sub(entry.last_activity()) > timeout_ms }) .map(|(addr, _)| *addr) .collect(); diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index bfd86aa..8223357 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -3,8 +3,8 @@ use super::*; use crate::node::session::EndToEndState; use crate::node::tests::spanning_tree::{ - TestNode, cleanup_nodes, generate_random_edges, process_available_packets, run_tree_test, - run_tree_test_with_mtus, verify_tree_convergence, + cleanup_nodes, generate_random_edges, process_available_packets, run_tree_test, + run_tree_test_with_mtus, verify_tree_convergence, TestNode, }; use crate::protocol::{SessionAck, SessionDatagram}; @@ -50,7 +50,10 @@ fn test_session_entry_new_initiating() { let identity_a = Identity::generate(); let identity_b = Identity::generate(); - let handshake = HandshakeState::new_initiator(identity_a.keypair(), identity_b.pubkey_full()); + let handshake = HandshakeState::new_initiator( + identity_a.keypair(), + identity_b.pubkey_full(), + ); let entry = crate::node::session::SessionEntry::new( *identity_b.node_addr(), @@ -74,7 +77,10 @@ fn test_session_entry_touch() { let identity_a = Identity::generate(); let identity_b = Identity::generate(); - let handshake = HandshakeState::new_initiator(identity_a.keypair(), identity_b.pubkey_full()); + let handshake = HandshakeState::new_initiator( + identity_a.keypair(), + identity_b.pubkey_full(), + ); let mut entry = crate::node::session::SessionEntry::new( *identity_b.node_addr(), @@ -96,8 +102,10 @@ fn test_session_table_operations() { let mut node = make_node(); let identity_b = Identity::generate(); - let handshake = - HandshakeState::new_initiator(node.identity().keypair(), identity_b.pubkey_full()); + let handshake = HandshakeState::new_initiator( + node.identity().keypair(), + identity_b.pubkey_full(), + ); let dest_addr = *identity_b.node_addr(); let entry = crate::node::session::SessionEntry::new( @@ -143,14 +151,12 @@ async fn test_session_direct_peer_handshake() { // Node 0 should have a session in Initiating state assert_eq!(nodes[0].node.session_count(), 1); - assert!( - nodes[0] - .node - .get_session(&node1_addr) - .unwrap() - .state() - .is_initiating() - ); + assert!(nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .state() + .is_initiating()); // Process packets: SessionSetup arrives at Node 1 tokio::time::sleep(Duration::from_millis(20)).await; @@ -159,14 +165,12 @@ async fn test_session_direct_peer_handshake() { // Node 1 should now have a session in AwaitingMsg3 state (XK: identity not yet known) assert_eq!(nodes[1].node.session_count(), 1); - assert!( - nodes[1] - .node - .get_session(&node0_addr) - .unwrap() - .state() - .is_awaiting_msg3() - ); + assert!(nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .state() + .is_awaiting_msg3()); // Process packets: SessionAck arrives at Node 0, Node 0 sends SessionMsg3 tokio::time::sleep(Duration::from_millis(20)).await; @@ -174,14 +178,12 @@ async fn test_session_direct_peer_handshake() { assert!(count > 0, "Expected SessionAck packet to arrive"); // Node 0 should now be Established (transitions after sending msg3) - assert!( - nodes[0] - .node - .get_session(&node1_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .state() + .is_established()); // Process packets: SessionMsg3 arrives at Node 1 tokio::time::sleep(Duration::from_millis(20)).await; @@ -189,14 +191,12 @@ async fn test_session_direct_peer_handshake() { assert!(count > 0, "Expected SessionMsg3 packet to arrive"); // Node 1 should now be Established (transitions after processing msg3) - assert!( - nodes[1] - .node - .get_session(&node0_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .state() + .is_established()); cleanup_nodes(&mut nodes).await; } @@ -226,22 +226,18 @@ async fn test_session_direct_peer_data_transfer() { tokio::time::sleep(Duration::from_millis(20)).await; process_available_packets(&mut nodes).await; // Msg3 → Node 1 - assert!( - nodes[0] - .node - .get_session(&node1_addr) - .unwrap() - .state() - .is_established() - ); - assert!( - nodes[1] - .node - .get_session(&node0_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .state() + .is_established()); + assert!(nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .state() + .is_established()); // Send data from Node 0 to Node 1 let test_data = b"Hello, FIPS session!"; @@ -295,14 +291,12 @@ async fn test_session_3node_forwarded_handshake() { nodes[2].node.get_session(&node0_addr).is_some(), "Node 2 should have a session entry for Node 0" ); - assert!( - nodes[2] - .node - .get_session(&node0_addr) - .unwrap() - .state() - .is_awaiting_msg3() - ); + assert!(nodes[2] + .node + .get_session(&node0_addr) + .unwrap() + .state() + .is_awaiting_msg3()); // Process: SessionAck: 2→1 (forwarded by transit B) tokio::time::sleep(Duration::from_millis(20)).await; @@ -313,14 +307,12 @@ async fn test_session_3node_forwarded_handshake() { process_available_packets(&mut nodes).await; // Node 0 should now be Established (transitions after sending msg3) - assert!( - nodes[0] - .node - .get_session(&node2_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[0] + .node + .get_session(&node2_addr) + .unwrap() + .state() + .is_established()); // Process: SessionMsg3: 0→1 (forwarded by transit B) tokio::time::sleep(Duration::from_millis(20)).await; @@ -331,14 +323,12 @@ async fn test_session_3node_forwarded_handshake() { process_available_packets(&mut nodes).await; // Node 2 should now be Established (transitions after processing msg3) - assert!( - nodes[2] - .node - .get_session(&node0_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[2] + .node + .get_session(&node0_addr) + .unwrap() + .state() + .is_established()); // Transit node B should NOT have a session assert_eq!( @@ -399,14 +389,12 @@ async fn test_session_3node_forwarded_data() { } // Node 2 should be Established (transitioned during XK handshake msg3) - assert!( - nodes[2] - .node - .get_session(&node0_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[2] + .node + .get_session(&node0_addr) + .unwrap() + .state() + .is_established()); cleanup_nodes(&mut nodes).await; } @@ -532,7 +520,12 @@ async fn test_session_100_nodes() { // Collect identities: (node_addr, pubkey) for all nodes let all_info: Vec<(NodeAddr, secp256k1::PublicKey)> = nodes .iter() - .map(|tn| (*tn.node.node_addr(), tn.node.identity().pubkey_full())) + .map(|tn| { + ( + *tn.node.node_addr(), + tn.node.identity().pubkey_full(), + ) + }) .collect(); // Each node picks one random target for its outbound session. @@ -647,7 +640,11 @@ async fn test_session_100_nodes() { // (Responder should already be Established after XK msg3) let rev_payload = format!("rev-{}", pair_idx).into_bytes(); let rev_ipv6 = build_ipv6_packet(&dst_fips, &src_fips, &rev_payload); - match nodes[dst].node.send_ipv6_packet(&src_addr, &rev_ipv6).await { + match nodes[dst] + .node + .send_ipv6_packet(&src_addr, &rev_ipv6) + .await + { Ok(()) => send_reverse_ok += 1, Err(_) => send_reverse_err += 1, } @@ -726,7 +723,10 @@ async fn test_session_100_nodes() { } } - let session_counts: Vec = nodes.iter().map(|tn| tn.node.session_count()).collect(); + let session_counts: Vec = nodes + .iter() + .map(|tn| tn.node.session_count()) + .collect(); let total_sessions: usize = session_counts.iter().sum(); let min_sessions = *session_counts.iter().min().unwrap(); let max_sessions = *session_counts.iter().max().unwrap(); @@ -770,8 +770,10 @@ async fn test_session_100_nodes() { }; // Coord cache stats - let coord_cache_sizes: Vec = - nodes.iter().map(|tn| tn.node.coord_cache().len()).collect(); + let coord_cache_sizes: Vec = nodes + .iter() + .map(|tn| tn.node.coord_cache().len()) + .collect(); let total_coord_entries: usize = coord_cache_sizes.iter().sum(); let min_coord = *coord_cache_sizes.iter().min().unwrap(); let max_coord = *coord_cache_sizes.iter().max().unwrap(); @@ -882,7 +884,10 @@ async fn test_session_100_nodes() { // === Assertions === - assert_eq!(send_forward_err, 0, "All forward sends should succeed"); + assert_eq!( + send_forward_err, 0, + "All forward sends should succeed" + ); assert_eq!( send_reverse_err, 0, "All reverse sends should succeed (responder Established after XK msg3)" @@ -910,11 +915,7 @@ async fn test_session_100_nodes() { // ============================================================================ /// Build a minimal valid IPv6 packet with given source and destination addresses. -fn build_ipv6_packet( - src: &crate::FipsAddress, - dst: &crate::FipsAddress, - payload: &[u8], -) -> Vec { +fn build_ipv6_packet(src: &crate::FipsAddress, dst: &crate::FipsAddress, payload: &[u8]) -> Vec { let payload_len = payload.len() as u16; let mut packet = vec![0u8; 40 + payload.len()]; // Version (6) + traffic class high nibble @@ -943,14 +944,17 @@ fn test_identity_cache_populated_on_promote() { let transport_id = TransportId::new(1); let link_id = LinkId::new(1); - let (conn, peer_identity) = make_completed_connection(&mut node, link_id, transport_id, 1000); + let (conn, peer_identity) = make_completed_connection( + &mut node, + link_id, + transport_id, + 1000, + ); node.add_connection(conn).unwrap(); // Promote - let result = node - .promote_connection(link_id, peer_identity, 2000) - .unwrap(); + let result = node.promote_connection(link_id, peer_identity, 2000).unwrap(); assert!(matches!(result, PromotionResult::Promoted(_))); // Identity cache should contain the peer @@ -958,10 +962,7 @@ fn test_identity_cache_populated_on_promote() { let mut prefix = [0u8; 15]; prefix.copy_from_slice(&peer_addr.as_bytes()[0..15]); let cached = node.lookup_by_fips_prefix(&prefix); - assert!( - cached.is_some(), - "Identity cache should contain promoted peer" - ); + assert!(cached.is_some(), "Identity cache should contain promoted peer"); let (cached_addr, cached_pk) = cached.unwrap(); assert_eq!(cached_addr, peer_addr); assert_eq!(cached_pk, peer_identity.pubkey_full()); @@ -985,11 +986,7 @@ async fn test_tun_outbound_established_session() { let dst_fips = crate::FipsAddress::from_node_addr(&node1_addr); // Establish session (XK: 3 messages — Setup, Ack, Msg3) - nodes[0] - .node - .initiate_session(node1_addr, node1_pubkey) - .await - .unwrap(); + nodes[0].node.initiate_session(node1_addr, node1_pubkey).await.unwrap(); tokio::time::sleep(Duration::from_millis(20)).await; process_available_packets(&mut nodes).await; // Setup → Node 1 tokio::time::sleep(Duration::from_millis(20)).await; @@ -997,14 +994,7 @@ async fn test_tun_outbound_established_session() { tokio::time::sleep(Duration::from_millis(20)).await; process_available_packets(&mut nodes).await; // Msg3 → Node 1 - assert!( - nodes[0] - .node - .get_session(&node1_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[0].node.get_session(&node1_addr).unwrap().state().is_established()); // Install TUN receiver on Node 1 let (tun_tx, tun_rx) = std::sync::mpsc::channel(); @@ -1023,10 +1013,7 @@ async fn test_tun_outbound_established_session() { // Verify plaintext arrived at Node 1's TUN let delivered: Vec> = std::iter::from_fn(|| tun_rx.try_recv().ok()).collect(); assert_eq!(delivered.len(), 1, "Exactly one packet should be delivered"); - assert_eq!( - delivered[0], ipv6_packet, - "Delivered packet should match original" - ); + assert_eq!(delivered[0], ipv6_packet, "Delivered packet should match original"); cleanup_nodes(&mut nodes).await; } @@ -1062,35 +1049,17 @@ async fn test_tun_outbound_triggers_session_initiation() { // Session should now be initiating assert_eq!(nodes[0].node.session_count(), 1); - assert!( - nodes[0] - .node - .get_session(&node1_addr) - .unwrap() - .state() - .is_initiating() - ); + assert!(nodes[0].node.get_session(&node1_addr).unwrap().state().is_initiating()); // Drain packets until session established and queued packet delivered drain_to_quiescence(&mut nodes).await; // Session should be established on Node 0 - assert!( - nodes[0] - .node - .get_session(&node1_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[0].node.get_session(&node1_addr).unwrap().state().is_established()); // Verify the queued packet was delivered to Node 1 let delivered: Vec> = std::iter::from_fn(|| tun_rx.try_recv().ok()).collect(); - assert_eq!( - delivered.len(), - 1, - "Queued packet should be delivered after handshake" - ); + assert_eq!(delivered.len(), 1, "Queued packet should be delivered after handshake"); assert_eq!(delivered[0], ipv6_packet); cleanup_nodes(&mut nodes).await; @@ -1118,19 +1087,12 @@ async fn test_tun_outbound_unknown_destination() { // Should receive ICMPv6 Destination Unreachable back on TUN let delivered: Vec> = std::iter::from_fn(|| tun_rx.try_recv().ok()).collect(); - assert_eq!( - delivered.len(), - 1, - "Should receive ICMPv6 Destination Unreachable" - ); + assert_eq!(delivered.len(), 1, "Should receive ICMPv6 Destination Unreachable"); // Verify it's an ICMPv6 Destination Unreachable (type 1, code 0) // ICMPv6 header starts at byte 40, type at byte 40, code at byte 41 assert!(delivered[0].len() >= 48, "ICMPv6 response too short"); assert_eq!(delivered[0][6], 58, "Next header should be ICMPv6 (58)"); - assert_eq!( - delivered[0][40], 1, - "ICMPv6 type should be Destination Unreachable (1)" - ); + assert_eq!(delivered[0][40], 1, "ICMPv6 type should be Destination Unreachable (1)"); assert_eq!(delivered[0][41], 0, "ICMPv6 code should be No Route (0)"); cleanup_nodes(&mut nodes).await; @@ -1169,14 +1131,7 @@ async fn test_tun_outbound_3node_forwarded() { drain_to_quiescence(&mut nodes).await; // Session should be established - assert!( - nodes[0] - .node - .get_session(&node2_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[0].node.get_session(&node2_addr).unwrap().state().is_established()); // Verify packet delivered to Node 2 let delivered: Vec> = std::iter::from_fn(|| tun_rx.try_recv().ok()).collect(); @@ -1215,34 +1170,16 @@ async fn test_tun_outbound_pending_queue_flush() { // First packet triggers session initiation, rest are queued assert_eq!(nodes[0].node.session_count(), 1); - assert!( - nodes[0] - .node - .get_session(&node1_addr) - .unwrap() - .state() - .is_initiating() - ); + assert!(nodes[0].node.get_session(&node1_addr).unwrap().state().is_initiating()); // Drain until session established and queued packets flushed drain_to_quiescence(&mut nodes).await; - assert!( - nodes[0] - .node - .get_session(&node1_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[0].node.get_session(&node1_addr).unwrap().state().is_established()); // All 5 packets should have been delivered let delivered: Vec> = std::iter::from_fn(|| tun_rx.try_recv().ok()).collect(); - assert_eq!( - delivered.len(), - 5, - "All 5 queued packets should be delivered" - ); + assert_eq!(delivered.len(), 5, "All 5 queued packets should be delivered"); for (i, pkt) in delivered.iter().enumerate() { assert_eq!(*pkt, packets[i], "Packet {} should match", i); } @@ -1261,8 +1198,10 @@ fn make_noise_session( ) -> crate::noise::NoiseSession { use crate::noise::HandshakeState; - let mut initiator = - HandshakeState::new_initiator(our_identity.keypair(), remote_identity.pubkey_full()); + let mut initiator = HandshakeState::new_initiator( + our_identity.keypair(), + remote_identity.pubkey_full(), + ); let mut responder = HandshakeState::new_responder(remote_identity.keypair()); // Set epochs for both sides (required for handshake message encryption) @@ -1331,11 +1270,7 @@ fn test_purge_idle_sessions_keeps_active() { let now_ms = 92_000; node.purge_idle_sessions(now_ms); - assert_eq!( - node.session_count(), - 1, - "Active session should survive purge" - ); + assert_eq!(node.session_count(), 1, "Active session should survive purge"); } #[test] @@ -1346,7 +1281,10 @@ fn test_purge_idle_sessions_ignores_initiating() { let remote = Identity::generate(); let remote_addr = *remote.node_addr(); - let handshake = HandshakeState::new_initiator(node.identity().keypair(), remote.pubkey_full()); + let handshake = HandshakeState::new_initiator( + node.identity().keypair(), + remote.pubkey_full(), + ); let entry = crate::node::session::SessionEntry::new( remote_addr, remote.pubkey_full(), @@ -1361,11 +1299,7 @@ fn test_purge_idle_sessions_ignores_initiating() { let now_ms = 1000 + 200_000; node.purge_idle_sessions(now_ms); - assert_eq!( - node.session_count(), - 1, - "Initiating session should not be purged by idle timeout" - ); + assert_eq!(node.session_count(), 1, "Initiating session should not be purged by idle timeout"); } #[test] @@ -1396,10 +1330,8 @@ fn test_purge_idle_sessions_cleans_pending_packets() { node.purge_idle_sessions(now_ms); assert_eq!(node.session_count(), 0); - assert!( - !node.pending_tun_packets.contains_key(&remote_addr), - "Pending packets should be cleaned up with idle session" - ); + assert!(!node.pending_tun_packets.contains_key(&remote_addr), + "Pending packets should be cleaned up with idle session"); } #[test] @@ -1425,11 +1357,7 @@ fn test_purge_idle_sessions_disabled_when_zero() { let now_ms = 1000 + 1_000_000; node.purge_idle_sessions(now_ms); - assert_eq!( - node.session_count(), - 1, - "Sessions should not be purged when idle timeout is disabled" - ); + assert_eq!(node.session_count(), 1, "Sessions should not be purged when idle timeout is disabled"); } #[test] @@ -1458,11 +1386,8 @@ fn test_purge_idle_sessions_mmp_activity_does_not_prevent_purge() { let now_ms = 92_000; node.purge_idle_sessions(now_ms); - assert_eq!( - node.session_count(), - 0, - "Session with MMP-only activity should be purged" - ); + assert_eq!(node.session_count(), 0, + "Session with MMP-only activity should be purged"); } // ============================================================================ @@ -1476,7 +1401,10 @@ fn test_coords_warmup_counter_default_zero_on_new() { let identity_a = Identity::generate(); let identity_b = Identity::generate(); - let handshake = HandshakeState::new_initiator(identity_a.keypair(), identity_b.pubkey_full()); + let handshake = HandshakeState::new_initiator( + identity_a.keypair(), + identity_b.pubkey_full(), + ); let entry = crate::node::session::SessionEntry::new( *identity_b.node_addr(), @@ -1486,11 +1414,8 @@ fn test_coords_warmup_counter_default_zero_on_new() { true, ); - assert_eq!( - entry.coords_warmup_remaining(), - 0, - "Counter should be 0 for non-Established sessions" - ); + assert_eq!(entry.coords_warmup_remaining(), 0, + "Counter should be 0 for non-Established sessions"); } #[test] @@ -1541,20 +1466,15 @@ fn test_coords_warmup_counter_decrement() { assert_eq!(entry.coords_warmup_remaining(), expected); } - assert_eq!( - entry.coords_warmup_remaining(), - 0, - "Counter should reach 0 after N decrements" - ); + assert_eq!(entry.coords_warmup_remaining(), 0, + "Counter should reach 0 after N decrements"); } #[test] fn test_coords_warmup_config_default() { let config = crate::config::Config::new(); - assert_eq!( - config.node.session.coords_warmup_packets, 5, - "Default coords_warmup_packets should be 5" - ); + assert_eq!(config.node.session.coords_warmup_packets, 5, + "Default coords_warmup_packets should be 5"); } // ============================================================================ @@ -1573,13 +1493,11 @@ fn test_identity_cache_lru_eviction() { // Insert first two with explicit timestamps to ensure deterministic ordering let mut prefix1 = [0u8; 15]; prefix1.copy_from_slice(&id1.node_addr().as_bytes()[0..15]); - node.identity_cache - .insert(prefix1, (*id1.node_addr(), id1.pubkey_full(), 1000)); + node.identity_cache.insert(prefix1, (*id1.node_addr(), id1.pubkey_full(), 1000)); let mut prefix2 = [0u8; 15]; prefix2.copy_from_slice(&id2.node_addr().as_bytes()[0..15]); - node.identity_cache - .insert(prefix2, (*id2.node_addr(), id2.pubkey_full(), 2000)); + node.identity_cache.insert(prefix2, (*id2.node_addr(), id2.pubkey_full(), 2000)); assert_eq!(node.identity_cache_len(), 2); @@ -1587,17 +1505,13 @@ fn test_identity_cache_lru_eviction() { node.register_identity(*id3.node_addr(), id3.pubkey_full()); assert_eq!(node.identity_cache_len(), 2); - assert!( - node.lookup_by_fips_prefix(&prefix1).is_none(), - "Oldest entry should have been evicted" - ); + assert!(node.lookup_by_fips_prefix(&prefix1).is_none(), + "Oldest entry should have been evicted"); let mut prefix3 = [0u8; 15]; prefix3.copy_from_slice(&id3.node_addr().as_bytes()[0..15]); - assert!( - node.lookup_by_fips_prefix(&prefix3).is_some(), - "Newest entry should be present" - ); + assert!(node.lookup_by_fips_prefix(&prefix3).is_some(), + "Newest entry should be present"); } #[test] @@ -1632,7 +1546,10 @@ fn test_session_entry_handshake_payload_storage() { let identity_a = Identity::generate(); let identity_b = Identity::generate(); - let handshake = HandshakeState::new_initiator(identity_a.keypair(), identity_b.pubkey_full()); + let handshake = HandshakeState::new_initiator( + identity_a.keypair(), + identity_b.pubkey_full(), + ); let mut entry = crate::node::session::SessionEntry::new( *identity_b.node_addr(), @@ -1664,7 +1581,10 @@ fn test_session_entry_resend_tracking() { let identity_a = Identity::generate(); let identity_b = Identity::generate(); - let handshake = HandshakeState::new_initiator(identity_a.keypair(), identity_b.pubkey_full()); + let handshake = HandshakeState::new_initiator( + identity_a.keypair(), + identity_b.pubkey_full(), + ); let mut entry = crate::node::session::SessionEntry::new( *identity_b.node_addr(), @@ -1695,7 +1615,10 @@ fn test_session_entry_clear_handshake_payload() { let identity_a = Identity::generate(); let identity_b = Identity::generate(); - let handshake = HandshakeState::new_initiator(identity_a.keypair(), identity_b.pubkey_full()); + let handshake = HandshakeState::new_initiator( + identity_a.keypair(), + identity_b.pubkey_full(), + ); let mut entry = crate::node::session::SessionEntry::new( *identity_b.node_addr(), @@ -1726,8 +1649,10 @@ async fn test_session_handshake_timeout() { let mut node = make_node(); let identity_b = Identity::generate(); - let handshake = - HandshakeState::new_initiator(node.identity.keypair(), identity_b.pubkey_full()); + let handshake = HandshakeState::new_initiator( + node.identity.keypair(), + identity_b.pubkey_full(), + ); let dest_addr = *identity_b.node_addr(); @@ -1747,18 +1672,12 @@ async fn test_session_handshake_timeout() { let timeout_secs = node.config.node.rate_limit.handshake_timeout_secs; let before_timeout = 1000 + timeout_secs * 1000 - 1; node.resend_pending_session_handshakes(before_timeout).await; - assert!( - node.sessions.contains_key(&dest_addr), - "Session should survive before timeout" - ); + assert!(node.sessions.contains_key(&dest_addr), "Session should survive before timeout"); // After timeout: session should be removed let after_timeout = 1000 + timeout_secs * 1000 + 1; node.resend_pending_session_handshakes(after_timeout).await; - assert!( - !node.sessions.contains_key(&dest_addr), - "Timed-out session should be removed" - ); + assert!(!node.sessions.contains_key(&dest_addr), "Timed-out session should be removed"); } /// Test that session handshake timeout removes stale AwaitingMsg3 sessions. @@ -1771,7 +1690,9 @@ async fn test_session_awaiting_msg3_timeout() { let identity_a = Identity::generate(); let identity_b = Identity::generate(); - let handshake = HandshakeState::new_xk_responder(identity_b.keypair()); + let handshake = HandshakeState::new_xx_responder( + identity_b.keypair(), + ); let src_addr = *identity_a.node_addr(); @@ -1791,10 +1712,7 @@ async fn test_session_awaiting_msg3_timeout() { let timeout_secs = node.config.node.rate_limit.handshake_timeout_secs; let after_timeout = 1000 + timeout_secs * 1000 + 1; node.resend_pending_session_handshakes(after_timeout).await; - assert!( - !node.sessions.contains_key(&src_addr), - "Timed-out AwaitingMsg3 session should be removed" - ); + assert!(!node.sessions.contains_key(&src_addr), "Timed-out AwaitingMsg3 session should be removed"); } #[tokio::test] @@ -1816,11 +1734,7 @@ async fn test_tun_outbound_path_mtu_generates_ptb() { let dst_fips = crate::FipsAddress::from_node_addr(&node1_addr); // Establish session (XK: 3 messages — Setup, Ack, Msg3) - nodes[0] - .node - .initiate_session(node1_addr, node1_pubkey) - .await - .unwrap(); + nodes[0].node.initiate_session(node1_addr, node1_pubkey).await.unwrap(); tokio::time::sleep(Duration::from_millis(20)).await; process_available_packets(&mut nodes).await; tokio::time::sleep(Duration::from_millis(20)).await; @@ -1828,14 +1742,7 @@ async fn test_tun_outbound_path_mtu_generates_ptb() { tokio::time::sleep(Duration::from_millis(20)).await; process_available_packets(&mut nodes).await; - assert!( - nodes[0] - .node - .get_session(&node1_addr) - .unwrap() - .state() - .is_established() - ); + assert!(nodes[0].node.get_session(&node1_addr).unwrap().state().is_established()); // Simulate receipt of MtuExceeded by reducing PathMtuState to a value // lower than the local transport MTU. @@ -1844,8 +1751,7 @@ async fn test_tun_outbound_path_mtu_generates_ptb() { { let entry = nodes[0].node.get_session_mut(&node1_addr).unwrap(); let mmp = entry.mmp_mut().unwrap(); - mmp.path_mtu - .apply_notification(reduced_mtu, std::time::Instant::now()); + mmp.path_mtu.apply_notification(reduced_mtu, std::time::Instant::now()); assert_eq!(mmp.path_mtu.current_mtu(), reduced_mtu); } @@ -1858,24 +1764,14 @@ async fn test_tun_outbound_path_mtu_generates_ptb() { let local_ipv6_mtu = nodes[0].node.effective_ipv6_mtu() as usize; let oversized_payload = vec![0u8; reduced_ipv6_mtu - 39]; // 40-byte hdr + payload > reduced MTU let ipv6_packet = build_ipv6_packet(&src_fips, &dst_fips, &oversized_payload); - assert!( - ipv6_packet.len() > reduced_ipv6_mtu, - "packet must exceed path MTU" - ); - assert!( - ipv6_packet.len() <= local_ipv6_mtu, - "packet must fit local MTU" - ); + assert!(ipv6_packet.len() > reduced_ipv6_mtu, "packet must exceed path MTU"); + assert!(ipv6_packet.len() <= local_ipv6_mtu, "packet must fit local MTU"); nodes[0].node.handle_tun_outbound(ipv6_packet).await; // Verify ICMPv6 Packet Too Big was generated let ptb_messages: Vec> = std::iter::from_fn(|| tun_rx.try_recv().ok()).collect(); - assert_eq!( - ptb_messages.len(), - 1, - "Should generate exactly one ICMPv6 PTB" - ); + assert_eq!(ptb_messages.len(), 1, "Should generate exactly one ICMPv6 PTB"); let ptb = &ptb_messages[0]; assert_eq!(ptb[0] >> 4, 6, "Should be IPv6"); @@ -1888,23 +1784,12 @@ async fn test_tun_outbound_path_mtu_generates_ptb() { // address, causing a PMTUD blackhole. let ptb_src = std::net::Ipv6Addr::from(<[u8; 16]>::try_from(&ptb[8..24]).unwrap()); let ptb_dst = std::net::Ipv6Addr::from(<[u8; 16]>::try_from(&ptb[24..40]).unwrap()); - assert_eq!( - ptb_src, - dst_fips.to_ipv6(), - "PTB source must be remote peer (original dst), not local node" - ); - assert_eq!( - ptb_dst, - src_fips.to_ipv6(), - "PTB destination must be local node (original src)" - ); + assert_eq!(ptb_src, dst_fips.to_ipv6(), "PTB source must be remote peer (original dst), not local node"); + assert_eq!(ptb_dst, src_fips.to_ipv6(), "PTB destination must be local node (original src)"); // Verify reported MTU (32-bit field at ICMPv6 header bytes 4-7) let reported_mtu = u32::from_be_bytes([ptb[44], ptb[45], ptb[46], ptb[47]]); - assert_eq!( - reported_mtu, reduced_ipv6_mtu as u32, - "Reported MTU should match path IPv6 MTU" - ); + assert_eq!(reported_mtu, reduced_ipv6_mtu as u32, "Reported MTU should match path IPv6 MTU"); // Verify a packet that fits within path MTU passes through (no PTB) let (tun_tx2, tun_rx2) = std::sync::mpsc::channel(); @@ -1917,11 +1802,7 @@ async fn test_tun_outbound_path_mtu_generates_ptb() { // No PTB should be generated for a fitting packet let ptb_messages2: Vec> = std::iter::from_fn(|| tun_rx2.try_recv().ok()).collect(); - assert_eq!( - ptb_messages2.len(), - 0, - "Should not generate PTB for fitting packet" - ); + assert_eq!(ptb_messages2.len(), 0, "Should not generate PTB for fitting packet"); cleanup_nodes(&mut nodes).await; } @@ -1964,19 +1845,10 @@ async fn test_multihop_pmtud_heterogeneous_mtu() { nodes[0].node.register_identity(node2_addr, node2_pubkey); // Establish session A→C via B (triggers routing through tree) - nodes[0] - .node - .initiate_session(node2_addr, node2_pubkey) - .await - .unwrap(); + nodes[0].node.initiate_session(node2_addr, node2_pubkey).await.unwrap(); drain_to_quiescence(&mut nodes).await; assert!( - nodes[0] - .node - .get_session(&node2_addr) - .unwrap() - .state() - .is_established(), + nodes[0].node.get_session(&node2_addr).unwrap().state().is_established(), "Session A→C should be established" ); @@ -1986,11 +1858,7 @@ async fn test_multihop_pmtud_heterogeneous_mtu() { // With coords (~66 extra), the wire could exceed B's recv buffer. for _ in 0..5 { let small = build_ipv6_packet(&src_fips, &dst_fips, &[0u8; 10]); - nodes[0] - .node - .send_ipv6_packet(&node2_addr, &small) - .await - .unwrap(); + nodes[0].node.send_ipv6_packet(&node2_addr, &small).await.unwrap(); } drain_to_quiescence(&mut nodes).await; @@ -2004,17 +1872,12 @@ async fn test_multihop_pmtud_heterogeneous_mtu() { assert!( ipv6_packet.len() <= local_effective_mtu, "packet ({}) must fit A's local MTU ({})", - ipv6_packet.len(), - local_effective_mtu + ipv6_packet.len(), local_effective_mtu ); // Send the oversized packet — B should fail to forward and send // MtuExceeded signal back. - nodes[0] - .node - .send_ipv6_packet(&node2_addr, &ipv6_packet) - .await - .unwrap(); + nodes[0].node.send_ipv6_packet(&node2_addr, &ipv6_packet).await.unwrap(); drain_to_quiescence(&mut nodes).await; // Verify PathMtuState was updated on A @@ -2039,8 +1902,7 @@ async fn test_multihop_pmtud_heterogeneous_mtu() { let ptb_messages: Vec> = std::iter::from_fn(|| tun_rx2.try_recv().ok()).collect(); assert_eq!( - ptb_messages.len(), - 1, + ptb_messages.len(), 1, "Should generate ICMPv6 PTB for oversized packet after PathMtuState update" ); @@ -2055,16 +1917,8 @@ async fn test_multihop_pmtud_heterogeneous_mtu() { // address, causing a PMTUD blackhole. let ptb_src = std::net::Ipv6Addr::from(<[u8; 16]>::try_from(&ptb[8..24]).unwrap()); let ptb_dst = std::net::Ipv6Addr::from(<[u8; 16]>::try_from(&ptb[24..40]).unwrap()); - assert_eq!( - ptb_src, - dst_fips.to_ipv6(), - "PTB source must be remote peer (original dst), not local node" - ); - assert_eq!( - ptb_dst, - src_fips.to_ipv6(), - "PTB destination must be local node (original src)" - ); + assert_eq!(ptb_src, dst_fips.to_ipv6(), "PTB source must be remote peer (original dst), not local node"); + assert_eq!(ptb_dst, src_fips.to_ipv6(), "PTB destination must be local node (original src)"); // Verify reported MTU is the path MTU (not local MTU) let reported_mtu = u32::from_be_bytes([ptb[44], ptb[45], ptb[46], ptb[47]]); @@ -2087,8 +1941,7 @@ async fn test_multihop_pmtud_heterogeneous_mtu() { let ptb_messages3: Vec> = std::iter::from_fn(|| tun_rx3.try_recv().ok()).collect(); assert_eq!( - ptb_messages3.len(), - 0, + ptb_messages3.len(), 0, "Should not generate PTB for packet fitting within path MTU" ); diff --git a/testing/static/configs/node.template.yaml b/testing/static/configs/node.template.yaml index 3d467d8..a07b510 100644 --- a/testing/static/configs/node.template.yaml +++ b/testing/static/configs/node.template.yaml @@ -6,6 +6,11 @@ node: identity: nsec: "{{NSEC}}" + discovery: + backoff_base_secs: 3 + rate_limit: + handshake_timeout_secs: 10 + tun: enabled: true name: fips0 diff --git a/testing/static/scripts/ping-test.sh b/testing/static/scripts/ping-test.sh index de0ede1..b34782a 100755 --- a/testing/static/scripts/ping-test.sh +++ b/testing/static/scripts/ping-test.sh @@ -109,7 +109,7 @@ elif [ "$PROFILE" = "mesh" ] || [ "$PROFILE" = "mesh-public" ]; then wait_for_peers fips-node-e 3 20 || true fi # Wait for FSP-level connectivity (discovery + session establishment) -wait_for_full_connectivity 30 || true +wait_for_full_connectivity 45 || true # Reset counters for the actual test PASSED=0 diff --git a/testing/static/scripts/rekey-test.sh b/testing/static/scripts/rekey-test.sh index 1c91842..377b09c 100755 --- a/testing/static/scripts/rekey-test.sh +++ b/testing/static/scripts/rekey-test.sh @@ -59,7 +59,7 @@ fi trap 'echo ""; echo "Test interrupted"; exit 130' INT # Wait times derived from rekey timer -BASELINE_CONVERGENCE_TIMEOUT=36 +BASELINE_CONVERGENCE_TIMEOUT=60 REKEY_SETTLE=5 # settle time after rekey for cutover to complete # First FMP rekey should follow shortly after the 35s interval once the mesh is # fully converged. Keep this bounded to preserve a meaningful scheduling check