Merge branch 'master' into next

Carries the connected-UDP source-pinning fix and the medium-change detection
that supplies its signal.

CHANGELOG.md was the only conflict: both sides offered content under
[Unreleased], so the two incoming entries were folded into next's existing
sections rather than either side being taken whole. The data-plane fix joins
next's Fixed section beside Packaging, and the medium-change detector joins its
Added section. Checked afterwards by diffing next's pre-merge changelog against
the result: no line of it was removed, and both incoming blocks are present
verbatim.

The version line did not conflict and next stays at 1.0.0-dev.
This commit is contained in:
Johnathan Corgan
2026-09-07 13:50:43 +00:00
27 changed files with 3097 additions and 3 deletions
+30
View File
@@ -522,6 +522,14 @@ jobs:
- suite: native-api
type: native-api
# Moves a multi-homed node's default route between two live paths
# while mesh traffic is in flight, and asserts the peering survives
# without a re-handshake. Includes a negative control that requires
# the outage with detection disabled, so a topology that stops
# exercising the bug fails loudly instead of passing green. ~6-8 min.
- suite: medium-change
type: medium-change
- suite: dns-resolver
type: dns-resolver
@@ -712,6 +720,28 @@ jobs:
docker compose -f testing/static/docker-compose.yml \
--profile gateway down --volumes --remove-orphans
# ── Transport-medium change ─────────────────────────────────────────
# Reads FIPS_TEST_IMAGE so it runs against the image this workflow
# built. Owns its own compose project and its own three bridges.
- name: Run medium-change test
if: matrix.type == 'medium-change'
timeout-minutes: 20
env:
FIPS_TEST_IMAGE: fips-test:latest
run: bash testing/medium-change/scripts/test.sh
- name: Collect logs on failure (medium-change)
if: matrix.type == 'medium-change' && failure()
run: |
docker compose -f testing/medium-change/docker-compose.yml \
logs --no-color || true
- name: Stop containers (medium-change)
if: matrix.type == 'medium-change' && always()
run: |
docker compose -f testing/medium-change/docker-compose.yml \
down --volumes --remove-orphans || true
# ── Native datagram API ─────────────────────────────────────────────
# Reads FIPS_TEST_IMAGE rather than defaulting to a name, so it runs
# against the image this workflow built. The two-node check creates and
+5
View File
@@ -44,3 +44,8 @@ __pycache__/
/fips.key
/fips.pub
/fips.yaml
# Per-run node configs written by the medium-change suite. `test.sh` renders
# them from the topology before compose starts and leaves them for post-mortem,
# and the run suffix means each run leaves its own directory behind.
/testing/medium-change/generated-configs*/
+58
View File
@@ -118,6 +118,31 @@ with v0.5.x or earlier peers.
additionally wired into the Noise XX handshake cluster
(msg1/msg2/msg3) and the rekey-initiator outbound sites on `next`.
#### Node lifecycle
- Transport-medium change detection, controlled by the new `node.netmon.*`
block (on by default). The node samples a coarse fingerprint of its network
attachment — the source addresses the routing table would pick for an off-link
destination, plus the set of up, non-loopback interface addresses — and
reports a change once the picture settles. A handover is not atomic (the old
address goes, briefly nothing has a route, the new one arrives), so a short
debounce coalesces the burst into one event and a fingerprint that settles
back where it started reports nothing. Linux subscribes to `NETLINK_ROUTE`
multicast (the groups `ip monitor` uses) and macOS and FreeBSD to a
`PF_ROUTE` socket, both reacting to the kernel event in milliseconds; every
other platform samples on a timer at `node.netmon.poll_interval_secs`, which
also runs underneath the kernel sources as a backstop, since a netlink socket
drops messages under memory pressure and the subscription can be refused in a
restricted sandbox. A backend only decides *when to look* — the fingerprint
comparison, the debounce and the settled-back suppression are shared — so the
remaining backends (`NotifyIpInterfaceChange` on Windows, an embedder push on
iOS) land behind the same seam without touching the reaction. Android takes
the netlink source, and falls back to the timer where policy refuses the
group bind. What the node does with the signal is the connected-socket
rebind described under Fixed above.
Bluetooth is not covered: an adapter's state is not an IP attachment and is
invisible to this detector.
### Changed
- `node.rekey.enabled` now means "initiate rekeys" and nothing else. The
@@ -302,6 +327,39 @@ with v0.5.x or earlier peers.
`fipstop` as "Own Loopback". `req_duplicate` returns to meaning only what it
says.
#### Data plane
- A per-peer `connect()`-ed UDP socket is no longer left pinned to an interface
the host has moved off. Established UDP peers get their own socket for the
send fast path; `open_connected_fd` binds the wildcard and then calls
`connect(2)`, which makes the kernel resolve the route once and auto-bind the
local source address to whichever interface was carrying it at that moment.
It never re-evaluates. So after the host changed transport medium — a laptop
moving between WLAN and LAN, a phone between Wi-Fi and cellular — every
established peer went on transmitting from an address the routing table had
abandoned, while the peer, which re-pins to whatever address it last heard
from, answered somewhere the node was no longer sending from. The peering
stayed marked connected and carried no traffic until the 30s liveness timeout
tore it down, roughly 60-90s of black-holed traffic per medium change,
followed by a full re-handshake and tree re-convergence. The mirror-image
case — the *peer* rotating its address — was already handled at the point the
rotation is observed; this is the local half, which had no signal to hang off
because a local move is invisible in the data plane. It now fires from the
medium-change detection added below, which is exactly that missing signal.
Dropping the sockets is self-healing rather than disruptive: the wildcard
listen socket resolves a route per packet, so sends keep working immediately,
and a correctly-bound connected socket is reinstalled on a later tick. Every
peer on a connectionless transport is also heartbeated at once, so the far
side re-pins to the new source address rather than waiting out its own
heartbeat interval. A peer on a connection-oriented transport keeps the
periodic heartbeat instead: that send awaits an unbounded `write_all` on a
stream the medium change has very likely just stranded, and this reaction
runs on the rx loop. Measured on a live
node, a WLAN/LAN switch in either direction now costs no reconnection at all —
the Noise session, tree position and routes survive it. Linux and macOS (the
platforms with the connected-socket fast path); elsewhere the heartbeat alone
carries the new address.
#### Packaging
- The Linux `.deb` and the systemd tarball now install and run on Debian 12 and
+51
View File
@@ -202,6 +202,53 @@ Auto-reconnect (triggered by MMP link-dead removal) uses the same backoff
parameters but bypasses `max_retries`, retrying indefinitely. See
`peers[].auto_reconnect` below.
### Medium-Change Detection (`node.netmon.*`)
Detects that the host moved between transport media — WLAN to LAN, WLAN to 5G,
an interface arriving or leaving — and rebinds the send path immediately.
| Parameter | Type | Default | Description |
|-----------|------|---------|-------------|
| `node.netmon.enabled` | bool | `true` | Whether medium-change detection runs |
| `node.netmon.poll_interval_secs` | u64 | `5` | How often the host's network attachment is sampled (backstop period where an event-driven backend exists) |
| `node.netmon.debounce_ms` | u64 | `250` | How long to wait for the picture to settle before acting (`0` disables) |
Established UDP peers use a per-peer `connect()`-ed socket for the send fast
path. `connect(2)` makes the kernel resolve the route once and pin the local
source address to whichever interface carried it then; it never re-evaluates.
Without detection, a medium change therefore leaves every peer transmitting
from an abandoned address while the peer answers where it last heard the node —
the peering reports itself connected and carries nothing until
`link_dead_timeout_secs` tears it down, typically 60–90s per switch.
On a detected change the node drops those sockets (the wildcard listen socket
resolves a route per packet, so sends keep working, and a correctly bound
connected socket is reinstalled on a later tick) and heartbeats every peer on a
connectionless transport at once so the far side re-pins to the new source
address. A peer reached over TCP, Tor, Nym or BLE is left to its periodic
heartbeat, since sending to it here would block the node's receive loop on a
stream the medium change has very likely just stranded; those transports
re-dial on send. No peering is torn down: sessions, tree positions and routes
survive the switch.
Detection uses the best backend the platform has:
| Platform | Backend | Latency |
|----------|---------|---------|
| Linux, Android | `NETLINK_ROUTE` multicast (as `ip monitor`) | kernel event, milliseconds |
| macOS, FreeBSD | `PF_ROUTE` socket | kernel event, milliseconds |
| Windows, iOS | timer | up to `poll_interval_secs` |
Where a backend exists, `poll_interval_secs` is only a backstop: a netlink
socket drops messages under memory pressure and the subscription can fail to
start in a restricted sandbox, so the timer keeps running underneath. A backend
that cannot start is logged once at `warn` and the node falls back to the timer.
None of this covers Bluetooth: a BLE adapter's state is not an IP attachment and
is invisible to this detector. The connected-socket fast path is Linux and macOS
only; elsewhere there are no pinned sockets to rebind, and the heartbeat alone
carries the new address.
### Cache Parameters (`node.cache.*`)
Controls caching of tree coordinates and identity mappings.
@@ -1122,6 +1169,10 @@ node:
max_retries: 5
base_interval_secs: 5
max_backoff_secs: 300
netmon:
enabled: true
poll_interval_secs: 5
debounce_ms: 250
cache:
coord_size: 50000
coord_ttl_secs: 300
+3 -2
View File
@@ -38,8 +38,9 @@ use zeroize::{Zeroize, Zeroizing};
pub use gateway::{ConntrackConfig, GatewayConfig, GatewayDnsConfig, PortForward, Proto};
pub use node::{
BloomConfig, BuffersConfig, CacheConfig, ControlConfig, LimitsConfig, LookupConfig, MmpConfig,
NativeApiConfig, NodeConfig, NostrRendezvousConfig, NostrRendezvousPolicy, RateLimitConfig,
RekeyConfig, RendezvousConfig, RetryConfig, SessionConfig, SessionMmpConfig, TreeConfig,
NativeApiConfig, NetmonConfig, NodeConfig, NostrRendezvousConfig, NostrRendezvousPolicy,
RateLimitConfig, RekeyConfig, RendezvousConfig, RetryConfig, SessionConfig, SessionMmpConfig,
TreeConfig,
};
pub use peer::{ConnectPolicy, PeerAddress, PeerConfig, TransportSpec};
pub use transport::{
+74
View File
@@ -212,6 +212,70 @@ impl RetryConfig {
}
}
/// Transport-medium change detection (`node.netmon.*`).
///
/// A node that moves between media (WLAN → LAN, WLAN → 5G) otherwise learns
/// about it only as silence: peers sit in the table until
/// `node.link_dead_timeout_secs` reaps them, and the reconnect then waits out
/// whatever backoff the *old* medium accumulated. The detector turns that into
/// an event; see [`crate::node::netmon`].
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct NetmonConfig {
/// Whether medium-change detection runs at all (`node.netmon.enabled`).
/// On by default: the node is strictly slower to recover without it.
#[serde(default = "NetmonConfig::default_enabled")]
pub enabled: bool,
/// How often the host's network attachment is sampled, in seconds
/// (`node.netmon.poll_interval_secs`).
///
/// On a platform with an event-driven backend (Linux, via netlink) changes
/// are acted on the moment the kernel reports them, and this is only the
/// backstop period — kept because a kernel event stream can drop messages
/// under memory pressure or stop altogether, and the node must not silently
/// revert to noticing nothing. Elsewhere it is the only signal, and so the
/// detection-latency floor. Not a correctness knob either way:
/// `link_dead_timeout_secs` remains the backstop behind it.
#[serde(default = "NetmonConfig::default_poll_interval_secs")]
pub poll_interval_secs: u64,
/// How long the detector waits for the picture to settle before reporting,
/// in milliseconds (`node.netmon.debounce_ms`). A handover is not atomic —
/// the old address goes, briefly nothing has a route, the new address
/// arrives — and acting mid-burst means acting on a state about to change
/// again. Zero disables the wait.
#[serde(default = "NetmonConfig::default_debounce_ms")]
pub debounce_ms: u64,
}
impl Default for NetmonConfig {
fn default() -> Self {
Self {
enabled: true,
poll_interval_secs: 5,
debounce_ms: 250,
}
}
}
impl NetmonConfig {
fn default_enabled() -> bool {
true
}
fn default_poll_interval_secs() -> u64 {
5
}
fn default_debounce_ms() -> u64 {
250
}
/// Whether this is the untouched default, so an absent `netmon:` block
/// stays absent on re-serialize.
pub(crate) fn is_default(&self) -> bool {
*self == Self::default()
}
}
/// Cache parameters (`node.cache.*`).
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CacheConfig {
@@ -1347,6 +1411,15 @@ pub struct NodeConfig {
#[serde(default)]
pub rekey: RekeyConfig,
/// Transport-medium change detection (`node.netmon.*`).
///
/// `skip_serializing_if` for the same reason `native_api` has it and
/// `drain_timeout_secs` is an `Option`: `NodeConfig` has no
/// `deny_unknown_fields`, so a plain `#[serde(default)]` add would write a
/// whole `netmon:` block into every deployed config on the next serialize.
#[serde(default, skip_serializing_if = "NetmonConfig::is_default")]
pub netmon: NetmonConfig,
/// Log level (`node.log_level`). Case-insensitive.
/// Valid values: trace, debug, info, warn, error. Default: info.
#[serde(default)]
@@ -1381,6 +1454,7 @@ impl Default for NodeConfig {
session_mmp: SessionMmpConfig::default(),
ecn: EcnConfig::default(),
rekey: RekeyConfig::default(),
netmon: NetmonConfig::default(),
log_level: None,
}
}
+20
View File
@@ -179,6 +179,18 @@ impl Node {
(rx, guard)
};
// Transport-medium change receiver, or a dummy channel when detection
// is disabled (or the node was seeded straight into Running without a
// start()). Same guard pattern as TUN outbound and DNS identity: the
// held sender keeps the channel open so the arm never sees it closed.
let (mut netmon_rx, _netmon_guard) = match self.supervisor.netmon_rx.take() {
Some(rx) => (rx, None),
None => {
let (tx, rx) = tokio::sync::mpsc::channel(1);
(rx, Some(tx))
}
};
// Decrypt-worker fallback receiver. The worker pushes each
// authenticated FMP plaintext here so rx_loop can finish the
// per-peer side-effects (stats, MMP, ECN, link dispatch).
@@ -356,6 +368,14 @@ impl Node {
);
self.register_identity(identity.node_addr, identity.pubkey);
}
// A transport medium change (WLAN -> LAN -> 5G, a BLE adapter
// arriving or leaving). Placed after the hot inbound path so
// the `biased` priority of packet processing is unchanged; the
// detector coalesces into a single-slot channel, so this arm
// never sees a burst and needs no drain loop.
Some(change) = netmon_rx.recv() => {
self.handle_net_change(change).await;
}
// Native API datagrams a client wrote to its descriptor. Drained
// in a burst like the TUN arm, for the same reason: one wake-up
// should clear what a client handed over, not one datagram.
+1
View File
@@ -5,6 +5,7 @@ pub(crate) mod lookup;
mod mmp;
mod native;
pub(in crate::node) use native::PendingNative;
pub(in crate::node) mod netmon;
pub(crate) mod probe;
// Widened from private by the rekey drain cap: `node::session` calls
// `rekey::drain_max_retention_ms` to bound how long a superseded epoch is
+176
View File
@@ -0,0 +1,176 @@
//! The node's reaction to a transport-medium change.
//!
//! [`crate::node::netmon`] detects that the host's network attachment moved and
//! publishes one [`NetChange`]; everything the node *does* about it lives here.
//! The split is deliberate — the per-OS backends feed the same channel, so none
//! of them has to restate this policy.
//!
//! # The problem
//!
//! Established UDP peers get a per-peer `connect()`-ed socket for the send fast
//! path. `open_connected_fd` binds the wildcard and then calls `connect(2)`,
//! which makes the kernel resolve the route **once** and auto-bind the local
//! source address to whichever interface was carrying it at that moment. It
//! never re-evaluates.
//!
//! So when the host changes medium — a laptop between WLAN and LAN, a phone
//! between Wi-Fi and cellular — every established peer goes on transmitting
//! from an address the routing table has abandoned. The peer, which re-pins to
//! whatever address it last heard from, answers somewhere the node is no longer
//! sending from. The peering stays marked connected and carries nothing until
//! `node.link_dead_timeout_secs` tears it down, and the reconnect then has to
//! redo the Noise handshake and the tree position. Measured on a live node
//! before this landed: 60–90s of black-holed traffic per switch.
//!
//! The mirror-image case — the *peer* rotating its address — is already handled
//! where the rotation is observed (`dataplane::encrypted`, on `address_changed`).
//! This is the local half, and it has no other signal to hang off: a medium
//! change is not visible anywhere in the data plane, which is precisely why it
//! went unhandled.
//!
//! # The reaction
//!
//! Two steps, both cheap enough to run on every detected change:
//!
//! 1. **Drop the stale sockets.** Self-healing rather than disruptive: the
//! wildcard listen socket resolves a route per packet, so sends keep working
//! immediately, and `activate_connected_udp_sessions` reinstalls a
//! correctly-bound connected socket on a later tick.
//! 2. **Heartbeat every peer whose send path cannot block.** The frame leaves
//! over the new path and carries the node's new source address, so the far
//! side re-pins on receipt instead of waiting out its own
//! `heartbeat_interval_secs`. Without it the forward direction is fixed but
//! the reverse still points at the old address until the node next happens
//! to send. This runs on the rx loop, so it covers the connectionless
//! transports only — see
//! [`Node::heartbeat_all_peers_after_net_change`] for why awaiting a
//! connection-oriented write here would hold the loop, and what a peer on
//! one of those gets instead.
//!
//! Nothing here tears a peering down. On a live node both WLAN→LAN and
//! LAN→WLAN now cost no reconnection at all — the Noise session, the tree
//! position and the routes all survive the switch. `link_dead_timeout_secs`
//! remains the backstop for a peer that genuinely cannot be reached on the new
//! medium.
use std::time::Instant;
use tracing::{debug, info, warn};
use crate::NodeAddr;
use crate::node::Node;
use crate::node::netmon::NetChange;
use crate::proto::link::LinkMessageType;
impl Node {
/// React to a settled transport-medium change.
pub(in crate::node) async fn handle_net_change(&mut self, change: NetChange) {
let peers = self.peers.len();
// Before the heartbeats: they must go out over a socket that resolves
// the route now, not one still pinned to the interface just left.
let sockets_rebound = self.drop_connected_sockets_after_net_change();
let heartbeated = self.heartbeat_all_peers_after_net_change().await;
info!(
generation = change.generation,
change = %change.summary,
peers,
sockets_rebound,
heartbeated,
"Transport medium changed; rebinding sends and re-pinning peers"
);
}
/// Drop every per-peer `connect()`-ed UDP socket, returning how many were
/// released.
///
/// See the module docs for why they are stale: `connect(2)` pins the local
/// source address to the interface that carried the route at connect time,
/// and never re-evaluates it.
#[cfg(any(target_os = "linux", target_os = "macos"))]
fn drop_connected_sockets_after_net_change(&mut self) -> usize {
let pinned: Vec<NodeAddr> = self
.peers
.iter()
.filter(|(_, peer)| peer.connected_udp().is_some())
.map(|(addr, _)| *addr)
.collect();
for addr in &pinned {
self.clear_connected_udp_for_peer(addr);
}
pinned.len()
}
/// No per-peer connected sockets on this platform, so nothing to rebind.
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
fn drop_connected_sockets_after_net_change(&mut self) -> usize {
0
}
/// Send one heartbeat to every peer whose send path cannot block, so each
/// learns the node's new source address in one RTT rather than at the next
/// due interval. Returns how many went out.
///
/// The filter is not an optimisation. A connectionless transport's send
/// completes without ever awaiting the wire: the UDP fast path hands the
/// frame to the encrypt workers and returns, and a raw datagram write does
/// not wait for a peer. A connection-oriented one awaits `write_all` on a
/// stream, unbounded — the connect above it is wrapped in a timeout, the
/// write is not — and a medium change is precisely the condition that
/// leaves a send window full against a path that has just gone away. This
/// runs on the rx loop, so that write would hold every other arm of the
/// select for as long as the stranded socket takes to fail.
///
/// Bounding it with a timeout is not the fix either: dropping a partial
/// `write_all` would leave a half-written frame on the stream, which the
/// peer cannot resynchronise from. Nor can the fan-out simply be spawned,
/// because the send needs `&mut self` for the session counter and the MMP
/// sender record.
///
/// So a peer on TCP, Tor, Nym or BLE keeps the periodic heartbeat it had
/// before this detector existed. It is not stranded by the omission: those
/// transports re-dial on send, and `link_dead_timeout_secs` remains the
/// backstop. Doing better for them means dropping the stale connection
/// rather than writing into it, which is a different change with a real
/// cost behind it — a Tor peer pays a fresh circuit — and is not this one.
async fn heartbeat_all_peers_after_net_change(&mut self) -> usize {
let now = Instant::now();
let heartbeat = [LinkMessageType::Heartbeat.to_byte()];
let targets: Vec<NodeAddr> = self
.peers
.iter()
.filter(|(_, peer)| {
peer.transport_id()
.and_then(|id| self.transports.get(&id))
.is_some_and(|t| !t.transport_type().connection_oriented)
})
.map(|(addr, _)| *addr)
.collect();
let sent = targets.len();
for addr in targets {
if let Some(peer) = self.peers.get_mut(&addr) {
peer.mark_heartbeat_sent(now);
}
if let Err(e) = self.send_encrypted_link_message(&addr, &heartbeat).await {
debug!(
peer = %self.peer_display_name(&addr),
error = %e,
"Failed to send post-medium-change heartbeat"
);
}
}
sent
}
}
/// Emitted once at startup when detection is configured off, so an operator
/// reading a slow recovery has something to find.
pub(in crate::node) fn warn_detection_disabled() {
warn!(
"node.netmon.enabled = false; a transport medium change will strand \
established peers on sockets bound to the old path until the link \
dead timeout"
);
}
+24
View File
@@ -2173,6 +2173,20 @@ impl Node {
});
}
// Transport-medium change detection. Spawned after bring-up so its
// first sample sees the transports the node actually came up with, and
// outside the FSM's child set for the reason documented on
// `Supervisor::netmon_rx`: a detector that dies costs recovery latency,
// not health.
let netmon_cfg = self.config().node.netmon.clone();
if netmon_cfg.enabled {
let (rx, task) = crate::node::netmon::spawn_detector(netmon_cfg);
self.supervisor.netmon_rx = Some(rx);
self.supervisor.netmon_task = Some(task);
} else {
crate::node::handlers::netmon::warn_detection_disabled();
}
info!("Node started:");
info!(" state: {}", self.supervisor.state);
info!(" transports: {}", self.transports.len());
@@ -2443,6 +2457,16 @@ impl Node {
self.supervisor.packet_tx.take();
self.packet_rx.take();
}
// The medium-change detector is a bare task with no teardown protocol
// of its own: it holds no node state and owns nothing but a timer, so
// aborting it is the whole shutdown. Dropping the receiver would also
// end it at its next send, but only after one more poll interval, and a
// stopped node should not still be sampling the network.
if let Some(task) = self.supervisor.netmon_task.take() {
task.abort();
}
self.supervisor.netmon_rx.take();
}
/// Retract anything a child published about itself, after it exited on its
+15
View File
@@ -731,6 +731,19 @@ pub(crate) struct Supervisor {
#[cfg(unix)]
pub(crate) decrypt_workers: Option<crate::node::decrypt_worker::DecryptWorkerPool>,
/// Transport-medium change detection: the receiver the rx loop drains and
/// the detector task behind it.
///
/// Deliberately **not** a [`Child`]. The FSM's children are the substrate
/// the node's health is defined against — a detector that dies makes
/// recovery slower, not the node unhealthy, and giving it a `Child` would
/// put it in the start plan, the teardown order and the N-of-M health
/// policy for no gain. It is spawned after bring-up and aborted in
/// teardown, like the runtime child-liveness monitor.
pub(in crate::node) netmon_rx: Option<crate::node::netmon::NetChangeRx>,
/// Handle for the detector task, aborted at teardown.
pub(in crate::node) netmon_task: Option<tokio::task::JoinHandle<()>>,
/// The sans-IO lifecycle FSM authoring spawn/teardown ordering.
pub(in crate::node) fsm: SupervisorFsm,
}
@@ -759,6 +772,8 @@ impl Supervisor {
encrypt_workers: None,
#[cfg(unix)]
decrypt_workers: None,
netmon_rx: None,
netmon_task: None,
fsm: SupervisorFsm::new(),
}
}
+1
View File
@@ -17,6 +17,7 @@ pub(crate) mod encrypt_worker;
mod handlers;
mod lifecycle;
pub(crate) mod metrics;
pub(crate) mod netmon;
mod peer_error_budget;
mod peering;
mod rate_limit;
+649
View File
@@ -0,0 +1,649 @@
//! Transport-medium change detection.
//!
//! A node that moves between media (WLAN → LAN, WLAN → 5G, a BLE adapter
//! coming or going) would otherwise learn about it only as *silence*: the peer
//! sits in the table until `node.link_dead_timeout_secs` reaps it, and the
//! reconnect then waits out whatever backoff the old medium had already
//! accumulated. The host kernel knew within milliseconds; the node would find
//! out half a minute later.
//!
//! This module closes that gap. It samples a coarse [`NetFingerprint`] of the
//! host's network attachment and publishes a [`NetChange`] on the channel the
//! rx loop drains whenever the fingerprint moves. The node's reaction lives in
//! [`crate::node::handlers::netmon`].
//!
//! # Backends
//!
//! Detection is split in two, and the split is what keeps a per-OS backend
//! small. A backend's whole job is to answer *when is it worth sampling* — see
//! [`WakeSource`]. Everything else, and in particular the decision about
//! whether a medium change actually happened, is the shared fingerprint
//! comparison below, so no backend parses kernel messages or owns its own
//! definition of a medium change.
//!
//! | Platform | Backend | Latency |
//! |---|---|---|
//! | Linux, Android | `NETLINK_ROUTE` multicast | kernel event, ~ms |
//! | macOS, FreeBSD | `PF_ROUTE` socket | kernel event, ~ms |
//! | everything else | timer, `node.netmon.poll_interval_secs` | up to one period |
//!
//! Both kernel sources are [`crate::transport::watcher::LinkWatcher`], shared
//! with the interface binder. It is asked for a wider set of netlink groups
//! here than the binder asks for: presence is a link question, but a default
//! route moving between two interfaces that both stay up emits nothing in the
//! link group, so a presence subscription would never fire for the change this
//! detector exists to catch. `PF_ROUTE` has no group selection and delivers
//! everything regardless.
//!
//! Still to come, behind the same seam and without touching the handler:
//! `NotifyIpInterfaceChange` on Windows, and an embedder push on iOS. Android
//! takes the netlink source above, which is the right backend when the policy
//! allows the group bind and degrades to the timer when it does not; a
//! `ConnectivityManager` push belongs there too, because a timer is not
//! reliable under Doze. Every platform runs the
//! timer regardless — as the only signal where there is no backend, and as a
//! backstop where there is one, since a kernel event stream can drop messages
//! or stop.
//!
//! # What the fingerprint captures
//!
//! Two independent signals, because neither alone is sufficient:
//!
//! - **The preferred source addresses.** A connected-but-never-sending UDP
//! socket makes the kernel run its route lookup and pick the source address
//! it *would* use to reach an off-link destination. That address changes
//! exactly when the default route moves between media, which is the
//! WLAN → 5G case. It costs two syscalls and no packets, and works
//! identically on every platform std supports.
//! - **The set of up, non-loopback interface addresses**, on unix, where
//! `getifaddrs(3)` is available through the `libc` dependency the crate
//! already carries. This catches a medium arriving or leaving without
//! displacing the default route. A peer on the same LAN is reached by the
//! on-link subnet route, and the probe above follows the *default* route by
//! construction, so it looks straight past that path: unplug a LAN cable on
//! a host whose default route is cellular and the source addresses do not
//! move, while every connected socket to that peer is stranded.
//!
//! It is a proxy for "could the local end of one of our peerings have
//! moved", and a broad one — it enumerates every address the host has,
//! rather than the local end of the peerings the node actually holds, so a
//! container bridge or a tunnel interface appearing moves it too.
//!
//! On platforms without `getifaddrs` (Windows) only the source-address probe
//! contributes, which still catches every default-route move. The native
//! backend is the answer there, not a richer sample.
//!
//! # What it deliberately does not capture
//!
//! A BLE adapter's state is invisible to both signals — it is not an IP
//! attachment at all. That signal comes from the radio (BlueZ properties, the
//! Android callback) and belongs on this same channel, pushed by the BLE
//! transport rather than sampled here.
use std::collections::BTreeSet;
use std::fmt;
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr, UdpSocket};
use std::time::Duration;
#[cfg(unix)]
use crate::transport::watcher::LinkWatcher;
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
use tracing::{debug, trace, warn};
use crate::config::NetmonConfig;
/// Off-link IPv4 probe destination — RFC 5737 TEST-NET-1, which is guaranteed
/// not to be routed anywhere. Nothing is ever sent to it; `connect(2)` on a UDP
/// socket only resolves the route and binds a source address.
const PROBE_V4: SocketAddr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(192, 0, 2, 1)), 9);
/// Off-link IPv6 probe destination — RFC 3849 documentation prefix. Same
/// no-packets-sent contract as [`PROBE_V4`].
const PROBE_V6: SocketAddr = SocketAddr::new(
IpAddr::V6(Ipv6Addr::new(0x2001, 0x0db8, 0, 0, 0, 0, 0, 1)),
9,
);
/// How many resample rounds the debounce will ride out before reporting
/// anyway. A handover emits a burst (address gone, address added, route
/// replaced), and reporting mid-burst would act on a picture that is about to
/// change again; riding it out coalesces the burst into one event. Bounded so
/// an interface that flaps continuously still produces events rather than
/// starving the handler forever.
const MAX_DEBOUNCE_ROUNDS: u32 = 8;
/// Minimum spacing between two reported changes.
///
/// The reaction is not free: it drops every peer's connected UDP socket (each
/// carrying a drain thread) and sends a heartbeat per peer. An interface that
/// flaps cleanly — settling between each transition, so the debounce reports
/// each one — could otherwise drive that several times a second across up to
/// `node.limits.max_peers` peers, which is thread churn rather than recovery.
///
/// A genuine change is delayed by at most this long, against a
/// `link_dead_timeout_secs` measured in tens of seconds, so the trade is
/// heavily one-sided.
const MIN_CHANGE_INTERVAL: Duration = Duration::from_secs(1);
/// Receiver the rx loop drains.
pub(crate) type NetChangeRx = mpsc::Receiver<NetChange>;
/// Sender held by a detection backend.
pub(crate) type NetChangeTx = mpsc::Sender<NetChange>;
/// A coarse fingerprint of how this host is attached to the network.
///
/// Equality is the whole point: the poller reports a change iff two
/// consecutive samples differ. The contents are only ever used for the
/// operator-facing description of what moved.
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct NetFingerprint {
/// Source address the routing table would pick for an off-link IPv4
/// destination. `None` when there is no IPv4 route at all — itself a
/// meaningful state, and distinct from any address.
v4_source: Option<IpAddr>,
/// IPv6 counterpart of `v4_source`.
v6_source: Option<IpAddr>,
/// Every up, non-loopback unicast address on the host. Always empty on
/// platforms with no `getifaddrs`, which makes the fingerprint degrade to
/// the source-address probe rather than to nothing.
local_addrs: BTreeSet<IpAddr>,
}
impl NetFingerprint {
/// Sample the host's current attachment.
///
/// Every call is a handful of non-blocking syscalls — a bind, a connect
/// that sends no packet, a `getifaddrs` walk — with no I/O wait, no name
/// resolution, and no allocation beyond the address set. It is called from
/// a dedicated task on a multi-second timer, so it runs inline rather than
/// through `spawn_blocking`.
pub(crate) fn sample() -> Self {
Self {
v4_source: preferred_source(PROBE_V4),
v6_source: preferred_source(PROBE_V6),
local_addrs: interface_addrs(),
}
}
/// Build a fingerprint directly, so a test can script a sequence of
/// samples instead of reading the host's real attachment.
#[cfg(test)]
pub(crate) fn for_test(v4_source: Option<IpAddr>, local_addrs: &[IpAddr]) -> Self {
Self {
v4_source,
v6_source: None,
local_addrs: local_addrs.iter().copied().collect(),
}
}
/// Describe the transition from `self` to `next` for the operator log.
fn diff(&self, next: &Self) -> NetChangeSummary {
NetChangeSummary {
added: next
.local_addrs
.difference(&self.local_addrs)
.copied()
.collect(),
removed: self
.local_addrs
.difference(&next.local_addrs)
.copied()
.collect(),
v4_source_moved: self.v4_source != next.v4_source,
v6_source_moved: self.v6_source != next.v6_source,
v4_source: next.v4_source,
v6_source: next.v6_source,
}
}
}
/// What moved between two fingerprints. Operator-facing only — the handler
/// re-evaluates everything regardless of which field changed, because it
/// cannot map an address back to the peers that were reaching over it.
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct NetChangeSummary {
/// Local addresses present now but not before.
pub added: Vec<IpAddr>,
/// Local addresses present before but not now.
pub removed: Vec<IpAddr>,
/// The preferred IPv4 source address changed (a default-route move).
pub v4_source_moved: bool,
/// The preferred IPv6 source address changed.
pub v6_source_moved: bool,
/// The preferred IPv4 source address as of this sample.
pub v4_source: Option<IpAddr>,
/// The preferred IPv6 source address as of this sample.
pub v6_source: Option<IpAddr>,
}
impl fmt::Display for NetChangeSummary {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let mut parts: Vec<String> = Vec::new();
if self.v4_source_moved {
parts.push(match self.v4_source {
Some(ip) => format!("v4 source -> {}", ip),
None => "v4 source lost".to_string(),
});
}
if self.v6_source_moved {
parts.push(match self.v6_source {
Some(ip) => format!("v6 source -> {}", ip),
None => "v6 source lost".to_string(),
});
}
if !self.added.is_empty() {
parts.push(format!("+{} addr", self.added.len()));
}
if !self.removed.is_empty() {
parts.push(format!("-{} addr", self.removed.len()));
}
if parts.is_empty() {
return write!(f, "no visible difference");
}
write!(f, "{}", parts.join(", "))
}
}
/// One settled transport-medium change, as delivered to the rx loop.
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct NetChange {
/// Monotonically increasing across the life of one detector, starting at 1.
/// Present so a log line can be tied to the handler's reaction, and so a
/// coalesced delivery is visibly a coalesced delivery.
pub generation: u64,
/// What moved.
pub summary: NetChangeSummary,
}
impl NetChange {
/// A synthetic change, for tests that exercise the node's *reaction* to a
/// medium change rather than its detection. The summary is empty because
/// the handler never reads it — it re-evaluates every peer regardless of
/// which address moved, having no way to map an address back to the peers
/// that were reaching over it.
#[cfg(test)]
pub(crate) fn for_test(generation: u64) -> Self {
let empty = NetFingerprint::default();
Self {
generation,
summary: empty.diff(&empty),
}
}
}
/// What tells the detector it is worth taking another sample.
///
/// The split between *being woken* and *deciding whether anything changed* is
/// the reason a per-OS backend stays small: a backend only has to say "something
/// happened", and the fingerprint comparison, the debounce and the
/// settled-back-unchanged suppression are shared by all of them. No backend
/// parses kernel messages or decides what a medium change is.
struct WakeSource {
/// What, besides the timer, can wake the detector.
source: Wake,
/// The timer. Without a backend it is the only signal, and its period is
/// the detection latency. With one it is a backstop, and not a
/// belt-and-braces backstop: a netlink socket drops messages under memory
/// pressure (`ENOBUFS`), and the backend task can exit on a socket error,
/// either of which would otherwise leave the node noticing nothing at all.
/// Keeping the period the poller would have used makes an event-driven
/// backend a strict latency improvement rather than a replacement that can
/// regress, for the cost of a few syscalls per period.
timer: tokio::time::Interval,
}
/// Where a wake-up can come from, besides the timer.
enum Wake {
/// Nothing but the timer. The platform has no event source, or one could
/// not be opened.
Timer,
/// Kernel link and route events, via the shared [`LinkWatcher`].
///
/// The watcher parks forever when it has no source and after it gives up
/// on a broken one, so selecting it against the timer degrades to the
/// timer without any bookkeeping here.
#[cfg(unix)]
Kernel(LinkWatcher),
/// An injected channel, so the tests can drive the detector on a paused
/// clock without a live network or a real interface to flap.
///
/// Single-slot upstream, so a burst coalesces into one wake-up rather
/// than queueing a wake-up per message.
#[cfg(test)]
Injected(mpsc::Receiver<()>),
}
impl WakeSource {
/// A wake source with no event-driven backend: the timer alone.
fn timer_only(period: Duration) -> Self {
Self {
source: Wake::Timer,
timer: Self::make_timer(period),
}
}
/// A wake source driven by the kernel, with the timer as backstop.
#[cfg(unix)]
fn kernel(watcher: LinkWatcher, period: Duration) -> Self {
Self {
source: Wake::Kernel(watcher),
timer: Self::make_timer(period),
}
}
/// A wake source driven by an injected channel, with the timer as backstop.
#[cfg(test)]
fn events(pings: mpsc::Receiver<()>, period: Duration) -> Self {
Self {
source: Wake::Injected(pings),
timer: Self::make_timer(period),
}
}
/// `Delay` rather than the default `Burst`: a debounced handover can hold
/// the loop for longer than one period, and catching up afterwards would
/// fire several immediate wake-ups to sample a picture that just settled.
///
/// `reset` drops the free first tick a fresh `Interval` hands out. Without
/// it the detector's opening `wait` returns instantly and re-samples an
/// attachment it read microseconds earlier — harmless, but it would also
/// mean the very first wake-up on a netlink host came from the backstop
/// rather than the backend.
fn make_timer(period: Duration) -> tokio::time::Interval {
let mut timer = tokio::time::interval(period);
timer.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
timer.reset();
timer
}
/// Wait until it is worth sampling again.
async fn wait(&mut self) {
let WakeSource { source, timer } = self;
// Only the injected source can stop: [`LinkWatcher`] parks forever
// once it gives up, so a kernel source that dies simply stops firing
// and the timer carries on underneath it with nothing to unwind here.
#[cfg(test)]
let mut backend_gone = false;
match source {
Wake::Timer => {
timer.tick().await;
}
#[cfg(unix)]
Wake::Kernel(watcher) => {
tokio::select! {
_ = watcher.changed() => {}
_ = timer.tick() => {}
}
}
#[cfg(test)]
Wake::Injected(pings) => {
tokio::select! {
ping = pings.recv() => {
// A closed channel is the sender giving up. Sampling
// once more on the way past is deliberate: it may have
// died part-way through a change.
backend_gone = ping.is_none();
}
_ = timer.tick() => {}
}
}
}
#[cfg(test)]
if backend_gone {
warn!("Network-change backend stopped; falling back to polling");
self.source = Wake::Timer;
}
}
}
/// Spawn the medium-change detector, using the best backend this platform has.
///
/// Returns the receiver the rx loop drains and the task handle the supervisor
/// aborts at teardown. The channel holds a single slot: a change already queued
/// and not yet handled makes a newer one redundant, because the handler's
/// reaction is "re-evaluate every peer and every backoff", which subsumes any
/// number of coalesced changes. A full channel therefore drops rather than
/// queues, and never applies backpressure to the detector.
pub(crate) fn spawn_detector(cfg: NetmonConfig) -> (NetChangeRx, JoinHandle<()>) {
let (tx, rx) = mpsc::channel(1);
let handle = tokio::spawn(async move {
let wake = build_wake_source(&cfg);
run_detector(tx, cfg, NetFingerprint::sample, wake).await;
});
(rx, handle)
}
/// Pick the wake source: the event-driven backend where one exists and starts,
/// the timer otherwise.
///
/// A backend that fails to start is a downgrade, not a failure — an unprivileged
/// container, a locked-down sandbox or a kernel without the socket all land
/// here, and the node keeps working with the detection latency the poller
/// gives. It is logged once at `warn` so a slow recovery is explicable.
///
/// Built inside the spawned task rather than by the caller because a backend
/// binds sockets and spawns tasks, which belongs on the runtime that will own
/// them.
fn build_wake_source(cfg: &NetmonConfig) -> WakeSource {
let period = Duration::from_secs(cfg.poll_interval_secs.max(1));
#[cfg(unix)]
{
// Where the backend is netlink the mask matters: the default route
// moving between two interfaces that stay up emits nothing in the
// link group, so a presence watcher would never fire for the change
// this detector exists to catch. Under `PF_ROUTE` every routing
// message is delivered regardless and the mask is ignored.
#[cfg(any(target_os = "linux", target_os = "android"))]
let watcher = LinkWatcher::with_groups(crate::transport::watcher::groups::EGRESS_PATH);
#[cfg(not(any(target_os = "linux", target_os = "android")))]
let watcher = LinkWatcher::new();
if watcher.is_event_driven() {
debug!(
backstop_secs = cfg.poll_interval_secs,
"Network-change detection: kernel events"
);
return WakeSource::kernel(watcher, period);
}
warn!("Kernel medium-change events unavailable; falling back to polling");
}
debug!(
poll_interval_secs = cfg.poll_interval_secs,
"Network-change detection: polling"
);
WakeSource::timer_only(period)
}
/// The detection loop, over an injected sampler and wake source.
///
/// `sample` is [`NetFingerprint::sample`] in production; the tests drive the
/// debounce and coalescing against a scripted one, so neither needs a live
/// network nor a real interface to flap.
async fn run_detector<F>(tx: NetChangeTx, cfg: NetmonConfig, sample: F, mut wake: WakeSource)
where
F: Fn() -> NetFingerprint,
{
let debounce = Duration::from_millis(cfg.debounce_ms);
let mut last = sample();
let mut generation: u64 = 0;
let mut last_emit: Option<tokio::time::Instant> = None;
debug!(
poll_interval_secs = cfg.poll_interval_secs,
debounce_ms = cfg.debounce_ms,
"Network-change detector started"
);
loop {
wake.wait().await;
let mut candidate = sample();
if candidate == last {
continue;
}
// The picture is moving. Ride out the burst: resample after the
// debounce window until two consecutive samples agree, so the reported
// change is against a settled state rather than a mid-handover one.
for _ in 0..MAX_DEBOUNCE_ROUNDS {
if debounce.is_zero() {
break;
}
tokio::time::sleep(debounce).await;
let resampled = sample();
if resampled == candidate {
break;
}
candidate = resampled;
}
// The burst may have settled back to where it started (an address that
// flapped away and returned). Nothing changed, so nothing is reported.
if candidate == last {
trace!("Network fingerprint settled back unchanged; no event");
continue;
}
// Space out reactions. Deliberately after the debounce and the
// settled-back check, so a burst that resolves to no change costs
// nothing here, and only a real report is paced.
if let Some(previous) = last_emit {
let since = previous.elapsed();
if since < MIN_CHANGE_INTERVAL {
tokio::time::sleep(MIN_CHANGE_INTERVAL - since).await;
}
}
last_emit = Some(tokio::time::Instant::now());
generation += 1;
let change = NetChange {
generation,
summary: last.diff(&candidate),
};
last = candidate;
match tx.try_send(change) {
Ok(()) => {}
Err(mpsc::error::TrySendError::Full(dropped)) => {
debug!(
generation = dropped.generation,
"Network change coalesced into the one already queued"
);
}
Err(mpsc::error::TrySendError::Closed(_)) => {
debug!("Network-change receiver gone; poller exiting");
return;
}
}
}
}
/// The source address the kernel would use to reach `probe`.
///
/// `connect(2)` on a UDP socket is a pure routing-table operation: it resolves
/// the route, binds a source address, and sends nothing. A failure — most often
/// `ENETUNREACH` with no route of that family — is itself a fingerprint value,
/// reported as `None` rather than swallowed.
fn preferred_source(probe: SocketAddr) -> Option<IpAddr> {
let bind: SocketAddr = match probe {
SocketAddr::V4(_) => SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0),
SocketAddr::V6(_) => SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 0),
};
let socket = UdpSocket::bind(bind).ok()?;
socket.connect(probe).ok()?;
let local = socket.local_addr().ok()?.ip();
// An unspecified local address means the kernel deferred the choice, which
// tells us nothing about the medium. Treat it as "no answer".
if local.is_unspecified() {
return None;
}
Some(local)
}
/// Every up, non-loopback unicast address on the host.
#[cfg(unix)]
fn interface_addrs() -> BTreeSet<IpAddr> {
let mut out = BTreeSet::new();
let mut head: *mut libc::ifaddrs = std::ptr::null_mut();
// SAFETY: `getifaddrs` either returns 0 and writes an owned linked list
// into `head`, or returns non-zero and leaves `head` untouched — so the
// list is only walked on success. Every node is read behind a null check,
// and `freeifaddrs` releases the list exactly once, after the walk.
if unsafe { libc::getifaddrs(&mut head) } != 0 {
return out;
}
let mut cursor = head;
while !cursor.is_null() {
// SAFETY: `cursor` is non-null here and points at a node of the list
// `getifaddrs` allocated, which stays valid until `freeifaddrs` below.
let entry = unsafe { &*cursor };
cursor = entry.ifa_next;
if entry.ifa_addr.is_null() {
continue;
}
let flags = entry.ifa_flags as i32;
let up = flags & libc::IFF_UP != 0 && flags & libc::IFF_RUNNING != 0;
if !up || flags & libc::IFF_LOOPBACK != 0 {
continue;
}
if let Some(ip) = sockaddr_ip(entry.ifa_addr) {
out.insert(ip);
}
}
// SAFETY: `head` is the list `getifaddrs` allocated above, freed once, and
// not read after this point (`cursor` is null by loop exit).
unsafe { libc::freeifaddrs(head) };
out
}
/// No portable interface enumeration without `getifaddrs`. The fingerprint
/// degrades to the source-address probe, which still catches every
/// default-route move; the native `NotifyIpInterfaceChange` backend is the
/// answer here rather than a richer poll.
#[cfg(not(unix))]
fn interface_addrs() -> BTreeSet<IpAddr> {
BTreeSet::new()
}
/// Read an `IpAddr` out of a kernel-supplied `sockaddr`, if it is one of the
/// two families we fingerprint.
#[cfg(unix)]
fn sockaddr_ip(sa: *const libc::sockaddr) -> Option<IpAddr> {
// SAFETY: `sa` is non-null (checked by the caller) and points at a
// kernel-supplied `sockaddr` whose `sa_family` selects the concrete layout
// that follows. Both branches copy out through `read_unaligned`, so nothing
// here assumes the pointer is aligned for the larger type.
let family = unsafe { std::ptr::addr_of!((*sa).sa_family).read_unaligned() } as i32;
match family {
libc::AF_INET => {
// SAFETY: family is AF_INET, so the allocation is at least a
// `sockaddr_in`.
let raw: libc::sockaddr_in = unsafe { std::ptr::read_unaligned(sa.cast()) };
// `s_addr` holds the octets in network order, so its native-endian
// bytes are the address octets in order.
Some(IpAddr::V4(Ipv4Addr::from(
raw.sin_addr.s_addr.to_ne_bytes(),
)))
}
libc::AF_INET6 => {
// SAFETY: family is AF_INET6, so the allocation is at least a
// `sockaddr_in6`.
let raw: libc::sockaddr_in6 = unsafe { std::ptr::read_unaligned(sa.cast()) };
Some(IpAddr::V6(Ipv6Addr::from(raw.sin6_addr.s6_addr)))
}
_ => None,
}
}
#[cfg(test)]
mod tests;
+441
View File
@@ -0,0 +1,441 @@
//! Poller tests.
//!
//! The debounce, coalescing and no-net-difference paths are driven through a
//! scripted sampler on tokio's paused clock, so none of them needs a real
//! interface to flap. The two live-sampling tests assert only what is true of
//! any host, including a CI container with a single interface.
use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
/// Virtual-time budget for the "nothing should arrive" assertions. On the
/// paused clock the runtime auto-advances whenever every task is idle, so this
/// covers several poll intervals without costing real wall time.
const QUIET_WINDOW: Duration = Duration::from_secs(30);
/// Await one change, failing rather than hanging if the poller never sends.
async fn expect_change(rx: &mut NetChangeRx) -> NetChange {
tokio::time::timeout(QUIET_WINDOW, rx.recv())
.await
.expect("the poller must report within the quiet window")
.expect("the channel must stay open")
}
/// Assert nothing arrives for a generous stretch of virtual time.
async fn expect_quiet(rx: &mut NetChangeRx, why: &str) {
assert!(
tokio::time::timeout(QUIET_WINDOW, rx.recv()).await.is_err(),
"{}",
why
);
}
fn v4(a: u8, b: u8, c: u8, d: u8) -> IpAddr {
IpAddr::V4(Ipv4Addr::new(a, b, c, d))
}
/// A timer-only wake source at the config's poll period — the portable
/// backend's behaviour, and the baseline the netlink tests compare against.
fn timer_wake(poll_secs: u64) -> WakeSource {
WakeSource::timer_only(Duration::from_secs(poll_secs))
}
fn cfg(poll_secs: u64, debounce_ms: u64) -> NetmonConfig {
NetmonConfig {
enabled: true,
poll_interval_secs: poll_secs,
debounce_ms,
}
}
/// A sampler that walks a script, holding on the last entry forever.
fn scripted(samples: Vec<NetFingerprint>) -> (impl Fn() -> NetFingerprint, Arc<AtomicUsize>) {
let calls = Arc::new(AtomicUsize::new(0));
let counter = calls.clone();
let sampler = move || {
let i = counter.fetch_add(1, Ordering::SeqCst);
samples[i.min(samples.len() - 1)].clone()
};
(sampler, calls)
}
#[tokio::test(start_paused = true)]
async fn steady_attachment_reports_nothing() {
let steady = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
let (sampler, _) = scripted(vec![steady]);
let (tx, mut rx) = mpsc::channel(1);
tokio::spawn(run_detector(tx, cfg(1, 0), sampler, timer_wake(1)));
expect_quiet(&mut rx, "an unchanging fingerprint must produce no events").await;
}
#[tokio::test(start_paused = true)]
async fn a_default_route_move_is_reported() {
// The WLAN → 5G shape: the interface set changes and the preferred source
// address moves with it.
let wlan = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
let cell = NetFingerprint::for_test(Some(v4(10, 40, 0, 7)), &[v4(10, 40, 0, 7)]);
let (sampler, _) = scripted(vec![wlan, cell]);
let (tx, mut rx) = mpsc::channel(1);
tokio::spawn(run_detector(tx, cfg(1, 0), sampler, timer_wake(1)));
let change = expect_change(&mut rx).await;
assert_eq!(change.generation, 1);
assert!(change.summary.v4_source_moved);
assert_eq!(change.summary.v4_source, Some(v4(10, 40, 0, 7)));
assert_eq!(change.summary.added, vec![v4(10, 40, 0, 7)]);
assert_eq!(change.summary.removed, vec![v4(192, 168, 1, 10)]);
}
#[tokio::test(start_paused = true)]
async fn a_handover_burst_coalesces_into_one_event() {
// A handover is not atomic: the old address goes, then briefly nothing has
// a route, then the new address arrives. Reporting each step would have the
// handler probing every peer three times against a picture still in
// motion. The debounce must ride the burst out and report once, against the
// settled state.
let wlan = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
let gone = NetFingerprint::for_test(None, &[]);
let cell = NetFingerprint::for_test(Some(v4(10, 40, 0, 7)), &[v4(10, 40, 0, 7)]);
let (sampler, _) = scripted(vec![wlan, gone, cell]);
let (tx, mut rx) = mpsc::channel(1);
tokio::spawn(run_detector(tx, cfg(1, 250), sampler, timer_wake(1)));
let change = expect_change(&mut rx).await;
assert_eq!(
change.generation, 1,
"the burst must report once, not per step"
);
assert_eq!(
change.summary.v4_source,
Some(v4(10, 40, 0, 7)),
"the reported state must be the settled one, not the mid-handover one"
);
expect_quiet(&mut rx, "no second event for the same handover").await;
}
#[tokio::test(start_paused = true)]
async fn a_flap_that_settles_back_reports_nothing() {
// An address that leaves and returns within the debounce window is not a
// medium change. Reporting it would have every peer probed for nothing,
// which on a host with churning routes is exactly the reconnect storm this
// is meant to avoid.
let steady = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
let gone = NetFingerprint::for_test(None, &[]);
let (sampler, _) = scripted(vec![steady.clone(), gone, steady]);
let (tx, mut rx) = mpsc::channel(1);
tokio::spawn(run_detector(tx, cfg(1, 250), sampler, timer_wake(1)));
expect_quiet(
&mut rx,
"a fingerprint that settles back where it started is not a change",
)
.await;
}
#[tokio::test(start_paused = true)]
async fn an_unread_change_coalesces_rather_than_queues() {
// The handler's reaction is "re-evaluate every peer and every backoff",
// which subsumes any number of changes. A second change arriving before the
// first is drained must therefore drop, not queue: the node must never work
// through a backlog of stale network states.
let a = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
let b = NetFingerprint::for_test(Some(v4(10, 40, 0, 7)), &[v4(10, 40, 0, 7)]);
let c = NetFingerprint::for_test(Some(v4(172, 16, 3, 2)), &[v4(172, 16, 3, 2)]);
let (sampler, _) = scripted(vec![a, b, c]);
let (tx, mut rx) = mpsc::channel(1);
tokio::spawn(run_detector(tx, cfg(1, 0), sampler, timer_wake(1)));
// Stay idle long enough for the poller to see both changes while nothing
// is draining, so the second meets a full channel.
tokio::time::sleep(Duration::from_secs(10)).await;
assert!(rx.try_recv().is_ok(), "the first change is delivered");
assert!(
rx.try_recv().is_err(),
"the second must have coalesced into the undrained first, not queued behind it"
);
}
#[tokio::test(start_paused = true)]
async fn a_closed_receiver_ends_the_poller() {
let a = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[]);
let b = NetFingerprint::for_test(Some(v4(10, 40, 0, 7)), &[]);
let (sampler, _) = scripted(vec![a, b]);
let (tx, rx) = mpsc::channel(1);
drop(rx);
let handle = tokio::spawn(run_detector(tx, cfg(1, 0), sampler, timer_wake(1)));
tokio::time::timeout(QUIET_WINDOW, handle)
.await
.expect("the poller must exit once nothing is listening")
.expect("and exit cleanly, not by panic");
}
#[test]
fn sampling_the_live_host_is_self_consistent() {
// Two samples taken back to back on an idle host describe the same
// attachment. This is the property the whole detector rests on: if plain
// sampling were noisy, every poll would look like a medium change.
let first = NetFingerprint::sample();
let second = NetFingerprint::sample();
assert_eq!(
first, second,
"consecutive samples of an unchanged host must agree"
);
}
#[test]
fn live_interface_addresses_exclude_loopback() {
// Loopback is present on every host and never changes, so including it
// would only add noise. Unix enumerates; elsewhere the set is empty by
// design and the assertion holds vacuously.
let sample = NetFingerprint::sample();
assert!(
!sample.local_addrs.iter().any(|ip| ip.is_loopback()),
"loopback must not contribute to the fingerprint: {:?}",
sample.local_addrs
);
}
#[test]
fn summary_of_an_empty_diff_is_legible() {
let same = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
assert_eq!(same.diff(&same).to_string(), "no visible difference");
}
#[test]
fn summary_names_the_new_source_address() {
let wlan = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
let cell = NetFingerprint::for_test(Some(v4(10, 40, 0, 7)), &[v4(10, 40, 0, 7)]);
let rendered = wlan.diff(&cell).to_string();
assert!(rendered.contains("v4 source -> 10.40.0.7"), "{}", rendered);
assert!(rendered.contains("+1 addr"), "{}", rendered);
assert!(rendered.contains("-1 addr"), "{}", rendered);
}
// === Wake source ===
/// The point of an event-driven backend: a change is acted on when the kernel
/// says so, not when the next poll happens to come round. The poll period here
/// is an hour, so only the ping can be what woke the detector.
#[tokio::test(start_paused = true)]
async fn an_event_ping_wakes_the_detector_before_the_timer_would() {
let wlan = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
let cell = NetFingerprint::for_test(Some(v4(10, 40, 0, 7)), &[v4(10, 40, 0, 7)]);
let (sampler, _) = scripted(vec![wlan, cell]);
let (tx, mut rx) = mpsc::channel(1);
let (pings, ping_rx) = mpsc::channel(1);
let wake = WakeSource::events(ping_rx, Duration::from_secs(3600));
tokio::spawn(run_detector(tx, cfg(3600, 0), sampler, wake));
pings.send(()).await.expect("the backend can ping");
let change = expect_change(&mut rx).await;
assert_eq!(change.summary.v4_source, Some(v4(10, 40, 0, 7)));
}
/// A netlink socket drops messages under memory pressure, and a backend can go
/// quiet without going away. The backstop timer must still get the node there,
/// so an event-driven backend is never worse than the poller it replaced.
#[tokio::test(start_paused = true)]
async fn the_backstop_still_fires_when_the_backend_says_nothing() {
let wlan = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
let cell = NetFingerprint::for_test(Some(v4(10, 40, 0, 7)), &[v4(10, 40, 0, 7)]);
let (sampler, _) = scripted(vec![wlan, cell]);
let (tx, mut rx) = mpsc::channel(1);
// Held, never sent on: the backend is alive but has missed the event.
let (_pings, ping_rx) = mpsc::channel(1);
let wake = WakeSource::events(ping_rx, Duration::from_secs(1));
tokio::spawn(run_detector(tx, cfg(1, 0), sampler, wake));
let change = expect_change(&mut rx).await;
assert_eq!(
change.summary.v4_source,
Some(v4(10, 40, 0, 7)),
"the backstop must reach the change the backend missed"
);
}
/// A backend that dies — socket error, sandbox revocation — must degrade the
/// node to polling, not stop detection. Spinning on the closed channel would be
/// worse still.
#[tokio::test(start_paused = true)]
async fn a_dead_backend_falls_back_to_the_timer() {
let wlan = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
let cell = NetFingerprint::for_test(Some(v4(10, 40, 0, 7)), &[v4(10, 40, 0, 7)]);
let (sampler, _) = scripted(vec![wlan, cell]);
let (tx, mut rx) = mpsc::channel(1);
let (pings, ping_rx) = mpsc::channel(1);
let wake = WakeSource::events(ping_rx, Duration::from_secs(1));
tokio::spawn(run_detector(tx, cfg(1, 0), sampler, wake));
// The backend gives up before the medium moves.
drop(pings);
let change = expect_change(&mut rx).await;
assert_eq!(
change.summary.v4_source,
Some(v4(10, 40, 0, 7)),
"detection must survive the backend it was using"
);
}
/// A wake source with no backend waits out its period rather than taking the
/// free first tick a fresh `Interval` hands out — otherwise the opening sample
/// is a duplicate of the one taken microseconds earlier.
#[tokio::test(start_paused = true)]
async fn the_first_wait_is_a_real_wait() {
let mut wake = WakeSource::timer_only(Duration::from_secs(60));
let start = tokio::time::Instant::now();
wake.wait().await;
assert!(
start.elapsed() >= Duration::from_secs(60),
"the first tick must not come free"
);
}
// === Kernel event source ===
/// The watcher this detector builds opens on a normal Linux host.
///
/// Distinct from the equivalent check in `transport::watcher`: that one pins
/// the *link* mask, this one pins the wider mask netmon actually asks for. A
/// group constant that was wrong only in the added bits would pass there and
/// fail here. A sandbox that refuses the subscription is a legitimate outcome
/// — it is why the fallback exists — so that case reports rather than fails.
#[cfg(target_os = "linux")]
#[tokio::test]
async fn the_egress_path_watcher_starts_or_cleanly_declines() {
use crate::transport::watcher::{LinkWatcher, groups};
let watcher = LinkWatcher::with_groups(groups::EGRESS_PATH);
if !watcher.is_event_driven() {
eprintln!("kernel events unavailable in this environment; fallback path applies");
}
}
/// A route change — with no link change alongside it — reaches the watcher.
///
/// This is the test that justifies the wider group mask, and the only one
/// that can fail if the mask is narrowed back. `RTMGRP_LINK` alone sees
/// nothing here: the interface does not appear, disappear or change state,
/// and yet the host's egress path has moved, which is exactly the event this
/// detector exists to catch. Every other test in this file drives the shared
/// decision logic through an injected channel and would pass against a
/// subscription that never fired for a route at all.
///
/// Ignored because it needs `CAP_NET_ADMIN` in a private network namespace —
/// it edits a routing table, which must not touch the developer's real
/// network. Run it with:
///
/// ```text
/// unshare -rn cargo test --lib netmon -- --ignored --nocapture
/// ```
#[cfg(target_os = "linux")]
#[tokio::test]
#[ignore = "needs CAP_NET_ADMIN in a private netns; run under `unshare -rn`"]
async fn a_route_change_alone_reaches_the_watcher() {
use crate::transport::watcher::{LinkWatcher, groups};
use futures::TryStreamExt;
use std::net::Ipv4Addr;
let (connection, handle, _) = rtnetlink::new_connection().expect("netlink connection");
tokio::spawn(connection);
// `lo` is down in a fresh namespace and a route needs a live interface,
// so bring it up first — before the watcher exists, so that this link
// change cannot be the thing the assertion below observes.
let index = handle
.link()
.get()
.match_name("lo".to_string())
.execute()
.try_next()
.await
.expect("link query")
.expect("lo exists")
.header
.index;
handle
.link()
.change(rtnetlink::LinkUnspec::new_with_index(index).up().build())
.execute()
.await
.expect("bringing lo up needs CAP_NET_ADMIN in this namespace");
tokio::time::sleep(Duration::from_millis(200)).await;
let watcher = LinkWatcher::with_groups(groups::EGRESS_PATH);
assert!(
watcher.is_event_driven(),
"this test cannot say anything without a live subscription"
);
// A route to TEST-NET-1 out of `lo`: no interface changes state, so the
// link group stays silent and only the route group can carry this.
handle
.route()
.add(
rtnetlink::RouteMessageBuilder::<Ipv4Addr>::new()
.destination_prefix(Ipv4Addr::new(192, 0, 2, 0), 24)
.output_interface(index)
.build(),
)
.execute()
.await
.expect("adding a route needs CAP_NET_ADMIN in this namespace");
tokio::time::timeout(Duration::from_secs(5), watcher.changed())
.await
.expect("a route change must reach the watcher well inside 5s");
}
/// Reactions are paced. Dropping every peer's connected socket and heartbeating
/// each of them is not free, so an interface that flaps cleanly — settling
/// between transitions, which defeats the debounce — must not drive that
/// several times a second across the whole peer set.
#[tokio::test(start_paused = true)]
async fn reports_are_spaced_out_under_clean_flapping() {
let a = NetFingerprint::for_test(Some(v4(192, 168, 1, 10)), &[v4(192, 168, 1, 10)]);
let b = NetFingerprint::for_test(Some(v4(10, 40, 0, 7)), &[v4(10, 40, 0, 7)]);
// Alternates every sample: each poll sees a settled but different picture.
let calls = Arc::new(AtomicUsize::new(0));
let counter = calls.clone();
let sampler = move || {
let i = counter.fetch_add(1, Ordering::SeqCst);
if i.is_multiple_of(2) {
a.clone()
} else {
b.clone()
}
};
let (tx, mut rx) = mpsc::channel(1);
// Poll far faster than the pacing floor, so only the floor can space these.
tokio::spawn(run_detector(tx, cfg(1, 0), sampler, timer_wake(1)));
let first = expect_change(&mut rx).await;
let started = tokio::time::Instant::now();
let second = expect_change(&mut rx).await;
assert_eq!(first.generation, 1);
assert_eq!(second.generation, 2);
assert!(
started.elapsed() >= MIN_CHANGE_INTERVAL,
"consecutive reports must be at least {:?} apart, got {:?}",
MIN_CHANGE_INTERVAL,
started.elapsed()
);
}
+1
View File
@@ -19,6 +19,7 @@ mod forwarding;
mod handshake;
mod heartbeat;
mod mmp_chartests;
mod netmon;
mod probe;
mod routing;
mod session;
+259
View File
@@ -0,0 +1,259 @@
//! What the node does when the transport medium changes.
//!
//! The detector itself is tested in `node::netmon::tests`; these drive the
//! reaction. The property that matters is that a medium change *rebinds* the
//! send path without disturbing the peering — the peer keeps its Noise session,
//! its tree position and its routes, and only the socket underneath it moves.
use super::spanning_tree::*;
use super::*;
use crate::config::PeerConfig;
use crate::config::TcpConfig;
use crate::node::netmon::NetChange;
use crate::transport::tcp::TcpTransport;
use crate::transport::{TransportAddr, TransportHandle, TransportId, packet_channel};
/// Add `peer` to the node's config as an auto-connect peer.
///
/// Node config is immutable after construction, so this goes through the same
/// copy-on-write context swap the heartbeat tests use.
fn configure_auto_peer(node: &mut Node, peer: &PeerIdentity) {
let peer_config = PeerConfig::new(peer.npub(), "udp", "127.0.0.1:1");
node.replace_context(|ctx| {
let mut cfg = (*ctx.config).clone();
cfg.peers.push(peer_config);
ctx.config = std::sync::Arc::new(cfg);
});
}
/// The peer identity node `j` presents to its peers.
fn identity_of(nodes: &[TestNode], j: usize) -> PeerIdentity {
PeerIdentity::from_pubkey_full(nodes[j].node.identity().pubkey_full())
}
/// Install a real `connect()`-ed UDP socket on a peer, the way the tick-driven
/// activation does.
///
/// The socket is opened against a discard port on loopback: nothing is ever
/// sent through it, and the test only cares whether the handle survives a
/// medium change.
#[cfg(target_os = "linux")]
fn install_connected_udp(node: &mut Node, addr: &NodeAddr, transport_id: TransportId) {
let local: std::net::SocketAddr = "0.0.0.0:0".parse().unwrap();
let peer_sa: std::net::SocketAddr = "127.0.0.1:9".parse().unwrap();
let owned = crate::transport::udp::open_connected_fd(local, peer_sa, 65_536, 65_536)
.expect("open a connected UDP socket");
let bound = crate::transport::udp::ConnectedPeerSocket::from_fd(owned, peer_sa, local);
let socket = std::sync::Arc::new(bound);
let (packet_tx, _packet_rx) = crate::transport::packet_channel(8);
let drain = crate::transport::udp::PeerRecvDrain::spawn(
socket.clone(),
transport_id,
peer_sa,
packet_tx,
)
.expect("spawn the peer recv drain");
node.get_peer_mut(addr)
.expect("peer present")
.set_connected_udp(socket, drain);
}
/// **The defect this feature exists for.**
///
/// Established UDP peers get a per-peer `connect()`-ed socket. `open_connected_fd`
/// binds the wildcard and then calls `connect(2)`, which — as its own comment
/// says — "locks in the per-packet kernel route": the kernel resolves the route
/// once and auto-binds the local source address to whichever interface was
/// carrying it at that moment. It never re-evaluates.
///
/// So when the host changes medium, every peer keeps transmitting from an
/// address the routing table has abandoned, on a socket pinned to the interface
/// the node has just moved off. The only other code that drops these sockets
/// fires when the *peer* rotates its address — the mirror-image case. Nothing
/// covered a local move, and a local move is invisible in the data plane, which
/// is why it went unhandled.
///
/// Observed in the field as a peering that carried exactly one packet after a
/// route change and then stalled until the 30s liveness timeout, reporting
/// itself connected the whole time.
#[cfg(target_os = "linux")]
#[tokio::test]
async fn a_medium_change_drops_connected_sockets_pinned_to_the_old_path() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
let peer_1 = identity_of(&nodes, 1);
configure_auto_peer(&mut nodes[0].node, &peer_1);
let transport_id = nodes[0].transport_id;
install_connected_udp(&mut nodes[0].node, &addr_1, transport_id);
assert!(
nodes[0]
.node
.get_peer(&addr_1)
.unwrap()
.connected_udp()
.is_some(),
"precondition: the peer holds a connected socket"
);
nodes[0]
.node
.handle_net_change(NetChange::for_test(1))
.await;
assert!(
nodes[0]
.node
.get_peer(&addr_1)
.unwrap()
.connected_udp()
.is_none(),
"a socket pinned to the old source address must not survive the change"
);
}
/// The rebind must not cost the peering. Everything above the socket — the
/// Noise session, the tree position, the routes — is unaffected by which local
/// address the node sends from, so a medium change that tore peers down would
/// be replacing a stall with a re-handshake for no reason.
#[tokio::test]
async fn a_medium_change_keeps_every_peering_intact() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
let peer_1 = identity_of(&nodes, 1);
configure_auto_peer(&mut nodes[0].node, &peer_1);
let link_before = nodes[0].node.get_peer(&addr_1).unwrap().link_id();
nodes[0]
.node
.handle_net_change(NetChange::for_test(1))
.await;
let peer = nodes[0]
.node
.get_peer(&addr_1)
.expect("the peering must survive a medium change");
assert_eq!(
peer.link_id(),
link_before,
"the same link, not a rebuilt one: no re-handshake"
);
cleanup_nodes(&mut nodes).await;
}
/// The far side has the same stale-address problem in reverse: it is still
/// sending to wherever it last heard us. One heartbeat over the new path
/// carries the node's new source address, so the peer re-pins on receipt rather
/// than waiting out its own heartbeat interval.
#[tokio::test]
async fn every_peer_is_heartbeated_so_the_far_side_re_pins() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
let peer_1 = identity_of(&nodes, 1);
configure_auto_peer(&mut nodes[0].node, &peer_1);
let before = nodes[0]
.node
.get_peer(&addr_1)
.unwrap()
.last_heartbeat_sent();
nodes[0]
.node
.handle_net_change(NetChange::for_test(1))
.await;
let after = nodes[0]
.node
.get_peer(&addr_1)
.unwrap()
.last_heartbeat_sent();
assert!(
after.is_some(),
"every peer is heartbeated on a medium change"
);
assert!(
before.is_none() || after > before,
"the heartbeat must go now, not at the next due interval"
);
cleanup_nodes(&mut nodes).await;
}
/// A peer on a connection-oriented transport is deliberately left out of the
/// immediate fan-out.
///
/// Its send would await `write_all` on a stream the medium change has very
/// likely just stranded — unbounded, and on the rx loop, where it would hold
/// every other arm of the select behind it. Such a peer keeps the periodic
/// heartbeat it had before this detector existed. If this ever starts passing
/// because the peer *was* heartbeated, the rx loop has a new way to stall.
#[tokio::test]
async fn a_peer_on_a_connection_oriented_transport_is_left_to_the_periodic_heartbeat() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
let peer_1 = identity_of(&nodes, 1);
configure_auto_peer(&mut nodes[0].node, &peer_1);
// Re-pin the peer onto a TCP transport. Nothing is connected on it, which
// is the point: the fan-out must decide from the transport's kind, before
// it ever reaches a send.
let tcp_id = TransportId::new(77);
let cfg = TcpConfig {
bind_addr: None,
..Default::default()
};
let (tx, _rx) = packet_channel(64);
nodes[0].node.transports.insert(
tcp_id,
TransportHandle::Tcp(TcpTransport::new(tcp_id, None, cfg, tx)),
);
nodes[0]
.node
.peers
.get_mut(&addr_1)
.expect("peer 1 is established")
.set_current_addr(tcp_id, TransportAddr::from_string("10.0.0.2:2121"));
let before = nodes[0]
.node
.get_peer(&addr_1)
.unwrap()
.last_heartbeat_sent();
nodes[0]
.node
.handle_net_change(NetChange::for_test(1))
.await;
let after = nodes[0]
.node
.get_peer(&addr_1)
.unwrap()
.last_heartbeat_sent();
assert_eq!(
before, after,
"a connection-oriented peer must not be heartbeated from the rx loop"
);
cleanup_nodes(&mut nodes).await;
}
/// A node with no peers has nothing to rebind and must not care.
#[tokio::test]
async fn a_change_with_no_peers_is_harmless() {
let mut node = make_node();
node.handle_net_change(NetChange::for_test(1)).await;
assert!(node.peers.is_empty());
}
+10
View File
@@ -15,6 +15,16 @@ pub mod udp;
#[cfg(any(target_os = "linux", target_os = "macos"))]
pub mod ethernet;
/// Kernel link-event notifications, shared by any transport or detector that
/// needs to react to interface state without polling for it.
///
/// Crate-internal on purpose. It is a mechanism the node's own subsystems
/// share, not a surface an embedder builds against, and publishing it would
/// commit the library to its shape before anything outside the crate has asked
/// for it.
#[cfg(unix)]
pub(crate) mod watcher;
#[cfg(ble_available)]
pub mod ble;
+453
View File
@@ -0,0 +1,453 @@
//! Kernel link-event sources.
//!
//! A watcher that resolves when the kernel reports that something about the
//! host's network links changed. Callers use it to react to interface state in
//! sub-second time instead of polling for it.
//!
//! | Platform | Source |
//! | -------- | ------ |
//! | Linux, Android | netlink `RTNLGRP_LINK` (`RTM_NEWLINK` / `RTM_DELLINK`) |
//! | macOS, FreeBSD | `PF_ROUTE` socket, `RTM_IFINFO` |
//! | Fallback | none — the watcher never fires, and callers poll |
//!
//! The messages themselves are deliberately **not parsed**. An event is a hint
//! to re-run whatever question the caller actually cares about, which is
//! cheap and authoritative; decoding `nlmsghdr`/`ifinfomsg` payloads to reach
//! the same answer would add a parser whose bugs would become the caller's
//! bugs. Any event on the socket wakes the caller, which then asks its own
//! question directly.
//!
//! Construction is best-effort, and that is the contract: a kernel or sandbox
//! that refuses the socket yields a watcher that never fires. `changed()` then
//! parks forever, which is what makes it safe to `select!` against a poll
//! ticker — the ticker simply always wins, and the caller degrades to polling
//! without a special case.
//!
//! Callers that need a *different* event group should extend
//! `open_link_socket` rather than opening a second socket beside this one:
//! `RTNLGRP_LINK` carries link state only, so a route change with both
//! interfaces up produces no event here.
use std::os::unix::io::{AsRawFd, RawFd};
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use std::time::Duration;
use tokio::io::unix::AsyncFd;
use tracing::{debug, warn};
/// Consecutive receive errors before the event source is abandoned for the
/// caller's poll.
const ERROR_GIVE_UP: u32 = 5;
/// An owned link-event socket. Closes its descriptor on drop.
struct LinkEventSocket {
fd: RawFd,
}
impl AsRawFd for LinkEventSocket {
fn as_raw_fd(&self) -> RawFd {
self.fd
}
}
impl Drop for LinkEventSocket {
fn drop(&mut self) {
unsafe { libc::close(self.fd) };
}
}
impl LinkEventSocket {
fn recv(&self, buf: &mut [u8]) -> std::io::Result<usize> {
let n = unsafe { libc::recv(self.fd, buf.as_mut_ptr() as *mut libc::c_void, buf.len(), 0) };
if n < 0 {
Err(std::io::Error::last_os_error())
} else {
Ok(n as usize)
}
}
}
/// Netlink multicast groups a watcher can subscribe to where the backend is
/// netlink.
///
/// Spelled as literals because the constants' names and availability differ
/// across libc versions; the values are ABI. These are the `RTMGRP_*` bitmask
/// form taken by `sockaddr_nl.nl_groups`, not the `RTNLGRP_*` ordinals.
#[cfg(any(target_os = "linux", target_os = "android"))]
pub mod groups {
/// Interfaces appearing, disappearing, or changing state.
pub const LINK: u32 = 0x1;
/// IPv4 addresses added to or removed from an interface.
pub const IPV4_IFADDR: u32 = 0x10;
/// IPv4 route table changes, including the default route moving.
pub const IPV4_ROUTE: u32 = 0x40;
/// IPv6 addresses added to or removed from an interface.
pub const IPV6_IFADDR: u32 = 0x100;
/// IPv6 route table changes.
pub const IPV6_ROUTE: u32 = 0x400;
/// Everything that can change which local address the host would use to
/// reach a given destination.
///
/// [`LINK`] alone does not cover it. A default route moving between two
/// interfaces that both stay up emits no link message at all — verified
/// with `ip monitor`, which reports zero events in the link group for
/// that change and two in the route group. A watcher that wants to hear
/// about egress-path changes rather than interface presence needs this.
pub const EGRESS_PATH: u32 = LINK | IPV4_IFADDR | IPV4_ROUTE | IPV6_IFADDR | IPV6_ROUTE;
}
/// Open the platform's link-event socket, non-blocking.
///
/// `groups` is ignored where the backend is `PF_ROUTE`: it has no group
/// selection and delivers every routing message to every reader, so a caller
/// that wants more than link events already has them there.
#[cfg(any(target_os = "linux", target_os = "android"))]
fn open_link_socket(groups: u32) -> std::io::Result<LinkEventSocket> {
let fd = unsafe {
libc::socket(
libc::AF_NETLINK,
libc::SOCK_RAW | libc::SOCK_NONBLOCK | libc::SOCK_CLOEXEC,
libc::NETLINK_ROUTE,
)
};
if fd < 0 {
return Err(std::io::Error::last_os_error());
}
let socket = LinkEventSocket { fd };
let mut sa: libc::sockaddr_nl = unsafe { std::mem::zeroed() };
sa.nl_family = libc::AF_NETLINK as u16;
sa.nl_groups = groups;
let ret = unsafe {
libc::bind(
fd,
&sa as *const libc::sockaddr_nl as *const libc::sockaddr,
std::mem::size_of::<libc::sockaddr_nl>() as libc::socklen_t,
)
};
if ret < 0 {
return Err(std::io::Error::last_os_error());
}
Ok(socket)
}
/// Open the platform's link-event socket, non-blocking.
#[cfg(not(any(target_os = "linux", target_os = "android")))]
fn open_link_socket(_groups: u32) -> std::io::Result<LinkEventSocket> {
// PF_ROUTE delivers RTM_IFINFO (and the rest of the routing messages) to
// every reader; no bind and no group selection exist for it.
let fd = unsafe { libc::socket(libc::PF_ROUTE, libc::SOCK_RAW, libc::AF_UNSPEC) };
if fd < 0 {
return Err(std::io::Error::last_os_error());
}
let socket = LinkEventSocket { fd };
let flags = unsafe { libc::fcntl(fd, libc::F_GETFL) };
if flags < 0 {
return Err(std::io::Error::last_os_error());
}
if unsafe { libc::fcntl(fd, libc::F_SETFL, flags | libc::O_NONBLOCK) } < 0 {
return Err(std::io::Error::last_os_error());
}
Ok(socket)
}
/// A source of "something about the links changed" wake-ups.
pub struct LinkWatcher {
/// `None` when no event source could be opened — the caller's poll is then
/// the whole mechanism, which is exactly the documented fallback.
inner: Option<AsyncFd<LinkEventSocket>>,
/// Consecutive receive errors. Reset by any successful read.
errors: AtomicU32,
/// Set once the source has been abandoned for good.
///
/// Abandonment has to outlive the future that decided it. `changed()` is
/// called fresh on every pass of the caller's `select!` and dropped
/// whenever the poll ticker wins, so a `pending()` inside that future
/// parks nothing beyond the current pass — without this flag the next pass
/// re-reads the dead socket, re-counts the error, and re-logs the
/// give-up warning, once per wake-up, forever.
given_up: AtomicBool,
}
impl Default for LinkWatcher {
fn default() -> Self {
Self::new()
}
}
impl LinkWatcher {
/// Open a watcher for interface presence.
///
/// Subscribes to link events only, which is what a caller asking "is this
/// interface here?" needs.
pub fn new() -> Self {
#[cfg(any(target_os = "linux", target_os = "android"))]
let groups = groups::LINK;
#[cfg(not(any(target_os = "linux", target_os = "android")))]
let groups = 0;
Self::with_groups(groups)
}
/// Open a watcher over an explicit set of netlink multicast groups.
///
/// Only meaningful on Linux, where the group mask decides what the kernel
/// sends; elsewhere `PF_ROUTE` delivers everything regardless and the mask
/// is ignored. See [`groups`] for the values, and `groups::EGRESS_PATH`
/// for the set that covers a change of egress path rather than of
/// interface presence.
pub fn with_groups(groups: u32) -> Self {
let inner = match open_link_socket(groups) {
Ok(socket) => match AsyncFd::new(socket) {
Ok(afd) => Some(afd),
Err(e) => {
debug!(error = %e, "Link event socket not registrable; polling instead");
None
}
},
Err(e) => {
debug!(error = %e, "No link event source available; polling instead");
None
}
};
Self {
inner,
errors: AtomicU32::new(0),
given_up: AtomicBool::new(false),
}
}
/// Whether an event source is actually backing this watcher.
pub fn is_event_driven(&self) -> bool {
self.inner.is_some()
}
/// Resolve when the kernel reports a link change.
///
/// Never resolves when no event source is available, which makes it safe
/// to `select!` against the poll ticker: the ticker simply always wins.
pub async fn changed(&self) {
let Some(afd) = &self.inner else {
std::future::pending::<()>().await;
unreachable!("pending never resolves")
};
// Already abandoned on an earlier pass. Park without touching the
// socket, so giving up costs one syscall in total rather than one per
// caller wake-up for the life of the process.
if self.given_up.load(Ordering::Relaxed) {
std::future::pending::<()>().await;
unreachable!("pending never resolves")
}
loop {
let Ok(mut guard) = afd.readable().await else {
// The registration died. Stop firing rather than spinning; the
// caller's poll continues to cover presence.
std::future::pending::<()>().await;
unreachable!("pending never resolves")
};
// Drain to WouldBlock so a burst of link messages is one wake-up
// and the socket buffer does not fill behind us.
let mut buf = [0u8; 4096];
let mut saw_event = false;
let mut failure = None;
loop {
match guard.try_io(|inner| inner.get_ref().recv(&mut buf)) {
Ok(Ok(n)) if n > 0 => saw_event = true,
// A zero-length read. Readiness is *not* cleared by
// `try_io` here — it clears only on `WouldBlock` — so
// breaking out plainly would leave `readable()` instantly
// ready with nothing to read, and this loop would spin
// without ever returning `Pending`. That starves the
// caller's `select!` of its poll ticker entirely, which
// takes presence detection down with it. Clear it by hand
// and treat it as a fault, so the give-up path applies.
Ok(Ok(_)) => {
guard.clear_ready();
failure = Some(std::io::Error::from(std::io::ErrorKind::UnexpectedEof));
break;
}
// A genuine socket error. Distinct from WouldBlock, and
// the distinction is the whole point: `try_io` clears
// readiness only on WouldBlock, so breaking out of a real
// error leaves `readable()` instantly ready, `recv`
// failing again, and the loop spinning a core flat with
// nothing logged. Clear it by hand and back off.
Ok(Err(e)) => {
guard.clear_ready();
failure = Some(e);
break;
}
// WouldBlock — readiness is cleared, drain complete.
Err(_) => break,
}
}
if saw_event {
self.errors.store(0, Ordering::Relaxed);
return;
}
if let Some(e) = failure {
let errors = self.errors.fetch_add(1, Ordering::Relaxed) + 1;
if errors == 1 {
// ENOBUFS is the realistic one: a burst of link events
// overflowed the socket buffer, so the kernel dropped some.
// Losing events is survivable — the caller polls — but the
// spin is not, and neither is doing it silently.
warn!(error = %e, "Link event source read failed");
}
if errors >= ERROR_GIVE_UP {
self.given_up.store(true, Ordering::Relaxed);
warn!(
errors,
"Link event source is not recoverable; falling back to \
polling for interface presence"
);
std::future::pending::<()>().await;
unreachable!("pending never resolves")
}
tokio::time::sleep(Duration::from_millis(100) * errors).await;
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
/// The watcher must construct on any host, with or without a usable event
/// source, because the binder builds one unconditionally.
#[tokio::test]
async fn watcher_constructs_and_reports_its_backing() {
let w = LinkWatcher::new();
// On Linux the source is a plain `AF_NETLINK` socket in the
// `RTNLGRP_LINK` group, which needs no capability and no privilege —
// so on this platform "a sandbox might refuse it" is not a licence to
// accept either answer. Discarding the result, which this test used
// to do, meant nothing anywhere asserted that the event path exists:
// the 1 s poll is a complete fallback, so the entire suite passed with
// the source unavailable and no test could tell.
#[cfg(target_os = "linux")]
assert!(
w.is_event_driven(),
"the netlink link-event source must open on Linux; \
falling back to the poll here is a silent loss of the fast path"
);
// Elsewhere both answers are legitimate, so pin only that asking is
// safe and that a watcher with no source parks rather than fires.
#[cfg(not(target_os = "linux"))]
{
let backed = w.is_event_driven();
assert!(
backed
|| tokio::time::timeout(Duration::from_millis(50), w.changed())
.await
.is_err(),
"a watcher with no source must never resolve"
);
}
}
/// The wider egress-path mask must open too.
///
/// Same no-privilege argument as the link group above: these are all
/// read-only `NETLINK_ROUTE` multicast groups. If this one cannot bind
/// while `new()` can, the mask is wrong rather than the environment
/// restricted — and the caller that needs it would silently fall back to
/// polling.
#[cfg(target_os = "linux")]
#[tokio::test]
async fn the_egress_path_mask_opens_a_source() {
let w = LinkWatcher::with_groups(groups::EGRESS_PATH);
assert!(
w.is_event_driven(),
"the egress-path group mask must bind on Linux"
);
}
/// The mask actually reaches the socket.
///
/// `EGRESS_PATH` is a superset of `LINK`, so a watcher built on it must
/// still be a watcher — this pins that widening the mask does not make
/// the bind fail in a way `is_event_driven` would report as a missing
/// source, which is how a wrong constant would present.
#[cfg(target_os = "linux")]
#[test]
fn the_egress_path_mask_is_a_superset_of_link() {
assert_eq!(groups::EGRESS_PATH & groups::LINK, groups::LINK);
assert_ne!(groups::EGRESS_PATH, groups::LINK);
}
/// A descriptor whose `recv` always fails must not become a busy loop.
///
/// `try_io` clears readiness only on `WouldBlock`. Breaking out of a real
/// error left `readable()` instantly ready, `recv` failing again, and the
/// loop spinning a core flat with nothing logged — the realistic trigger
/// being `ENOBUFS` when a burst of link events overflows the socket
/// buffer. A pipe stands in for that here: `recv` on one answers
/// `ENOTSOCK`, every time, which is exactly the shape of a persistent
/// error.
#[tokio::test]
async fn a_persistently_failing_source_gives_up_instead_of_spinning() {
let mut fds = [0i32; 2];
assert_eq!(unsafe { libc::pipe(fds.as_mut_ptr()) }, 0, "pipe()");
let (read_fd, write_fd) = (fds[0], fds[1]);
// AsyncFd requires a non-blocking descriptor.
let flags = unsafe { libc::fcntl(read_fd, libc::F_GETFL) };
assert!(unsafe { libc::fcntl(read_fd, libc::F_SETFL, flags | libc::O_NONBLOCK) } >= 0);
let watcher = LinkWatcher {
inner: Some(AsyncFd::new(LinkEventSocket { fd: read_fd }).expect("register")),
errors: AtomicU32::new(0),
given_up: AtomicBool::new(false),
};
// Keep producing readiness edges. The error arm calls `clear_ready`,
// and a descriptor that was already readable before re-registration
// may never deliver another edge on its own — which would stall the
// loop at one error and hide whether the give-up path works. A steady
// trickle stands in for the burst of link events that provokes the
// real failure.
let writer = tokio::task::spawn_blocking(move || {
for _ in 0..200 {
if unsafe { libc::write(write_fd, b"x".as_ptr().cast(), 1) } < 0 {
break;
}
std::thread::sleep(Duration::from_millis(25));
}
unsafe { libc::close(write_fd) };
});
// Never resolves — there is no event to report — but it must reach the
// give-up state rather than burn until the timeout.
let fired = tokio::time::timeout(Duration::from_secs(5), watcher.changed()).await;
assert!(fired.is_err(), "a failing source must not report an event");
assert!(
watcher.errors.load(Ordering::Relaxed) >= ERROR_GIVE_UP,
"the error path must count, back off and stop, not spin silently"
);
writer.abort();
}
/// A watcher with no event source must never resolve, so a `select!`
/// against the poll ticker degrades cleanly instead of spinning.
#[tokio::test]
async fn a_sourceless_watcher_never_fires() {
let w = LinkWatcher {
inner: None,
errors: AtomicU32::new(0),
given_up: AtomicBool::new(false),
};
let fired = tokio::time::timeout(std::time::Duration::from_millis(50), w.changed()).await;
assert!(fired.is_err(), "sourceless watcher resolved");
}
}
+12
View File
@@ -97,6 +97,18 @@ Checks the experimental native datagram API: a client process opens a
flow to a remote pubkey over a Unix socket, receives a file descriptor,
and exchanges datagrams on it with no TUN device and no IPv6 emulation.
### [medium-change/](medium-change/) -- Transport-Medium Change
A multi-homed node whose default route moves between two live access paths
while mesh traffic is in flight, with the far peer reachable only through a
router so the path to it actually follows that default route. Asserts the
peering survives without a re-handshake (`link_id` and `authenticated_at_ms`
unchanged) and that the far side re-pins to the new source address.
Includes a negative control that runs the same move with
`node.netmon.enabled: false` and requires the outage, so a topology that
has stopped exercising the bug fails rather than passing quietly.
### [dns-resolver/](dns-resolver/) -- `fips-dns-setup` Backends
Runs `fips-dns-setup` against each supported Linux resolver backend in
+25 -1
View File
@@ -31,7 +31,7 @@
# nat-lan, nostr-publish-consume, stun-faults,
# chaos-churn-mixed-10, chaos-ethernet-mesh,
# chaos-ethernet-only, chaos-tcp-mesh, chaos-congestion-stress,
# sidecar, dns-resolver, deb-install
# sidecar, dns-resolver, deb-install, medium-change
#
# Opt-in (require --with-tor; depend on live Tor network):
# tor-socks5, tor-directory
@@ -213,6 +213,7 @@ NOSTR_RELAY_SUITES=(nostr-publish-consume)
STUN_FAULTS_SUITES=(stun-faults)
DNS_RESOLVER_SUITES=(dns-resolver)
NATIVE_API_SUITES=(native-api)
MEDIUM_CHANGE_SUITES=(medium-change)
DEB_INSTALL_SUITES=(deb-install)
TOR_SUITES=(tor-socks5 tor-directory)
@@ -267,6 +268,9 @@ list_suites() {
echo " Native API:"
for s in "${NATIVE_API_SUITES[@]}"; do echo " $s"; done
echo ""
echo " Medium change:"
for s in "${MEDIUM_CHANGE_SUITES[@]}"; do echo " $s"; done
echo ""
echo " DNS resolver:"
for s in "${DNS_RESOLVER_SUITES[@]}"; do echo " $s"; done
echo ""
@@ -1002,6 +1006,19 @@ run_native_api() {
fi
}
# Run the transport-medium change suite.
#
# Owns its own compose project and its own three bridges, so it neither
# shares container names with the NAT lab nor has to run after it.
run_medium_change() {
info "[medium-change] Running transport-medium change test"
if FIPS_TEST_IMAGE="$CI_IMAGE_TEST" bash testing/medium-change/scripts/test.sh 2>&1; then
record "medium-change" 0
else
record "medium-change" 1
fi
}
# Run dns-resolver harness (multi-distro + e2e scenarios)
run_dns_resolver() {
info "[dns-resolver] Running multi-distro test (slow — builds per-distro images)"
@@ -1165,6 +1182,11 @@ run_integration() {
run_stun_faults
done
# Transport-medium change (sequential — owns its own compose project)
for _suite in "${MEDIUM_CHANGE_SUITES[@]}"; do
run_medium_change
done
# Chaos scenarios (parallel, throttled)
if [[ "$SKIP_CHAOS" != true ]]; then
info "Running ${#CHAOS_SUITES[@]} chaos scenarios (max $PARALLEL_JOBS parallel)"
@@ -1291,6 +1313,8 @@ run_suite() {
run_dns_resolver ;;
native-api)
run_native_api ;;
medium-change)
run_medium_change ;;
deb-install)
run_deb_install ;;
tor-socks5)
+74
View File
@@ -0,0 +1,74 @@
# Transport-Medium Change Lab
A node whose network attachment moves under it — WLAN to LAN, Wi-Fi to
cellular — while its peers stay where they are.
```
node-a ──┬── mc-primary ───┐
│ ├── router ── mc-far ── node-b
└── mc-secondary ─┘
```
`node-a` is multi-homed with two equally usable paths to the router. `node-b`
sits beyond the router and is reachable **only** through it. That last part
carries the whole design: because `node-b` is off-link, the route to it
follows `node-a`'s *default* route, which is what the suite moves. Put
`node-b` on a bridge shared with `node-a` and the directly-connected route
wins, the source address never changes, and there is nothing left to test.
Both of `node-a`'s interfaces stay **up** throughout. Nothing is unplugged.
The only thing that changes is which of them the default route points at,
which is what makes this a medium change rather than a link failure — and
which is why link-state watching alone does not see it.
## Running
```bash
./testing/medium-change/scripts/test.sh
# or
./testing/ci-local.sh --only medium-change
```
## What it asserts
Traffic returning after the move is a weak signal: a peering that was torn
down by the liveness timeout and rebuilt by a re-dial also ends with traffic
flowing. The suite therefore checks *continuity*, from `fipsctl show peers` on
both nodes:
| Observation | Meaning |
| ----------- | ------- |
| `link_id` unchanged on node-a | the link was never rebuilt |
| `authenticated_at_ms` unchanged on node-a | no second handshake ran |
| `transport_addr` changed on node-b | the far side re-pinned to the new source |
| longest ping gap within budget | the data plane genuinely carried through |
The third is what stops the first two from passing vacuously on a topology
where nothing actually moved.
Phases 1 and 2 move the route in each direction, since the two are not
symmetric — one direction leaves the old interface holding an address the
routing table has abandoned, the other returns to it.
## The negative control
Phase 3 repeats the move with `node.netmon.enabled: false` and **requires**
the outage. If traffic survives with detection off, this topology is not
exercising the code path and every assertion above is vacuous — so the suite
fails rather than passing quietly.
This is deliberate. A regression test that has never been seen to fail is a
claim, not a test, and the claim is cheap to make and expensive to trust.
## Knobs
| Variable | Default | Meaning |
| -------- | ------- | ------- |
| `MC_MAX_GAP_SECS` | `5` | longest tolerated break in traffic across a move |
| `MC_CONTROL_DARK_SECS` | `12` | how long the control must stay dark |
| `MC_PRIMARY_PREFIX` | `172.31.60` | first access path `/24` |
| `MC_SECONDARY_PREFIX` | `172.31.61` | second access path `/24` |
| `MC_FAR_PREFIX` | `172.31.62` | far segment `/24` |
The gap budget sits far below the 30 s liveness timeout on purpose: a pass
must mean the move was absorbed, not that the reaper was quick.
+130
View File
@@ -0,0 +1,130 @@
# Transport-medium change lab.
#
# node-a ──┬── mc-primary ───┐
# │ ├── router ── mc-far ── node-b
# └── mc-secondary ─┘
#
# node-a is multi-homed with two equally usable paths to the router; node-b
# sits beyond it and is reachable only through the router. That last part is
# the whole design: node-b has to be off-link so the route to it follows
# node-a's *default* route, which is what the suite moves. Put node-b on a
# shared bridge instead and the directly-connected route wins, the source
# address never changes, and the bug under test cannot reproduce.
networks:
mc-primary:
driver: bridge
labels:
- "com.corganlabs.fips-ci=1"
ipam:
config:
- subnet: ${MC_PRIMARY_PREFIX:-172.31.60}.0/24
mc-secondary:
driver: bridge
labels:
- "com.corganlabs.fips-ci=1"
ipam:
config:
- subnet: ${MC_SECONDARY_PREFIX:-172.31.61}.0/24
mc-far:
driver: bridge
labels:
- "com.corganlabs.fips-ci=1"
ipam:
config:
- subnet: ${MC_FAR_PREFIX:-172.31.62}.0/24
x-fips-common: &fips-common
image: ${FIPS_TEST_IMAGE:-fips-test:latest}
cap_add:
- NET_ADMIN
devices:
- /dev/net/tun:/dev/net/tun
sysctls:
- net.ipv6.conf.all.disable_ipv6=0
restart: "no"
entrypoint:
- /usr/local/bin/mc-node-entrypoint.sh
environment:
- RUST_LOG=info,fips::node::netmon=debug,fips::node::handlers::netmon=debug
- PRIMARY_PREFIX=${MC_PRIMARY_PREFIX:-172.31.60}
- SECONDARY_PREFIX=${MC_SECONDARY_PREFIX:-172.31.61}
- FAR_PREFIX=${MC_FAR_PREFIX:-172.31.62}
- ROUTER_OCTET=254
services:
router:
build:
context: ./router
container_name: fips-mc-router${FIPS_CI_NAME_SUFFIX:-}
cap_add:
- NET_ADMIN
sysctls:
- net.ipv4.ip_forward=1
# Strict reverse-path filtering, and the suite does not work without it.
#
# It is what makes a stale source address *hurt*. Both of node-a's
# interfaces stay up and both stay routable, so a packet still sourced
# from the old path is otherwise forwarded and answered quite happily —
# the pin is stale but harmless, the bug does not bite, and the negative
# control passes, which would make every assertion in this suite vacuous.
#
# A real gateway drops that packet as spoofed, because the reverse route
# for its source points out a different interface. That is the actual
# reason a medium change black-holes traffic in the field, so it is the
# thing the lab has to model.
- net.ipv4.conf.all.rp_filter=1
- net.ipv4.conf.default.rp_filter=1
restart: "no"
networks:
mc-primary:
ipv4_address: ${MC_PRIMARY_PREFIX:-172.31.60}.254
mc-secondary:
ipv4_address: ${MC_SECONDARY_PREFIX:-172.31.61}.254
mc-far:
ipv4_address: ${MC_FAR_PREFIX:-172.31.62}.254
# The node under test. Two paths up at once, default route on the primary,
# and the suite moves it to the secondary mid-traffic.
node-a:
<<: *fips-common
container_name: fips-mc-node-a${FIPS_CI_NAME_SUFFIX:-}
hostname: fips-mc-node-a
depends_on:
- router
environment:
- RUST_LOG=info,fips::node::netmon=debug,fips::node::handlers::netmon=debug
- PRIMARY_PREFIX=${MC_PRIMARY_PREFIX:-172.31.60}
- SECONDARY_PREFIX=${MC_SECONDARY_PREFIX:-172.31.61}
- ROUTER_OCTET=254
- DEFAULT_VIA=primary
volumes:
- ../docker/resolv.conf:/etc/resolv.conf:ro
- ./node/entrypoint.sh:/usr/local/bin/mc-node-entrypoint.sh:ro
- ./generated-configs${FIPS_CI_NAME_SUFFIX:-}/node-a.yaml:/etc/fips/fips.yaml:ro
networks:
mc-primary:
ipv4_address: ${MC_PRIMARY_PREFIX:-172.31.60}.10
mc-secondary:
ipv4_address: ${MC_SECONDARY_PREFIX:-172.31.61}.10
# The far peer. Single-homed and stationary — it never moves, so anything
# the suite observes at this end is a consequence of node-a's move.
node-b:
<<: *fips-common
container_name: fips-mc-node-b${FIPS_CI_NAME_SUFFIX:-}
hostname: fips-mc-node-b
depends_on:
- router
environment:
- RUST_LOG=info
- FAR_PREFIX=${MC_FAR_PREFIX:-172.31.62}
- ROUTER_OCTET=254
- DEFAULT_VIA=far
volumes:
- ../docker/resolv.conf:/etc/resolv.conf:ro
- ./node/entrypoint.sh:/usr/local/bin/mc-node-entrypoint.sh:ro
- ./generated-configs${FIPS_CI_NAME_SUFFIX:-}/node-b.yaml:/etc/fips/fips.yaml:ro
networks:
mc-far:
ipv4_address: ${MC_FAR_PREFIX:-172.31.62}.20
+74
View File
@@ -0,0 +1,74 @@
#!/bin/bash
# Pin this node's routing before the daemon starts.
#
# Interfaces are resolved by the subnet they carry, never by name. Docker
# assigns eth0/eth1 in an order that is not the order the networks appear in
# compose, so a name-based rule silently binds the wrong path on some hosts
# and the suite then measures nothing — the default route would already be on
# the interface the test is about to "switch" to.
set -euo pipefail
WAIT_TIMEOUT_SECS="${WAIT_TIMEOUT_SECS:-30}"
# "<prefix>=<router-host-octet>" for every path this node sits on.
PRIMARY_PREFIX="${PRIMARY_PREFIX:-}"
SECONDARY_PREFIX="${SECONDARY_PREFIX:-}"
FAR_PREFIX="${FAR_PREFIX:-}"
ROUTER_OCTET="${ROUTER_OCTET:-254}"
# Which path the default route starts on: primary, secondary, or far.
DEFAULT_VIA="${DEFAULT_VIA:-primary}"
iface_for_subnet() {
local prefix="$1"
ip -4 -oneline addr show \
| awk -v p="${prefix}." '$4 ~ "^"p {print $2; exit}'
}
wait_for_subnet() {
local prefix="$1" name="$2" deadline=$((SECONDS + WAIT_TIMEOUT_SECS))
while [ "$SECONDS" -lt "$deadline" ]; do
if [ -n "$(iface_for_subnet "$prefix")" ]; then
return 0
fi
sleep 0.5
done
echo "Timed out waiting for an address on ${prefix}.0/24 (${name})" >&2
ip -4 -brief addr show >&2 || true
return 1
}
ip link set lo up
for spec in "primary:$PRIMARY_PREFIX" "secondary:$SECONDARY_PREFIX" "far:$FAR_PREFIX"; do
name="${spec%%:*}"
prefix="${spec#*:}"
[ -n "$prefix" ] || continue
wait_for_subnet "$prefix" "$name"
ip link set "$(iface_for_subnet "$prefix")" up
done
# The default route. Docker installs one of its own per attached bridge; on a
# multi-homed container which one wins is not something the suite can depend
# on, so it is replaced outright rather than adjusted.
case "$DEFAULT_VIA" in
primary) via_prefix="$PRIMARY_PREFIX" ;;
secondary) via_prefix="$SECONDARY_PREFIX" ;;
far) via_prefix="$FAR_PREFIX" ;;
*) echo "Unknown DEFAULT_VIA: $DEFAULT_VIA" >&2; exit 1 ;;
esac
via_if="$(iface_for_subnet "$via_prefix")"
# The default route, and deliberately nothing more specific.
#
# Every off-link segment in this lab is reachable through the router, so the
# default covers them all. Adding a per-subnet route as well would be worse
# than redundant: a /24 to the far segment outranks the default, so moving the
# default would leave the path to the far node exactly where it was. The suite
# would then detect a medium change, drop the sockets, and assert against a
# peer that never actually moved.
ip route replace default via "${via_prefix}.${ROUTER_OCTET}" dev "$via_if"
echo "node: addresses"
ip -4 -brief addr show | sed 's/^/ /'
echo "node: routes"
ip -4 route | sed 's/^/ /'
exec /usr/local/bin/entrypoint.sh
+11
View File
@@ -0,0 +1,11 @@
FROM debian:trixie-slim
RUN apt-get update && \
apt-get install -y --no-install-recommends \
iproute2 iputils-ping procps tcpdump && \
rm -rf /var/lib/apt/lists/*
COPY entrypoint.sh /usr/local/bin/entrypoint.sh
RUN chmod +x /usr/local/bin/entrypoint.sh
ENTRYPOINT ["/usr/local/bin/entrypoint.sh"]
+26
View File
@@ -0,0 +1,26 @@
#!/bin/sh
# Plain forwarder between the two access paths and the far segment.
#
# Deliberately no NAT and no firewall. This lab is about which *source
# address* a node picks, so anything that rewrites one would hide the very
# thing under test — a masquerading router would make both paths look
# identical to the far node and the bug would not reproduce.
#
# Forwarding is enabled by compose's `sysctls:`, not here: docker mounts
# /proc/sys read-only in an unprivileged container, so `sysctl -w` fails and,
# under `set -e`, takes the router down with it — leaving a topology that is
# wired correctly and forwards nothing.
set -eu
forwarding="$(cat /proc/sys/net/ipv4/ip_forward)"
if [ "$forwarding" != "1" ]; then
echo "router: ip_forward is '$forwarding', expected 1" >&2
echo "router: compose must set net.ipv4.ip_forward=1 for this container" >&2
exit 1
fi
echo "router: interfaces"
ip -4 -brief addr show | sed 's/^/ /'
echo "router: forwarding enabled, no NAT"
exec sleep infinity
+114
View File
@@ -0,0 +1,114 @@
#!/bin/bash
# Render the two node configs for the medium-change lab.
#
# Usage: generate-configs.sh [mesh-name] [netmon-enabled]
#
# `netmon-enabled` re-renders node-a with detection off. The test script uses
# that for its negative control: the same topology and the same move, with the
# only difference being the mechanism under test. A regression test that has
# never been seen to fail is a claim, not a test, so the suite proves the
# failure rather than asserting it from the changelog.
set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
MC_DIR="$(cd "$SCRIPT_DIR/.." && pwd)"
ROOT_DIR="$(cd "$MC_DIR/../.." && pwd)"
DERIVE_KEYS="$ROOT_DIR/testing/lib/derive_keys.py"
# Per-run, matching every other generator in the tree: compose bind-mounts
# this directory and the test script reads the npubs back after the containers
# are up, so a shared path lets a second run overwrite what a first is about
# to ping.
OUTPUT_DIR="$MC_DIR/generated-configs${FIPS_CI_NAME_SUFFIX:-}"
MESH_NAME="${1:-medium-change-$(date +%s)-$$}"
NETMON_ENABLED="${2:-true}"
primary="${MC_PRIMARY_PREFIX:-172.31.60}"
far="${MC_FAR_PREFIX:-172.31.62}"
mkdir -p "$OUTPUT_DIR"
keys_a="$(python3 "$DERIVE_KEYS" "$MESH_NAME" "a")"
keys_b="$(python3 "$DERIVE_KEYS" "$MESH_NAME" "b")"
nsec_a="$(echo "$keys_a" | awk -F= '/^nsec=/{print $2}')"
npub_a="$(echo "$keys_a" | awk -F= '/^npub=/{print $2}')"
nsec_b="$(echo "$keys_b" | awk -F= '/^nsec=/{print $2}')"
npub_b="$(echo "$keys_b" | awk -F= '/^npub=/{print $2}')"
write_config() {
local output_file="$1" nsec="$2" peer_npub="$3" peer_alias="$4" \
peer_addr="$5" netmon_block="$6"
cat > "$output_file" <<EOF
# Generated by testing/medium-change/scripts/generate-configs.sh
node:
identity:
nsec: "$nsec"
$netmon_block
retry:
max_retries: 5
base_interval_secs: 2
tun:
enabled: true
name: fips0
mtu: 1280
dns:
enabled: true
transports:
udp:
# Wildcard on purpose. A per-peer connected socket is only opened over a
# wildcard-bound transport, and that socket is what this lab exercises.
bind_addr: "0.0.0.0:2121"
mtu: 1400
peers:
- npub: "$peer_npub"
alias: "$peer_alias"
addresses:
- transport: udp
addr: "$peer_addr"
priority: 1
EOF
}
if [ "$NETMON_ENABLED" = "true" ]; then
netmon_a=$(cat <<'EOF'
netmon:
enabled: true
poll_interval_secs: 5
debounce_ms: 250
EOF
)
else
netmon_a=$(cat <<'EOF'
# Negative control: the detector is off, so nothing tells the node its
# egress path moved and the connected socket stays pinned to the old one.
netmon:
enabled: false
EOF
)
fi
# node-a dials the far node through the router, so the route to it follows
# node-a's default route — the thing the suite moves.
write_config "$OUTPUT_DIR/node-a.yaml" "$nsec_a" "$npub_b" "node-b" \
"${far}.20:2121" "$netmon_a"
# node-b is given node-a's *primary* address only. After the move that address
# is stale, which is the realistic shape: a configured peer whose path went
# away. Recovery by re-dial is therefore still possible here — and the suite
# distinguishes it from the fix by asserting session continuity, not merely
# that traffic eventually returns.
write_config "$OUTPUT_DIR/node-b.yaml" "$nsec_b" "$npub_a" "node-a" \
"${primary}.10:2121" " netmon:
enabled: true"
cat > "$OUTPUT_DIR/npubs.env" <<EOF
NPUB_A=$npub_a
NPUB_B=$npub_b
EOF
echo "Generated configs in $OUTPUT_DIR (netmon on node-a: $NETMON_ENABLED)"
+360
View File
@@ -0,0 +1,360 @@
#!/bin/bash
# Transport-medium change suite.
#
# Moves node-a's default route between two live access paths while mesh
# traffic is in flight, and asserts the peering survives it intact.
#
# What makes the assertions meaningful rather than "the mesh still works":
#
# link_id / authenticated_at_ms unchanged → no re-handshake happened
# node-b's transport_addr changed → the far side re-pinned
# ping gap bounded → the data plane really carried
#
# The first two are what separate "the fix worked" from "the liveness reaper
# tore it down and a re-dial rebuilt it", which look identical if you only
# check that traffic eventually returns.
#
# Phase 3 runs the same move with detection disabled and requires the outage,
# so the suite demonstrates the regression rather than asserting it from a
# changelog entry.
set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
MC_DIR="$(cd "$SCRIPT_DIR/.." && pwd)"
ROOT_DIR="$(cd "$MC_DIR/../.." && pwd)"
BUILD_SCRIPT="$ROOT_DIR/testing/scripts/build.sh"
GENERATE_SCRIPT="$SCRIPT_DIR/generate-configs.sh"
WAIT_LIB="$ROOT_DIR/testing/lib/wait-converge.sh"
CONFIG_DIR="$MC_DIR/generated-configs${FIPS_CI_NAME_SUFFIX:-}"
PRIMARY="${MC_PRIMARY_PREFIX:-172.31.60}"
SECONDARY="${MC_SECONDARY_PREFIX:-172.31.61}"
ROUTER_OCTET=254
NODE_A="fips-mc-node-a${FIPS_CI_NAME_SUFFIX:-}"
NODE_B="fips-mc-node-b${FIPS_CI_NAME_SUFFIX:-}"
COMPOSE=(docker compose -f "$MC_DIR/docker-compose.yml")
# Ping cadence during a move. 5/s is fast enough to resolve a sub-second gap
# without the send loop itself becoming the thing under test.
PING_INTERVAL=0.2
# The move must cost less than this. Generous next to the ~82s the unpatched
# daemon took, and still far below the 30s liveness timeout, so a pass here
# cannot be the reaper doing the work.
MAX_GAP_SECS="${MC_MAX_GAP_SECS:-5}"
# How long the negative control must stay dark before the outage is believed.
CONTROL_DARK_SECS="${MC_CONTROL_DARK_SECS:-12}"
source "$WAIT_LIB"
PASS=0
FAIL=0
# Build the test image only when nobody has handed us one.
#
# Building here is right for a hand run and wrong under a harness: when
# FIPS_TEST_IMAGE is set the caller has already built the image it named, so a
# miss means something upstream is broken and building a substitute would hide
# that behind a green run of binaries nobody asked for.
require_test_image() {
local img="${FIPS_TEST_IMAGE:-fips-test:latest}"
if docker image inspect "$img" >/dev/null 2>&1; then
echo "Using test image $img"
return 0
fi
if [ -n "${FIPS_TEST_IMAGE:-}" ]; then
echo "ERROR: $img not present, and FIPS_TEST_IMAGE names the caller's own image" >&2
echo "The harness that set it is expected to have built it." >&2
exit 1
fi
echo "$img not found; building test image"
"$BUILD_SCRIPT"
}
cleanup() {
"${COMPOSE[@]}" down -v --remove-orphans >/dev/null 2>&1 || true
}
trap 'echo ""; echo "medium-change interrupted"; cleanup; exit 130' INT TERM
ok() { echo " ✓ $1"; PASS=$((PASS + 1)); }
bad() { echo " ✗ $1"; FAIL=$((FAIL + 1)); }
dump_state() {
echo ""
echo "--- $NODE_A routes ---"
docker exec "$NODE_A" ip -4 route 2>&1 | sed 's/^/ /' || true
for c in "$NODE_A" "$NODE_B"; do
echo "--- $c: last 60 log lines ---"
docker logs "$c" 2>&1 | tail -60 | sed 's/^/ /' || true
done
}
# One field out of `fipsctl show peers`, for the single peer that is present.
# Prints nothing when the container does not answer, which every caller
# treats as a failed read rather than as a value.
peer_field() {
local container="$1" field="$2"
docker exec "$container" fipsctl show peers 2>/dev/null \
| python3 -c "
import sys, json
try:
peers = json.load(sys.stdin).get('peers', [])
except Exception:
sys.exit(0)
if peers:
v = peers[0].get('$field')
if v is not None:
print(v)
" 2>/dev/null || true
}
# Move node-a's default route to the named path, leaving both interfaces up.
#
# Interfaces are resolved by subnet for the same reason the entrypoint does
# it: docker's eth0/eth1 ordering is not the compose ordering.
move_default_route() {
local prefix="$1"
docker exec "$NODE_A" sh -c "
set -e
dev=\$(ip -4 -oneline addr show | awk -v p='${prefix}.' '\$4 ~ \"^\"p {print \$2; exit}')
test -n \"\$dev\"
ip route replace default via ${prefix}.${ROUTER_OCTET} dev \$dev
"
}
# Start a timestamped ping in the background inside node-a, writing to a file
# in the container. Returns immediately.
ping_start() {
local target="$1"
docker exec "$NODE_A" sh -c "rm -f /tmp/mc-ping.log"
docker exec -d "$NODE_A" sh -c \
"ping -D -i $PING_INTERVAL '$target' > /tmp/mc-ping.log 2>&1"
}
# Stop the ping and remember when. The stop time is what closes an outage
# that never ended — see ping_max_gap.
ping_stop() {
docker exec "$NODE_A" sh -c "pkill -f 'ping -D' || true" >/dev/null 2>&1 || true
PING_STOPPED_AT="$(date +%s.%N)"
}
# Longest interval without a successful reply, in seconds.
#
# Measured from `ping -D` timestamps rather than from the loss count, because
# loss alone cannot distinguish twenty scattered drops from one twenty-packet
# blackout — and only the second is the failure this suite is about.
#
# The observation window's end counts as a boundary. Without it an outage that
# never recovers scores *zero*: `ping -D` writes a line only for a reply, so a
# permanent blackout simply stops producing lines and the largest interval
# between two surviving replies stays one ping apart. That reads as perfect
# continuity and passes — the exact failure this suite exists to catch. The
# containers share the host's clock, so the two timebases are comparable.
ping_max_gap() {
docker exec "$NODE_A" cat /tmp/mc-ping.log 2>/dev/null \
| python3 -c "
import re, sys
end = float(sys.argv[1])
stamps = []
for line in sys.stdin:
m = re.match(r'\[(\d+\.\d+)\].*bytes from', line)
if m:
stamps.append(float(m.group(1)))
if not stamps:
print('-1')
else:
gaps = [b - a for a, b in zip(stamps, stamps[1:])]
gaps.append(end - stamps[-1])
print('%.2f' % max(gaps))
" "${PING_STOPPED_AT:-$(date +%s.%N)}"
}
ping_reply_count() {
docker exec "$NODE_A" sh -c "grep -c 'bytes from' /tmp/mc-ping.log 2>/dev/null || echo 0"
}
wait_for_mesh() {
wait_for_peers "$NODE_A" 1 60 || return 1
wait_for_peers "$NODE_B" 1 60 || return 1
}
# ── One move, fully asserted ────────────────────────────────────────────────
#
# Records the peering identity on both sides, moves the route under live
# traffic, and checks continuity against what was recorded.
assert_move_survives() {
local label="$1" to_prefix="$2" expect_src="$3"
echo ""
echo "── $label ──"
local link_before auth_before addr_before
link_before="$(peer_field "$NODE_A" link_id)"
auth_before="$(peer_field "$NODE_A" authenticated_at_ms)"
addr_before="$(peer_field "$NODE_B" transport_addr)"
if [ -z "$link_before" ] || [ -z "$auth_before" ]; then
bad "$label: could not read the peering before the move"
return
fi
echo " before: link_id=$link_before authenticated_at_ms=$auth_before"
echo " before: node-b sees node-a at $addr_before"
ping_start "${NPUB_B}.fips"
sleep 3
local baseline
baseline="$(ping_reply_count)"
if [ "$baseline" -lt 5 ]; then
bad "$label: traffic was not flowing before the move ($baseline replies)"
ping_stop
return
fi
echo " moving default route to ${to_prefix}.${ROUTER_OCTET} ..."
move_default_route "$to_prefix"
# Long enough for detection, the socket drop and the heartbeat to land,
# and for the far side to re-pin — but well short of the liveness timeout,
# so a pass cannot be the reaper's doing.
sleep 10
ping_stop
local gap replies
gap="$(ping_max_gap)"
replies="$(ping_reply_count)"
local link_after auth_after addr_after
link_after="$(peer_field "$NODE_A" link_id)"
auth_after="$(peer_field "$NODE_A" authenticated_at_ms)"
addr_after="$(peer_field "$NODE_B" transport_addr)"
echo " after: link_id=$link_after authenticated_at_ms=$auth_after"
echo " after: node-b sees node-a at $addr_after"
echo " traffic: $replies replies, longest gap ${gap}s"
# 1. The data plane carried through the move.
if [ "$gap" = "-1" ]; then
bad "$label: no replies at all — the mesh never carried traffic"
elif awk "BEGIN{exit !($gap <= $MAX_GAP_SECS)}"; then
ok "$label: traffic continuous, longest gap ${gap}s (limit ${MAX_GAP_SECS}s)"
else
bad "$label: ${gap}s outage exceeds the ${MAX_GAP_SECS}s limit"
fi
# 2. No re-handshake. This is the assertion that distinguishes the fix
# from a reconnect that merely happened fast enough.
if [ "$link_after" = "$link_before" ] && [ "$auth_after" = "$auth_before" ]; then
ok "$label: peering survived intact (same link_id, same authenticated_at_ms)"
else
bad "$label: peering was rebuilt — link_id $link_before→$link_after, authenticated_at_ms $auth_before→$auth_after"
fi
# 3. The far side actually re-pinned to the new source address. Without
# this the first two could pass on a topology where nothing moved.
if [ "$addr_after" = "${expect_src}.10:2121" ]; then
ok "$label: node-b re-pinned to ${expect_src}.10:2121"
else
bad "$label: node-b still sees node-a at ${addr_after:-<none>}, expected ${expect_src}.10:2121"
fi
}
# ── Negative control ────────────────────────────────────────────────────────
#
# The same move with detection off. The point is not to characterise the bug
# precisely, only to prove this suite can see it: if traffic survives here,
# the topology is not exercising the code path and every pass above is
# vacuous.
assert_control_fails() {
echo ""
echo "── Phase 3: negative control (netmon disabled) ──"
"$GENERATE_SCRIPT" "$MESH_NAME" false >/dev/null
"${COMPOSE[@]}" restart node-a >/dev/null 2>&1
if ! wait_for_mesh; then
bad "control: mesh did not re-form after restarting node-a with detection off"
return
fi
# Start from the primary path again, whatever the previous phase left.
move_default_route "$PRIMARY"
sleep 5
ping_start "${NPUB_B}.fips"
sleep 3
local baseline
baseline="$(ping_reply_count)"
if [ "$baseline" -lt 5 ]; then
bad "control: traffic was not flowing before the move ($baseline replies)"
ping_stop
return
fi
echo " moving default route to ${SECONDARY}.${ROUTER_OCTET} with detection off ..."
move_default_route "$SECONDARY"
sleep "$CONTROL_DARK_SECS"
ping_stop
local gap
gap="$(ping_max_gap)"
echo " traffic: longest gap ${gap}s over a ${CONTROL_DARK_SECS}s observation"
if [ "$gap" = "-1" ]; then
bad "control: no replies at all, so nothing was demonstrated"
elif awk "BEGIN{exit !($gap > $MAX_GAP_SECS)}"; then
ok "control: the move black-holed traffic for ${gap}s with detection off — the suite can see the regression"
else
bad "control: traffic survived a ${gap}s gap without detection; this topology does not exercise the bug, so the passes above prove nothing"
fi
}
main() {
echo "=============================================="
echo " FIPS transport-medium change suite"
echo "=============================================="
trap cleanup EXIT
cleanup
MESH_NAME="medium-change-$(date +%s)-$$"
echo ""
require_test_image
echo "Generating configs ..."
"$GENERATE_SCRIPT" "$MESH_NAME" true
# shellcheck disable=SC1090
source "$CONFIG_DIR/npubs.env"
echo "Starting topology ..."
"${COMPOSE[@]}" up -d
if ! wait_for_mesh; then
bad "mesh never converged"
dump_state
echo ""
echo "Result: $PASS passed, $FAIL failed"
exit 1
fi
ok "mesh converged with node-a on the primary path"
assert_move_survives "Phase 1: primary → secondary" "$SECONDARY" "$SECONDARY"
sleep 5
assert_move_survives "Phase 2: secondary → primary" "$PRIMARY" "$PRIMARY"
assert_control_fails
echo ""
echo "=============================================="
echo " Result: $PASS passed, $FAIL failed"
echo "=============================================="
if [ "$FAIL" -ne 0 ]; then
dump_state
exit 1
fi
}
main "$@"