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()