diff --git a/src/node/dataplane/rx_loop.rs b/src/node/dataplane/rx_loop.rs index e4574860..7e81268d 100644 --- a/src/node/dataplane/rx_loop.rs +++ b/src/node/dataplane/rx_loop.rs @@ -359,20 +359,27 @@ impl Node { // latch. maybe_presence = presence_rx.recv() => { if let Some(edge) = maybe_presence { - let child = crate::node::lifecycle::supervisor::Child::Transport( - edge.transport_id, - ); - let event = if edge.present { - crate::node::lifecycle::supervisor::Event::ChildPresent { child } - } else { - crate::node::lifecycle::supervisor::Event::ChildAbsent { child } - }; - let actions = self.supervisor.fsm.step(event); - for action in actions { - if let crate::node::lifecycle::supervisor::Action::PublishState(ns) = - action - { - self.supervisor.state = ns; + // Health is policy-filtered; the MTU floor below is + // not. An `optional` interface's absence is normal and + // must not move the node off `Full`, but it changes + // the bound set all the same. + if edge.health_relevant { + let child = crate::node::lifecycle::supervisor::Child::Transport( + edge.transport_id, + ); + let event = if edge.present { + crate::node::lifecycle::supervisor::Event::ChildPresent { child } + } else { + crate::node::lifecycle::supervisor::Event::ChildAbsent { child } + }; + let actions = self.supervisor.fsm.step(event); + for action in actions { + if let crate::node::lifecycle::supervisor::Action::PublishState( + ns, + ) = action + { + self.supervisor.state = ns; + } } } // The bound set just changed, so the node's egress MTU diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index 37039acf..8afefb13 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -2290,6 +2290,11 @@ impl Node { let mut published = None; let saw_edge = !edges.is_empty(); for edge in edges { + // Health is policy-filtered; the MTU refresh below is not. See + // `saw_edge`. + if !edge.health_relevant { + continue; + } let child = Child::Transport(edge.transport_id); let event = if edge.present { Event::ChildPresent { child } diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index e2aa36e4..92cf8d01 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -1929,6 +1929,7 @@ async fn a_presence_edge_refreshes_the_tun_mss_ceiling_without_being_asked() { .send(crate::transport::TransportPresence { transport_id: TransportId::new(1), present: true, + health_relevant: true, }) .await .expect("presence edge queued"); diff --git a/src/transport/ethernet/io.rs b/src/transport/ethernet/io.rs index 39d091c8..fd7e30df 100644 --- a/src/transport/ethernet/io.rs +++ b/src/transport/ethernet/io.rs @@ -39,6 +39,17 @@ pub const ETHERNET_BROADCAST: [u8; 6] = [0xff; 6]; /// spelled the same on Linux and the BSDs. #[cfg(unix)] pub fn interface_present(interface: &str) -> bool { + interface_present_probe(interface).unwrap_or(false) +} + +/// [`interface_present`], keeping "the probe failed" distinct from "absent". +/// +/// `None` means the kernel would not answer. A caller deciding whether to +/// *bind* can treat that as absence and retry on the next tick, which is what +/// [`interface_present`] does. A caller deciding whether to *unbind* must not: +/// see [`interface_has_flags`]. +#[cfg(unix)] +pub fn interface_present_probe(interface: &str) -> Option { interface_has_flags(interface, libc::IFF_UP as u32) } @@ -67,19 +78,31 @@ pub fn interface_index(interface: &str) -> Option { /// from being told apart here; presence answers that. #[cfg(unix)] pub fn interface_carrier(interface: &str) -> bool { - interface_has_flags(interface, (libc::IFF_UP | libc::IFF_RUNNING) as u32) + // Report-only, so a probe failure reads the same as no carrier. + interface_has_flags(interface, (libc::IFF_UP | libc::IFF_RUNNING) as u32).unwrap_or(false) } -/// Whether the named interface exists and has every flag in `wanted` set. +/// Whether the named interface exists and has every flag in `wanted` set, or +/// `None` if the question could not be asked. +/// +/// The `None` matters. `getifaddrs` is a netlink dump on Linux and it does +/// fail for reasons that have nothing to do with the interface — `ENOBUFS` +/// under memory pressure or a busy netlink socket, `EMFILE`/`ENFILE` under fd +/// exhaustion, since it opens a socket of its own. Answering `false` there +/// reports a present interface as gone, and a caller holding a live binding +/// would tear a working socket down over a transient syscall failure. Callers +/// that can tell the two apart should. #[cfg(unix)] -fn interface_has_flags(interface: &str, wanted: u32) -> bool { +fn interface_has_flags(interface: &str, wanted: u32) -> Option { let Ok(c_name) = std::ffi::CString::new(interface) else { - return false; + // An interior NUL is not a probe failure — no such interface can + // exist, and no retry will change that. + return Some(false); }; let mut addrs: *mut libc::ifaddrs = std::ptr::null_mut(); if unsafe { libc::getifaddrs(&mut addrs) } != 0 { - return false; + return None; } let mut matched = false; @@ -97,7 +120,7 @@ fn interface_has_flags(interface: &str, wanted: u32) -> bool { } unsafe { libc::freeifaddrs(addrs) }; - matched + Some(matched) } // Platform-specific PacketSocket implementation. @@ -284,6 +307,12 @@ mod async_impl { /// A received frame: (payload, source_mac). type Frame = (Vec, [u8; 6]); + /// Consecutive failed BPF reads before the reader thread gives up. + /// + /// Mirrors the receive loop's own error threshold: the point is not to + /// tolerate errors but to end the task so the binder can rebind. + const READ_ERROR_EXIT_THRESHOLD: u32 = 5; + pub struct AsyncPacketSocket { inner: Arc, /// `None` once shutdown has taken the receiver, which is what makes @@ -310,6 +339,9 @@ mod async_impl { let mut parse_buf = vec![0u8; bpf_buflen]; let mut parse_offset: usize = 0; let mut parse_len: usize = 0; + // Consecutive failed reads, to bound a socket whose + // interface went away underneath it. + let mut read_errors: u32 = 0; let nfds = bpf_fd.max(shutdown_fd) + 1; loop { @@ -377,11 +409,31 @@ mod async_impl { if err.raw_os_error() == Some(libc::EBADF) { break; } + if err.kind() == std::io::ErrorKind::Interrupted { + continue; + } + } + // Anything else — `ENXIO` is the one that matters, + // which is what BPF answers once the interface it + // was attached to is torn away — used to loop here + // forever. That mattered beyond the spin: the + // binder's detach check asks whether this thread is + // still running, so a thread that never returns + // reports a dead socket as a live one, and the + // transport sits `present` and deaf until the name + // or index happens to change too. Give up after a + // streak and let the return close the channel, + // which fails `recv_from`, which ends the tokio + // task the binder is actually watching. + read_errors += 1; + if read_errors >= READ_ERROR_EXIT_THRESHOLD { + break; } parse_len = 0; parse_offset = 0; continue; } + read_errors = 0; parse_len = ret as usize; parse_offset = 0; } diff --git a/src/transport/ethernet/mod.rs b/src/transport/ethernet/mod.rs index cebcc4ff..64667500 100644 --- a/src/transport/ethernet/mod.rs +++ b/src/transport/ethernet/mod.rs @@ -30,7 +30,9 @@ use super::{ TransportId, TransportPresence, TransportState, TransportType, }; use crate::config::EthernetConfig; -use io::{AsyncPacketSocket, ETHERNET_BROADCAST, PacketSocket, interface_present}; +use io::{ + AsyncPacketSocket, ETHERNET_BROADCAST, PacketSocket, interface_present, interface_present_probe, +}; use neighbor::{FRAME_TYPE_BEACON, FRAME_TYPE_DATA, NeighborBuffer, build_beacon, parse_beacon}; use presence::{ABSENCE_ERROR_AFTER, ChurnGuard, PresenceState, bind_backoff}; use stats::EthernetStats; @@ -369,6 +371,10 @@ impl EthernetTransport { // previous run before the binder can observe it. self.shutdown.store(false, Ordering::SeqCst); + // The bring-up window is measured from here, not from whenever this + // object happened to be constructed. + self.presence.mark_starting(); + let ctx = Arc::new(BinderContext { shutdown: self.shutdown.clone(), transport_id: self.transport_id, @@ -388,14 +394,22 @@ impl EthernetTransport { // First attempt inline, so the common case (interface already there) // keeps its ordering: the transport is bound and logged before // `start_async` returns, exactly as before this mechanism existed. + // An edge the channel refused here is handed to the binder to retry, + // rather than dropped: it is the only edge either consumer will ever + // see for this transport until the interface next changes state. + let mut initial_unpublished: Option = None; match bind_and_spawn(&ctx).await { Ok(()) => { - publish_presence(&ctx, true); + if !publish_presence(&ctx, true) { + initial_unpublished = Some(true); + } } Err(TransportError::InterfaceUnavailable { .. }) => { // Absence is a state, not a start failure. Come up and wait. log_initial_absence(&ctx); - publish_presence(&ctx, false); + if !publish_presence(&ctx, false) { + initial_unpublished = Some(false); + } } // Anything else — no CAP_NET_RAW, no free BPF device, a buffer // the kernel refused — is a fault, not a state, and it will not @@ -411,7 +425,7 @@ impl EthernetTransport { let watcher_ctx = ctx.clone(); self.binder_task = Some(tokio::spawn(async move { - binder_loop(watcher_ctx).await; + binder_loop(watcher_ctx, initial_unpublished).await; })); self.state = TransportState::Up; @@ -441,9 +455,21 @@ impl EthernetTransport { // aborted tasks are deliberately not awaited: on a current_thread // runtime an aborted task cannot be polled while we are blocked on its // `JoinHandle`, which deadlocks. + let was_present = self.presence.is_present(); self.binding.tear_down(); self.presence.transition(Presence::Absent); + // Retract the edge the binder left standing. Every other transition + // pairs its edges; a transport stopped while `Present` would otherwise + // leave health reading `present: true` for a socket that is gone. + if was_present && let Some(tx) = &self.presence_tx { + let _ = tx.try_send(TransportPresence { + transport_id: self.transport_id, + present: false, + health_relevant: !self.policy.is_optional(), + }); + } + self.state = TransportState::Down; info!( @@ -805,20 +831,25 @@ fn log_initial_absence(ctx: &Arc) { /// it carries. The caller retries a refused edge on its next tick, so a /// momentarily full channel costs latency rather than correctness. /// -/// An `optional` interface publishes nothing, and reports success for it. +/// An `optional` interface publishes its edge with `health_relevant: false`. /// That is the entire health half of the policy: absence of an interface -/// whose absence is normal must not move the node off `Full`, and since -/// nothing was ever published, its return has nothing to clear. +/// whose absence is normal must not move the node off `Full`. +/// +/// It is deliberately *not* silence. The edge still goes out, because the +/// consumer derives two things from it and only one of them is about health: +/// the node's egress MTU floor is a function of the bound set, and an +/// `optional` interface binding or unbinding changes that set exactly as a +/// `required` one does. Filtering the edge here — which is what this used to +/// do — left the TUN MSS clamp stale for every `optional` transport, which on +/// the shipped OpenWrt config is five of seven. fn publish_presence(ctx: &Arc, present: bool) -> bool { - if ctx.policy.is_optional() { - return true; - } let Some(tx) = &ctx.presence_tx else { return true; }; match tx.try_send(TransportPresence { transport_id: ctx.transport_id, present, + health_relevant: !ctx.policy.is_optional(), }) { Ok(()) => true, Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => { @@ -837,7 +868,7 @@ fn publish_presence(ctx: &Arc, present: bool) -> bool { /// fallback; on platforms with a link-event source the same loop is driven by /// events instead (see `docs`/the netlink watcher), and the interval below is /// what "immediate" degrades to without one. -async fn binder_loop(ctx: Arc) { +async fn binder_loop(ctx: Arc, initial_unpublished: Option) { let watcher = LinkWatcher::new(); let mut ticker = tokio::time::interval(WATCH_INTERVAL); ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); @@ -857,10 +888,22 @@ async fn binder_loop(ctx: Arc) { let mut absence_reported = false; // Damping for bindings that keep succeeding and then dying young. let mut churn = ChurnGuard::new(); + // `start_async` binds inline and publishes that edge itself, outside this + // guard. Seed the guard from that bind, or the guard believes it has + // announced nothing: `detached` takes `announced` to decide whether an + // edge is owed, so the first detach after a clean start would retract + // nothing and health would keep reading `Full` for as long as the + // interface stayed away. `stabilized` cannot repair it either — it + // early-returns while `bound_at` is `None`, which it is for a bind this + // loop did not perform. + if ctx.presence.is_present() { + churn.bound(Instant::now()); + } // A presence edge the channel refused, to retry on the next tick. Health // is a level, so only the most recent value matters — an older pending - // edge is superseded rather than queued. - let mut unpublished: Option = None; + // edge is superseded rather than queued. Seeded with the start edge if + // that one was refused. + let mut unpublished: Option = initial_unpublished; // When the presence probe last ran, for the coalescing floor. let mut last_probe = Instant::now(); @@ -890,7 +933,23 @@ async fn binder_loop(ctx: Arc) { if ctx.presence.is_present() { // Detach is either the interface going away or the socket under it // dying while the name stays (a recreated veth, a reloaded phy). - let gone = !interface_present(&ctx.interface); + // + // A probe that could not run is not an interface that went away. + // Hold the binding and re-ask next tick, rather than tearing a + // working socket down because `getifaddrs` hit `ENOBUFS` or the + // process ran out of descriptors — the latter being a state the + // rebind could not recover from anyway. + let gone = match interface_present_probe(&ctx.interface) { + Some(present) => !present, + None => { + debug!( + transport_id = %ctx.transport_id, + interface = %ctx.interface, + "Interface presence probe failed; holding the binding" + ); + false + } + }; let replaced = !gone && ctx.binding.device_replaced(&ctx.interface); let dead = !ctx.binding.tasks_alive(); if !gone && !replaced && !dead { @@ -991,12 +1050,20 @@ async fn binder_loop(ctx: Arc) { } Err(e) => { // `bind_and_spawn` has already reverted presence to Absent. - let attempts = ctx.presence.record_attempt(); if matches!(e, TransportError::InterfaceUnavailable { .. }) { // Raced with a detach between the probe and the bind. No // backoff and no log: the next tick re-probes. + // + // And no `record_attempt` either — the counter below gates + // the only log a genuine bind fault ever gets, on + // `attempts == 1`. Charging absence races to it means a + // flap followed by a real fault (CAP_NET_RAW dropped, no + // free BPF device) is never reported at all, and the + // operator gets `report_sustained_absence`'s "still + // missing" instead — for an interface that is present. continue; } + let attempts = ctx.presence.record_attempt(); if attempts == 1 { error!( transport_id = %ctx.transport_id, @@ -1493,24 +1560,55 @@ mod tests { } #[tokio::test] - async fn an_optional_interface_publishes_no_health_edge() { + async fn an_optional_interface_publishes_an_edge_that_does_not_move_health() { // `optional: true` is a statement about presence, and its whole health // effect is this: a dock adapter that is not plugged in must not make // the node report Degraded. + // + // It is not a statement about the *edge*. The edge still goes out, + // carrying `health_relevant: false`, because the consumer derives the + // node's egress MTU floor from the same channel and an optional + // interface changes the bound set exactly as a required one does. + // Suppressing the edge here — which is what this used to do — left the + // TUN MSS clamp stale for every optional transport. let (mut eth, _rx) = absent_transport(true); let (tx, mut presence_rx) = tokio::sync::mpsc::channel(4); eth.set_presence_tx(tx); eth.start_async().await.expect("start"); + let edge = presence_rx + .try_recv() + .expect("an optional interface still publishes its edge"); + assert_eq!(edge.transport_id, TransportId::new(1)); + assert!(!edge.present); assert!( - presence_rx.try_recv().is_err(), + !edge.health_relevant, "an optional interface must not report absence to node health" ); eth.stop_async().await.expect("stop"); } + #[tokio::test] + async fn a_required_interface_publishes_an_edge_that_moves_health() { + // The other half of the policy, pinned alongside it so the pair cannot + // drift into agreeing. + let (mut eth, _rx) = absent_transport(false); + let (tx, mut presence_rx) = tokio::sync::mpsc::channel(4); + eth.set_presence_tx(tx); + + eth.start_async().await.expect("start"); + + let edge = presence_rx + .try_recv() + .expect("an absence edge is published"); + assert!(!edge.present); + assert!(edge.health_relevant); + + eth.stop_async().await.expect("stop"); + } + #[tokio::test] async fn the_binder_stops_with_the_transport() { // Teardown must abort the binder first: a rebind racing a stop would @@ -1887,6 +1985,7 @@ mod tests { tx.try_send(TransportPresence { transport_id: TransportId::new(1), present: true, + health_relevant: true, }) .expect("first slot"); diff --git a/src/transport/ethernet/presence.rs b/src/transport/ethernet/presence.rs index 33269e41..c79c738c 100644 --- a/src/transport/ethernet/presence.rs +++ b/src/transport/ethernet/presence.rs @@ -198,6 +198,18 @@ impl PresenceState { read(&self.since).elapsed() } + /// Restart the episode clock, for a transport that is about to start. + /// + /// `new()` stamps the clock at construction, but construction and + /// `start_async` need not be adjacent — config load and supervisor staging + /// sit between them. Left alone, a transport staged for longer than + /// [`ABSENCE_ERROR_AFTER`] logs the sustained-absence error on its very + /// first binder tick, having given the interface no bring-up window at + /// all. The window is supposed to absorb exactly that race. + pub fn mark_starting(&self) { + *write(&self.since) = Instant::now(); + } + /// Successful binds since creation (`1` after a clean start). pub fn binds(&self) -> u64 { self.binds.load(Ordering::Relaxed) @@ -584,6 +596,48 @@ mod tests { assert_eq!(g.streak(), 0); } + #[test] + fn an_unseeded_guard_retracts_nothing_and_never_repairs_itself() { + // Why `binder_loop` seeds the guard when it inherits a binding from + // `start_async`, rather than leaving it fresh. + // + // A guard that was never told about a bind believes it has announced + // nothing, so it asks for no retraction — and health, which learned + // `present: true` from the inline bind, would keep reading `Full` with + // the interface gone. `stabilized` cannot rescue it either: with no + // `bound_at` there is nothing for it to judge stable. + let mut g = ChurnGuard::new(); + let t0 = Instant::now(); + + assert!( + !g.stabilized(t0 + MIN_STABLE_BINDING + Duration::from_secs(60)), + "a guard with no recorded bind has nothing to stabilize" + ); + + let out = g.detached(t0 + Duration::from_secs(60)); + assert!( + !out.retract, + "an unseeded guard retracts nothing — which is exactly why the \ + binder must seed it from the inline bind" + ); + } + + #[test] + fn a_seeded_guard_retracts_the_edge_the_inline_bind_published() { + // The fix, from the binder's angle: seeding with `bound` is what makes + // the first detach after a clean start reach node health. + let mut g = ChurnGuard::new(); + let t0 = Instant::now(); + g.bound(t0); + + let out = g.detached(t0 + MIN_STABLE_BINDING + Duration::from_secs(60)); + assert!( + out.retract, + "the edge `start_async` published must be retracted on detach" + ); + assert!(out.log_edge); + } + #[test] fn short_lived_bindings_back_off() { // The failure this guards: a receive loop that gives up on a diff --git a/src/transport/ethernet/watcher.rs b/src/transport/ethernet/watcher.rs index 8dbeacc2..5e00b6bf 100644 --- a/src/transport/ethernet/watcher.rs +++ b/src/transport/ethernet/watcher.rs @@ -22,7 +22,7 @@ //! yields a watcher that never fires, and the binder degrades to its poll. use std::os::unix::io::{AsRawFd, RawFd}; -use std::sync::atomic::{AtomicU32, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; use std::time::Duration; use tokio::io::unix::AsyncFd; @@ -123,6 +123,15 @@ pub(crate) struct LinkWatcher { inner: Option>, /// Consecutive receive errors. Reset by any successful read. errors: AtomicU32, + /// Set once the source has been abandoned for good. + /// + /// Abandonment has to outlive the future that decided it. `changed()` is + /// called fresh on every pass of the caller's `select!` and dropped + /// whenever the poll ticker wins, so a `pending()` inside that future + /// parks nothing beyond the current pass — without this flag the next pass + /// re-reads the dead socket, re-counts the error, and re-logs the + /// give-up warning, once per wake-up, forever. + given_up: AtomicBool, } impl LinkWatcher { @@ -144,6 +153,7 @@ impl LinkWatcher { Self { inner, errors: AtomicU32::new(0), + given_up: AtomicBool::new(false), } } @@ -162,6 +172,14 @@ impl LinkWatcher { unreachable!("pending never resolves") }; + // Already abandoned on an earlier pass. Park without touching the + // socket, so giving up costs one syscall in total rather than one per + // caller wake-up for the life of the process. + if self.given_up.load(Ordering::Relaxed) { + std::future::pending::<()>().await; + unreachable!("pending never resolves") + } + loop { let Ok(mut guard) = afd.readable().await else { // The registration died. Stop firing rather than spinning; the @@ -178,8 +196,19 @@ impl LinkWatcher { loop { match guard.try_io(|inner| inner.get_ref().recv(&mut buf)) { Ok(Ok(n)) if n > 0 => saw_event = true, - // A zero-length read: nothing more to take this round. - Ok(Ok(_)) => break, + // A zero-length read. Readiness is *not* cleared by + // `try_io` here — it clears only on `WouldBlock` — so + // breaking out plainly would leave `readable()` instantly + // ready with nothing to read, and this loop would spin + // without ever returning `Pending`. That starves the + // caller's `select!` of its poll ticker entirely, which + // takes presence detection down with it. Clear it by hand + // and treat it as a fault, so the give-up path applies. + Ok(Ok(_)) => { + guard.clear_ready(); + failure = Some(std::io::Error::from(std::io::ErrorKind::UnexpectedEof)); + break; + } // A genuine socket error. Distinct from WouldBlock, and // the distinction is the whole point: `try_io` clears // readiness only on WouldBlock, so breaking out of a real @@ -211,6 +240,7 @@ impl LinkWatcher { warn!(error = %e, "Link event source read failed"); } if errors >= ERROR_GIVE_UP { + self.given_up.store(true, Ordering::Relaxed); warn!( errors, "Link event source is not recoverable; falling back to \ @@ -261,6 +291,7 @@ mod tests { let watcher = LinkWatcher { inner: Some(AsyncFd::new(LinkEventSocket { fd: read_fd }).expect("register")), errors: AtomicU32::new(0), + given_up: AtomicBool::new(false), }; // Keep producing readiness edges. The error arm calls `clear_ready`, @@ -298,6 +329,7 @@ mod tests { let w = LinkWatcher { inner: None, errors: AtomicU32::new(0), + given_up: AtomicBool::new(false), }; let fired = tokio::time::timeout(std::time::Duration::from_millis(50), w.changed()).await; assert!(fired.is_err(), "sourceless watcher resolved"); diff --git a/src/transport/mod.rs b/src/transport/mod.rs index 9217450b..87a5745a 100644 --- a/src/transport/mod.rs +++ b/src/transport/mod.rs @@ -142,6 +142,14 @@ pub struct TransportPresence { pub transport_id: TransportId, /// Whether the interface is now bound. pub present: bool, + /// Whether this edge should move node health. + /// + /// `false` for an `optional` interface, whose absence is normal and must + /// not take the node off `Full`. The edge is still published, because the + /// bound set changed either way and the node's egress MTU floor is derived + /// from it — health and MTU are two different questions riding one + /// channel, and only the first one is policy-filtered. + pub health_relevant: bool, } /// Channel sender for transport presence edges.