fix(transport): pair every presence edge the binder publishes

Six defects, all in the wiring around the presence machine rather than in
the machine itself. The state machines were well tested; what went untested
was how the binder feeds them, and every one of these lives there.

**The first detach after a clean start never reached node health.**
start_async binds inline and publishes that edge itself, before binder_loop
exists. The loop then built a fresh ChurnGuard, which had therefore recorded
no bind — and detached() takes `announced` to decide whether an edge is
owed, so it asked for no retraction. A node that booted with its interface
present and then lost it kept reporting Full for as long as the interface
stayed away, while show_transports said absent and the log said detached.
stabilized() could not repair it either: it early-returns while bound_at is
None, which it is for a bind the loop did not perform, so the guard stayed
unseeded until a second detach happened to fix it. The loop now seeds the
guard from the binding it inherits.

This is not the cable-unplug case. Presence is IFF_UP, so a carrier loss is
correctly not a detach at all; it takes an admin down, a netdev delete, or a
device removal — a wifi reload, a hostapd restart, a dongle pulled.

**An optional interface published nothing, and the MTU floor rode the same
channel.** Filtering the edge at the source conflated two questions that
happen to share a transport: whether node health should move, and whether
the bound set changed. Only the first is policy. The second determines the
node's egress MTU floor, and refresh_tun_mss_ceiling has no other trigger —
so an optional transport binding or unbinding at runtime left the TUN MSS
clamp derived from a bound set that no longer existed, reporting one
effective MTU and clamping to another. That is the defect the MssCeiling
work was written to remove, reintroduced for exactly the transports the
shipped OpenWrt config marks optional: five of seven. TransportPresence
gains health_relevant; the edge always goes out and only health is filtered.

**A permanent bind fault went unlogged if any absence race preceded it.**
record_attempt ran ahead of the InterfaceUnavailable arm, so a probe that
won a race the bind then lost burned the counter that gates the only error
a real fault ever gets — emitted on attempts == 1. A flapping interface
followed by CAP_NET_RAW being dropped or /dev/bpf* exhausting therefore
reported nothing at all, for the life of the process, and the operator got
report_sustained_absence's "still missing" instead: wrong, for an interface
that is present. An absence race is not a bind attempt and no longer counts
as one.

**The watcher's give-up did not stick.** Abandoning the event source parked
on pending() *inside the changed() future*, and that future is constructed
fresh on every pass of the binder's select! and dropped whenever the poll
ticker wins. So the next pass re-read the dead socket, re-counted the error
and re-logged "not recoverable" — once per wake-up, forever, which at a 1 s
tick across the shipped seven-transport config is seven warnings and seven
failing syscalls a second on flash-backed logging. The 100 ms backoff lived
in the dropped future too and never applied across passes. Abandonment is
now state on the watcher.

**A zero-length read livelocked the binder.** try_io clears readiness only
on WouldBlock, which the sibling error arm handles by hand and this one did
not — so breaking out left readable() instantly ready with nothing to read,
and the loop never returned Pending. That starves the select! of its ticker
entirely and takes presence detection down with it. Cleared and treated as a
fault so the give-up path applies.

**A failed presence probe read as an absent interface.** getifaddrs is a
netlink dump and fails for reasons that have nothing to do with the
interface: ENOBUFS under memory pressure, EMFILE or ENFILE under fd
exhaustion, since it opens a socket of its own. Answering "not present"
there tore down a working socket and degraded health over a transient
syscall failure, undiagnosably, and under fd exhaustion the rebind could not
have succeeded anyway. interface_has_flags now distinguishes the two; the
detach gate holds its binding when the kernel will not answer, while the
bind gate still treats it as absence and retries.

**On macOS the reader thread could not die, so a dead socket read as live.**
Any read() failure other than EBADF — ENXIO being the one that matters,
which is what BPF answers once the interface it was attached to is torn away
— reset the parse buffer and continued. The thread never returned, so the
channel never closed, so recv_from never failed, so the tokio task never
exited, so tasks_alive() reported a dead socket as a live one. Detach
detection on macOS reduced to the name and the index, and an interface reset
in place left the transport present and deaf. It now gives up after a
streak, and the return closes the channel the binder is actually watching.

Two smaller pairings while here: stop_async retracts the edge it would
otherwise leave standing for a socket that is gone, and start_async hands a
refused edge to the binder to retry rather than dropping the only edge
either consumer will see until the interface next moves. The absence
deadline is stamped at start rather than at construction, so a transport
staged for longer than the bring-up window no longer reports sustained
absence on its first tick having given the interface no window at all.
This commit is contained in:
Arjen
2026-09-01 13:16:42 +01:00
parent cc3c7c4588
commit 6ad450b951
8 changed files with 298 additions and 40 deletions
+21 -14
View File
@@ -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
+5
View File
@@ -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 }
+1
View File
@@ -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");
+58 -6
View File
@@ -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<bool> {
interface_has_flags(interface, libc::IFF_UP as u32)
}
@@ -67,19 +78,31 @@ pub fn interface_index(interface: &str) -> Option<u32> {
/// 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<bool> {
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>, [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<PacketSocket>,
/// `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;
}
+116 -17
View File
@@ -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<bool> = 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<BinderContext>) {
/// 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<BinderContext>, 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<BinderContext>, 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<BinderContext>) {
async fn binder_loop(ctx: Arc<BinderContext>, initial_unpublished: Option<bool>) {
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<BinderContext>) {
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<bool> = None;
// edge is superseded rather than queued. Seeded with the start edge if
// that one was refused.
let mut unpublished: Option<bool> = 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<BinderContext>) {
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<BinderContext>) {
}
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");
+54
View File
@@ -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
+35 -3
View File
@@ -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<AsyncFd<LinkEventSocket>>,
/// 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");
+8
View File
@@ -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.