From 957e05b28328a7ed4dc6c5800fa61798d38597e6 Mon Sep 17 00:00:00 2001 From: fr34aky <162515565+fr34aky@users.noreply.github.com> Date: Wed, 9 Sep 2026 14:48:22 +0000 Subject: [PATCH] feat(netmon): fingerprint the path to each peer, not the host MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The medium-change detector sampled two host-wide signals: the source address the routing table would pick for an off-link destination, and the set of up, non-loopback interface addresses. The second one was the problem. It enumerated every address the host had, so a docker bridge coming up, a VPN connecting or a container network appearing moved the fingerprint with no peering affected at all — and the reaction to a moved fingerprint is to drop every connected UDP socket and heartbeat every peer. Self-healing, so it cost work rather than connectivity, but on any host running containers it could fire repeatedly for nothing. Ask the question per peer instead. For each peer whose transport address is a numeric IP endpoint, `connect(2)` a UDP socket to it and read back the local address — the same no-packets operation, aimed at the peers we actually hold rather than at a documentation prefix. The fingerprint becomes the set of local addresses the kernel would use to reach our peers. That is the quantity the reaction cares about. The stale `connect(2)` this subsystem exists to repair pinned a local source address chosen for one destination, so measuring the same thing for the same destinations asks the kernel the question the bug is about rather than a proxy for it. Three things follow: - Interfaces the node does not peer over cannot move it, by construction rather than by a filter guessing which interface names are infrastructure. `docker compose up` moves nothing. - On-link peers become visible. A peer on the same LAN is reached by its subnet route, and the old probe followed the *default* route by construction, so it looked straight past that path. - A more specific route moving under one peer is representable at all, which no single host-wide sample could be. Samples are compared over the *intersection* of their peer sets, never the union, so peers joining and leaving cannot fire the fan-out on their own. The sample is still adopted on the no-change path, or `last` would freeze on the peer set the detector started with. A peer whose probe stops answering is a move to "no route" and does count: that peer is exactly the one now stranded. **A peer's first sample is judged against its socket, not against nothing.** The intersection rule skips a peer present in only one sample, which is right for peer churn and wrong for the sample in which a peer first appears, because that sample may already be the post-change one. `last` gains a peer only at the first wake after it shows up in the entity snapshot, so a medium change inside that window is consumed rather than delayed: the peer's connected socket stays pinned to the path the host has just left, and the peering black-holes until the liveness timeout tears it down. A first-seen peer is therefore compared against the source its own connected socket is bound to, where it has one. One residual on that rule, stated exactly rather than understated: the tick publishes the entity snapshot *before* it installs connected sockets, so a socket installed on tick N is first visible on tick N+1, and a peer whose path moves inside that window is first seen with `bound` still `None` while genuinely holding a pinned socket. About one `tick_interval_secs` per join. A hole, not a harmless skip. **The probe binds the way the send path binds.** `open_connected_fd` binds `local_addr` verbatim before connecting, so the socket keeps a configured address whatever the route says, while the probe took the kernel's choice. Under a non-wildcard `bind_addr` the two answered different questions and every first-seen peer reported a phantom move. The probe now binds what the transport binds, address only and port 0; under the default wildcard bind nothing changes. Operator-facing corrections in the same surface. `PeerSourceMove::before` was recorded on every move and then wildcarded away by the only thing that read it, so the log said where a peer moved to but not where from — and the address it moved *from* is the one the stale `connect(2)` had pinned. Both ends are rendered now. The peer id used a private four-byte hex helper that duplicated `NodeAddr::short_hex` and dropped the `...` suffix every other operator surface prints; the duplicate is gone. The cost note said three syscalls per target: `UdpSocket::bind` is a `socket(2)` and a `bind(2)`, and the socket takes a `close(2)` on drop, so it is five. In the same place, "bounded by `node.limits.max_peers`" does not hold when that value is 0, which the configuration defines as unlimited. This deletes `interface_addrs()`'s only call site, and with it the `getifaddrs` walk and its `sockaddr` decoding. It therefore absorbs the Android `getifaddrs` issue rather than leaving it to be fixed separately. The peer table is reached through `entities_snapshot` rather than a new sharing primitive, so the detector stays a detached task holding no node state and taking no node lock. `PeerRow` gains a typed `probe_target` rather than having the detector re-parse the display string next to it, so a change to that string's rendering cannot silently leave the detector with an empty table and no way to notice. Peers that are not probeable IP destinations contribute nothing and need no per-transport special-casing here: a MAC on Ethernet or BLE, a .onion or Nym recipient behind a local proxy, a scoped IPv6 literal, and a peer still carrying its configured hostname all arrive as `None`. The last is deliberate — resolving one would put a DNS lookup with its timeouts on the sample path — and the window is small, since the address is replaced by the observed numeric source the first time an authenticated packet arrives. A node with no peers detects nothing, which is right: nothing is bound to the old path. Also corrects two doc claims that did not match the code, both in the text being rewritten: `transport::watcher` has no consumer besides this module, so it is not "shared with the interface binder"; and a connection-oriented transport does not "re-dial on send" in the case that matters, because `send_async` only dials when the pool holds no connection for the address and a connection stranded by a medium change is still in the pool — it is evicted after a write to it fails. Tests, each run against the defect it guards rather than only against the fix: - The churn rules are mutation-checked: iterating the union instead of the intersection fails four tests, and dropping the sample-adoption fails the one that pins a newly joined peer entering the comparison. - The snapshot seam is table-driven over six address shapes and fails if the publish site stops populating `probe_target` — nothing renders that field, so nothing else would have caught it. - A namespace test brings up a dummy interface with its own subnet and asserts the fingerprint does not move, then puts a more specific route to the peer out of that same interface and asserts it does, so the negative half cannot pass because sampling quietly stopped working. The same test now asserts that an unconstrained probe answers with the carrier and a constrained one with the other address, and that the two differ — the disagreement that would otherwise report a phantom move on every first-seen peer. Ignoring the constraint in `preferred_source` fails it. - The two links carrying the pinned source were asserted by nothing. Substituting `None` where `probe_targets` reads the row, or where `sample` writes the fingerprint, left the whole suite green while silently restoring the bug the first-sight rule exists to fix. Both are asserted now, and both mutations fail. - Every live-probe test passed if `preferred_source` returned `None` for everything: two all-`None` samples are self-consistent, the recorded-keys test never inspected a value, and the loopback assertion skipped through its `if let`. A probe to loopback must now answer with loopback, which holds on any host that can run the suite, including one started with `--network none`. - `reports_are_spaced_out_under_clean_flapping` polled at one second against a one-second pacing floor, so the two were indistinguishable and deleting the pacing block still passed. The wake is 100ms now. - The pinned-source publish test was Linux-gated though `open_connected_fd` and the field it asserts are available on macOS too. Widened, along with the two sibling tests on the same helper. - The detector's netlink subscription asserts the route groups rather than logging a decline, so a wrong group mask reds instead of passing quietly. --- CHANGELOG.md | 33 +- docs/reference/configuration.md | 41 +- src/control/snapshot.rs | 38 + src/node/handlers/netmon.rs | 30 +- src/node/lifecycle/mod.rs | 3 +- src/node/mod.rs | 17 + src/node/netmon/mod.rs | 649 +++++++++++------ src/node/netmon/tests.rs | 845 +++++++++++++++++++++-- src/node/tests/netmon.rs | 190 ++++- src/transport/udp/io/connected/socket.rs | 53 +- src/transport/watcher.rs | 40 ++ 11 files changed, 1632 insertions(+), 307 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6da2a06d..e5ee08bd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -67,13 +67,32 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 #### 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` + block (on by default). For each peer whose transport address is a numeric IP + endpoint, the node asks the kernel which local address it would use to reach + *that peer* — a `connect(2)` on a UDP socket, which resolves the route and + sends nothing — and reports a change once some peer held across two + consecutive samples is reached from a different local address, or has stopped + being reachable at all. Asking the question per peer rather than about the + host is what keeps it quiet: a container bridge, a VPN, a `veth` pair or a + tunnel appearing is not the route to any peer and cannot move the + fingerprint, while a peer on the same LAN — reached by its subnet route, not + the default route — is covered, as is a more specific route moving under a + single peer. Peers joining and leaving are ignored on their own, being + ordinary node behaviour rather than a statement about the medium — except + that a peer seen for the first time is checked against its own + `connect()`-ed socket, and reported if that socket is pinned to a source the + routing table would no longer choose, so a medium change in the window + between a peer authenticating and the next sample is not adopted silently + while that peer sits stranded on the old path. A peer + addressed by MAC, by `.onion` or Nym recipient, by a scoped IPv6 literal, or + by a hostname it has not yet been heard from on, has no route to ask about + and contributes nothing; a node with no peers detects nothing, having nothing + bound to the old path to repair. The peer table is read through the node's + existing lock-free entity snapshot, so the detector stays a detached task + holding no node state. A handover is not atomic (the route goes, briefly + there is none, 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 diff --git a/docs/reference/configuration.md b/docs/reference/configuration.md index 708888a5..29d40274 100644 --- a/docs/reference/configuration.md +++ b/docs/reference/configuration.md @@ -202,7 +202,7 @@ 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.poll_interval_secs` | u64 | `5` | How often the path to each peer 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 @@ -219,9 +219,42 @@ 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. +stream the medium change has very likely just stranded. No peering is torn +down: sessions, tree positions and routes survive the switch. + +**What counts as a change.** For each peer whose transport address is a numeric +IP endpoint, the node asks the kernel which local address it would use to reach +*that peer* — a `connect(2)` on a UDP socket, which resolves the route and sends +nothing. A change is reported when a peer present in two consecutive samples is +now reached from a different local address, or has stopped being reachable at +all. + +Because the question is asked per peer, an interface the node does not peer over +cannot trigger anything: a container bridge, a VPN, a `veth` pair or a tunnel +appearing is not the route to any peer, so it does not enter the sample. A peer +on the same LAN, reached by its subnet route rather than the default route, is +covered as well as one across the internet, and so is a more specific route +moving under a single peer. + +Peers appearing and leaving are ignored on their own — that is ordinary node +behaviour and says nothing about the medium. A peer seen for the first time is +the one exception, and it is not judged against history but against its own +send path: if its `connect()`-ed socket is pinned to a source the routing table +would no longer choose, it is reported. Without that, a medium change in the +window between a peer authenticating and the detector's next sample would be +the detector's first sight of that peer, and would be adopted silently while +the peer's socket stayed pinned to the path the host had just left. A peer +joining onto a path that has not moved has its socket pinned exactly where its +traffic goes, so it still reports nothing. A peer whose address is not a +probeable IP endpoint contributes nothing: a MAC on Ethernet or BLE, a `.onion` +or Nym recipient reached through a local proxy, an IPv6 literal with a scope +suffix, or a peer still carrying the hostname it was configured with (resolving +one would put a DNS lookup on the sample path; the address becomes numeric as +soon as an authenticated packet arrives from the peer). A node holding no peers +detects nothing, which is correct — it has nothing bound to the old path. + +The cost is three syscalls per peer per sample, bounded by +`node.limits.max_peers`, with no packets sent and no name resolution. Detection uses the best backend the platform has: diff --git a/src/control/snapshot.rs b/src/control/snapshot.rs index bea9d8d7..bf7c4074 100644 --- a/src/control/snapshot.rs +++ b/src/control/snapshot.rs @@ -18,6 +18,7 @@ //! also advances only on the tick. use std::collections::HashMap; +use std::net::{IpAddr, SocketAddr}; use std::sync::Arc; use crate::identity::NodeAddr; @@ -629,6 +630,43 @@ pub(crate) struct PeerRow { pub is_parent: bool, pub is_child: bool, pub transport_addr: Option, + /// The peer's current transport address as a numeric IP endpoint, when it + /// is one. Not rendered anywhere: this is the medium-change detector's + /// read of the peer table (see [`crate::node::netmon`]), carried here + /// because the detector is a detached task and this snapshot is the + /// node's existing lock-free read side. + /// + /// `None` covers everything that is not a probeable IP destination — a + /// MAC on Ethernet or BLE, a `.onion` or Nym recipient, a peer still + /// carrying the hostname it was configured with, an IPv6 literal with a + /// scope suffix. Typed rather than re-parsed from `transport_addr` above + /// so a change to that string's rendering cannot silently leave the + /// detector with nothing to probe. + pub probe_target: Option, + /// Source address this peer's per-peer `connect()`-ed UDP socket was + /// pinned to by `connect(2)`, when it has one. Also not rendered, and read by the same detector: + /// it is what the send path is *actually* using, as against the + /// `probe_target` lookup's answer for what the kernel would choose now. + /// + /// `None` where there is no such socket — every platform but Linux and + /// macOS, a peer on another transport, and a peer whose socket has not + /// been installed yet or was just released — and also where the kernel + /// declined to name a source, which is not an address and must not be + /// compared as one. + pub bound_source: Option, + /// Address this peer's transport is bound to, when that bind is not the + /// wildcard. Read by the same detector, which has to put its probe the + /// same constrained question the send path answers. + /// + /// `open_connected_fd` binds the transport's configured address verbatim + /// and only then connects, so a non-wildcard `bind_addr` pins the source + /// whatever the routing table says, while an unconstrained probe takes the + /// kernel's choice. Left unequal, those two answers differ permanently and + /// every first-seen peer reports a move that never happened. + /// + /// `None` for the wildcard bind, which is the default and the case where + /// the kernel chooses on both sides. + pub probe_bind: Option, pub link_info: Option, pub tree_depth: Option, /// `effective_depth = tree_depth + link_cost` — the same quantity diff --git a/src/node/handlers/netmon.rs b/src/node/handlers/netmon.rs index 89ae4e0e..cddd1326 100644 --- a/src/node/handlers/netmon.rs +++ b/src/node/handlers/netmon.rs @@ -52,6 +52,23 @@ //! 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. +//! +//! # Why the reaction is still node-wide +//! +//! [`NetChange`] now names the peers whose local source address moved, because +//! the fingerprint is keyed on them. The reaction deliberately does not use +//! that yet: it drops every connected socket and heartbeats every +//! connectionless peer, exactly as it did when the detector could only say +//! "something about this host moved". +//! +//! That is over-broad and known to be. It is left node-wide here because +//! narrowing it is a behavioural change with its own failure mode — a peer +//! left un-rebound because it was absent from the moved set is stranded for +//! `link_dead_timeout_secs`, which is the bug this subsystem exists to close — +//! and it wants its own tests rather than a free ride on a change to the +//! fingerprint. The cost of staying broad is now small: a change is only +//! reported when a peering's own path moved, so the fan-out no longer fires +//! for a container bridge appearing. use std::time::Instant; @@ -64,6 +81,9 @@ use crate::proto::link::LinkMessageType; impl Node { /// React to a settled transport-medium change. + /// + /// `change.summary` names the peers that moved; see the module docs for + /// why the reaction is node-wide regardless. 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 @@ -132,9 +152,13 @@ impl Node { /// 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 + /// before this detector existed, and `link_dead_timeout_secs` remains the + /// backstop. Note that it does *not* recover by redialling: `send_async` + /// only dials when the pool holds no connection for the address, and a + /// connection stranded by a medium change is still in the pool. It is + /// evicted after a write to it fails, so the redial happens on the send + /// after the failure, not on the first one. 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 { diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index 1810bf76..41a3a67f 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -1944,7 +1944,8 @@ impl Node { // not health. let netmon_cfg = self.config().node.netmon.clone(); if netmon_cfg.enabled { - let (rx, task) = crate::node::netmon::spawn_detector(netmon_cfg); + let (rx, task) = + crate::node::netmon::spawn_detector(netmon_cfg, self.entities_snapshot.clone()); self.supervisor.netmon_rx = Some(rx); self.supervisor.netmon_task = Some(task); } else { diff --git a/src/node/mod.rs b/src/node/mod.rs index 67b2af38..adc7dc84 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -2204,6 +2204,23 @@ impl Node { is_parent, is_child, transport_addr: peer.current_addr().map(|a| format!("{}", a)), + probe_target: peer + .current_addr() + .and_then(|a| a.as_str()) + .and_then(|s| s.parse::().ok()), + #[cfg(any(target_os = "linux", target_os = "macos"))] + bound_source: peer.connected_udp().and_then(|s| s.pinned_source()), + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + bound_source: None, + probe_bind: peer + .transport_id() + .and_then(|id| self.transports.get(&id)) + .and_then(|t| match t { + TransportHandle::Udp(u) => u.local_addr(), + _ => None, + }) + .map(|sa| sa.ip()) + .filter(|ip| !ip.is_unspecified()), link_info, tree_depth: peer.coords().map(|c| c.depth()), effective_depth, diff --git a/src/node/netmon/mod.rs b/src/node/netmon/mod.rs index 1f5969da..7a02f36e 100644 --- a/src/node/netmon/mod.rs +++ b/src/node/netmon/mod.rs @@ -27,13 +27,14 @@ //! | 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. +//! Both kernel sources are [`crate::transport::watcher::LinkWatcher`], which +//! this module is currently the only consumer of. It is asked for +//! [`groups::EGRESS_PATH`](crate::transport::watcher::groups::EGRESS_PATH) +//! rather than the watcher's default link-presence mask: a 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 @@ -47,42 +48,99 @@ //! //! # What the fingerprint captures //! -//! Two independent signals, because neither alone is sufficient: +//! One local source address per peer: for every peer whose transport address is +//! a numeric IP endpoint, the address the kernel would pick to reach *that +//! peer*. A connected-but-never-sending UDP socket makes the kernel run its +//! route lookup and bind the source address it would use; three syscalls, no +//! packets, no name resolution, and it works identically on every platform std +//! supports. //! -//! - **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. +//! That set is exactly the quantity the reaction cares about. The stale +//! `connect(2)` this whole subsystem exists to repair pinned a *local source +//! address chosen for one destination*, so measuring the same thing for the +//! same destinations asks the kernel the question the bug is about, rather +//! than a proxy for it. //! -//! 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. +//! Two consequences fall out of aiming the probe at peers rather than at the +//! host: //! -//! 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. +//! - **Interfaces the node does not peer over cannot move it.** A container +//! bridge, a VPN coming up, a `veth` pair, a tunnel — none of them is the +//! route to any peer, so none of them enters the fingerprint. This is by +//! construction rather than by a filter that has to keep guessing which +//! interface names are infrastructure, which is what the host-wide address +//! set this replaced could never get right: on a host running containers it +//! moved, and the node dropped every connected socket and heartbeated every +//! peer, for a `docker compose up`. +//! - **On-link peers are visible.** A peer on the same LAN is reached by its +//! subnet route, not the default route, so a host-wide probe at an off-link +//! destination looked straight past it: unplug the LAN cable on a host whose +//! default route is cellular and nothing host-wide moved while every socket +//! to that peer was stranded. Its own probe follows its own route and moves. +//! +//! It is also the right granularity for a *more specific* route changing under +//! one peer while the rest of the host is untouched, which no single host-wide +//! sample can represent at all. //! //! # 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 +//! **A peer whose address is not a probeable IP destination** contributes +//! nothing — an Ethernet or BLE peer addressed by MAC, a `.onion` or a Nym +//! recipient reached through a local proxy, an IPv6 literal carrying a scope +//! suffix. None of them is IP-attached in the way this detector reasons about, +//! and the connected-UDP pinning it repairs cannot happen to them. +//! +//! **A peer still carrying the hostname it was configured with**, because +//! resolving one on the sample path would put a DNS lookup, with its timeouts, +//! inside the detector's tick. In practice the window is small: the address is +//! replaced by the observed numeric source the first time an authenticated +//! packet arrives (`dataplane::encrypted`), and a peer that has never been +//! heard from has no established peering to strand. +//! +//! **A node with no peers** has an empty fingerprint and detects nothing, which +//! is correct — there is nothing bound to the old path to repair. +//! +//! **A BLE adapter's state** is invisible here, as it was before: 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; +//! +//! # Where the peer list comes from +//! +//! [`crate::control::snapshot::EntitySnapshot`], the node's existing lock-free +//! read side, republished from the tick. The detector is deliberately a +//! detached task holding no node state and taking no node lock, so it reads +//! the peer table the same way the off-loop `show_peers` renderer does. The +//! view is at most one `node.tick_interval_secs` stale, which does not matter: +//! a peer that has just appeared is absorbed on the next sample (see below), +//! and one that has just left is dropped from the comparison rather than +//! reported. +//! +//! # Comparing two samples +//! +//! Over the **intersection** of the two peer sets, never their union: a change +//! is reported when some peer present in both samples is now reached from a +//! different local address. Peers joining and leaving is ordinary node +//! behaviour and says nothing about the medium, so on its own it must not fire +//! a reaction that drops every connected socket. A peer whose probe stops +//! answering entirely — the route to it is gone — is a move to "no source +//! address" and does count, because that peer is exactly the one now stranded. +//! +//! A peer seen for the *first* time is the one case the intersection cannot +//! decide, and it cannot simply be skipped: the sample in which a peer first +//! appears may already be the post-change one, and adopting it silently would +//! swallow the event while that peer's socket stayed pinned to the path the +//! host has just left. Such a peer is judged against its own send path +//! instead — the source `connect(2)` pinned its socket to, which needs no +//! history and asks directly whether that socket is already stale. Churn still +//! fires nothing on its own, because a peer joining onto a path that has not +//! moved is pinned exactly where its traffic goes. See +//! [`NetFingerprint::moved`] for the residual this leaves. +//! +use std::collections::BTreeMap; use std::fmt; use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr, UdpSocket}; +use std::sync::Arc; use std::time::Duration; #[cfg(unix)] @@ -92,18 +150,8 @@ 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, -); +use crate::control::snapshot::EntitySnapshot; +use crate::identity::NodeAddr; /// How many resample rounds the debounce will ride out before reporting /// anyway. A handover emits a burst (address gone, address added, route @@ -131,117 +179,272 @@ 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. +/// Where this host sits relative to the peers it holds: one local source +/// address per peer, as the routing table would choose it right now. /// -/// 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. +/// Each peer carries both the answer to that question and the source its +/// connected socket is already pinned to, because a peer seen for the first +/// time has no earlier sample to be compared against and is judged against its +/// own socket instead. [`NetFingerprint::moved`] is the whole definition of +/// "the medium changed" and the only thing the detector asks of a sample. #[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, + /// Peer → where its traffic leaves from. A peer with no probeable address + /// never appears at all. + sources: BTreeMap, +} + +/// One peer's local end, from two directions. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct PeerPath { + /// The local address the kernel would choose to reach this peer right now. + /// `None` when the route lookup fails, which is a real value rather than a + /// missing one: the peer is still ours, and having no route to it is + /// precisely the state worth reacting to. + current: Option, + /// The local address this peer's connected UDP socket is bound to, if it + /// has one. Unlike `current` this is not a question put to the kernel — it + /// is what the send path is already doing, and it is the only thing that + /// gives a peer the detector has not seen before a baseline to be judged + /// against. See [`NetFingerprint::moved`]. + bound: Option, } impl NetFingerprint { - /// Sample the host's current attachment. + /// Probe every target and record the local address the kernel picks. /// - /// 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 { + /// Five non-blocking syscalls per target: `socket(2)` and `bind(2)` behind + /// `UdpSocket::bind`, a `connect(2)` that sends no packet, a + /// `getsockname(2)`, and the `close(2)` the socket takes on drop. No I/O + /// wait, no name resolution, and no allocation beyond the map. + /// + /// The count matters because a debounced handover resamples: up to + /// `MAX_DEBOUNCE_ROUNDS` rounds plus the settled sample, times the peers + /// held. `node.limits.max_peers` bounds that only where it is set — + /// the value 0 means unlimited, and there the cost tracks the live peer + /// count instead. It runs inline in the detector's own task rather than + /// through `spawn_blocking`, which is what keeps it off every other task + /// regardless. + pub(in crate::node) fn sample(targets: &[ProbeTarget]) -> Self { Self { - v4_source: preferred_source(PROBE_V4), - v6_source: preferred_source(PROBE_V6), - local_addrs: interface_addrs(), + sources: targets + .iter() + .map(|t| { + ( + t.peer, + PeerPath { + current: preferred_source(t.dest, t.bind), + bound: t.bound, + }, + ) + }) + .collect(), } } /// Build a fingerprint directly, so a test can script a sequence of - /// samples instead of reading the host's real attachment. + /// samples instead of probing real peers. No peer has a connected socket; + /// [`Self::for_test_bound`] is the variant that gives one. #[cfg(test)] - pub(crate) fn for_test(v4_source: Option, local_addrs: &[IpAddr]) -> Self { + pub(crate) fn for_test(sources: &[(NodeAddr, Option)]) -> Self { Self { - v4_source, - v6_source: None, - local_addrs: local_addrs.iter().copied().collect(), + sources: sources + .iter() + .map(|(peer, current)| { + ( + *peer, + PeerPath { + current: *current, + bound: None, + }, + ) + }) + .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() + /// As [`Self::for_test`], with each peer's connected socket bound where the + /// third element says. + #[cfg(test)] + pub(crate) fn for_test_bound(sources: &[(NodeAddr, Option, Option)]) -> Self { + Self { + sources: sources + .iter() + .map(|(peer, current, bound)| { + ( + *peer, + PeerPath { + current: *current, + bound: *bound, + }, + ) + }) .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, } } + + /// Which peers are now leaving from somewhere other than where their + /// traffic is actually going out. + /// + /// Two rules, because there are two ways to know: + /// + /// **A peer in both samples** is judged on whether its probe answer moved. + /// The comparison is over the *intersection* of the two peer sets, never + /// the union: a peer that has only just been authenticated, or one that has + /// just been reaped, differs between the samples for reasons that have + /// nothing to do with the host's attachment, and the reaction — drop every + /// connected socket, heartbeat every peer — is far too blunt to fire on + /// ordinary peer churn. + /// + /// **A peer only in the newer sample** has no previous probe answer to be + /// compared against, and skipping it outright leaves a hole this detector + /// cannot afford. `last` gains a peer only at the first wake *after* it + /// appears, so a medium change in that window is the detector's first + /// sight of that peer, and adopting it silently would swallow the very + /// event being adopted — while the peer's connected socket stays pinned to + /// the path the host has just left. The window is up to one + /// `poll_interval_secs` after every peer that authenticates, and a medium + /// change wakes the detector, so the two coincide readily rather than + /// rarely. + /// + /// So such a peer is judged against `bound` instead: the address its + /// connected socket is *actually* using. That needs no history — it asks + /// whether the send path is already stale, which is the question the whole + /// subsystem exists to answer, and it is exactly the peer that would + /// otherwise be left stranded. Churn still cannot fire anything on its own: + /// a peer joining onto a path that has not moved has `bound == current` and + /// reports nothing. + /// + /// A first-seen peer with no `bound` is still skipped, and that is the + /// residual. It covers three groups, and they are not equally harmless. + /// + /// Where there is genuinely no connected socket — every platform but Linux + /// and macOS, and every peer on a stream or proxied transport on those two + /// — nothing is pinned to repair, because the wildcard socket resolves a + /// route per packet. Such a peer is not stranded; it loses only the + /// immediate heartbeat that would have told the far side to re-pin, and + /// notices at its next `heartbeat_interval_secs`. + /// + /// The third group is a real hole rather than a harmless one, and it is one + /// tick wide. The tick publishes the entity snapshot before it installs + /// connected sockets (`record_stats_history` then + /// `activate_connected_udp_sessions`, in that order), so a socket installed + /// on tick N is first visible to this detector in the snapshot published on + /// tick N+1. A peer that joins and has its socket installed, and whose path + /// then moves before that next publish, is first seen with `bound` still + /// `None` and is skipped — and it *does* hold a pinned socket. About one + /// `tick_interval_secs` per join against a `poll_interval_secs` five times + /// longer, and the peer is recovered by the ordinary diff on the sample + /// after, so it is bounded rather than permanent. + /// + /// A non-wildcard `transports.udp.bind_addr` is not part of the residual; + /// see [`ProbeTarget::bind`], which keeps both sides answering the same + /// question rather than skipping the peer. + /// + /// An empty result means nothing moved; it is the detector's entire + /// definition of "no change". + pub(in crate::node) fn moved(&self, next: &Self) -> Vec { + next.sources + .iter() + .filter_map(|(peer, now)| { + let before = match self.sources.get(peer) { + // Seen before: its own previous probe answer. + Some(then) => then.current, + // First sight: what its socket is bound to, if it has one. + None => match now.bound { + Some(bound) => Some(bound), + None => return None, + }, + }; + (before != now.current).then_some(PeerSourceMove { + peer: *peer, + before, + after: now.current, + }) + }) + .collect() + } } -/// 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. +/// One peer to probe: where to aim, and what its send path is already using. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(in crate::node) struct ProbeTarget { + /// The peer this is about. + pub peer: NodeAddr, + /// Its transport address — the destination the route lookup is run for. + pub dest: SocketAddr, + /// The local address its connected UDP socket is bound to, if it has one. + pub bound: Option, + /// The address to bind the probe to, when the peer's transport binds a + /// specific one rather than the wildcard. `None` means bind unspecified + /// and let the kernel choose, which is the default posture. + /// + /// This exists so the probe asks the same question the send path answers. + /// `open_connected_fd` binds the transport's configured address verbatim + /// and only then connects, so under a non-wildcard `transports.udp.bind_addr` + /// the socket's source is that address whatever the routing table says. An + /// unconstrained probe would answer with the kernel's choice instead, and + /// the two would disagree permanently — reporting a first-sight move, on + /// every peer, forever, with nothing having moved. + pub bind: Option, +} + +/// One peer whose local source address changed between two samples. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) struct PeerSourceMove { + /// The peer that is now reached from somewhere else. + pub peer: NodeAddr, + /// The local address it was reached from, or `None` if there was no route. + pub before: Option, + /// The local address it is reached from now, or `None` if the route is gone. + pub after: Option, +} + +/// What moved between two fingerprints. +/// +/// Operator-facing, and — unlike the host-wide summary this replaced — it now +/// names the peers affected, because the fingerprint is keyed on them. The +/// handler still re-evaluates every peer regardless; narrowing the reaction to +/// exactly [`Self::moved`] is a separate change. #[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, + /// Every peer whose local source address changed, in `NodeAddr` order. + pub moved: Vec, + /// How many peers were probed in the newer of the two samples, so a log + /// line shows what fraction of the table moved. + pub probed: usize, } 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() { + if self.moved.is_empty() { return write!(f, "no visible difference"); } - write!(f, "{}", parts.join(", ")) + write!(f, "{}/{} peers: ", self.moved.len(), self.probed)?; + // Bounded: an operator needs the shape of the change, and a full table + // moving at once is the common case rather than the interesting one. + const NAMED: usize = 3; + for (i, m) in self.moved.iter().take(NAMED).enumerate() { + if i > 0 { + write!(f, ", ")?; + } + // Both ends: the address the stale `connect(2)` had pinned is what + // an operator correlates against route history, so a line naming + // only the destination leaves out the half being diagnosed. + let before = match m.before { + Some(ip) => ip.to_string(), + None => "no route".to_string(), + }; + let after = match m.after { + Some(ip) => ip.to_string(), + None => "no route".to_string(), + }; + write!(f, "{} {} -> {}", m.peer.short_hex(), before, after)?; + } + if self.moved.len() > NAMED { + write!(f, ", +{} more", self.moved.len() - NAMED)?; + } + Ok(()) } } @@ -259,15 +462,16 @@ pub(crate) struct NetChange { 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. + /// the handler does not read it — it re-evaluates every peer regardless of + /// which peer's source address moved. #[cfg(test)] pub(crate) fn for_test(generation: u64) -> Self { - let empty = NetFingerprint::default(); Self { generation, - summary: empty.diff(&empty), + summary: NetChangeSummary { + moved: Vec::new(), + probed: 0, + }, } } } @@ -357,6 +561,17 @@ impl WakeSource { timer } + /// The netlink groups this wake source is subscribed to, or `None` if it + /// is not a live netlink source. For the group-mask assertion in the + /// tests — see `the_detector_subscribes_to_the_route_groups_not_just_link`. + #[cfg(all(test, any(target_os = "linux", target_os = "android")))] + fn subscribed_groups(&self) -> Option { + match &self.source { + Wake::Kernel(watcher) => watcher.subscribed_groups(), + _ => None, + } + } + /// Wait until it is worth sampling again. async fn wait(&mut self) { let WakeSource { source, timer } = self; @@ -407,15 +622,42 @@ impl WakeSource { /// 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<()>) { +pub(crate) fn spawn_detector( + cfg: NetmonConfig, + peers: Arc>, +) -> (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; + let sample = move || NetFingerprint::sample(&probe_targets(&peers.load())); + run_detector(tx, cfg, sample, wake).await; }); (rx, handle) } +/// The peers worth probing, read off the node's published entity snapshot. +/// +/// Every peer carrying a numeric IP endpoint, paired with it. The filter is +/// [`crate::control::snapshot::PeerRow::probe_target`] being `Some`, which is +/// already exactly "this peer is an IP destination we could `connect(2)` to": +/// a MAC, a `.onion`, a Nym recipient and an unresolved hostname all arrive +/// here as `None` and are skipped, with no per-transport special-casing in +/// this module. +pub(in crate::node) fn probe_targets(snapshot: &EntitySnapshot) -> Vec { + snapshot + .peers + .iter() + .filter_map(|row| { + row.probe_target.map(|dest| ProbeTarget { + peer: row.node_addr, + dest, + bound: row.bound_source, + bind: row.probe_bind, + }) + }) + .collect() +} + /// Pick the wake source: the event-driven backend where one exists and starts, /// the timer otherwise. /// @@ -484,29 +726,42 @@ where wake.wait().await; let mut candidate = sample(); - if candidate == last { + if last.moved(&candidate).is_empty() { + // Nothing the node is peering over moved. Adopt the sample anyway: + // it is how a peer that has just joined enters the comparison, and + // one that has left leaves it. Skipping this would freeze `last` on + // the peer set the detector started with, and a peer authenticated + // later would never be compared against anything. + last = candidate; 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. + // Settling is judged the same way — no peer moved since the previous + // round — so peers joining or leaving mid-handover cannot extend the + // debounce on their own either. for _ in 0..MAX_DEBOUNCE_ROUNDS { if debounce.is_zero() { break; } tokio::time::sleep(debounce).await; let resampled = sample(); - if resampled == candidate { + let settled = candidate.moved(&resampled).is_empty(); + candidate = resampled; + if settled { 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 { + // The burst may have settled back to where it started (a route that + // flapped away and returned). Nothing moved, so nothing is reported — + // but the sample is still adopted, for the reason above. + let moved = last.moved(&candidate); + if moved.is_empty() { trace!("Network fingerprint settled back unchanged; no event"); + last = candidate; continue; } @@ -524,7 +779,10 @@ where generation += 1; let change = NetChange { generation, - summary: last.diff(&candidate), + summary: NetChangeSummary { + moved, + probed: candidate.sources.len(), + }, }; last = candidate; @@ -544,16 +802,30 @@ where } } -/// The source address the kernel would use to reach `probe`. +/// The source address the kernel would use to reach `probe`, from `bind_to` if +/// the transport constrains it. /// /// `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), +/// the route, binds a source address, and sends nothing. The socket is dropped +/// here and never written to, so a peer is probed without a single packet +/// reaching it. +/// +/// A failure — most often `ENETUNREACH`, no route to that peer at all — is +/// itself a fingerprint value, reported as `None` rather than swallowed. It is +/// the state a stranded peer is in, so losing it would blind the detector to +/// the case it most needs to see. +fn preferred_source(probe: SocketAddr, bind_to: Option) -> Option { + // Port 0 always: the probe wants the transport's *address* constraint, not + // its port, and binding the live port would collide with the socket the + // transport already holds there. A family mismatch between the configured + // bind and this peer is not an error to report — the transport could not + // have reached the peer from it either — so fall back to unspecified and + // let the connect below fail on its own terms. + let bind: SocketAddr = match (probe, bind_to) { + (SocketAddr::V4(_), Some(ip @ IpAddr::V4(_))) => SocketAddr::new(ip, 0), + (SocketAddr::V6(_), Some(ip @ IpAddr::V6(_))) => SocketAddr::new(ip, 0), + (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()?; @@ -566,84 +838,5 @@ fn preferred_source(probe: SocketAddr) -> Option { 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 index 957123b3..b079a322 100644 --- a/src/node/netmon/tests.rs +++ b/src/node/netmon/tests.rs @@ -35,6 +35,19 @@ fn v4(a: u8, b: u8, c: u8, d: u8) -> IpAddr { IpAddr::V4(Ipv4Addr::new(a, b, c, d)) } +/// A distinct peer address. Only identity matters here, so the byte pattern is +/// arbitrary as long as two peers differ. +fn peer(n: u8) -> NodeAddr { + NodeAddr::from_bytes([n; 16]) +} + +/// A fingerprint in which `peers` are all reached from `source`. The common +/// shape: every peering rides one medium, so they move together. +fn all_from(peers: &[NodeAddr], source: Option) -> NetFingerprint { + let sources: Vec<_> = peers.iter().map(|p| (*p, source)).collect(); + NetFingerprint::for_test(&sources) +} + /// 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 { @@ -62,7 +75,7 @@ fn scripted(samples: Vec) -> (impl Fn() -> NetFingerprint, Arc ProbeTarget { + ProbeTarget { + peer: p, + dest, + bound: None, + bind: None, + } +} + #[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(); + let targets = [target(peer(1), OFF_LINK)]; + let first = NetFingerprint::sample(&targets); + let second = NetFingerprint::sample(&targets); assert_eq!( first, second, "consecutive samples of an unchanged host must agree" ); + assert!( + first.moved(&second).is_empty(), + "and must show no peer as having moved" + ); } +/// The live-probe tests above are all satisfied by a `preferred_source` that +/// returns `None` for everything: two all-`None` samples are self-consistent, +/// the recorded-keys test never inspects a value, and the loopback assertion +/// skips through its `if let`. This one pins that the probe actually answers. +/// +/// Loopback is the destination because it is routable on any host that can run +/// this suite, including a container started with `--network none`, and the +/// source for it is loopback itself. #[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 +fn a_probe_to_loopback_answers_with_loopback() { + let dest = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 9); + let sample = NetFingerprint::sample(&[target(peer(1), dest)]); + assert_eq!( + sample.sources.get(&peer(1)).map(|p| p.current), + Some(Some(IpAddr::V4(Ipv4Addr::LOCALHOST))), + "the probe must return the kernel's source address, not None" ); } #[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"); +fn every_target_is_recorded_whether_or_not_it_has_a_route() { + // A peer with no route must stay in the map as `None` rather than dropping + // out of it. Dropping it would make "the route to this peer just vanished" + // indistinguishable from "this peer was reaped", and the intersection rule + // would then discard exactly the event the detector exists to catch. The + // assertion holds on a CI container with no route at all. + let targets = [target(peer(1), OFF_LINK), target(peer(2), OFF_LINK)]; + let sample = NetFingerprint::sample(&targets); + assert_eq!(sample.sources.len(), 2); + assert!(sample.sources.contains_key(&peer(1))); + assert!(sample.sources.contains_key(&peer(2))); } #[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); +fn a_probe_never_yields_a_loopback_or_unspecified_source() { + // Either of those would be the kernel declining to choose, not an answer, + // and treating one as an address would make the fingerprint move whenever + // the route lookup failed differently. + let sample = NetFingerprint::sample(&[target(peer(1), OFF_LINK)]); + if let Some(Some(ip)) = sample.sources.get(&peer(1)).map(|p| p.current) { + assert!(!ip.is_loopback(), "loopback source for an off-link probe"); + assert!( + !ip.is_unspecified(), + "unspecified source reported as an answer" + ); + } +} + +// === Summary rendering === + +#[test] +fn summary_of_nothing_moving_is_legible() { + let summary = NetChangeSummary { + moved: Vec::new(), + probed: 4, + }; + assert_eq!(summary.to_string(), "no visible difference"); +} + +#[test] +fn summary_names_the_peer_and_its_new_source() { + let wlan = all_from(&[peer(0xab)], Some(v4(192, 168, 1, 10))); + let cell = all_from(&[peer(0xab)], Some(v4(10, 40, 0, 7))); + let summary = NetChangeSummary { + moved: wlan.moved(&cell), + probed: 1, + }; + let rendered = summary.to_string(); + assert!(rendered.contains("1/1 peers"), "{}", rendered); + // Both ends, and the peer id in the same shape every other operator + // surface prints, so a line found in `show_peers` matches here. + assert!( + rendered.contains("abababab... 192.168.1.10 -> 10.40.0.7"), + "{}", + rendered + ); +} + +#[test] +fn summary_says_when_a_peer_lost_its_route() { + let up = all_from(&[peer(1)], Some(v4(192, 168, 1, 10))); + let down = all_from(&[peer(1)], None); + let summary = NetChangeSummary { + moved: up.moved(&down), + probed: 1, + }; + assert!(summary.to_string().contains("no route"), "{}", summary); +} + +#[test] +fn summary_truncates_a_whole_table_moving_at_once() { + // The common case is every peer moving together, and a log line naming a + // hundred of them is not a log line. + let peers: Vec = (1..=10).map(peer).collect(); + let before = all_from(&peers, Some(v4(192, 168, 1, 10))); + let after = all_from(&peers, Some(v4(10, 40, 0, 7))); + let summary = NetChangeSummary { + moved: before.moved(&after), + probed: peers.len(), + }; + let rendered = summary.to_string(); + assert!(rendered.contains("10/10 peers"), "{}", rendered); + assert!(rendered.contains("+7 more"), "{}", rendered); } // === Wake source === @@ -228,8 +512,8 @@ fn summary_names_the_new_source_address() { /// 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 wlan = all_from(&[peer(1)], Some(v4(192, 168, 1, 10))); + let cell = all_from(&[peer(1)], Some(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); @@ -240,7 +524,7 @@ async fn an_event_ping_wakes_the_detector_before_the_timer_would() { 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))); + assert_eq!(change.summary.moved[0].after, Some(v4(10, 40, 0, 7))); } /// A netlink socket drops messages under memory pressure, and a backend can go @@ -248,8 +532,8 @@ async fn an_event_ping_wakes_the_detector_before_the_timer_would() { /// 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 wlan = all_from(&[peer(1)], Some(v4(192, 168, 1, 10))); + let cell = all_from(&[peer(1)], Some(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. @@ -260,7 +544,7 @@ async fn the_backstop_still_fires_when_the_backend_says_nothing() { let change = expect_change(&mut rx).await; assert_eq!( - change.summary.v4_source, + change.summary.moved[0].after, Some(v4(10, 40, 0, 7)), "the backstop must reach the change the backend missed" ); @@ -271,8 +555,8 @@ async fn the_backstop_still_fires_when_the_backend_says_nothing() { /// 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 wlan = all_from(&[peer(1)], Some(v4(192, 168, 1, 10))); + let cell = all_from(&[peer(1)], Some(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); @@ -285,7 +569,7 @@ async fn a_dead_backend_falls_back_to_the_timer() { let change = expect_change(&mut rx).await; assert_eq!( - change.summary.v4_source, + change.summary.moved[0].after, Some(v4(10, 40, 0, 7)), "detection must survive the backend it was using" ); @@ -408,8 +692,8 @@ async fn a_route_change_alone_reaches_the_watcher() { /// 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)]); + let a = all_from(&[peer(1)], Some(v4(192, 168, 1, 10))); + let b = all_from(&[peer(1)], Some(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(); @@ -424,7 +708,11 @@ async fn reports_are_spaced_out_under_clean_flapping() { 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))); + // Not `timer_wake(1)`: a one-second poll against a one-second floor makes + // the two indistinguishable, and deleting the pacing block would still + // produce one-second spacing and still pass. + let wake = WakeSource::timer_only(Duration::from_millis(100)); + tokio::spawn(run_detector(tx, cfg(1, 0), sampler, wake)); let first = expect_change(&mut rx).await; let started = tokio::time::Instant::now(); @@ -439,3 +727,440 @@ async fn reports_are_spaced_out_under_clean_flapping() { started.elapsed() ); } + +/// The claim this whole shape exists to make: an interface appearing that is +/// not the route to any peer does not move the fingerprint. +/// +/// This is the case that made the host-wide address set unusable — a container +/// bridge, a VPN, a `veth` pair or a tunnel coming up moved it, and the node +/// answered by dropping every connected socket and heartbeating every peer, for +/// a `docker compose up`. Here a second interface arrives with an address and a +/// subnet of its own, carrying no route to the peer, and nothing is reported. +/// +/// Structurally this cannot fail while the fingerprint holds only per-peer +/// probe results — there is no host-wide enumeration left in the module to go +/// wrong. The test is here so that a future signal added back into +/// `NetFingerprint::sample` has to answer to it. +/// +/// Like the watcher test above, this needs `CAP_NET_ADMIN` in a namespace it +/// may reconfigure: +/// +/// ```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 an_interface_no_peer_is_reached_through_does_not_move_the_fingerprint() { + use futures::TryStreamExt; + use std::net::Ipv4Addr; + + // Its own destination, in RFC 5737 TEST-NET-2 rather than the TEST-NET-1 + // that [`OFF_LINK`] uses: `unshare -rn` gives the whole test binary one + // namespace, so the routes the netlink test above installs are still there + // and a shared prefix collides with EEXIST depending on the order they run. + const NET: Ipv4Addr = Ipv4Addr::new(198, 51, 100, 0); + const CARRIER: Ipv4Addr = Ipv4Addr::new(10, 99, 0, 1); + const BRIDGE: Ipv4Addr = Ipv4Addr::new(172, 30, 0, 1); + const DEST: SocketAddr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(198, 51, 100, 1)), 9); + + /// Bring up a dummy interface carrying `addr/24`, returning its index. + async fn dummy_up(handle: &rtnetlink::Handle, name: &str, addr: Ipv4Addr) -> u32 { + handle + .link() + .add(rtnetlink::LinkDummy::new(name).build()) + .execute() + .await + .expect("creating a dummy link needs CAP_NET_ADMIN in this namespace"); + let index = handle + .link() + .get() + .match_name(name.to_string()) + .execute() + .try_next() + .await + .expect("link query") + .expect("the link just created exists") + .header + .index; + handle + .address() + .add(index, std::net::IpAddr::V4(addr), 24) + .execute() + .await + .expect("adding an address"); + handle + .link() + .set(rtnetlink::LinkUnspec::new_with_index(index).up().build()) + .execute() + .await + .expect("bringing the link up"); + index + } + + let (connection, handle, _) = rtnetlink::new_connection().expect("netlink connection"); + tokio::spawn(connection); + + // The medium the peer is actually reached over: an interface plus the + // default route out of it. + let carrier = dummy_up(&handle, "mc-carrier", CARRIER).await; + handle + .route() + .add( + rtnetlink::RouteMessageBuilder::::new() + .output_interface(carrier) + .build(), + ) + .execute() + .await + .expect("adding a default route"); + tokio::time::sleep(Duration::from_millis(200)).await; + + let targets = [target(peer(1), DEST)]; + let before = NetFingerprint::sample(&targets); + assert_eq!( + before.sources.get(&peer(1)).map(|p| p.current), + Some(Some(IpAddr::V4(CARRIER))), + "the peer must be reached over the carrier before anything else appears, \ + or this test proves nothing about what happens next" + ); + + // The interloper: a bridge-shaped interface with its own subnet, exactly + // what `docker compose up` leaves behind. It is up, it is not loopback, and + // it carries an address — every property the old host-wide set keyed on — + // but no peer is reached through it. + dummy_up(&handle, "mc-bridge", BRIDGE).await; + tokio::time::sleep(Duration::from_millis(200)).await; + + let after = NetFingerprint::sample(&targets); + assert_eq!( + before.moved(&after), + Vec::new(), + "an interface carrying no route to any peer must not be a medium change" + ); + + // While two addresses exist here, measure the reason the probe carries a + // bind constraint at all. + // + // `open_connected_fd` binds the transport's configured address verbatim + // and only then connects, so under a non-wildcard `transports.udp.bind_addr` + // the socket's source is that address whatever the routing table says. An + // unconstrained probe answers with the kernel's choice instead. Where the + // two differ, the first-sight rule would compare them and report a move on + // every peer, permanently, with nothing having moved. + // + // The interloper's address is reachable-from but is not what the route to + // DEST would pick, so it is exactly that disagreement, made concrete: the + // unconstrained probe answers with the carrier, the constrained one with + // what it was told to bind. + let unconstrained = preferred_source(DEST, None); + let constrained = preferred_source(DEST, Some(IpAddr::V4(BRIDGE))); + assert_eq!( + unconstrained, + Some(IpAddr::V4(CARRIER)), + "an unconstrained probe follows the route" + ); + assert_eq!( + constrained, + Some(IpAddr::V4(BRIDGE)), + "a constrained probe answers from the address it was told to bind, which \ + is what a non-wildcard bind_addr makes the send path do" + ); + assert_ne!( + unconstrained, constrained, + "if these agreed the constraint would be untested, and the phantom-move \ + case it exists for could not arise" + ); + + // The other half, in the same namespace and against the same live sampler: + // put a more specific route to the peer out of the interloper, and the + // fingerprint must move with it. Two things ride on this. It stops the + // assertion above passing because sampling had quietly stopped working, + // which is the failure mode a negative assertion is worst at catching. And + // it is the per-peer route case in its own right — the default route never + // moves here, nothing about the host's attachment changes, and no + // host-wide sample could represent this at all. + let bridge = handle + .link() + .get() + .match_name("mc-bridge".to_string()) + .execute() + .try_next() + .await + .expect("link query") + .expect("mc-bridge exists") + .header + .index; + handle + .route() + .add( + rtnetlink::RouteMessageBuilder::::new() + .destination_prefix(NET, 24) + .output_interface(bridge) + .build(), + ) + .execute() + .await + .expect("adding a more specific route to the peer"); + tokio::time::sleep(Duration::from_millis(200)).await; + + let moved = after.moved(&NetFingerprint::sample(&targets)); + assert_eq!( + moved, + vec![PeerSourceMove { + peer: peer(1), + before: Some(IpAddr::V4(Ipv4Addr::new(10, 99, 0, 1))), + after: Some(IpAddr::V4(Ipv4Addr::new(172, 30, 0, 1))), + }], + "the route to the peer moving is exactly what must be reported" + ); +} + +/// An interface going down and coming back up, which the docker suite does not +/// cover: it moves the default route with both interfaces held up throughout, +/// deliberately, so that it tests a medium change rather than a link failure. +/// This is the other shape — the interface carrying a peer is taken away and +/// given back. +/// +/// Three transitions, and the third is the one worth having. Downing the +/// interface a peer is reached over must report; bringing it back must report; +/// downing an interface no peer is reached over must not. That last case is the +/// down-direction counterpart of the container test above, and it is the one a +/// link-state watcher gets wrong — the kernel emits exactly the same link event +/// for all three. +/// +/// Routes here are specific to this test's own prefix rather than defaults, so +/// it shares the `unshare -rn` namespace with the tests above without fighting +/// them over the default route. +/// +/// ```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 an_interface_going_down_is_reported_only_when_a_peer_was_reached_over_it() { + use futures::TryStreamExt; + use std::net::Ipv4Addr; + + // RFC 5737 TEST-NET-3, so this test's routes cannot collide with either of + // the two above in the shared namespace. + const NET: Ipv4Addr = Ipv4Addr::new(203, 0, 113, 0); + const DEST: SocketAddr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(203, 0, 113, 1)), 9); + const PRIMARY: Ipv4Addr = Ipv4Addr::new(10, 77, 0, 1); + const BACKUP: Ipv4Addr = Ipv4Addr::new(10, 78, 0, 1); + const IDLE: Ipv4Addr = Ipv4Addr::new(10, 79, 0, 1); + + let (connection, handle, _) = rtnetlink::new_connection().expect("netlink connection"); + tokio::spawn(connection); + + async fn index_of(handle: &rtnetlink::Handle, name: &str) -> u32 { + handle + .link() + .get() + .match_name(name.to_string()) + .execute() + .try_next() + .await + .expect("link query") + .expect("the link exists") + .header + .index + } + + async fn dummy_up(handle: &rtnetlink::Handle, name: &str, addr: Ipv4Addr) -> u32 { + handle + .link() + .add(rtnetlink::LinkDummy::new(name).build()) + .execute() + .await + .expect("creating a dummy link needs CAP_NET_ADMIN in this namespace"); + let index = index_of(handle, name).await; + handle + .address() + .add(index, std::net::IpAddr::V4(addr), 24) + .execute() + .await + .expect("adding an address"); + handle + .link() + .set(rtnetlink::LinkUnspec::new_with_index(index).up().build()) + .execute() + .await + .expect("bringing the link up"); + index + } + + async fn set_link(handle: &rtnetlink::Handle, index: u32, up: bool) { + let msg = if up { + rtnetlink::LinkUnspec::new_with_index(index).up().build() + } else { + rtnetlink::LinkUnspec::new_with_index(index).down().build() + }; + handle + .link() + .set(msg) + .execute() + .await + .expect("changing link state"); + tokio::time::sleep(Duration::from_millis(200)).await; + } + + async fn route_to_peer(handle: &rtnetlink::Handle, index: u32, metric: u32) { + handle + .route() + .add( + rtnetlink::RouteMessageBuilder::::new() + .destination_prefix(NET, 24) + .output_interface(index) + .priority(metric) + .build(), + ) + .execute() + .await + .expect("adding a route to the peer"); + } + + // Two paths to the peer, primary preferred, plus an interface carrying no + // route to it at all. + let primary = dummy_up(&handle, "mc-updn-a", PRIMARY).await; + let backup = dummy_up(&handle, "mc-updn-b", BACKUP).await; + let idle = dummy_up(&handle, "mc-updn-c", IDLE).await; + route_to_peer(&handle, primary, 100).await; + route_to_peer(&handle, backup, 200).await; + tokio::time::sleep(Duration::from_millis(200)).await; + + let targets = [target(peer(1), DEST)]; + let on_primary = NetFingerprint::sample(&targets); + assert_eq!( + on_primary.sources.get(&peer(1)).map(|p| p.current), + Some(Some(IpAddr::V4(PRIMARY))), + "the peer must start out on the primary, or nothing below means anything" + ); + + // 1. The interface the peer is reached over goes down. The route with it, + // so the kernel falls back to the higher-metric path. + set_link(&handle, primary, false).await; + let on_backup = NetFingerprint::sample(&targets); + assert_eq!( + on_primary.moved(&on_backup), + vec![PeerSourceMove { + peer: peer(1), + before: Some(IpAddr::V4(PRIMARY)), + after: Some(IpAddr::V4(BACKUP)), + }], + "losing the interface a peer was reached over is a medium change" + ); + + // 2. And back. Not symmetric with the above: this direction returns to an + // interface that has been holding a stale address throughout. + // + // The route has to be re-added by hand, because the kernel deleted it + // when the link went down and does not restore it when the link returns + // — verified in a namespace, not assumed. On a real host that re-add is + // what the DHCP client or the network manager does on carrier-up, so + // re-adding it here is modelling the real sequence rather than working + // around it. The address, by contrast, does survive, which is exactly + // the trap the detector exists for: the interface is up and addressed + // again the instant the link returns, and only the route says whether + // anything is reached over it. + set_link(&handle, primary, true).await; + route_to_peer(&handle, primary, 100).await; + tokio::time::sleep(Duration::from_millis(200)).await; + let back_on_primary = NetFingerprint::sample(&targets); + assert_eq!( + on_backup.moved(&back_on_primary), + vec![PeerSourceMove { + peer: peer(1), + before: Some(IpAddr::V4(BACKUP)), + after: Some(IpAddr::V4(PRIMARY)), + }], + "the interface returning and reclaiming the route is a medium change too" + ); + + // 3. An interface no peer is reached over goes down. Same kernel link + // event as case 1, and it must report nothing. + set_link(&handle, idle, false).await; + assert_eq!( + back_on_primary.moved(&NetFingerprint::sample(&targets)), + Vec::new(), + "an interface no peer was reached over going down is not a medium change" + ); + + // 4. Every path this test installed goes away at once. Where the peer + // lands afterwards is deliberately not asserted: it depends on what + // else the host offers, and in this shared namespace it falls back to + // the default route another `--ignored` test installed. What must hold + // either way is that the peer moved off the primary and that the move + // is reported — a peer resolving to nothing at all is pinned by + // `every_target_is_recorded_whether_or_not_it_has_a_route`, which runs + // on a CI container with no route to fall back to. + set_link(&handle, primary, false).await; + set_link(&handle, backup, false).await; + let stranded = NetFingerprint::sample(&targets); + assert_ne!( + stranded.sources.get(&peer(1)).map(|p| p.current), + Some(Some(IpAddr::V4(PRIMARY))), + "the peer cannot still be reached over an interface that is down" + ); + assert_eq!( + back_on_primary.moved(&stranded).len(), + 1, + "losing every path this test installed is a medium change" + ); +} + +/// The group mask the detector actually ends up subscribed to, read back from +/// the kernel. +/// +/// This is the one property of the netlink backend that CI could not check. +/// `a_route_change_alone_reaches_the_watcher` discriminates a wrong mask by +/// provoking a real route change, but it needs `CAP_NET_ADMIN` and is skipped +/// everywhere CI runs — so a regression to link-events-only would have passed +/// every gate while the detector silently stopped seeing the default route +/// move, which is the change it exists to catch. +/// +/// A wrong mask cannot be caught by watching the bind: subscribing to the +/// wrong groups succeeds exactly like subscribing to the right ones, and only +/// differs in what never arrives afterwards. So this asks the kernel what the +/// socket is subscribed to instead, which needs no privileges at all. +/// +/// It goes through `build_wake_source` rather than constructing a watcher +/// directly, so it is the production path being asserted on and not a second +/// copy of the same constant. +#[cfg(any(target_os = "linux", target_os = "android"))] +#[tokio::test] +async fn the_detector_subscribes_to_the_route_groups_not_just_link() { + use crate::transport::watcher::groups; + + let wake = build_wake_source(&cfg(5, 250)); + // Asserted, not skipped. An early return here would make this test green in + // exactly the environment that differs from a real check — a sandbox where + // the bind is refused — so a mask regression would pass everywhere the + // subscription could not be inspected. `the_egress_path_mask_opens_a_source` + // already holds the same line: the bind is expected to work on Linux. + let subscribed = wake.subscribed_groups().expect( + "the detector must have a live netlink subscription to inspect; without one \ + this test cannot say anything about the group mask", + ); + + assert_eq!( + subscribed, + groups::EGRESS_PATH, + "the detector must be subscribed to the egress-path groups it asked for" + ); + for (name, group) in [ + ("IPV4_ROUTE", groups::IPV4_ROUTE), + ("IPV6_ROUTE", groups::IPV6_ROUTE), + ("IPV4_IFADDR", groups::IPV4_IFADDR), + ("IPV6_IFADDR", groups::IPV6_IFADDR), + ] { + assert_ne!( + subscribed & group, + 0, + "{name} is missing: a default route moving between two interfaces that both \ + stay up emits nothing in the link group, so without this the detector would \ + never fire for the change it exists to catch" + ); + } +} diff --git a/src/node/tests/netmon.rs b/src/node/tests/netmon.rs index b76e629c..ef6f35a8 100644 --- a/src/node/tests/netmon.rs +++ b/src/node/tests/netmon.rs @@ -9,7 +9,7 @@ use super::spanning_tree::*; use super::*; use crate::config::PeerConfig; use crate::config::TcpConfig; -use crate::node::netmon::NetChange; +use crate::node::netmon::{NetChange, NetFingerprint, ProbeTarget}; use crate::transport::tcp::TcpTransport; use crate::transport::{TransportAddr, TransportHandle, TransportId, packet_channel}; @@ -37,7 +37,7 @@ fn identity_of(nodes: &[TestNode], j: usize) -> PeerIdentity { /// 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")] +#[cfg(any(target_os = "linux", target_os = "macos"))] 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(); @@ -78,7 +78,7 @@ fn install_connected_udp(node: &mut Node, addr: &NodeAddr, transport_id: Transpo /// 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")] +#[cfg(any(target_os = "linux", target_os = "macos"))] #[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; @@ -257,3 +257,187 @@ async fn a_change_with_no_peers_is_harmless() { node.handle_net_change(NetChange::for_test(1)).await; assert!(node.peers.is_empty()); } + +/// The detector reads the peer table through the published entity snapshot, +/// and this is the seam: what a peer's transport address is determines whether +/// it arrives on the other side as something to probe. Nothing else in the +/// tree exercises `PeerRow::probe_target`, because nothing renders it — so if +/// the publish site stopped populating it, every other test here would still +/// pass while the detector silently probed an empty table and never reported +/// anything again. +/// +/// Each case re-pins the same established peer, because the address is the +/// only variable that matters: the projection is a property of the address, +/// not of the transport it was learned on. (The harness's own peers sit on a +/// synthetic `loopback:1` transport, which is itself correctly unprobeable.) +#[tokio::test] +async fn only_a_peer_with_an_ip_endpoint_reaches_the_probe() { + // (address as the peer carries it, the destination the detector should + // probe, why) + let cases: [(TransportAddr, Option<&str>, &str); 6] = [ + ( + TransportAddr::from_string("10.0.0.2:2121"), + Some("10.0.0.2:2121"), + "an ordinary IPv4 peer is the whole point", + ), + ( + TransportAddr::from_string("[2001:db8::1]:2121"), + Some("[2001:db8::1]:2121"), + "IPv6 literals round-trip through the row", + ), + ( + TransportAddr::from_bytes(&[0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff]), + None, + "a MAC has no IP destination to ask the routing table about", + ), + ( + TransportAddr::from_string("example.com:2121"), + None, + "resolving a hostname would put DNS on the detector's sample path", + ), + ( + TransportAddr::from_string("abcdefghij234567.onion:2121"), + None, + "a .onion is reached through a local proxy, not a route", + ), + ( + TransportAddr::from_string("[fe80::1%eth0]:2121"), + None, + "a scoped link-local literal is not a parseable SocketAddr", + ), + ]; + + let mut nodes = run_tree_test(2, &[(0, 1)], false).await; + verify_tree_convergence(&nodes); + + let addr_1 = *nodes[1].node.node_addr(); + let transport_id = nodes[0] + .node + .peers + .get(&addr_1) + .and_then(|p| p.transport_id()) + .expect("peer 1 has a transport"); + + for (addr, expected, why) in cases { + nodes[0] + .node + .peers + .get_mut(&addr_1) + .expect("peer 1 is established") + .set_current_addr(transport_id, addr.clone()); + + // The snapshot is published from the tick, which is its only writer. + nodes[0].node.record_stats_history(); + let snapshot = nodes[0].node.entities_snapshot.load_full(); + + let got = crate::node::netmon::probe_targets(&snapshot) + .into_iter() + .find(|t| t.peer == addr_1) + .map(|t| t.dest); + + let want = expected.map(|s| s.parse::().unwrap()); + assert_eq!(got, want, "{}: {}", addr, why); + } + + cleanup_nodes(&mut nodes).await; +} + +/// The seed the join-window fix rests on, wired end to end. +/// +/// A peer the detector has not seen before is judged against the source its +/// connected socket was pinned to, and that value has to be the address +/// `connect(2)` actually chose — not the wildcard the bind was requested with. +/// `ConnectedPeerSocket::local_addr()` is the wildcard (`0.0.0.0:port`), and +/// reading *that* would compare an unspecified address against a real one for +/// every peer, so every peer joining would report a medium change: precisely +/// the "peer churn fires the fan-out" behaviour the intersection rule exists to +/// prevent. Nothing renders `bound_source`, so no other test would notice. +#[cfg(any(target_os = "linux", target_os = "macos"))] +#[tokio::test] +async fn a_peers_connected_socket_publishes_the_source_it_was_pinned_to() { + let mut nodes = run_tree_test(2, &[(0, 1)], false).await; + verify_tree_convergence(&nodes); + + let addr_1 = *nodes[1].node.node_addr(); + let transport_id = nodes[0] + .node + .peers + .get(&addr_1) + .and_then(|p| p.transport_id()) + .expect("peer 1 has a transport"); + + // No socket yet: nothing to seed from, and the peer must say so rather + // than offering the wildcard. + nodes[0].node.record_stats_history(); + let row = |n: &Node| { + n.entities_snapshot + .load_full() + .peers + .iter() + .find(|r| r.node_addr == addr_1) + .expect("peer 1 has a row") + .clone() + }; + assert_eq!( + row(&nodes[0].node).bound_source, + None, + "a peer with no connected socket has no pinned source to be judged against" + ); + + // The helper connects to 127.0.0.1:9, so the kernel pins the loopback + // source — a real address, and demonstrably not the `0.0.0.0` the bind was + // requested with. + install_connected_udp(&mut nodes[0].node, &addr_1, transport_id); + // The harness peers sit on a synthetic `loopback:1` address, which is + // correctly not probeable. Re-pin to a numeric endpoint on the same + // transport so the row reaches the probe at all — the pinned source is a + // property of the socket, not of the address, and survives this. + nodes[0] + .node + .peers + .get_mut(&addr_1) + .expect("peer 1 is established") + .set_current_addr(transport_id, TransportAddr::from_string("10.0.0.2:2121")); + nodes[0].node.record_stats_history(); + + assert_eq!( + row(&nodes[0].node).bound_source, + Some(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)), + "the published source must be what connect(2) pinned, not the wildcard bind" + ); + + // Publishing it is only half the wiring. Nothing else asserts that + // `probe_targets` carries `bound_source` through to the target, so + // substituting `None` there leaves the whole suite green while silently + // restoring the bug the first-sight rule exists to fix — the same failure + // class as reading the wildcard `local_addr()`, one layer further on. + let snapshot = nodes[0].node.entities_snapshot.load_full(); + let target = crate::node::netmon::probe_targets(&snapshot) + .into_iter() + .find(|t| t.peer == addr_1) + .expect("an established UDP peer must be probeable"); + assert_eq!( + target.bound, + Some(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)), + "the pinned source must reach the probe target, not stop at the row" + ); + + // And the last link: `sample()` has to carry it into the fingerprint, or a + // first-seen peer is judged against nothing again. An empty previous + // fingerprint is exactly the first-sight case, and the peer's socket is + // pinned to loopback while the probe answers for a routable destination, + // so the two disagree and a move must be reported. + let sampled = NetFingerprint::sample(&[ProbeTarget { + peer: addr_1, + dest: "192.0.2.1:9".parse().unwrap(), + bound: Some(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)), + bind: None, + }]); + assert!( + !NetFingerprint::default().moved(&sampled).is_empty(), + "sample() must carry the pinned source into the fingerprint, or first \ + sight has nothing to judge against" + ); + + cleanup_nodes(&mut nodes).await; +} diff --git a/src/transport/udp/io/connected/socket.rs b/src/transport/udp/io/connected/socket.rs index 0404e59a..ac540443 100644 --- a/src/transport/udp/io/connected/socket.rs +++ b/src/transport/udp/io/connected/socket.rs @@ -4,7 +4,7 @@ //! and closes it on drop. See that function's docs for why established //! peers get their own connected socket. -use std::net::SocketAddr; +use std::net::{IpAddr, SocketAddr}; use std::os::unix::io::{AsRawFd, OwnedFd, RawFd}; /// A `connect()`-ed UDP socket for one established peer. @@ -28,6 +28,9 @@ pub(crate) struct ConnectedPeerSocket { fd: OwnedFd, peer_addr: SocketAddr, local_addr: SocketAddr, + /// The source address `connect(2)` actually pinned, read back once at + /// construction. See [`ConnectedPeerSocket::pinned_source`]. + pinned_source: Option, } impl ConnectedPeerSocket { @@ -35,10 +38,12 @@ impl ConnectedPeerSocket { /// `crate::transport::udp::open_connected_fd`) into an owning /// handle. Takes ownership of the fd; the `OwnedFd` closes it on drop. pub(crate) fn from_fd(fd: OwnedFd, peer_addr: SocketAddr, local_addr: SocketAddr) -> Self { + let pinned_source = pinned_source_of(fd.as_raw_fd()); Self { fd, peer_addr, local_addr, + pinned_source, } } @@ -50,6 +55,52 @@ impl ConnectedPeerSocket { pub fn local_addr(&self) -> SocketAddr { self.local_addr } + + /// The source address the kernel bound when this socket was + /// `connect(2)`-ed, as opposed to [`Self::local_addr`], which is the + /// wildcard the bind was *requested* with and carries no interface + /// information at all. + /// + /// This is the quantity the whole connected-socket fast path turns on: the + /// kernel resolves the route once at connect time and pins the source + /// address to whichever interface was carrying it then, and never + /// re-evaluates. Read back once here rather than per call, because it + /// cannot change for the life of the socket — that being exactly the + /// problem. `crate::node::netmon` compares it against the address the + /// routing table would choose now, which is how a peer whose socket is + /// already stale is recognised without any earlier sample to compare + /// against. + /// + /// `None` if `getsockname` fails or reports a family this does not decode, + /// which is treated as "no answer" rather than guessed at. + pub(crate) fn pinned_source(&self) -> Option { + self.pinned_source + } +} + +/// `getsockname` on a connected UDP socket, reduced to the local IP. +/// +/// Returns `None` on any failure: the caller's contract is that an unknown +/// pinned source is indistinguishable from not having one, and both mean "do +/// not draw a conclusion from this socket". +fn pinned_source_of(fd: RawFd) -> Option { + let mut storage: libc::sockaddr_storage = unsafe { std::mem::zeroed() }; + let mut len = std::mem::size_of::() as libc::socklen_t; + // SAFETY: `fd` is the socket this handle owns, and `storage` / `len` are a + // correctly sized and initialised out-parameter pair for `getsockname`, + // which writes at most `len` bytes and updates `len` to what it wrote. + let rc = + unsafe { libc::getsockname(fd, &mut storage as *mut _ as *mut libc::sockaddr, &mut len) }; + if rc < 0 { + return None; + } + let addr = super::super::unix::sockaddr_to_socket_addr(&storage).ok()?; + // An unspecified source means the kernel declined to choose — no route of + // that family — which is not an address and must not be compared as one. + if addr.ip().is_unspecified() { + return None; + } + Some(addr.ip()) } impl AsRawFd for ConnectedPeerSocket { diff --git a/src/transport/watcher.rs b/src/transport/watcher.rs index 9b519e35..724fde66 100644 --- a/src/transport/watcher.rs +++ b/src/transport/watcher.rs @@ -57,6 +57,30 @@ impl Drop for LinkEventSocket { } impl LinkEventSocket { + /// The multicast group mask this socket is actually subscribed to, read + /// back from the kernel rather than remembered from the bind. + /// + /// `getsockname` on a netlink socket fills `sockaddr_nl.nl_groups` with the + /// legacy 32-bit subscription mask, which covers every group in + /// [`groups`]. Reading it back is the only way to tell a watcher that + /// *asked* for the right groups from one that got them: a bind with a + /// wrong mask succeeds just as happily as a bind with the right one, and + /// then silently never delivers the messages the caller subscribed for. + #[cfg(all(test, any(target_os = "linux", target_os = "android")))] + fn bound_groups(&self) -> Option { + let mut sa: libc::sockaddr_nl = unsafe { std::mem::zeroed() }; + let mut len = std::mem::size_of::() as libc::socklen_t; + // SAFETY: `self.fd` is the netlink socket this struct owns, and `sa` / + // `len` are a correctly sized out-parameter pair for `getsockname`. + let rc = unsafe { + libc::getsockname(self.fd, &mut sa as *mut _ as *mut libc::sockaddr, &mut len) + }; + if rc < 0 { + return None; + } + Some(sa.nl_groups) + } + fn recv(&self, buf: &mut [u8]) -> std::io::Result { let n = unsafe { libc::recv(self.fd, buf.as_mut_ptr() as *mut libc::c_void, buf.len(), 0) }; if n < 0 { @@ -218,6 +242,22 @@ impl LinkWatcher { } } + /// The netlink multicast groups this watcher is actually subscribed to, as + /// the kernel reports them. + /// + /// `None` when there is no live source, and on every platform whose backend + /// is `PF_ROUTE`, which has no group selection to report. + /// + /// This exists to be asserted on. A bind with the wrong group mask succeeds + /// exactly like a bind with the right one and then silently never delivers + /// what the caller subscribed for, so nothing short of reading the + /// subscription back can tell the two apart without provoking a real + /// kernel event — which needs privileges CI does not have. + #[cfg(all(test, any(target_os = "linux", target_os = "android")))] + pub(crate) fn subscribed_groups(&self) -> Option { + self.inner.as_ref()?.get_ref().bound_groups() + } + /// Whether an event source is actually backing this watcher. pub fn is_event_driven(&self) -> bool { self.inner.is_some()