diff --git a/CHANGELOG.md b/CHANGELOG.md index 64017af9..87f52319 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -133,6 +133,30 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 #### Library surface and internals +- `Node::netmon_trigger()` and `NetmonTrigger::poke()`, an embedder push into + the medium-change detector + ([#146](https://github.com/jmcorgan/fips/issues/146)). The detector's + channel was private end to end, so an embedder whose platform announces a + network change had nowhere to put it. That is Android's position: the + netlink group bind is refused to an app there, which leaves detection on the + poll timer, while the app's `ConnectivityManager` callback knows the moment + the default network moved. `poke()` is a bare wake-up with no payload — the + detector still samples, and the fingerprint comparison, the debounce, the + settled-back suppression and the minimum spacing between reports apply to a + pushed wake-up exactly as they do to a kernel event, so a spurious callback + costs two samples at the default debounce rather than a socket rebind. It + is synchronous and callable from any thread, before or after `start()`; a + poke that lands while the detector is sampling is held rather than lost, + and a burst in that window coalesces into the one held wake-up (a burst + that finds the detector waiting costs two: the one that woke it and the one + held). A platform callback is not ordered against the routing table it + reports on, so a pushed wake-up whose sample shows nothing moved looks once + more after `node.netmon.debounce_ms` instead of leaving the change to the + timer; wake-ups from the kernel and the timer are unchanged. The push is + selected beside the platform's own source and the timer backstop, never + instead of them. The platform wiring (`registerNetworkCallback`, or + `NWPathMonitor` on iOS) stays with the embedder, since the crate has no JNI + layer. - `TransportError::InterfaceUnavailable { interface }`. A missing interface and a typo'd interface name were previously the same flat `StartFailed(String)`; nothing downstream could branch on absence. diff --git a/src/lib.rs b/src/lib.rs index 8afea9c3..d4675d42 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -103,4 +103,4 @@ pub use peer::{ActivePeer, ConnectivityState, PeerError}; // Re-export node types #[cfg(unix)] pub use node::AppOwnedUdpSocket; -pub use node::{Node, NodeError, NodeState, UpdatePeersOutcome}; +pub use node::{NetmonTrigger, Node, NodeError, NodeState, UpdatePeersOutcome}; diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index d9811dcb..261f1600 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -2002,8 +2002,11 @@ 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, self.entities_snapshot.clone()); + let (rx, task) = crate::node::netmon::spawn_detector( + netmon_cfg, + self.entities_snapshot.clone(), + self.netmon_trigger.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 cc015bf5..d5e9cbcd 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -18,6 +18,7 @@ mod handlers; mod lifecycle; pub(crate) mod metrics; pub(crate) mod netmon; +pub use netmon::NetmonTrigger; mod peer_error_budget; mod peering; mod rate_limit; @@ -439,6 +440,10 @@ pub struct Node { /// removal is one of its callers, so a peer that never comes back leaves /// nothing behind. path_mtu_seeded_by: Arc>>, + /// The embedder's wake-up for the medium-change detector, handed out by + /// [`Node::netmon_trigger`] and into the detector at `start()`. Always + /// allocated so the trigger works whenever it is fetched. + netmon_trigger: netmon::NetmonTrigger, // === Transports & Links === /// Active transports (owned by Node). @@ -951,6 +956,7 @@ impl Node { crate::upper::icmp::mss_ceiling(crate::upper::tun::IPV6_MIN_MTU), )), path_mtu_seeded_by: Arc::new(std::sync::RwLock::new(HashMap::new())), + netmon_trigger: netmon::NetmonTrigger::new(), #[cfg(unix)] decrypt_registered_sessions: std::collections::HashSet::new(), #[cfg(unix)] @@ -1121,6 +1127,7 @@ impl Node { crate::upper::icmp::mss_ceiling(crate::upper::tun::IPV6_MIN_MTU), )), path_mtu_seeded_by: Arc::new(std::sync::RwLock::new(HashMap::new())), + netmon_trigger: netmon::NetmonTrigger::new(), #[cfg(unix)] decrypt_registered_sessions: std::collections::HashSet::new(), #[cfg(unix)] @@ -3689,6 +3696,25 @@ impl Node { self.supervisor.dns_local_addr } + /// A handle that wakes the medium-change detector (`node.netmon.*`) now + /// rather than at its next poll — see [`NetmonTrigger`]. + /// + /// For an embedder whose platform tells it when the network moved but + /// refuses the node its kernel event source: an Android `VpnService` gets + /// a `ConnectivityManager` callback the moment the default network + /// changes, while the netlink group bind the detector would otherwise use + /// is denied to apps, leaving the detector on its poll timer. Poking from + /// the callback turns a poll-period latency into a debounce-period one. + /// Callable before or after [`Self::start`], from any thread, and the + /// same handle keeps working across a [`Self::stop`] and another + /// [`Self::start`]; a poke while no detector is running is held for the + /// next one — and with `node.netmon.enabled: false` there is no next one, + /// so the poke does nothing. A new `Node` has a new trigger — an embedder + /// that rebuilds the node must fetch it again. + pub fn netmon_trigger(&self) -> NetmonTrigger { + self.netmon_trigger.clone() + } + // === Sending === /// Encrypt and send a link-layer message to an authenticated peer. diff --git a/src/node/netmon/mod.rs b/src/node/netmon/mod.rs index 8e3956f8..b67d8f16 100644 --- a/src/node/netmon/mod.rs +++ b/src/node/netmon/mod.rs @@ -36,15 +36,18 @@ //! to catch. `PF_ROUTE` has no group selection and delivers everything //! regardless. //! -//! Still to come, behind the same seam and without touching the handler: -//! `NotifyIpInterfaceChange` on Windows, and an embedder push on iOS. Android -//! takes the netlink source above, which is the right backend when the policy -//! allows the group bind and degrades to the timer when it does not; a -//! `ConnectivityManager` push belongs there too, because a timer is not -//! reliable under Doze. Every platform runs the -//! timer regardless — as the only signal where there is no backend, and as a -//! backstop where there is one, since a kernel event stream can drop messages -//! or stop. +//! An embedder can push a wake-up of its own through [`NetmonTrigger`] +//! ([`Node::netmon_trigger`](crate::Node::netmon_trigger)), which sits +//! *beside* whatever backend the platform has rather than replacing it. That +//! is the Android path: the policy there refuses the netlink group bind for +//! an app, so the kernel source degrades to the timer, while the app's +//! `ConnectivityManager` callback knows the exact moment the default network +//! moved — and a timer is not reliable under Doze. iOS would use the same +//! seam from its path monitor. Still to come, behind the same seam and +//! without touching the handler: `NotifyIpInterfaceChange` on Windows. Every +//! platform runs the timer regardless — as the only signal where there is no +//! backend, and as a backstop where there is one, since a kernel event stream +//! can drop messages or stop. //! //! # What the fingerprint captures //! @@ -552,6 +555,69 @@ struct WakeSource { /// regress, for the cost of one sample per period: five syscalls per probed /// peer, or 640 at the default of 128 peers. timer: tokio::time::Interval, + /// An embedder's push, beside the source above rather than instead of it. + /// Always present: a trigger nobody holds can never fire, so a wake source + /// built without one simply carries one that stays silent. + push: NetmonTrigger, +} + +/// Which arm of [`WakeSource::wait`] returned. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum Woke { + /// The platform source or the timer. + Source, + /// An embedder's [`NetmonTrigger::poke`]. + Push, +} + +/// An embedder's handle for waking the medium-change detector now, rather +/// than at its next poll. +/// +/// Obtained from [`Node::netmon_trigger`](crate::Node::netmon_trigger). +/// [`poke`](Self::poke) is synchronous, cheap, and safe from any thread — it +/// is meant to be called straight from a platform network callback (an +/// Android `ConnectivityManager` one, say). A poke while the detector is +/// mid-sample is not lost: one wake-up is held until the detector next waits, +/// and further pokes in that window coalesce into the one held. A poke that +/// finds the detector waiting wakes it at once, so a burst that starts there +/// costs two wake-ups — the one that woke it and the one held — and never +/// more. A poke while no detector is running — before `start()`, or between +/// a `stop()` and the next `start()` — is held the same way and spent at the +/// next start, where it costs the two samples below, since nothing will have +/// moved since the baseline taken moments earlier. With +/// `node.netmon.enabled: false` no detector ever runs and a poke does +/// nothing. +/// +/// A poke costs a sample — five syscalls per probed peer — and, when that +/// sample shows nothing moved, one more a debounce period later, because a +/// platform callback can run ahead of the routing table it reports on; so a +/// spurious poke costs two samples at the default debounce. Wire it to the +/// callbacks that mean the attachment moved (on Android: `onAvailable`, +/// `onLost` and +/// `onLinkPropertiesChanged` of the default network) rather than to +/// `onCapabilitiesChanged`, which fires every few seconds on cellular for +/// signal and bandwidth estimates and would turn the push into a faster poll. +#[derive(Clone, Debug)] +pub struct NetmonTrigger { + inner: Arc, +} + +impl NetmonTrigger { + pub(crate) fn new() -> Self { + Self { + inner: Arc::new(tokio::sync::Notify::new()), + } + } + + /// Wake the detector: sample the path to every peer now. + pub fn poke(&self) { + self.inner.notify_one(); + } + + /// Wait for a poke, consuming the one held if there is one. + pub(in crate::node) async fn poked(&self) { + self.inner.notified().await; + } } /// Where a wake-up can come from, besides the timer. @@ -581,6 +647,7 @@ impl WakeSource { Self { source: Wake::Timer, timer: Self::make_timer(period), + push: NetmonTrigger::new(), } } @@ -590,6 +657,7 @@ impl WakeSource { Self { source: Wake::Kernel(watcher), timer: Self::make_timer(period), + push: NetmonTrigger::new(), } } @@ -599,6 +667,7 @@ impl WakeSource { Self { source: Wake::Injected(pings), timer: Self::make_timer(period), + push: NetmonTrigger::new(), } } @@ -618,6 +687,12 @@ impl WakeSource { timer } + /// Add an embedder push beside the existing source. See [`NetmonTrigger`]. + fn with_push(mut self, push: NetmonTrigger) -> Self { + self.push = push; + self + } + /// 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`. @@ -629,9 +704,22 @@ impl WakeSource { } } - /// Wait until it is worth sampling again. - async fn wait(&mut self) { - let WakeSource { source, timer } = self; + /// Wait until it is worth sampling again, and say what it was. + async fn wait(&mut self) -> Woke { + // The push is selected beside the platform source, never instead of + // it: an embedder that pokes is a latency improvement on top of the + // timer backstop, and a `Notify` holds one permit for a poke that + // lands while the detector is busy sampling, so it cannot be missed. + let push = self.push.clone(); + tokio::select! { + _ = push.poked() => Woke::Push, + _ = self.wait_source() => Woke::Source, + } + } + + /// Wait on the platform source and the timer alone. + async fn wait_source(&mut self) { + let WakeSource { source, timer, .. } = self; // Only the injected source can stop: [`LinkWatcher`] parks forever // once it gives up, so a kernel source that dies simply stops firing // and the timer carries on underneath it with nothing to unwind here. @@ -689,10 +777,11 @@ impl WakeSource { pub(crate) fn spawn_detector( cfg: NetmonConfig, peers: Arc>, + push: NetmonTrigger, ) -> (NetChangeRx, JoinHandle<()>) { let (tx, rx) = mpsc::channel(1); let handle = tokio::spawn(async move { - let wake = build_wake_source(&cfg); + let wake = build_wake_source(&cfg).with_push(push); let sample = move || NetFingerprint::sample(&probe_targets(&peers.load())); run_detector(tx, cfg, sample, wake).await; }); @@ -787,9 +876,21 @@ where ); loop { - wake.wait().await; + let woke = wake.wait().await; let mut candidate = sample(); + // A push gets one second look. A kernel source re-arms itself — every + // route message in a handover's burst is another wake-up — but an + // embedder's callback is a single shot, and the platform can deliver + // it ahead of the routing table it is about: on Android the + // `ConnectivityManager` callback is not ordered against netd + // finishing the default-route swap. Without this the poke would be + // spent on a sample of the old picture and the change left to the + // timer, which is the latency the push exists to remove. + if woke == Woke::Push && !debounce.is_zero() && last.moved(&candidate).is_empty() { + tokio::time::sleep(debounce).await; + candidate = sample(); + } 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 diff --git a/src/node/netmon/tests.rs b/src/node/netmon/tests.rs index e2ae35dd..36a1e234 100644 --- a/src/node/netmon/tests.rs +++ b/src/node/netmon/tests.rs @@ -78,6 +78,28 @@ fn scripted(samples: Vec) -> (impl Fn() -> NetFingerprint, Arc (impl Fn() -> NetFingerprint, Arc) { + let calls = Arc::new(AtomicUsize::new(0)); + let counter = calls.clone(); + let epoch = tokio::time::Instant::now(); + let sampler = move || { + counter.fetch_add(1, Ordering::SeqCst); + if epoch.elapsed() < at { + before.clone() + } else { + after.clone() + } + }; + (sampler, calls) +} + #[tokio::test(start_paused = true)] async fn steady_attachment_reports_nothing() { let steady = all_from(&[peer(1), peer(2)], Some(v4(192, 168, 1, 10))); @@ -598,6 +620,163 @@ async fn an_event_ping_wakes_the_detector_before_the_timer_would() { assert_eq!(change.summary.moved[0].after, Some(v4(10, 40, 0, 7))); } +/// An embedder's push is the Android path: the netlink group bind is refused +/// for an app there, so the platform source is the timer alone, and the app's +/// own network callback is what knows the moment the medium moved. The poll +/// period here is an hour, so only the poke can be what woke the detector. +#[tokio::test(start_paused = true)] +async fn an_embedder_poke_wakes_the_detector_before_the_timer_would() { + 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 trigger = NetmonTrigger::new(); + + let wake = WakeSource::timer_only(Duration::from_secs(3600)).with_push(trigger.clone()); + tokio::spawn(run_detector(tx, cfg(3600, 0), sampler, wake)); + // Let the detector take its baseline and park on `wait()`; a poke that + // lands before that must also work (the permit is held), but the common + // case is the parked one. + tokio::task::yield_now().await; + + trigger.poke(); + + let change = expect_change(&mut rx).await; + assert_eq!(change.summary.moved[0].after, Some(v4(10, 40, 0, 7))); +} + +/// A poke that lands while the detector is not waiting — here, before it has +/// even started — is held as a permit rather than dropped, so the embedder +/// never has to sequence its callback against the detector's own loop. +#[tokio::test(start_paused = true)] +async fn a_poke_before_the_detector_waits_is_not_lost() { + 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, calls) = scripted(vec![wlan, cell]); + let (tx, mut rx) = mpsc::channel(1); + let trigger = NetmonTrigger::new(); + trigger.poke(); + trigger.poke(); // nothing is waiting, so the burst is held as one wake-up + + let wake = WakeSource::timer_only(Duration::from_secs(3600)).with_push(trigger.clone()); + tokio::spawn(run_detector(tx, cfg(3600, 0), sampler, wake)); + + let change = expect_change(&mut rx).await; + assert_eq!(change.summary.moved[0].after, Some(v4(10, 40, 0, 7))); + + // The baseline and the one sample the burst bought. A second wake-up for + // the second poke would show up here as a third. + expect_quiet(&mut rx, "the second poke of a burst must not wake it again").await; + assert_eq!(calls.load(Ordering::SeqCst), 2); +} + +/// A platform callback is not ordered against the routing table it reports +/// on: the poke can land while the kernel still answers with the old source. +/// A kernel source would be woken again by the route message itself, but a +/// push is a single shot, so the detector looks once more a debounce period +/// later rather than spending the poke on the old picture and leaving the +/// change to an hour-long timer. The swap lands on the paused clock 100 ms +/// after the poke, so this passes only if the second look really waits: a +/// second look taken at once would still see the old source. +#[tokio::test(start_paused = true)] +async fn a_poke_that_runs_ahead_of_the_route_change_still_catches_it() { + 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, calls) = switching_at(wlan, cell, Duration::from_millis(100)); + let (tx, mut rx) = mpsc::channel(1); + let trigger = NetmonTrigger::new(); + + let wake = WakeSource::timer_only(Duration::from_secs(3600)).with_push(trigger.clone()); + tokio::spawn(run_detector(tx, cfg(3600, 250), sampler, wake)); + tokio::task::yield_now().await; + + trigger.poke(); + + let change = expect_change(&mut rx).await; + assert_eq!(change.summary.moved[0].after, Some(v4(10, 40, 0, 7))); + // The baseline, the sample the poke bought (still the old picture), the + // second look that caught the swap, and the one debounce round that + // found it settled. + assert_eq!(calls.load(Ordering::SeqCst), 4); +} + +/// A poke when nothing moved is the spurious-callback case, and the sampler +/// underneath the push is what keeps it from becoming a spurious socket +/// rebind: the detector samples, takes its second look, and reports nothing. +/// Those two samples are the whole cost. +#[tokio::test(start_paused = true)] +async fn a_poke_when_nothing_moved_reports_nothing() { + let wlan = all_from(&[peer(1)], Some(v4(192, 168, 1, 10))); + let (sampler, calls) = scripted(vec![wlan]); + let (tx, mut rx) = mpsc::channel(1); + let trigger = NetmonTrigger::new(); + + let wake = WakeSource::timer_only(Duration::from_secs(3600)).with_push(trigger.clone()); + tokio::spawn(run_detector(tx, cfg(3600, 250), sampler, wake)); + tokio::task::yield_now().await; + + trigger.poke(); + + expect_quiet( + &mut rx, + "a poke that finds nothing moved must report nothing", + ) + .await; + // The baseline, then the sample the poke bought and its second look. + assert_eq!(calls.load(Ordering::SeqCst), 3); +} + +/// The push sits beside the platform source, not instead of it: on a host +/// whose kernel source did open, the kernel's own event and an embedder's +/// poke must each still wake the detector. +#[tokio::test(start_paused = true)] +async fn a_push_beside_an_event_source_leaves_both_live() { + 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 wired = all_from(&[peer(1)], Some(v4(172, 16, 0, 5))); + let (sampler, _) = scripted(vec![wlan, cell, wired]); + let (tx, mut rx) = mpsc::channel(1); + let (pings, ping_rx) = mpsc::channel(1); + let trigger = NetmonTrigger::new(); + + let wake = WakeSource::events(ping_rx, Duration::from_secs(3600)).with_push(trigger.clone()); + tokio::spawn(run_detector(tx, cfg(3600, 0), sampler, wake)); + tokio::task::yield_now().await; + + pings.send(()).await.expect("the backend can ping"); + let change = expect_change(&mut rx).await; + assert_eq!( + change.summary.moved[0].after, + Some(v4(10, 40, 0, 7)), + "the kernel's event must still wake it" + ); + + trigger.poke(); + let change = expect_change(&mut rx).await; + assert_eq!( + change.summary.moved[0].after, + Some(v4(172, 16, 0, 5)), + "the embedder's poke must wake it too" + ); +} + +/// The second look belongs to a push alone. A timer tick that finds nothing +/// moved must go back to sleep, or every quiet poll period would cost two +/// samples instead of one. +#[tokio::test(start_paused = true)] +async fn a_timer_wake_that_finds_nothing_does_not_take_a_second_look() { + let wlan = all_from(&[peer(1)], Some(v4(192, 168, 1, 10))); + let (sampler, calls) = scripted(vec![wlan]); + let (tx, mut rx) = mpsc::channel(1); + + let wake = WakeSource::timer_only(Duration::from_secs(20)); + tokio::spawn(run_detector(tx, cfg(20, 250), sampler, wake)); + + // One tick inside the quiet window: the baseline plus one sample. + expect_quiet(&mut rx, "nothing moved").await; + assert_eq!(calls.load(Ordering::SeqCst), 2); +} + /// A netlink socket drops messages under memory pressure, and a backend can go /// quiet without going away. The backstop timer must still get the node there, /// so an event-driven backend is never worse than the poller it replaced. diff --git a/src/node/tests/netmon.rs b/src/node/tests/netmon.rs index d9884d0f..c7c5c6c8 100644 --- a/src/node/tests/netmon.rs +++ b/src/node/tests/netmon.rs @@ -635,3 +635,44 @@ async fn a_transports_bind_address_reaches_the_probe_target() { cleanup_nodes(&mut nodes).await; } + +/// The detector's own tests attach a trigger to a wake source by hand, so they +/// would all still pass if `start()` handed the detector some other trigger +/// than the one [`Node::netmon_trigger`] gives out. The observable here is the +/// held poke: a running detector wired to this trigger consumes it, and +/// anything else leaves it sitting there for the test to collect. Twice, +/// because the same handle has to reach the detector a `stop()` and another +/// `start()` later as well. +#[tokio::test] +async fn the_nodes_trigger_reaches_the_running_detector_across_a_restart() { + let mut node = make_healthy_node(); + // Fetched before `start()`, which is when an embedder registering its + // platform callback early would fetch it. + let trigger = node.netmon_trigger(); + + for round in ["first start", "restart"] { + node.start().await.unwrap(); + trigger.poke(); + + let consumed = tokio::time::timeout(Duration::from_secs(5), async { + loop { + tokio::time::sleep(Duration::from_millis(20)).await; + let still_held = + tokio::time::timeout(Duration::from_millis(1), trigger.poked()).await; + match still_held { + // Nobody took it, and collecting it just consumed it: put + // it back and give the detector another turn. + Ok(()) => trigger.poke(), + Err(_) => break, + } + } + }) + .await; + assert!( + consumed.is_ok(), + "{round}: the detector never took the poke, so it is not listening on the node's trigger" + ); + + node.stop().await.unwrap(); + } +}