diff --git a/CHANGELOG.md b/CHANGELOG.md index f87f6c51..52395e92 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -267,6 +267,50 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 seed or at the clamp. Legitimate narrow paths are unaffected: adaptation to hops well below the IPv6 minimum, which the mesh does use, continues to work. +- A session setup message naming an already-established peer no longer replaces + that peer's session. The handler did this whenever `node.rekey.enabled` was + false: it ran a fresh responder handshake and overwrote the entry, discarding + the live keys. The message carries no authenticator and its source address is + an envelope field, so anyone able to reach a node could name an established + peer and take that session down, repeatedly, and hold it down by repeating + the message. The established case now always arms the handshake alongside the + running session and adopts the new keys only after a msg3 whose authenticated + static key matches the key the session was opened with, which is the check + the rekey path already applied; a peer that genuinely restarted still + re-establishes, and a forged setup leaves the session carrying traffic. This + changes no wire format and adds no configuration: a node with rekey disabled + already answered such a message, it simply destroyed the session afterwards. + +- The session drain sweep and the cut-over that retires an old key epoch now + run whether or not periodic rekey is enabled. Both sat behind the + periodic-rekey gate, so a node with rekey disabled that adopted new keys held + the superseded ones for the life of the session. + +- A session rekey armed by a peer's setup message is now abandoned if the + matching msg3 never arrives, rather than persisting for the life of the + session. A stuck one made the node treat a later genuine setup message as a + simultaneous initiation and drop it, which would otherwise have turned the + fix above into a lasting block on re-establishment for roughly half of peer + pairs. Only the armed handshake expires, and only it: a rekey that completed + is the key epoch the peer has already moved to, since it exists only because + a msg3 carrying that peer's authenticated key arrived and the sender of that + msg3 promotes the new epoch on an unconditional two-second timer. Expiring + those keys on any timer would drop every later frame from that peer, so they + are now held until the peer's own frame promotes them, a newer completed + rekey replaces them, or the session goes away. What the wait does bound is + precedence, not the keys: a completed rekey outranks a fresh setup message + from that peer only until it has waited a full idle timeout, after which the + setup is answered normally, so a peer that restarted while we held such a + session is no longer refused for as long as our own sends keep the session + from idling out. The handshake timeout logs at INFO, since it costs nothing, + and a completed session displaced by a newer one at WARN, since that does + throw away keys the peer may hold. Session counters record the arming of a + handshake by a setup message, each of the three ways such a message is + refused, and each displaced session, so a node under a sustained spray of + setup messages shows a rate rather than nothing; the per-message log lines + stay at DEBUG because an unauthenticated sender can drive them at line rate. + These counters are not yet readable through the control socket. + - The FSP session address is now bound to the peer key the Noise handshake authenticated, on both the initial and the rekey path. The responder recorded a session under the source address carried in the datagram without ever diff --git a/docs/reference/configuration.md b/docs/reference/configuration.md index f1796237..738ff592 100644 --- a/docs/reference/configuration.md +++ b/docs/reference/configuration.md @@ -326,7 +326,7 @@ cutover. | Parameter | Type | Default | Description | |-----------|------|---------|-------------| -| `node.rekey.enabled` | bool | `true` | Enable periodic Noise rekey on all links and sessions | +| `node.rekey.enabled` | bool | `true` | Initiate periodic Noise rekey on links and sessions. A peer-driven session rekey is still answered when this is off, so session keys can still rotate | | `node.rekey.after_secs` | u64 | `120` | Initiate rekey after this many seconds on a session | | `node.rekey.after_messages` | u64 | `65536` | Initiate rekey after this many messages sent on a session | diff --git a/src/node/handlers/rekey.rs b/src/node/handlers/rekey.rs index 73e367aa..fc2c22fb 100644 --- a/src/node/handlers/rekey.rs +++ b/src/node/handlers/rekey.rs @@ -10,7 +10,7 @@ use crate::node::Node; use crate::node::wire::build_msg1; use crate::noise::HandshakeState; use crate::protocol::{SessionDatagram, SessionSetup}; -use tracing::{debug, trace, warn}; +use tracing::{debug, info, trace, warn}; /// Keep previous session alive for this long after cutover. const DRAIN_WINDOW_SECS: u64 = 10; @@ -419,15 +419,22 @@ impl Node { /// timer, perform the K-bit cutover (overlapping-epoch decrypt /// makes this safe on any schedule — see `FSP_CUTOVER_DELAY_MS`) /// - If the drain window has expired, clean up the previous session + /// - If a responder-side handshake the peer never finished has aged + /// out, abandon it (the handshake only — a completed rekey session + /// is never discarded on a timer, see below) /// - If the rekey timer/counter fires, initiate a new XK handshake + /// (this last one only when `node.rekey.enabled`) /// /// msg3 retransmission is handled separately by /// `resend_pending_session_msg3`; its lifetime is tied to the /// responder receiving msg3, not to this initiator's cutover. pub(in crate::node) async fn check_session_rekey(&mut self) { - if !self.config().node.rekey.enabled { - return; - } + // The cutover, drain and abandoned-rekey sweeps run whether or not + // periodic rekey is enabled: a peer's setup message is answered in + // either configuration, so both a superseded key epoch and an + // abandoned handshake can exist with rekey disabled. Only the + // trigger that starts a rekey of our own is gated. + let rekey_enabled = self.config().node.rekey.enabled; let rekey_after_secs = self.config().node.rekey.after_secs; let rekey_after_messages = self.config().node.rekey.after_messages; @@ -438,6 +445,26 @@ impl Node { let mut sessions_to_cutover: Vec = Vec::new(); let mut sessions_to_drain: Vec = Vec::new(); let mut sessions_to_rekey: Vec = Vec::new(); + let mut handshakes_to_abandon: Vec<(NodeAddr, u64)> = Vec::new(); + + // Bound for a responder-side handshake the peer armed and never + // finished. A parked handshake blocks a later genuine setup message + // through the dual-initiation tie-break, so it clears on the + // handshake timeout. Nothing is lost with it: an armed handshake + // holds no key material either side can be using. + // + // A *completed* rekey has no such bound, and must not acquire one. + // A responder-side pending session exists only because a msg3 + // authenticated by the session's own peer key arrived, and the + // initiator that sent that msg3 promotes the new epoch on an + // unconditional timer (`FSP_CUTOVER_DELAY_MS`, branch 1 below). + // The pending slot is therefore the epoch the peer has already + // moved to, not key material it abandoned, and discarding it on a + // timer makes every later frame from that peer undecryptable. It is + // released only by the events that supersede it: the peer's own + // frame promoting it (`handle_peer_kbit_flip`), a newer completed + // rekey replacing it, or the session going away. + let stale_handshake_ms = self.config().node.rate_limit.handshake_timeout_secs * 1000; for (node_addr, entry) in &self.sessions { if !entry.is_established() { @@ -466,7 +493,25 @@ impl Node { sessions_to_drain.push(*node_addr); } - // 3. Rekey trigger + // 3. Abandon a responder-side handshake the peer never finished. + // Anchored on the peer's last accepted setup message, which is + // the only stamp this path writes. A pending session alongside + // it survives: only the handshake is dropped. + if !entry.is_rekey_initiator() + && entry.has_rekey_in_progress() + && entry.last_peer_rekey_ms() != 0 + { + let age = now_ms.saturating_sub(entry.last_peer_rekey_ms()); + if age > stale_handshake_ms { + handshakes_to_abandon.push((*node_addr, age)); + continue; + } + } + + // 4. Rekey trigger + if !rekey_enabled { + continue; + } if entry.has_rekey_in_progress() { continue; } @@ -519,6 +564,20 @@ impl Node { } } + // Abandon a handshake the peer armed and never finished. Cheap: no + // key material is lost, and the slot was blocking re-establishment. + for (node_addr, age_ms) in handshakes_to_abandon { + if let Some(entry) = self.sessions.get_mut(&node_addr) { + entry.abandon_handshake(); + self.stats_mut().session.rekey_expired += 1; + info!( + peer = %self.peer_display_name(&node_addr), + age_ms, + "FSP rekey armed by peer expired without msg3, session retained" + ); + } + } + // Initiate new rekeys for node_addr in sessions_to_rekey { self.initiate_session_rekey(&node_addr).await; diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index bb70c0d9..ac5e94d2 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -501,85 +501,116 @@ impl Node { } return; } else if existing.is_established() { - // Rekey: if rekey enabled, treat as rekey for key rotation. - // The existing established session remains active for traffic. - if self.config().node.rekey.enabled { - let rekey_in_progress = existing.has_rekey_in_progress(); - let has_pending = existing.pending_new_session().is_some(); + // A SessionSetup naming an already-established peer is + // unauthenticated: msg1 is a bare ephemeral and the source + // address is an envelope field, so anyone able to reach us + // can claim it. It may therefore only arm a handshake + // alongside the running session, never replace it. The new + // keys are adopted in `handle_session_msg3` and only when the + // authenticated static key matches the key this session was + // opened with, and the cut-over waits for a frame that + // authenticates against the pending epoch. Every exit below + // returns, so an established entry never reaches the + // re-establishment path that replaces it. + let rekey_in_progress = existing.has_rekey_in_progress(); + // A completed rekey outranks a fresh setup message while the + // cut-over it is waiting for can still arrive. Once it has + // waited a full idle timeout for a peer that never appeared + // on the new epoch, it stops vetoing: a peer that restarted, + // or one whose own cycle lapsed, would otherwise be refused + // for as long as our own sends kept the session from idling + // out. The pending keys are not discarded here either way — + // only an authenticated msg3 replaces them. + let pending_outranks = existing.pending_new_session().is_some() + && !existing.pending_stale( + Self::now_ms(), + self.config().node.session.idle_timeout_secs * 1000, + ); - // Dual-initiation detection: both sides sent SessionSetup - // simultaneously. Apply tie-breaker — smaller NodeAddr - // wins as initiator (same as initial session setup). - if rekey_in_progress { - if self.identity().node_addr() < src_addr { - // We win as initiator — drop their msg1. - debug!( - src = %self.peer_display_name(src_addr), - "Dual FSP rekey initiation: we win (smaller addr), dropping their msg1" - ); - return; - } - // We lose — abandon our rekey, become responder below. + // Dual-initiation detection: both sides sent SessionSetup + // simultaneously. Apply tie-breaker — smaller NodeAddr + // wins as initiator (same as initial session setup). + if rekey_in_progress { + if self.identity().node_addr() < src_addr { + // We win as initiator — drop their msg1. debug!( src = %self.peer_display_name(src_addr), - "Dual FSP rekey initiation: we lose (larger addr), abandoning ours" - ); - let entry = self.sessions.get_mut(src_addr).unwrap(); - entry.abandon_rekey(); - } else if has_pending { - // Guard: already have a pending session waiting for K-bit cutover - debug!( - src = %self.peer_display_name(src_addr), - "FSP rekey msg1 received but already have pending session, dropping" + "Dual FSP rekey initiation: we win (smaller addr), dropping their msg1" ); + self.stats_mut() + .record_reject(RejectReason::Session(SessionReject::RekeyTiebreak)); return; } - let our_keypair = self.identity().keypair(); - let mut handshake = HandshakeState::new_xk_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"); - return; - } - - // Generate msg2 - let msg2 = match handshake.write_xk_message_2() { - Ok(m) => m, - Err(e) => { - debug!(error = %e, "Failed to generate rekey XK msg2"); - return; - } - }; - - // Build and send SessionAck - let our_coords = self.tree_state.my_coords().clone(); - let ack = SessionAck::new(our_coords, setup.src_coords).with_handshake(msg2); - let ack_payload = ack.encode(); - let my_addr = *self.node_addr(); - let mut datagram = SessionDatagram::new(my_addr, *src_addr, ack_payload) - .with_ttl(self.config().node.session.default_ttl); - - if let Err(e) = self.send_session_datagram(&mut datagram).await { - debug!(error = %e, dest = %self.peer_display_name(src_addr), "Failed to send rekey SessionAck"); - return; - } - - // Store rekey state on the existing entry - let now_ms = Self::now_ms(); - let entry = self.sessions.get_mut(src_addr).unwrap(); - entry.set_rekey_state(handshake, false); - entry.record_peer_rekey(now_ms); - + // We lose — abandon our rekey, become responder below. debug!( src = %self.peer_display_name(src_addr), - "FSP rekey: processed peer's msg1, sent msg2, awaiting msg3" + "Dual FSP rekey initiation: we lose (larger addr), abandoning ours" + ); + let entry = self.sessions.get_mut(src_addr).unwrap(); + entry.abandon_rekey(); + self.stats_mut() + .record_reject(RejectReason::Session(SessionReject::RekeyYielded)); + } else if pending_outranks { + // Guard: already have a pending session waiting for K-bit cutover + debug!( + src = %self.peer_display_name(src_addr), + "FSP rekey msg1 received but already have pending session, dropping" + ); + self.stats_mut() + .record_reject(RejectReason::Session(SessionReject::RekeyPending)); + return; + } + let our_keypair = self.identity().keypair(); + let mut handshake = HandshakeState::new_xk_responder(our_keypair); + handshake.set_local_epoch(self.startup_epoch()); + + if let Err(e) = handshake.read_xk_message_1(&setup.handshake_payload) { + debug!( + src = %self.peer_display_name(src_addr), + error = %e, + "Failed to process rekey XK msg1" ); return; } - // Re-establishment: replace existing session below - debug!(src = %self.peer_display_name(src_addr), "Session re-establishment from peer"); + // Generate msg2 + let msg2 = match handshake.write_xk_message_2() { + Ok(m) => m, + Err(e) => { + debug!( + src = %self.peer_display_name(src_addr), + error = %e, + "Failed to generate rekey XK msg2" + ); + return; + } + }; + + // Build and send SessionAck + let our_coords = self.tree_state.my_coords().clone(); + let ack = SessionAck::new(our_coords, setup.src_coords).with_handshake(msg2); + let ack_payload = ack.encode(); + let my_addr = *self.node_addr(); + let mut datagram = SessionDatagram::new(my_addr, *src_addr, ack_payload) + .with_ttl(self.config().node.session.default_ttl); + + if let Err(e) = self.send_session_datagram(&mut datagram).await { + debug!(error = %e, dest = %self.peer_display_name(src_addr), "Failed to send rekey SessionAck"); + return; + } + + // Store rekey state on the existing entry + let now_ms = Self::now_ms(); + let entry = self.sessions.get_mut(src_addr).unwrap(); + entry.set_rekey_state(handshake, false); + entry.record_peer_rekey(now_ms); + self.stats_mut().session.rekey_armed += 1; + + debug!( + src = %self.peer_display_name(src_addr), + "FSP rekey: processed peer's msg1, sent msg2, awaiting msg3" + ); + return; } } @@ -858,7 +889,11 @@ 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"); + debug!( + src = %self.peer_display_name(src_addr), + error = %e, + "Failed to process rekey XK msg3" + ); entry.abandon_rekey(); self.sessions.insert(*src_addr, entry); return; @@ -900,8 +935,22 @@ impl Node { } }; + // A pending session already held for this peer is superseded + // here rather than by any timer: only a msg3 carrying the + // session's own peer key can replace the epoch that peer moved + // to. The keys it displaces may still be in use, so the event is + // counted and logged rather than silent. + let superseded = entry.pending_new_session().is_some(); entry.set_pending_session(session); + entry.set_rekey_completed_ms(Self::now_ms()); self.sessions.insert(*src_addr, entry); + if superseded { + self.stats_mut().session.pending_replaced += 1; + warn!( + src = %self.peer_display_name(src_addr), + "FSP rekey: newly completed session replaced one still awaiting cutover" + ); + } debug!( src = %self.peer_display_name(src_addr), diff --git a/src/node/reject.rs b/src/node/reject.rs index 89451f37..d156b249 100644 --- a/src/node/reject.rs +++ b/src/node/reject.rs @@ -209,6 +209,23 @@ pub enum SessionReject { /// Tracked via /// [`SessionStats::rekey_key_mismatch`](crate::node::stats::SessionStats). RekeyKeyMismatch, + /// A setup message named an established peer while our own rekey of + /// that session was in flight, and our address sorted smaller, so the + /// tie-break kept us as initiator and their msg1 was dropped. Tracked + /// via [`SessionStats::rekey_tiebreak`](crate::node::stats::SessionStats). + RekeyTiebreak, + /// A setup message named an established peer while our own rekey of + /// that session was in flight, and our address sorted larger, so we + /// abandoned our rekey and answered as responder. The message carries + /// no authenticator, so a sustained rate here means local key rotation + /// is being suppressed. Tracked via + /// [`SessionStats::rekey_yielded`](crate::node::stats::SessionStats). + RekeyYielded, + /// A setup message named an established peer that already holds a + /// completed rekey awaiting cut-over, so the message was dropped + /// rather than arming a second handshake. Tracked via + /// [`SessionStats::rekey_pending`](crate::node::stats::SessionStats). + RekeyPending, } /// MMP rejection reasons. diff --git a/src/node/session.rs b/src/node/session.rs index 956d3ea3..4af0d463 100644 --- a/src/node/session.rs +++ b/src/node/session.rs @@ -153,11 +153,15 @@ pub(crate) struct SessionEntry { rekey_initiator: bool, /// Dampening: last time peer sent us a rekey msg1 (Unix ms). last_peer_rekey_ms: u64, - /// When the FSP rekey handshake completed (initiator sent msg3, Unix ms). - /// Drives the initiator's liveness-bound cutover timer. Cleared on - /// cutover. The timer is no longer safety-critical: overlapping-epoch - /// trial-decrypt covers any cutover skew. It only bounds how long the - /// initiator advertises the old K-bit. + /// When this side's FSP rekey handshake completed and produced the + /// `pending` session (Unix ms): the initiator sending msg3, or the + /// responder accepting it. Cleared on cutover. + /// + /// On the initiator it drives the liveness-bound cutover timer, which + /// is no longer safety-critical: overlapping-epoch trial-decrypt covers + /// any cutover skew, so it only bounds how long the initiator + /// advertises the old K-bit. On the responder it dates the wait for the + /// peer's cut-over (`pending_stale`) and never expires the keys. rekey_completed_ms: u64, /// Encoded SessionMsg3 payload retained for retransmission (initiator). /// Set when the rekey initiator sends msg3; cleared once the responder @@ -435,6 +439,14 @@ impl SessionEntry { now_ms.saturating_sub(self.last_peer_rekey_ms) < dampening_ms } + /// When the peer last initiated a rekey, or 0 if it never has. + /// + /// Bounds the age of any rekey state armed by the peer's setup message, + /// including a pending session derived from it. + pub(crate) fn last_peer_rekey_ms(&self) -> u64 { + self.last_peer_rekey_ms + } + /// Record that the peer initiated a rekey (for dampening). pub(crate) fn record_peer_rekey(&mut self, now_ms: u64) { self.last_peer_rekey_ms = now_ms; @@ -453,7 +465,8 @@ impl SessionEntry { } } - /// When the FSP rekey handshake completed (initiator sent msg3). + /// When this side's FSP rekey handshake completed: the initiator + /// sending msg3, or the responder accepting it. pub(crate) fn rekey_completed_ms(&self) -> u64 { self.rekey_completed_ms } @@ -468,7 +481,8 @@ impl SessionEntry { self.rekey_jitter_secs } - /// Record when the FSP rekey handshake completed (initiator side). + /// Record when this side's FSP rekey handshake completed and the + /// `pending` session appeared. pub(crate) fn set_rekey_completed_ms(&mut self, ms: u64) { self.rekey_completed_ms = ms; } @@ -742,6 +756,39 @@ impl SessionEntry { self.previous_last_used_ms = 0; } + /// Whether a completed rekey session has waited longer than `stale_ms` + /// for the cut-over that would consume it. + /// + /// This is a staleness question about the *wait*, never a licence to + /// drop the keys: the peer may be on this epoch and simply silent, so + /// the pending slot is retained regardless. It answers only whether a + /// fresh setup message from that peer still has to yield to the + /// cut-over that has not arrived. + /// + /// False when no pending session is held, and false while the entry + /// carries no completion stamp, so an unstamped pending is treated as + /// freshly completed. + pub(crate) fn pending_stale(&self, now_ms: u64, stale_ms: u64) -> bool { + self.pending_new_session.is_some() + && self.rekey_completed_ms != 0 + && now_ms.saturating_sub(self.rekey_completed_ms) > stale_ms + } + + /// Abandon an in-progress rekey handshake, keeping any completed + /// session already waiting for cut-over. + /// + /// Used when a handshake the peer armed times out without its msg3. + /// The handshake holds no key material either endpoint can be using, + /// so dropping it costs nothing; a `pending` session alongside it is + /// the epoch the peer may already have moved to and must survive. + /// + /// `rekey_initiator` is deliberately left alone: it describes whichever + /// rekey artefact the entry still holds, and every reader gates on a + /// live handshake or a pending session first. + pub(crate) fn abandon_handshake(&mut self) { + self.rekey_state = None; + } + /// Abandon an in-progress rekey. /// /// Drops the in-flight handshake state, the pending session, and any diff --git a/src/node/stats.rs b/src/node/stats.rs index 1a89e468..0d69c4ce 100644 --- a/src/node/stats.rs +++ b/src/node/stats.rs @@ -40,6 +40,36 @@ pub struct SessionStats { /// the key the session was established with. The rekey is /// abandoned and the existing session is left intact. pub rekey_key_mismatch: u64, + /// A setup message naming an already-established peer armed a + /// responder-side handshake alongside the running session and a + /// SessionAck was sent. The message carries no authenticator, so this + /// counts genuine peer restarts and forged setups alike; its rate is + /// the signal that something is spraying setup messages, which the + /// per-message DEBUG line cannot carry safely at line rate. + pub rekey_armed: u64, + /// A setup message named an established peer while our own rekey of + /// that session was in flight and our address sorted smaller, so the + /// tie-break dropped their msg1 and kept us as initiator. + pub rekey_tiebreak: u64, + /// A setup message named an established peer while our own rekey of + /// that session was in flight and our address sorted larger, so we + /// abandoned our own rekey and answered as responder. A sustained + /// rate here means local key rotation is being suppressed. + pub rekey_yielded: u64, + /// A setup message named an established peer that already holds a + /// completed rekey awaiting cut-over, so the message was dropped + /// rather than arming a second handshake. + pub rekey_pending: u64, + /// A responder-side handshake armed by a peer's setup message passed + /// the handshake timeout without a msg3 and was discarded. The + /// established session is retained. + pub rekey_expired: u64, + /// A completed rekey session still waiting for the peer's cut-over was + /// replaced by a newer one, completed from a msg3 carrying the same + /// authenticated peer key. This drops key material the peer may + /// already have adopted, so a sustained rate means one side keeps + /// rekeying while the other never appears on the new epoch. + pub pending_replaced: u64, } impl SessionStats { @@ -49,6 +79,12 @@ impl SessionStats { bad_state: self.bad_state, addr_mismatch: self.addr_mismatch, rekey_key_mismatch: self.rekey_key_mismatch, + rekey_armed: self.rekey_armed, + rekey_tiebreak: self.rekey_tiebreak, + rekey_yielded: self.rekey_yielded, + rekey_pending: self.rekey_pending, + rekey_expired: self.rekey_expired, + pending_replaced: self.pending_replaced, } } @@ -58,6 +94,9 @@ impl SessionStats { SessionReject::BadState => self.bad_state += 1, SessionReject::AddrMismatch => self.addr_mismatch += 1, SessionReject::RekeyKeyMismatch => self.rekey_key_mismatch += 1, + SessionReject::RekeyTiebreak => self.rekey_tiebreak += 1, + SessionReject::RekeyYielded => self.rekey_yielded += 1, + SessionReject::RekeyPending => self.rekey_pending += 1, } } } @@ -319,6 +358,12 @@ pub struct SessionStatsSnapshot { pub bad_state: u64, pub addr_mismatch: u64, pub rekey_key_mismatch: u64, + pub rekey_armed: u64, + pub rekey_tiebreak: u64, + pub rekey_yielded: u64, + pub rekey_pending: u64, + pub rekey_expired: u64, + pub pending_replaced: u64, } #[derive(Clone, Debug, Default, Serialize)] @@ -400,6 +445,33 @@ mod tests { assert_eq!(stats.addr_mismatch, 0); } + #[test] + fn session_stats_record_reject_separates_the_three_rekey_arming_refusals() { + let mut stats = SessionStats::default(); + stats.record_reject(SessionReject::RekeyTiebreak); + stats.record_reject(SessionReject::RekeyYielded); + stats.record_reject(SessionReject::RekeyYielded); + stats.record_reject(SessionReject::RekeyPending); + assert_eq!(stats.rekey_tiebreak, 1); + assert_eq!(stats.rekey_yielded, 2); + assert_eq!(stats.rekey_pending, 1); + assert_eq!(stats.rekey_armed, 0); + } + + #[test] + fn session_stats_snapshot_carries_the_rekey_arming_counters() { + let mut stats = SessionStats::default(); + stats.record_reject(SessionReject::RekeyTiebreak); + stats.rekey_armed = 7; + stats.rekey_expired = 3; + stats.pending_replaced = 2; + let snap = stats.snapshot(); + assert_eq!(snap.rekey_tiebreak, 1); + assert_eq!(snap.rekey_armed, 7); + assert_eq!(snap.rekey_expired, 3); + assert_eq!(snap.pending_replaced, 2); + } + #[test] fn session_stats_snapshot_carries_identity_binding_counters() { let mut stats = SessionStats::default(); diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index ff61a8ed..9642bc8a 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -4,7 +4,8 @@ use super::*; use crate::node::session::EndToEndState; use crate::node::tests::spanning_tree::{ TestNode, cleanup_nodes, generate_random_edges, lock_large_network_test, - process_available_packets, run_tree_test, run_tree_test_with_mtus, verify_tree_convergence, + process_available_packets, run_tree_test, run_tree_test_with_configs, run_tree_test_with_mtus, + verify_tree_convergence, }; use crate::protocol::{CoordsRequired, PathBroken, SessionAck, SessionDatagram, SessionMsg3}; @@ -3363,3 +3364,716 @@ async fn test_rekey_msg3_accepts_odd_parity_peer_stored_as_even() { ); assert_eq!(node.stats().session.rekey_key_mismatch, 0); } + +// ============================================================================ +// Integration tests: a setup message naming an established peer +// ============================================================================ + +/// Build a two-node routable mesh whose nodes both have periodic rekey off. +async fn make_rekey_disabled_pair() -> Vec { + let configs = (0..2) + .map(|_| { + let mut config = Config::new(); + config.node.rekey.enabled = false; + config + }) + .collect(); + let mut nodes = run_tree_test_with_configs(configs, &[(0, 1)]).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + nodes +} + +/// Establish an FSP session from nodes[0] to nodes[1] and assert both sides +/// reached Established. +async fn establish_pair_session(nodes: &mut [TestNode]) { + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let node1_pubkey = nodes[1].node.identity().pubkey_full(); + + nodes[0] + .node + .initiate_session(node1_addr, node1_pubkey) + .await + .expect("initiate_session failed"); + + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(nodes).await; + } + + assert!( + nodes[0] + .node + .get_session(&node1_addr) + .expect("initiator session present") + .is_established(), + "initiator session must be established before the test body" + ); + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder session present") + .is_established(), + "responder session must be established before the test body" + ); +} + +/// Forge a SessionSetup carrying an unrelated ephemeral but claiming +/// `nodes[0]`'s coordinates, as an attacker able to reach `nodes[1]` would. +fn forge_setup_from_stranger(nodes: &[TestNode]) -> Vec { + use crate::noise::HandshakeState; + use crate::protocol::SessionSetup; + + let attacker = Identity::generate(); + let mut handshake = HandshakeState::new_xk_initiator( + attacker.keypair(), + nodes[1].node.identity().pubkey_full(), + ); + handshake.set_local_epoch([0xA5; 8]); + let msg1 = handshake + .write_xk_message_1() + .expect("attacker msg1 must build"); + + SessionSetup::new( + nodes[0].node.tree_state().my_coords().clone(), + nodes[1].node.tree_state().my_coords().clone(), + ) + .with_handshake(msg1) + .encode() +} + +#[tokio::test] +async fn test_forged_setup_naming_established_peer_leaves_session_carrying_traffic_rekey_disabled() +{ + let mut nodes = make_rekey_disabled_pair().await; + establish_pair_session(&mut nodes).await; + + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let recv_before = nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .traffic_counters() + .1; + + let forged = forge_setup_from_stranger(&nodes); + nodes[1] + .node + .handle_session_payload(&node0_addr, &forged, 1280, false) + .await; + + let entry = nodes[1] + .node + .get_session(&node0_addr) + .expect("the established entry must survive an unauthenticated setup"); + assert!( + entry.is_established(), + "an unauthenticated setup must not replace the established session" + ); + assert!( + entry.has_rekey_in_progress(), + "the forged setup must have been observed as a side handshake, \ + not dropped for an unrelated reason" + ); + assert_eq!( + nodes[1].node.stats().session.rekey_armed, + 1, + "arming a handshake from an unauthenticated setup must be counted, \ + since the DEBUG line at that site is invisible at the default level" + ); + + // The session must still decrypt the real peer's next frame. + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"after the forgery") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + + let recv_after = nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .traffic_counters() + .1; + assert!( + recv_after > recv_before, + "the real peer's frame must still decrypt: received {} packets before, {} after", + recv_before, + recv_after + ); + + cleanup_nodes(&mut nodes).await; +} + +#[tokio::test] +async fn test_genuine_peer_restart_reestablishes_session_with_rekey_disabled() { + let mut nodes = make_rekey_disabled_pair().await; + establish_pair_session(&mut nodes).await; + + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let node1_pubkey = nodes[1].node.identity().pubkey_full(); + + // Simulate node 0 restarting: it loses its session state but keeps its + // identity, so its setup message names an address node 1 still holds an + // established session for. + nodes[0].node.remove_session(&node1_addr); + nodes[0] + .node + .initiate_session(node1_addr, node1_pubkey) + .await + .expect("re-initiate_session failed"); + + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + } + + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder entry present") + .pending_new_session() + .is_some(), + "the restarted peer's msg3 must have produced a pending session" + ); + + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"after the restart") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + + let entry = nodes[1].node.get_session(&node0_addr).unwrap(); + assert!( + entry.pending_new_session().is_none(), + "the first frame on the new epoch must complete the cutover" + ); + assert!( + entry.traffic_counters().1 > 0, + "node 1 must have decrypted the restarted peer's frame" + ); + + cleanup_nodes(&mut nodes).await; +} + +#[tokio::test] +async fn test_genuine_peer_restart_reestablishes_session_with_rekey_enabled() { + let mut nodes = run_tree_test(2, &[(0, 1)], false).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + establish_pair_session(&mut nodes).await; + + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let node1_pubkey = nodes[1].node.identity().pubkey_full(); + + nodes[0].node.remove_session(&node1_addr); + nodes[0] + .node + .initiate_session(node1_addr, node1_pubkey) + .await + .expect("re-initiate_session failed"); + + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + } + + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder entry present") + .pending_new_session() + .is_some(), + "the restarted peer's msg3 must have produced a pending session" + ); + + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"after the restart") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + + let entry = nodes[1].node.get_session(&node0_addr).unwrap(); + assert!( + entry.pending_new_session().is_none(), + "the first frame on the new epoch must complete the cutover" + ); + assert!( + entry.traffic_counters().1 > 0, + "node 1 must have decrypted the restarted peer's frame" + ); + + cleanup_nodes(&mut nodes).await; +} + +// ============================================================================ +// Tick-loop maintenance with periodic rekey disabled +// ============================================================================ + +/// Build a node with periodic rekey disabled holding one established session +/// with `peer`, returning the node and the peer's address. +fn make_node_with_established_peer( + rekey_enabled: bool, + peer: &Identity, +) -> (Node, crate::NodeAddr) { + let mut config = Config::new(); + config.node.rekey.enabled = rekey_enabled; + let mut node = make_node_with(config); + let peer_addr = *peer.node_addr(); + + let session = make_noise_session(node.identity(), peer); + let mut entry = crate::node::session::SessionEntry::new( + peer_addr, + peer.pubkey_full(), + EndToEndState::Established(session), + 1000, + true, + ); + entry.mark_established(1000); + node.sessions.insert(peer_addr, entry); + (node, peer_addr) +} + +/// Wall-clock milliseconds, matching the clock the tick loop reads. +fn wall_clock_ms() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis() as u64 +} + +#[tokio::test] +async fn test_superseded_key_epoch_is_drained_with_rekey_disabled() { + let peer = Identity::generate(); + let (mut node, peer_addr) = make_node_with_established_peer(false, &peer); + + let old = make_noise_session(node.identity(), &peer); + node.sessions + .get_mut(&peer_addr) + .unwrap() + .set_previous_session_for_test(old, 1000); + assert!(node.sessions.get(&peer_addr).unwrap().is_draining()); + + node.check_session_rekey().await; + + assert!( + !node.sessions.get(&peer_addr).unwrap().is_draining(), + "a drain window that expired long ago must be completed even when \ + periodic rekey is disabled" + ); +} + +#[tokio::test] +async fn test_abandoned_peer_armed_rekey_expires_with_rekey_disabled() { + let peer = Identity::generate(); + let attacker = Identity::generate(); + let (mut node, peer_addr) = make_node_with_established_peer(false, &peer); + + // A setup message armed a responder-side handshake whose msg3 never came. + let (responder, _msg3) = drive_xk_to_msg3(&attacker, node.identity()); + let now_ms = wall_clock_ms(); + let entry = node.sessions.get_mut(&peer_addr).unwrap(); + entry.set_rekey_state(responder, false); + entry.record_peer_rekey(now_ms - 31_000); + + node.check_session_rekey().await; + + let entry = node.sessions.get(&peer_addr).unwrap(); + assert!( + !entry.has_rekey_in_progress(), + "rekey state armed by a peer that never sent msg3 must not persist \ + for the life of the session" + ); + assert!( + entry.is_established(), + "expiring the abandoned handshake must leave the session intact" + ); + assert_eq!( + node.stats().session.rekey_expired, + 1, + "the expired handshake must be counted as a handshake timeout" + ); + assert_eq!( + node.stats().session.pending_replaced, + 0, + "no pending session was replaced, so that counter must not move" + ); +} + +/// Forge a SessionSetup addressed to `node` from an unrelated identity, +/// as an off-path sender naming an established peer would send it. +fn forge_setup_for(node: &Node) -> Vec { + use crate::noise::HandshakeState; + use crate::protocol::SessionSetup; + + let stranger = Identity::generate(); + let mut handshake = + HandshakeState::new_xk_initiator(stranger.keypair(), node.identity().pubkey_full()); + handshake.set_local_epoch([0x5A; 8]); + let msg1 = handshake + .write_xk_message_1() + .expect("stranger msg1 must build"); + + let coords = node.tree_state().my_coords().clone(); + SessionSetup::new(coords.clone(), coords) + .with_handshake(msg1) + .encode() +} + +#[tokio::test] +async fn test_setup_naming_peer_with_pending_session_is_dropped_and_counted() { + let peer = Identity::generate(); + let (mut node, peer_addr) = make_node_with_established_peer(false, &peer); + + // A completed rekey is already waiting for the peer to cut over. + let pending = make_noise_session(node.identity(), &peer); + node.sessions + .get_mut(&peer_addr) + .unwrap() + .set_pending_session(pending); + + let forged = forge_setup_for(&node); + node.handle_session_payload(&peer_addr, &forged, 1280, false) + .await; + + let entry = node.sessions.get(&peer_addr).unwrap(); + assert!( + !entry.has_rekey_in_progress(), + "a setup message must not arm a second handshake while a pending \ + session is still waiting for cutover" + ); + assert!( + entry.pending_new_session().is_some(), + "the pending session must survive the dropped setup message" + ); + assert_eq!( + node.stats().session.rekey_pending, + 1, + "the dropped setup message must be counted, since its DEBUG line is \ + invisible at the default log level" + ); + assert_eq!( + node.stats().session.rekey_armed, + 0, + "nothing was armed, so the arming counter must not move" + ); +} + +#[tokio::test] +async fn test_completed_peer_rekey_session_is_never_expired_by_the_tick_loop() { + let peer = Identity::generate(); + let (mut node, peer_addr) = make_node_with_established_peer(false, &peer); + + // A rekey the peer armed completed long ago, and the peer has not yet + // appeared on the new epoch. The keys are the epoch that peer cut over + // to, so no amount of waiting may discard them. + let pending = make_noise_session(node.identity(), &peer); + let idle_ms = node.config().node.session.idle_timeout_secs * 1000; + let now_ms = wall_clock_ms(); + let entry = node.sessions.get_mut(&peer_addr).unwrap(); + entry.set_pending_session(pending); + entry.set_rekey_completed_ms(now_ms - idle_ms - 60_000); + entry.record_peer_rekey(now_ms - idle_ms - 60_000); + + for _ in 0..3 { + node.check_session_rekey().await; + } + + let entry = node.sessions.get(&peer_addr).unwrap(); + assert!( + entry.pending_new_session().is_some(), + "a completed rekey session must survive any wait for the peer's \ + cutover: discarding it makes the peer's next frame undecryptable" + ); + assert!( + entry.is_established(), + "the running session must be left intact alongside it" + ); + assert_eq!( + node.stats().session.rekey_expired, + 0, + "no armed handshake timed out, so that counter must not move" + ); + assert_eq!( + node.stats().session.pending_replaced, + 0, + "nothing replaced the pending session, so that counter must not move" + ); +} + +#[tokio::test] +async fn test_expiring_an_armed_handshake_keeps_the_completed_session_beside_it() { + let peer = Identity::generate(); + let (mut node, peer_addr) = make_node_with_established_peer(false, &peer); + + // The peer's earlier rekey completed and is still waiting for its + // cutover; a later setup message armed a handshake whose msg3 never + // came. Expiring the handshake must not take the keys with it. + let pending = make_noise_session(node.identity(), &peer); + let (responder, _msg3) = drive_xk_to_msg3(&peer, node.identity()); + let now_ms = wall_clock_ms(); + let entry = node.sessions.get_mut(&peer_addr).unwrap(); + entry.set_pending_session(pending); + entry.set_rekey_completed_ms(now_ms - 120_000); + entry.set_rekey_state(responder, false); + entry.record_peer_rekey(now_ms - 31_000); + + node.check_session_rekey().await; + + let entry = node.sessions.get(&peer_addr).unwrap(); + assert!( + !entry.has_rekey_in_progress(), + "the armed handshake must still expire on the handshake timeout" + ); + assert!( + entry.pending_new_session().is_some(), + "expiring the armed handshake must leave the completed session that \ + the peer may already have cut over to" + ); + assert_eq!( + node.stats().session.rekey_expired, + 1, + "the expired handshake must be counted as a handshake timeout" + ); +} + +#[tokio::test] +async fn test_fresh_peer_armed_rekey_is_not_expired() { + let peer = Identity::generate(); + let (mut node, peer_addr) = make_node_with_established_peer(false, &peer); + + let (responder, _msg3) = drive_xk_to_msg3(&peer, node.identity()); + let now_ms = wall_clock_ms(); + let entry = node.sessions.get_mut(&peer_addr).unwrap(); + entry.set_rekey_state(responder, false); + entry.record_peer_rekey(now_ms); + + node.check_session_rekey().await; + + assert!( + node.sessions + .get(&peer_addr) + .unwrap() + .has_rekey_in_progress(), + "a handshake still within the timeout must not be expired out from \ + under a peer whose msg3 is in flight" + ); +} + +// ============================================================================ +// Integration tests: a peer's cutover after a long silence +// ============================================================================ + +#[tokio::test] +async fn test_silent_peers_cutover_still_lands_after_the_idle_timeout_has_passed() { + // Both nodes rekey after a single message, so one data frame drives a + // full FSP rekey cycle with node 0 as initiator. + let configs = (0..2) + .map(|_| { + let mut config = Config::new(); + config.node.rekey.after_messages = 1; + config + }) + .collect(); + let mut nodes = run_tree_test_with_configs(configs, &[(0, 1)]).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + establish_pair_session(&mut nodes).await; + + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + + // One frame arms node 0's rekey trigger; the cycle then runs to + // completion, leaving node 1 holding the new epoch as `pending`. + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"before the rekey") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + + nodes[0].node.check_session_rekey().await; + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + } + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder entry present") + .pending_new_session() + .is_some(), + "the rekey cycle must have left node 1 holding a pending session" + ); + + // Node 0 emits nothing for longer than the idle timeout — with MMP in + // minimal mode and traffic flowing one way, nothing authenticates + // against node 1's pending slot in that time. + let now_ms = wall_clock_ms(); + let idle_ms = nodes[1].node.config().node.session.idle_timeout_secs * 1000; + let stamp = now_ms - idle_ms - 10_000; + nodes[1] + .node + .sessions + .get_mut(&node0_addr) + .unwrap() + .set_rekey_completed_ms(stamp); + nodes[1] + .node + .sessions + .get_mut(&node0_addr) + .unwrap() + .record_peer_rekey(stamp); + for _ in 0..3 { + nodes[1].node.check_session_rekey().await; + } + + // Node 0 now cuts over on its own liveness timer and speaks again. + nodes[0] + .node + .sessions + .get_mut(&node1_addr) + .unwrap() + .set_rekey_completed_ms(now_ms - 10_000); + nodes[0].node.check_session_rekey().await; + let recv_before = nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .traffic_counters() + .1; + + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"after the long silence") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + + let entry = nodes[1].node.get_session(&node0_addr).unwrap(); + assert_eq!( + entry.traffic_counters().1, + recv_before + 1, + "the peer's first frame on the epoch it cut over to must still \ + decrypt: received {} packets before the frame, {} after", + recv_before, + entry.traffic_counters().1 + ); + assert!( + entry.pending_new_session().is_none(), + "that frame must also complete node 1's cutover" + ); + + cleanup_nodes(&mut nodes).await; +} + +#[tokio::test] +async fn test_peer_restart_reestablishes_through_a_pending_session_that_waited_out_the_timeout() { + // One-second idle timeout, and a rekey after a single message, so a + // genuine rekey cycle leaves a pending session that ages out of the + // veto within the test rather than after a minute and a half. + let configs = (0..2) + .map(|_| { + let mut config = Config::new(); + config.node.rekey.after_messages = 1; + config.node.session.idle_timeout_secs = 1; + config + }) + .collect(); + let mut nodes = run_tree_test_with_configs(configs, &[(0, 1)]).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + establish_pair_session(&mut nodes).await; + + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let node1_pubkey = nodes[1].node.identity().pubkey_full(); + + // A real rekey cycle leaves node 1 holding a completed session whose + // cutover never comes, stamped by the handler that completed it. + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"before the rekey") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + nodes[0].node.check_session_rekey().await; + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + } + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder entry present") + .pending_new_session() + .is_some(), + "the rekey cycle must have left node 1 holding a pending session" + ); + tokio::time::sleep(Duration::from_millis(1200)).await; + + // Node 0 restarts and re-initiates. Its setup message must not be + // refused indefinitely on account of that pending session, or node 1's + // own sends keep the session alive and the peer is locked out for good. + nodes[0].node.remove_session(&node1_addr); + nodes[0] + .node + .initiate_session(node1_addr, node1_pubkey) + .await + .expect("re-initiate_session failed"); + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + } + + assert_eq!( + nodes[1].node.stats().session.rekey_pending, + 0, + "a pending session that has waited out the idle timeout must stop \ + vetoing the peer's setup message" + ); + assert_eq!( + nodes[1].node.stats().session.pending_replaced, + 1, + "the restarted peer's authenticated msg3 must be what replaces the \ + waiting session, and the replacement must be counted" + ); + + nodes[0] + .node + .send_session_data(&node1_addr, 0, 0, b"after the restart") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + + let entry = nodes[1].node.get_session(&node0_addr).unwrap(); + assert!( + entry.pending_new_session().is_none(), + "the restarted peer's first frame must complete the cutover" + ); + assert!( + entry.traffic_counters().1 > 0, + "node 1 must have decrypted the restarted peer's frame" + ); + + cleanup_nodes(&mut nodes).await; +} diff --git a/src/node/tests/spanning_tree.rs b/src/node/tests/spanning_tree.rs index 76cab547..468893fb 100644 --- a/src/node/tests/spanning_tree.rs +++ b/src/node/tests/spanning_tree.rs @@ -78,7 +78,16 @@ pub(super) async fn make_test_node() -> TestNode { /// mirroring UDP, so heterogeneous-MTU / PMTUD tests still exercise the /// forward-path bottleneck. pub(super) async fn make_test_node_with_mtu(mtu: u16) -> TestNode { - let mut node = make_node(); + make_test_node_with_config(Config::new(), mtu).await +} + +/// Create a test node with a specific `Config` and transport MTU. +/// +/// Node configuration is immutable after construction (see `make_node_with`), +/// so a test that needs a non-default setting must supply the `Config` here +/// rather than poking the node afterwards. +pub(super) async fn make_test_node_with_config(config: Config, mtu: u16) -> TestNode { + let mut node = make_node_with(config); let transport_id = TransportId::new(1); let (tx, rx) = tokio::sync::mpsc::unbounded_channel::(); @@ -754,6 +763,29 @@ pub(super) async fn run_tree_test_with_mtus( nodes.push(make_test_node_with_mtu(mtu).await); } + converge_nodes(nodes, edges).await +} + +/// Like `run_tree_test` but with a per-node `Config`. +/// +/// `configs` must have one entry per node. Used by tests that need a +/// non-default setting on a node in a routable mesh, which cannot be produced +/// any other way because node config is immutable after construction. +pub(super) async fn run_tree_test_with_configs( + configs: Vec, + edges: &[(usize, usize)], +) -> Vec { + let mut nodes = Vec::new(); + for config in configs { + nodes.push(make_test_node_with_config(config, 1280).await); + } + + converge_nodes(nodes, edges).await +} + +/// Drive the given nodes to convergence over `edges` and assert every edge +/// established a bidirectional peer. +async fn converge_nodes(mut nodes: Vec, edges: &[(usize, usize)]) -> Vec { for &(i, j) in edges { initiate_handshake(&mut nodes, i, j).await; }