diff --git a/src/control/queries.rs b/src/control/queries.rs index a813462e..b85dd3bf 100644 --- a/src/control/queries.rs +++ b/src/control/queries.rs @@ -393,10 +393,8 @@ pub fn show_peers(node: &Node) -> Value { if let Some(smoothed_etx) = mmp.metrics.smoothed_etx() { mmp_json["smoothed_etx"] = json!(smoothed_etx); } - if let Some(srtt) = mmp.metrics.srtt_ms() - && let Some(setx) = mmp.metrics.smoothed_etx() - { - mmp_json["lqi"] = json!(setx * (1.0 + srtt / 100.0)); + if let Some(qi) = mmp.metrics.quality_index() { + mmp_json["lqi"] = json!(qi); } peer_json["mmp"] = mmp_json; } @@ -813,10 +811,8 @@ pub fn show_sessions(node: &Node) -> Value { if let Some(smoothed_etx) = mmp.metrics.smoothed_etx() { mmp_json["smoothed_etx"] = json!(smoothed_etx); } - if let Some(srtt) = mmp.metrics.srtt_ms() - && let Some(setx) = mmp.metrics.smoothed_etx() - { - mmp_json["sqi"] = json!(setx * (1.0 + srtt / 100.0)); + if let Some(qi) = mmp.metrics.quality_index() { + mmp_json["sqi"] = json!(qi); } session_json["mmp"] = mmp_json; } @@ -1013,9 +1009,9 @@ pub fn show_mmp(node: &Node) -> Value { } if let Some(srtt) = metrics.srtt_ms() { link_layer["srtt_ms"] = json!(srtt); - if let Some(setx) = metrics.smoothed_etx() { - link_layer["lqi"] = json!(setx * (1.0 + srtt / 100.0)); - } + } + if let Some(qi) = metrics.quality_index() { + link_layer["lqi"] = json!(qi); } // Trend indicators @@ -1065,9 +1061,9 @@ pub fn show_mmp(node: &Node) -> Value { } if let Some(srtt) = metrics.srtt_ms() { session_layer["srtt_ms"] = json!(srtt); - if let Some(setx) = metrics.smoothed_etx() { - session_layer["sqi"] = json!(setx * (1.0 + srtt / 100.0)); - } + } + if let Some(qi) = metrics.quality_index() { + session_layer["sqi"] = json!(qi); } // Session-layer trend indicators (srtt / loss / etx), mirroring the diff --git a/src/node/mod.rs b/src/node/mod.rs index 92356bed..264be9f0 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -2523,10 +2523,7 @@ impl Node { let metrics = &mmp.metrics; let srtt_ms = metrics.srtt_ms(); let smoothed_etx = metrics.smoothed_etx(); - let lqi = match (srtt_ms, smoothed_etx) { - (Some(srtt), Some(setx)) => Some(setx * (1.0 + srtt / 100.0)), - _ => None, - }; + let lqi = metrics.quality_index(); let trend = |dual: &crate::proto::mmp::DualEwma| { dual.initialized() .then(|| crate::control::queries::trend_label(dual.short(), dual.long())) @@ -2564,10 +2561,7 @@ impl Node { let metrics = &mmp.metrics; let srtt_ms = metrics.srtt_ms(); let smoothed_etx = metrics.smoothed_etx(); - let sqi = match (srtt_ms, smoothed_etx) { - (Some(srtt), Some(setx)) => Some(setx * (1.0 + srtt / 100.0)), - _ => None, - }; + let sqi = metrics.quality_index(); let trend = |dual: &crate::proto::mmp::DualEwma| { dual.initialized() .then(|| crate::control::queries::trend_label(dual.short(), dual.long())) @@ -4033,10 +4027,7 @@ fn project_entity_mmp( ) -> crate::control::snapshot::EntityMmp { let srtt_ms = metrics.srtt_ms(); let smoothed_etx = metrics.smoothed_etx(); - let quality_index = match (srtt_ms, smoothed_etx) { - (Some(srtt), Some(setx)) => Some(setx * (1.0 + srtt / 100.0)), - _ => None, - }; + let quality_index = metrics.quality_index(); crate::control::snapshot::EntityMmp { mode, srtt_ms, diff --git a/src/node/tests/multi_path.rs b/src/node/tests/multi_path.rs index a585d9f0..9c27e25d 100644 --- a/src/node/tests/multi_path.rs +++ b/src/node/tests/multi_path.rs @@ -225,3 +225,26 @@ fn set_current_addr_roams_inside_the_bound_transport_only() { Some(&TransportAddr::from_string("10.0.0.7:1")) ); } + +#[test] +fn a_promoted_peer_holds_one_path_and_rebind_repoints_it() { + let cable = TransportId::new(1); + let wifi = TransportId::new(2); + let (node, node_addr, _our_index, _far_side) = promoted_peer_with_the_far_side_session(cable); + let peer = node.get_peer(&node_addr).unwrap(); + assert_eq!(peer.paths().len(), 1, "promotion binds exactly one path"); + assert_eq!(peer.active_path().map(|p| p.transport_id()), Some(cable)); + assert_eq!( + peer.active_path().map(|p| p.addr()), + Some(&TransportAddr::from_string(PROMOTED_ADDR)) + ); + + // Until the probe exchange adds paths, a rebind re-points the single + // path rather than growing the set. + let mut peer = crate::peer::ActivePeer::new(make_peer_identity(), LinkId::new(1), 0); + assert!(peer.paths().is_empty()); + assert!(peer.rebind_transport(cable, TransportAddr::from_string("10.0.0.1:1"))); + assert!(peer.rebind_transport(wifi, TransportAddr::from_string("10.0.0.7:1"))); + assert_eq!(peer.paths().len(), 1); + assert_eq!(peer.transport_id(), Some(wifi)); +} diff --git a/src/peer/active.rs b/src/peer/active.rs index 9a425236..c4cb9e91 100644 --- a/src/peer/active.rs +++ b/src/peer/active.rs @@ -54,6 +54,74 @@ impl fmt::Display for ConnectivityState { } } +/// One transport-level path to a peer. +/// +/// Path identity is the transport instance: one transport holds at most one +/// path to a given peer, and an address roams *inside* a path. Everything +/// that is per session (Noise slots, K-bit, indices, rekey state) stays on +/// the peer; a path carries only what is bound to the medium it runs over. +/// See `docs/design/fips-multi-path-switchover.md` §1–2. +/// +/// Today a peer holds at most one path, added at promotion. The probe +/// exchange that adds further paths under the existing session is the next +/// step of that design; nothing here assumes a single path. +#[derive(Debug)] +pub struct PeerPath { + /// The transport instance this path runs over. + transport_id: TransportId, + /// The peer's current address on that transport (roams). + addr: TransportAddr, + + /// Unix UDP fast-path: per-path `connect()`-ed socket (paired with + /// the listen socket via `SO_REUSEPORT`). The kernel demux prefers + /// the connected 5-tuple, so inbound packets land here; the + /// encrypt-worker send path sends with `msg_name = NULL`, skipping + /// per-packet sockaddr handling + route lookup. Behind an `Arc` so + /// in-flight worker jobs survive rekey/address-change rotations. + #[cfg(any(target_os = "linux", target_os = "macos"))] + connected_udp: Option>, + + /// Recv drain thread for `connected_udp`. Always paired with it: the + /// kernel routes inbound packets from this peer to the connected + /// socket, so it *must* be drained or the kernel recv buffer fills. + /// Drop signals shutdown via self-pipe. + #[cfg(any(target_os = "linux", target_os = "macos"))] + peer_recv_drain: Option, +} + +impl PeerPath { + fn new(transport_id: TransportId, addr: TransportAddr) -> Self { + Self { + transport_id, + addr, + #[cfg(any(target_os = "linux", target_os = "macos"))] + connected_udp: None, + #[cfg(any(target_os = "linux", target_os = "macos"))] + peer_recv_drain: None, + } + } + + /// The transport instance this path runs over. + pub fn transport_id(&self) -> TransportId { + self.transport_id + } + + /// The peer's current address on this path. + pub fn addr(&self) -> &TransportAddr { + &self.addr + } + + /// Drop the connected socket and its drain. The drain goes first so + /// its last fd reference is released cleanly; the kernel fd closes on + /// the last `Arc` drop, so in-flight worker jobs holding the old `Arc` + /// stay valid until they complete. + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn clear_connected_udp(&mut self) { + self.peer_recv_drain = None; + self.connected_udp = None; + } +} + /// Published active-send-state for a peer (the two-tier boundary). /// /// This is the send-critical subset of an `ActivePeer` that the data plane @@ -104,30 +172,15 @@ struct PeerSendState { session_start: Instant, // === Transport target === - /// Transport ID for this peer's link. - transport_id: Option, - /// Current transport address (for roaming support). - current_addr: Option, + /// The paths this peer is reachable over, one per transport instance. + /// Empty for a peer that has not been bound to a transport yet. + paths: Vec, + /// Index into `paths` of the path *our* frames go out on. `None` only + /// while `paths` is empty. + active: Option, /// Link used to reach this peer. link_id: LinkId, - // === Connected-UDP handles === - /// Unix UDP fast-path: per-peer `connect()`-ed socket (paired with - /// the listen socket via `SO_REUSEPORT`). The kernel demux prefers - /// the connected 5-tuple, so inbound packets land here; the - /// encrypt-worker send path sends with `msg_name = NULL`, skipping - /// per-packet sockaddr handling + route lookup. Behind an `Arc` so - /// in-flight worker jobs survive rekey/address-change rotations. - #[cfg(any(target_os = "linux", target_os = "macos"))] - connected_udp: Option>, - - /// Per-peer recv drain thread. Always paired with `connected_udp`: - /// the kernel routes inbound packets from this peer to the - /// connected socket, so it *must* be drained or the kernel recv - /// buffer fills. Drop signals shutdown via self-pipe. - #[cfg(any(target_os = "linux", target_os = "macos"))] - peer_recv_drain: Option, - // === Hot counters === /// Link statistics. link_stats: LinkStats, @@ -157,13 +210,9 @@ impl PeerSendState { pending_their_index: None, current_k_bit: false, session_start, - transport_id: None, - current_addr: None, + paths: Vec::new(), + active: None, link_id, - #[cfg(any(target_os = "linux", target_os = "macos"))] - connected_udp: None, - #[cfg(any(target_os = "linux", target_os = "macos"))] - peer_recv_drain: None, link_stats: LinkStats::new(), last_seen, replay_suppressed_count: 0, @@ -171,6 +220,15 @@ impl PeerSendState { mmp: None, } } + + /// The path our frames go out on, if any. + fn active_path(&self) -> Option<&PeerPath> { + self.active.and_then(|i| self.paths.get(i)) + } + + fn active_path_mut(&mut self) -> Option<&mut PeerPath> { + self.active.and_then(|i| self.paths.get_mut(i)) + } } /// A fully authenticated remote FIPS node. @@ -360,8 +418,8 @@ impl ActivePeer { send.noise_session = Some(noise_session); send.our_index = Some(our_index); send.their_index = Some(their_index); - send.transport_id = Some(transport_id); - send.current_addr = Some(current_addr); + send.paths.push(PeerPath::new(transport_id, current_addr)); + send.active = Some(0); send.link_stats = link_stats; send.mmp = Some(MmpPeerState::new( mmp_config.mode, @@ -404,42 +462,48 @@ impl ActivePeer { // === Connected-UDP fast path === - /// Refcount the per-peer `connect()`-ed UDP socket if installed. + /// Refcount the active path's `connect()`-ed UDP socket if installed. /// Encrypt-worker send path uses this to bypass the wildcard /// listen socket's per-packet sockaddr handling. #[cfg(any(target_os = "linux", target_os = "macos"))] pub(crate) fn connected_udp( &self, ) -> Option> { - self.send.connected_udp.clone() + self.send + .active_path() + .and_then(|path| path.connected_udp.clone()) } - /// Install a per-peer `connect()`-ed UDP socket with its paired - /// recv drain thread. The two own each other's lifetime: the drain - /// is the only consumer of packets on this socket. + /// Install a `connect()`-ed UDP socket with its paired recv drain + /// thread on the active path. The two own each other's lifetime: the + /// drain is the only consumer of packets on this socket. A no-op on a + /// peer with no path: there is no address to have connected to. #[cfg(any(target_os = "linux", target_os = "macos"))] pub(crate) fn set_connected_udp( &mut self, socket: std::sync::Arc, drain: crate::transport::udp::PeerRecvDrain, ) { + let Some(path) = self.send.active_path_mut() else { + return; + }; // Drop the old drain BEFORE the old socket so its last fd // reference is released cleanly. - self.send.peer_recv_drain = None; - self.send.connected_udp = None; - self.send.connected_udp = Some(socket); - self.send.peer_recv_drain = Some(drain); + path.clear_connected_udp(); + path.connected_udp = Some(socket); + path.peer_recv_drain = Some(drain); } - /// Clear the per-peer connected UDP socket + drain. The drain + /// Clear the active path's connected UDP socket + drain. The drain /// exits via self-pipe signal; the kernel fd closes on last `Arc` /// drop (any in-flight worker jobs holding the old `Arc` stay /// valid until they complete). #[cfg(any(target_os = "linux", target_os = "macos"))] #[allow(dead_code)] // called from session-deregister + rekey follow-up pub(crate) fn clear_connected_udp(&mut self) { - self.send.peer_recv_drain = None; - self.send.connected_udp = None; + if let Some(path) = self.send.active_path_mut() { + path.clear_connected_udp(); + } } // === Identity Accessors === @@ -554,53 +618,106 @@ impl ActivePeer { old_our_index } - /// Get the transport ID for this peer. + /// The transport our frames to this peer go out on: the active path's. pub fn transport_id(&self) -> Option { - self.send.transport_id + self.send.active_path().map(|path| path.transport_id) } - /// Get the current transport address. + /// The address our frames to this peer go to: the active path's. pub fn current_addr(&self) -> Option<&TransportAddr> { - self.send.current_addr.as_ref() + self.send.active_path().map(|path| &path.addr) + } + + /// Every path this peer is reachable over. The active one is + /// [`active_path`](Self::active_path). + pub fn paths(&self) -> &[PeerPath] { + &self.send.paths + } + + /// The path our frames go out on, if the peer has one. + pub fn active_path(&self) -> Option<&PeerPath> { + self.send.active_path() } /// Update the current address (for roaming support). /// /// Called when we receive a valid authenticated packet from a new address. - /// An address roams only *inside* the transport the peer is bound to. A - /// frame that arrives on another transport is still delivered (the demux - /// is by index alone) but does not move the peer: an authentic frame - /// proves the peer produced it, not that it came from where it claims, - /// so an on-path relay rewriting the source (a rogue AP, anyone on a - /// shared L2) could otherwise move the whole send side onto another - /// transport, undamped and unprobed. Only a deliberate - /// [`rebind_transport`](Self::rebind_transport) changes the transport. + /// An address roams only *inside* a path: the frame updates the address + /// of the path on `transport_id` if the peer has one. A frame that + /// arrives on a transport the peer has no path on is still delivered + /// (the demux is by index alone) but creates nothing and moves nothing: + /// an authentic frame proves the peer produced it, not that it came from + /// where it claims, so an on-path relay rewriting the source (a rogue + /// AP, anyone on a shared L2) could otherwise move the whole send side + /// onto another transport, undamped and unprobed. A peer with no path + /// at all is bound by its first authentic frame, as before. Only a + /// deliberate [`rebind_transport`](Self::rebind_transport) changes which + /// transport the peer sends on. /// - /// Returns `true` if the address actually changed — callers use this to - /// invalidate per-peer `connect(2)`-ed UDP sockets whose 5-tuple just - /// went stale. A frame refused for being on another transport returns - /// `false`: nothing moved. + /// Returns `true` if the *active* path's address changed — callers use + /// this to invalidate the `connect(2)`-ed UDP socket whose 5-tuple just + /// went stale. A roam on a standby path clears that path's own socket + /// here and returns `false`; a frame refused for being on an unknown + /// transport returns `false` too: nothing moved. pub fn set_current_addr(&mut self, transport_id: TransportId, addr: TransportAddr) -> bool { - if let Some(bound) = self.send.transport_id - && bound != transport_id - { + if self.send.paths.is_empty() { + return self.rebind_transport(transport_id, addr); + } + let Some(idx) = self + .send + .paths + .iter() + .position(|path| path.transport_id == transport_id) + else { + return false; + }; + let path = &mut self.send.paths[idx]; + if path.addr == addr { return false; } - self.rebind_transport(transport_id, addr) + path.addr = addr; + let on_active = self.send.active == Some(idx); + // A standby's connected socket is its own to drop; the active one + // is the caller's, on the `true` return. + #[cfg(any(target_os = "linux", target_os = "macos"))] + if !on_active { + path.clear_connected_udp(); + } + on_active } - /// Bind the peer to `(transport_id, addr)` outright, whatever it was on. + /// Bind the peer's send side to `(transport_id, addr)` outright. /// /// The deliberate counterpart of [`set_current_addr`](Self::set_current_addr): - /// that one is the roaming rule and refuses to cross transports; this one - /// is a path change and does not. Returns `true` if either the transport - /// or the address changed. + /// that one is the roaming rule and never crosses transports; this one + /// is a path change. If the peer already has a path on `transport_id` + /// it becomes the active one at `addr`; otherwise the active path is + /// re-pointed at the new transport, or created if there was none. + /// Returns `true` if either the active transport or its address changed. pub fn rebind_transport(&mut self, transport_id: TransportId, addr: TransportAddr) -> bool { - let changed = self.send.transport_id != Some(transport_id) - || self.send.current_addr.as_ref() != Some(&addr); - self.send.transport_id = Some(transport_id); - self.send.current_addr = Some(addr); - changed + let changed = + self.transport_id() != Some(transport_id) || self.current_addr() != Some(&addr); + if !changed { + return false; + } + if let Some(idx) = self + .send + .paths + .iter() + .position(|path| path.transport_id == transport_id) + { + self.send.paths[idx].addr = addr; + self.send.active = Some(idx); + } else if let Some(path) = self.send.active_path_mut() { + #[cfg(any(target_os = "linux", target_os = "macos"))] + path.clear_connected_udp(); + path.transport_id = transport_id; + path.addr = addr; + } else { + self.send.paths.push(PeerPath::new(transport_id, addr)); + self.send.active = Some(0); + } + true } // === Handshake Resend === @@ -749,7 +866,8 @@ impl ActivePeer { /// Link cost for routing decisions. /// /// Returns a scalar cost where lower is better (1.0 = ideal). - /// Computed as RTT-weighted ETX: `etx * (1.0 + srtt_ms / 100.0)`. + /// The [`quality_index`](crate::proto::mmp::quality_index) of the link's + /// ETX and smoothed RTT. /// /// Returns 1.0 (optimistic default) when MMP metrics are not yet /// available, matching depth-only parent selection behavior. @@ -758,7 +876,7 @@ impl ActivePeer { Some(mmp) => { let etx = mmp.metrics.etx; match mmp.metrics.srtt_ms() { - Some(srtt_ms) => etx * (1.0 + srtt_ms / 100.0), + Some(srtt_ms) => crate::proto::mmp::quality_index(etx, srtt_ms), None => 1.0, } } diff --git a/src/peer/mod.rs b/src/peer/mod.rs index bd186174..3f493fe3 100644 --- a/src/peer/mod.rs +++ b/src/peer/mod.rs @@ -8,7 +8,7 @@ mod active; pub(crate) mod machine; -pub use active::{ActivePeer, ConnectivityState}; +pub use active::{ActivePeer, ConnectivityState, PeerPath}; use crate::NodeAddr; use crate::transport::LinkId; diff --git a/src/proto/mmp/algorithms.rs b/src/proto/mmp/algorithms.rs index 27e837ba..e13253d4 100644 --- a/src/proto/mmp/algorithms.rs +++ b/src/proto/mmp/algorithms.rs @@ -270,6 +270,18 @@ pub fn compute_etx(d_forward: f64, d_reverse: f64) -> f64 { (1.0 / product).clamp(1.0, 100.0) } +/// The quality index of a link, path or session: `etx × (1 + rtt_ms / 100)`, +/// lower is better, `1.0` ideal. +/// +/// The one place the RTT weighting lives. Path selection scores a path by +/// it, the tree costs a link by it, and the control socket reports it as +/// `lqi`/`sqi`; if they computed it apart they could disagree about which +/// of two links is better, and a path switch would then move traffic onto +/// a link the tree calls worse. +pub fn quality_index(etx: f64, rtt_ms: f64) -> f64 { + etx * (1.0 + rtt_ms / 100.0) +} + // ============================================================================ // Spin Bit // ============================================================================ diff --git a/src/proto/mmp/metrics.rs b/src/proto/mmp/metrics.rs index 0803ce2e..2e18b65e 100644 --- a/src/proto/mmp/metrics.rs +++ b/src/proto/mmp/metrics.rs @@ -342,6 +342,15 @@ impl MmpMetrics { } /// Smoothed ETX (long-term EWMA), or `None` if not yet initialized. + /// The quality index, [`quality_index`](super::quality_index) of the + /// smoothed ETX and the smoothed RTT. `None` until both are measured. + pub fn quality_index(&self) -> Option { + match (self.srtt_ms(), self.smoothed_etx()) { + (Some(srtt), Some(setx)) => Some(super::quality_index(setx, srtt)), + _ => None, + } + } + pub fn smoothed_etx(&self) -> Option { if self.etx_trend.initialized() { Some(self.etx_trend.long()) diff --git a/src/proto/mmp/mod.rs b/src/proto/mmp/mod.rs index d5f9a9b8..14b14f7e 100644 --- a/src/proto/mmp/mod.rs +++ b/src/proto/mmp/mod.rs @@ -38,7 +38,7 @@ mod wire; #[cfg(test)] mod tests; -pub(crate) use algorithms::DualEwma; +pub(crate) use algorithms::{DualEwma, quality_index}; pub(crate) use core::{ BackoffUpdate, LinkReportKind, LinkReportSnapshot, MmpAction, PeerLivenessSnapshot, SendResult, SessionReportKind, SessionReportSnapshot, diff --git a/src/proto/mmp/tests/algorithms.rs b/src/proto/mmp/tests/algorithms.rs index 898b3194..d8b73617 100644 --- a/src/proto/mmp/tests/algorithms.rs +++ b/src/proto/mmp/tests/algorithms.rs @@ -2,6 +2,7 @@ use crate::proto::mmp::algorithms::{ DualEwma, JitterEstimator, OwdTrendDetector, SpinBitState, SrttEstimator, compute_etx, + quality_index, }; #[test] @@ -173,3 +174,19 @@ fn test_spin_bit_responder_counter_guard() { responder.rx_observe(false, 3, 0); assert!(responder.tx_bit()); // unchanged } + +#[test] +fn quality_index_weights_etx_by_rtt() { + assert!( + (quality_index(1.0, 0.0) - 1.0).abs() < f64::EPSILON, + "ideal link" + ); + assert!( + (quality_index(1.0, 100.0) - 2.0).abs() < f64::EPSILON, + "100 ms doubles it" + ); + assert!((quality_index(2.0, 50.0) - 3.0).abs() < f64::EPSILON); + // Lower is better, and both inputs move it the same way. + assert!(quality_index(1.0, 10.0) < quality_index(1.0, 20.0)); + assert!(quality_index(1.0, 10.0) < quality_index(1.5, 10.0)); +}