mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-05 19:18:25 +00:00
refactor(peer): hold the transport binding in a PeerPath
`PeerSendState`'s `transport_id`, `current_addr` and the connected-UDP socket + drain move into a `PeerPath`; the send state holds `paths` and an `active` index. Promotion adds exactly one path, `transport_id()`/`current_addr()` read the active one, and the connected-UDP accessors operate on it. No behaviour change: a peer still has at most one path, and `rebind_transport` re-points it rather than growing the set until the probe exchange exists. `set_current_addr` is written for the general case already: an authentic frame roams the address of the path on its transport, whether that path is active or a standby (a standby roam drops that path's own connected socket and reports no change to the caller), and a frame on a transport the peer has no path on creates nothing. `MmpPeerState` stays per peer. The design sketches it per path, but its detection section keeps the receiver report per peer and counts per-path loss from heartbeat sequence gaps; the per-path liveness and RTT fields land with that step. `etx × (1 + rtt_ms / 100)` was typed out at eight sites: the tree's link cost, three snapshot builders and four control-socket fields. Tuning the RTT weighting at one without the others would have had the tree and the views disagree about the same link. One `quality_index(etx, rtt_ms)` next to `compute_etx`, and `MmpMetrics::quality_index()` for the seven sites that read the smoothed pair. Path selection, which lands later, scores a path with the same function, so selection cannot move traffic onto a link the tree calls worse. Refs docs/design/fips-multi-path-switchover.md
This commit is contained in:
+10
-14
@@ -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
|
||||
|
||||
+3
-12
@@ -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,
|
||||
|
||||
@@ -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));
|
||||
}
|
||||
|
||||
+190
-72
@@ -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<std::sync::Arc<crate::transport::udp::ConnectedPeerSocket>>,
|
||||
|
||||
/// 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<crate::transport::udp::PeerRecvDrain>,
|
||||
}
|
||||
|
||||
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<TransportId>,
|
||||
/// Current transport address (for roaming support).
|
||||
current_addr: Option<TransportAddr>,
|
||||
/// 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<PeerPath>,
|
||||
/// Index into `paths` of the path *our* frames go out on. `None` only
|
||||
/// while `paths` is empty.
|
||||
active: Option<usize>,
|
||||
/// 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<std::sync::Arc<crate::transport::udp::ConnectedPeerSocket>>,
|
||||
|
||||
/// 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<crate::transport::udp::PeerRecvDrain>,
|
||||
|
||||
// === 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<std::sync::Arc<crate::transport::udp::ConnectedPeerSocket>> {
|
||||
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<crate::transport::udp::ConnectedPeerSocket>,
|
||||
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<TransportId> {
|
||||
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,
|
||||
}
|
||||
}
|
||||
|
||||
+1
-1
@@ -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;
|
||||
|
||||
@@ -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
|
||||
// ============================================================================
|
||||
|
||||
@@ -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<f64> {
|
||||
match (self.srtt_ms(), self.smoothed_etx()) {
|
||||
(Some(srtt), Some(setx)) => Some(super::quality_index(setx, srtt)),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn smoothed_etx(&self) -> Option<f64> {
|
||||
if self.etx_trend.initialized() {
|
||||
Some(self.etx_trend.long())
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user