diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 212b7a80..dbb365c7 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -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 diff --git a/.gitignore b/.gitignore index d3bbb5e6..6c3b8621 100644 --- a/.gitignore +++ b/.gitignore @@ -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*/ diff --git a/CHANGELOG.md b/CHANGELOG.md index 194f5509..1a598f33 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,8 +7,67 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] -Nothing yet. Everything previously staged here is folded into -`[0.5.1]` below. +### Fixed + +#### 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. + +### Added + +#### 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. ## [0.5.1] - 2026-09-06 diff --git a/docs/reference/configuration.md b/docs/reference/configuration.md index b18652f8..708888a5 100644 --- a/docs/reference/configuration.md +++ b/docs/reference/configuration.md @@ -194,6 +194,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. @@ -1113,6 +1160,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 diff --git a/src/config/mod.rs b/src/config/mod.rs index 67491aa6..6dcb35ea 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -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::{ diff --git a/src/config/node.rs b/src/config/node.rs index 20969af4..25e6695e 100644 --- a/src/config/node.rs +++ b/src/config/node.rs @@ -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 { @@ -1332,6 +1396,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)] @@ -1365,6 +1438,7 @@ impl Default for NodeConfig { session_mmp: SessionMmpConfig::default(), ecn: EcnConfig::default(), rekey: RekeyConfig::default(), + netmon: NetmonConfig::default(), log_level: None, } } diff --git a/src/node/dataplane/rx_loop.rs b/src/node/dataplane/rx_loop.rs index 8dee06ad..4fec454d 100644 --- a/src/node/dataplane/rx_loop.rs +++ b/src/node/dataplane/rx_loop.rs @@ -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. diff --git a/src/node/handlers/mod.rs b/src/node/handlers/mod.rs index dea42f47..607f160b 100644 --- a/src/node/handlers/mod.rs +++ b/src/node/handlers/mod.rs @@ -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 diff --git a/src/node/handlers/netmon.rs b/src/node/handlers/netmon.rs new file mode 100644 index 00000000..452b46d2 --- /dev/null +++ b/src/node/handlers/netmon.rs @@ -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 = 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 = 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" + ); +} diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index 1e40b272..1810bf76 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -1937,6 +1937,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()); @@ -2207,6 +2221,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 diff --git a/src/node/lifecycle/supervisor.rs b/src/node/lifecycle/supervisor.rs index af5197da..31a3bb76 100644 --- a/src/node/lifecycle/supervisor.rs +++ b/src/node/lifecycle/supervisor.rs @@ -731,6 +731,19 @@ pub(crate) struct Supervisor { #[cfg(unix)] pub(crate) decrypt_workers: Option, + /// 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, + /// Handle for the detector task, aborted at teardown. + pub(in crate::node) netmon_task: Option>, + /// 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(), } } diff --git a/src/node/mod.rs b/src/node/mod.rs index ec93cb04..67b2af38 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -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; diff --git a/src/node/netmon/mod.rs b/src/node/netmon/mod.rs new file mode 100644 index 00000000..1f5969da --- /dev/null +++ b/src/node/netmon/mod.rs @@ -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; +/// Sender held by a detection backend. +pub(crate) type NetChangeTx = mpsc::Sender; + +/// 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, + /// IPv6 counterpart of `v4_source`. + v6_source: Option, + /// 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, +} + +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, 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, + /// Local addresses present before but not now. + pub removed: Vec, + /// 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, + /// The preferred IPv6 source address as of this sample. + pub v6_source: Option, +} + +impl fmt::Display for NetChangeSummary { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + let mut parts: Vec = 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(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 = 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 { + 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 { + 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 { + 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 { + // 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; diff --git a/src/node/netmon/tests.rs b/src/node/netmon/tests.rs new file mode 100644 index 00000000..957123b3 --- /dev/null +++ b/src/node/netmon/tests.rs @@ -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) -> (impl Fn() -> NetFingerprint, Arc) { + 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::::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() + ); +} diff --git a/src/node/tests/mod.rs b/src/node/tests/mod.rs index dd56ed0d..7031f7f4 100644 --- a/src/node/tests/mod.rs +++ b/src/node/tests/mod.rs @@ -20,6 +20,7 @@ mod forwarding; mod handshake; mod heartbeat; mod mmp_chartests; +mod netmon; mod probe; mod routing; mod session; diff --git a/src/node/tests/netmon.rs b/src/node/tests/netmon.rs new file mode 100644 index 00000000..b76e629c --- /dev/null +++ b/src/node/tests/netmon.rs @@ -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()); +} diff --git a/src/transport/watcher.rs b/src/transport/watcher.rs index 2997b607..9b519e35 100644 --- a/src/transport/watcher.rs +++ b/src/transport/watcher.rs @@ -6,7 +6,7 @@ //! //! | Platform | Source | //! | -------- | ------ | -//! | Linux | netlink `RTNLGRP_LINK` (`RTM_NEWLINK` / `RTM_DELLINK`) | +//! | 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 | //! @@ -67,13 +67,43 @@ impl LinkEventSocket { } } -/// Open the platform's link-event socket, non-blocking. -#[cfg(target_os = "linux")] -fn open_link_socket() -> std::io::Result { - // RTMGRP_LINK. Spelled as a literal because the constant's name and - // availability differ across libc versions; the value is ABI. - const RTMGRP_LINK: u32 = 1; +/// 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 { let fd = unsafe { libc::socket( libc::AF_NETLINK, @@ -88,7 +118,7 @@ fn open_link_socket() -> std::io::Result { let mut sa: libc::sockaddr_nl = unsafe { std::mem::zeroed() }; sa.nl_family = libc::AF_NETLINK as u16; - sa.nl_groups = RTMGRP_LINK; + sa.nl_groups = groups; let ret = unsafe { libc::bind( fd, @@ -103,8 +133,8 @@ fn open_link_socket() -> std::io::Result { } /// Open the platform's link-event socket, non-blocking. -#[cfg(not(target_os = "linux"))] -fn open_link_socket() -> std::io::Result { +#[cfg(not(any(target_os = "linux", target_os = "android")))] +fn open_link_socket(_groups: u32) -> std::io::Result { // 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) }; @@ -148,9 +178,27 @@ impl Default for LinkWatcher { } impl LinkWatcher { - /// Open the platform link-event source, falling back to nothing. + /// 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 { - let inner = match open_link_socket() { + #[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) => { @@ -307,6 +355,36 @@ mod tests { } } + /// 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 diff --git a/testing/README.md b/testing/README.md index 9658b5c6..03ff2384 100644 --- a/testing/README.md +++ b/testing/README.md @@ -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 diff --git a/testing/ci-local.sh b/testing/ci-local.sh index 9b097be5..b2bc39cb 100755 --- a/testing/ci-local.sh +++ b/testing/ci-local.sh @@ -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 @@ -199,6 +199,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) @@ -253,6 +254,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 "" @@ -988,6 +992,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)" @@ -1151,6 +1168,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)" @@ -1277,6 +1299,8 @@ run_suite() { run_dns_resolver ;; native-api) run_native_api ;; + medium-change) + run_medium_change ;; deb-install) run_deb_install ;; tor-socks5) diff --git a/testing/medium-change/README.md b/testing/medium-change/README.md new file mode 100644 index 00000000..6e11be4e --- /dev/null +++ b/testing/medium-change/README.md @@ -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. diff --git a/testing/medium-change/docker-compose.yml b/testing/medium-change/docker-compose.yml new file mode 100644 index 00000000..e68a0cb9 --- /dev/null +++ b/testing/medium-change/docker-compose.yml @@ -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 diff --git a/testing/medium-change/node/entrypoint.sh b/testing/medium-change/node/entrypoint.sh new file mode 100755 index 00000000..a7184565 --- /dev/null +++ b/testing/medium-change/node/entrypoint.sh @@ -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}" +# "=" 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 diff --git a/testing/medium-change/router/Dockerfile b/testing/medium-change/router/Dockerfile new file mode 100644 index 00000000..1b197c44 --- /dev/null +++ b/testing/medium-change/router/Dockerfile @@ -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"] diff --git a/testing/medium-change/router/entrypoint.sh b/testing/medium-change/router/entrypoint.sh new file mode 100755 index 00000000..73121cdb --- /dev/null +++ b/testing/medium-change/router/entrypoint.sh @@ -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 diff --git a/testing/medium-change/scripts/generate-configs.sh b/testing/medium-change/scripts/generate-configs.sh new file mode 100755 index 00000000..5c493f2f --- /dev/null +++ b/testing/medium-change/scripts/generate-configs.sh @@ -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" < "$OUTPUT_DIR/npubs.env" </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:-}, 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 "$@"