mirror of
https://github.com/jmcorgan/fips.git
synced 2026-08-09 16:24:45 +00:00
Switch FSP handshake from Noise XK to XX
Replace the 3-message XK handshake with XX for FSP session establishment. XX requires no prior knowledge of the peer's static key — the responder's identity is revealed in msg2, the initiator's in msg3. Key changes: - session.rs: XX initiator/responder, post-handshake identity verification using x-only key comparison (parity-independent for npub compatibility), negotiation payload in msg2/msg3 (FSP version [0,0], features=0) - Rekey: switched from XK to XX for FSP rekey handshake - timeout.rs: suppress msg1 resends when target peer is already promoted, preventing cross-connection session mismatch from duplicate handshakes - Test template: discovery backoff 3s and handshake timeout 10s for faster convergence in integration tests - Integration test timeouts restored to 45s (ping) and 60s (rekey) Squashed commits: - Switch FSP handshake from Noise XK to XX - Fix integration test convergence by reducing discovery backoff - Fix cross-connection session mismatch from msg1 resend - Fix FSP identity verification parity mismatch
This commit is contained in:
@@ -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;
|
||||
}
|
||||
|
||||
+221
-214
@@ -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
|
||||
}
|
||||
|
||||
@@ -22,9 +22,7 @@ impl Node {
|
||||
.unwrap_or(0);
|
||||
let timeout_ms = self.config.node.rate_limit.handshake_timeout_secs * 1000;
|
||||
|
||||
let stale: Vec<LinkId> = self
|
||||
.connections
|
||||
.iter()
|
||||
let stale: Vec<LinkId> = 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<u8>)> = 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<u8>)> = 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<crate::NodeAddr> = self
|
||||
.sessions
|
||||
.iter()
|
||||
let timed_out: Vec<crate::NodeAddr> = 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<u8>)> = self
|
||||
.sessions
|
||||
.iter()
|
||||
let candidates: Vec<(crate::NodeAddr, Vec<u8>)> = 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();
|
||||
|
||||
+190
-337
@@ -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<usize> = nodes.iter().map(|tn| tn.node.session_count()).collect();
|
||||
let session_counts: Vec<usize> = 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<usize> =
|
||||
nodes.iter().map(|tn| tn.node.coord_cache().len()).collect();
|
||||
let coord_cache_sizes: Vec<usize> = 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<u8> {
|
||||
fn build_ipv6_packet(src: &crate::FipsAddress, dst: &crate::FipsAddress, payload: &[u8]) -> Vec<u8> {
|
||||
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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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"
|
||||
);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user