refactor(transport): lift presence out of the ethernet module

The presence machinery was written inside `transport::ethernet` because that is
where it was needed first, not because it belongs to ethernet. `presence.rs`
imports nothing but `std` — its only `super::` references were a doc link and
its own test module's `use super::*`. It is a standalone state machine that
happened to live in a transport's directory.

So this is a move, not a generalisation. `PresenceState`, `Presence`,
`AbsencePolicy`, `ChurnGuard` and `bind_backoff` go up to `transport::presence`;
ethernet re-exports what it used to own, and the duplicate `watcher.rs` is
deleted in favour of the shared one this branch now starts from.

What deliberately does *not* move is the binder loop. `BinderContext` holds an
`EthernetConfig`, a `NeighborBuffer`, `EthernetStats` and a pubkey for beacons,
and `bind_now` spawns the ethernet receive loop and beacon sender. Making that
generic needs associated types or `dyn`, and every implementer would still
supply its own bind-and-spawn body — abstracting thirty lines of control flow
while leaving two hundred lines of substance per transport, with an indirection
sitting between the reader and the platform-specific unsafe. The shared part
really is just the state and the decisions.

The line matters because ethernet is not the only transport bound to something
that can disappear. BLE binds an HCI adapter — a USB dongle that unplugs,
rfkills, or resets — and today has no presence handling at all: a missing
`hci0` at start returns an error and the transport is skipped for the life of
the process, which is precisely the boot race this branch exists to fix. A
future BLE binder writes its own loop and reuses the state machine, which is
the half that took a review to get right: the churn damping, the
announce/retract pairing, the episode clock, the detach classification.

The move itself is behaviour-neutral: `presence.rs` changes by a single
doc-link line and nothing else. The one suite change is the deletion of the
branch's duplicate `src/transport/ethernet/watcher.rs`, which takes its three
tests with it. `watcher_constructs_and_reports_its_backing`,
`a_persistently_failing_source_gives_up_instead_of_spinning` and
`a_sourceless_watcher_never_fires` are byte-identical to three of the five in
the shared `src/transport/watcher.rs`, which this commit does not touch, so
what goes is three duplicates and no coverage.

The module is narrowed to `pub(crate)` while it is being moved. `pub` here was
inherited from `ethernet::presence` rather than chosen, nothing outside the
crate names it, and it matches what `transport::watcher` already is. The two
types an embedder could reach stay reachable: `ethernet` re-exports
`AbsencePolicy` and `Presence`, so `transport::ethernet::{AbsencePolicy,
Presence}` is unchanged. Everything else under the module — `PresenceState`,
`ChurnGuard`, `BindOutcome`, `DetachOutcome`, `bind_backoff` and the three
constants — becomes crate-only, which is a public-surface narrowing riding
inside a move and is named here for that reason. It carries ethernet's target
gate, since ethernet is its only consumer and a crate-private module with no
consumer is dead code on a target where ethernet is cfg'd out.
This commit is contained in:
Arjen
2026-09-10 19:18:09 +00:00
committed by Johnathan Corgan
parent 40b24cc9df
commit a354501514
4 changed files with 23 additions and 369 deletions
+5 -6
View File
@@ -13,30 +13,29 @@
//! for the interface, binds when it appears, tears down when it goes away, and
//! rebinds when it returns. Start-time absence and runtime detach are the same
//! transition, so a node that boots before wifi and a node whose wifi reloads
//! at 03:00 take one code path. See [`presence`] for the state machine.
//! at 03:00 take one code path. See [`crate::transport::presence`] for the
//! state machine.
pub mod addr;
pub mod io;
pub mod neighbor;
pub mod presence;
pub mod stats;
mod watcher;
pub use crate::transport::presence::{AbsencePolicy, Presence};
pub use addr::parse_mac_string;
pub use presence::{AbsencePolicy, Presence};
use super::{
DiscoveredPeer, PacketTx, PresenceTx, ReceivedPacket, Transport, TransportAddr, TransportError,
TransportId, TransportPresence, TransportState, TransportType,
};
use crate::config::EthernetConfig;
use crate::transport::presence::{ABSENCE_ERROR_AFTER, ChurnGuard, PresenceState, bind_backoff};
use crate::transport::watcher::LinkWatcher;
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;
use watcher::LinkWatcher;
use secp256k1::XOnlyPublicKey;
use std::sync::atomic::{AtomicBool, AtomicU16, AtomicU32, Ordering};
-362
View File
@@ -1,362 +0,0 @@
//! Link-event sources for the interface presence watcher.
//!
//! The presence machine works on a 1-second poll alone. This module removes
//! the latency: where the kernel offers a link-event source, the binder blocks
//! on it and reacts in sub-second time, and the poll stays underneath as a
//! backstop rather than as the mechanism.
//!
//! | Platform | Source |
//! | -------- | ------ |
//! | Linux | netlink `RTNLGRP_LINK` (`RTM_NEWLINK` / `RTM_DELLINK`) |
//! | macOS, FreeBSD | `PF_ROUTE` socket, `RTM_IFINFO` |
//! | Fallback | poll `getifaddrs` + flags, 1 s |
//!
//! The messages themselves are deliberately **not parsed**. A link event is a
//! hint to re-run the presence probe, which is cheap and authoritative;
//! decoding `nlmsghdr`/`ifinfomsg` payloads to reach the same answer would add
//! a parser whose bugs would be presence bugs. Any event on the socket wakes
//! the binder, which then asks
//! [`interface_present`](super::io::interface_present).
//!
//! Construction is best-effort. A kernel or sandbox that refuses the socket
//! yields a watcher that never fires, and the binder degrades to its poll.
use std::os::unix::io::{AsRawFd, RawFd};
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use std::time::Duration;
use tokio::io::unix::AsyncFd;
use tracing::{debug, warn};
/// Consecutive receive errors before the event source is abandoned for the
/// caller's poll.
const ERROR_GIVE_UP: u32 = 5;
/// An owned link-event socket. Closes its descriptor on drop.
struct LinkEventSocket {
fd: RawFd,
}
impl AsRawFd for LinkEventSocket {
fn as_raw_fd(&self) -> RawFd {
self.fd
}
}
impl Drop for LinkEventSocket {
fn drop(&mut self) {
unsafe { libc::close(self.fd) };
}
}
impl LinkEventSocket {
fn recv(&self, buf: &mut [u8]) -> std::io::Result<usize> {
let n = unsafe { libc::recv(self.fd, buf.as_mut_ptr() as *mut libc::c_void, buf.len(), 0) };
if n < 0 {
Err(std::io::Error::last_os_error())
} else {
Ok(n as usize)
}
}
}
/// Open the platform's link-event socket, non-blocking.
#[cfg(target_os = "linux")]
fn open_link_socket() -> std::io::Result<LinkEventSocket> {
// RTMGRP_LINK. Spelled as a literal because the constant's name and
// availability differ across libc versions; the value is ABI.
const RTMGRP_LINK: u32 = 1;
let fd = unsafe {
libc::socket(
libc::AF_NETLINK,
libc::SOCK_RAW | libc::SOCK_NONBLOCK | libc::SOCK_CLOEXEC,
libc::NETLINK_ROUTE,
)
};
if fd < 0 {
return Err(std::io::Error::last_os_error());
}
let socket = LinkEventSocket { fd };
let mut sa: libc::sockaddr_nl = unsafe { std::mem::zeroed() };
sa.nl_family = libc::AF_NETLINK as u16;
sa.nl_groups = RTMGRP_LINK;
let ret = unsafe {
libc::bind(
fd,
&sa as *const libc::sockaddr_nl as *const libc::sockaddr,
std::mem::size_of::<libc::sockaddr_nl>() as libc::socklen_t,
)
};
if ret < 0 {
return Err(std::io::Error::last_os_error());
}
Ok(socket)
}
/// Open the platform's link-event socket, non-blocking.
#[cfg(not(target_os = "linux"))]
fn open_link_socket() -> std::io::Result<LinkEventSocket> {
// PF_ROUTE delivers RTM_IFINFO (and the rest of the routing messages) to
// every reader; no bind and no group selection exist for it.
let fd = unsafe { libc::socket(libc::PF_ROUTE, libc::SOCK_RAW, libc::AF_UNSPEC) };
if fd < 0 {
return Err(std::io::Error::last_os_error());
}
let socket = LinkEventSocket { fd };
let flags = unsafe { libc::fcntl(fd, libc::F_GETFL) };
if flags < 0 {
return Err(std::io::Error::last_os_error());
}
if unsafe { libc::fcntl(fd, libc::F_SETFL, flags | libc::O_NONBLOCK) } < 0 {
return Err(std::io::Error::last_os_error());
}
Ok(socket)
}
/// A source of "something about the links changed" wake-ups.
pub(crate) struct LinkWatcher {
/// `None` when no event source could be opened — the caller's poll is then
/// the whole mechanism, which is exactly the documented fallback.
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 {
/// Open the platform link-event source, falling back to nothing.
pub(crate) fn new() -> Self {
let inner = match open_link_socket() {
Ok(socket) => match AsyncFd::new(socket) {
Ok(afd) => Some(afd),
Err(e) => {
debug!(error = %e, "Link event socket not registrable; polling instead");
None
}
},
Err(e) => {
debug!(error = %e, "No link event source available; polling instead");
None
}
};
Self {
inner,
errors: AtomicU32::new(0),
given_up: AtomicBool::new(false),
}
}
/// Whether an event source is actually backing this watcher.
pub(crate) fn is_event_driven(&self) -> bool {
self.inner.is_some()
}
/// Resolve when the kernel reports a link change.
///
/// Never resolves when no event source is available, which makes it safe
/// to `select!` against the poll ticker: the ticker simply always wins.
pub(crate) async fn changed(&self) {
let Some(afd) = &self.inner else {
std::future::pending::<()>().await;
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
// caller's poll continues to cover presence.
std::future::pending::<()>().await;
unreachable!("pending never resolves")
};
// Drain to WouldBlock so a burst of link messages is one wake-up
// and the socket buffer does not fill behind us.
let mut buf = [0u8; 4096];
let mut saw_event = false;
let mut failure = None;
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. 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
// error leaves `readable()` instantly ready, `recv`
// failing again, and the loop spinning a core flat with
// nothing logged. Clear it by hand and back off.
Ok(Err(e)) => {
guard.clear_ready();
failure = Some(e);
break;
}
// WouldBlock — readiness is cleared, drain complete.
Err(_) => break,
}
}
if saw_event {
self.errors.store(0, Ordering::Relaxed);
return;
}
if let Some(e) = failure {
let errors = self.errors.fetch_add(1, Ordering::Relaxed) + 1;
if errors == 1 {
// ENOBUFS is the realistic one: a burst of link events
// overflowed the socket buffer, so the kernel dropped some.
// Losing events is survivable — the caller polls — but the
// spin is not, and neither is doing it silently.
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 \
polling for interface presence"
);
std::future::pending::<()>().await;
unreachable!("pending never resolves")
}
tokio::time::sleep(Duration::from_millis(100) * errors).await;
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
/// The watcher must construct on any host, with or without a usable event
/// source, because the binder builds one unconditionally.
#[tokio::test]
async fn watcher_constructs_and_reports_its_backing() {
let w = LinkWatcher::new();
// On Linux the source is a plain `AF_NETLINK` socket in the
// `RTNLGRP_LINK` group, which needs no capability and no privilege —
// so on this platform "a sandbox might refuse it" is not a licence to
// accept either answer. Discarding the result, which this test used
// to do, meant nothing anywhere asserted that the event path exists:
// the 1 s poll is a complete fallback, so the entire suite passed with
// the source unavailable and no test could tell.
#[cfg(target_os = "linux")]
assert!(
w.is_event_driven(),
"the netlink link-event source must open on Linux; \
falling back to the poll here is a silent loss of the fast path"
);
// Elsewhere both answers are legitimate, so pin only that asking is
// safe and that a watcher with no source parks rather than fires.
#[cfg(not(target_os = "linux"))]
{
let backed = w.is_event_driven();
assert!(
backed
|| tokio::time::timeout(Duration::from_millis(50), w.changed())
.await
.is_err(),
"a watcher with no source must never resolve"
);
}
}
/// A descriptor whose `recv` always fails must not become a busy loop.
///
/// `try_io` clears readiness only on `WouldBlock`. Breaking out of a real
/// error left `readable()` instantly ready, `recv` failing again, and the
/// loop spinning a core flat with nothing logged — the realistic trigger
/// being `ENOBUFS` when a burst of link events overflows the socket
/// buffer. A pipe stands in for that here: `recv` on one answers
/// `ENOTSOCK`, every time, which is exactly the shape of a persistent
/// error.
#[tokio::test]
async fn a_persistently_failing_source_gives_up_instead_of_spinning() {
let mut fds = [0i32; 2];
assert_eq!(unsafe { libc::pipe(fds.as_mut_ptr()) }, 0, "pipe()");
let (read_fd, write_fd) = (fds[0], fds[1]);
// AsyncFd requires a non-blocking descriptor.
let flags = unsafe { libc::fcntl(read_fd, libc::F_GETFL) };
assert!(unsafe { libc::fcntl(read_fd, libc::F_SETFL, flags | libc::O_NONBLOCK) } >= 0);
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`,
// and a descriptor that was already readable before re-registration
// may never deliver another edge on its own — which would stall the
// loop at one error and hide whether the give-up path works. A steady
// trickle stands in for the burst of link events that provokes the
// real failure.
let writer = tokio::task::spawn_blocking(move || {
for _ in 0..200 {
if unsafe { libc::write(write_fd, b"x".as_ptr().cast(), 1) } < 0 {
break;
}
std::thread::sleep(Duration::from_millis(25));
}
unsafe { libc::close(write_fd) };
});
// Never resolves — there is no event to report — but it must reach the
// give-up state rather than burn until the timeout.
let fired = tokio::time::timeout(Duration::from_secs(5), watcher.changed()).await;
assert!(fired.is_err(), "a failing source must not report an event");
assert!(
watcher.errors.load(Ordering::Relaxed) >= ERROR_GIVE_UP,
"the error path must count, back off and stop, not spin silently"
);
writer.abort();
}
/// A watcher with no event source must never resolve, so a `select!`
/// against the poll ticker degrades cleanly instead of spinning.
#[tokio::test]
async fn a_sourceless_watcher_never_fires() {
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");
}
}
+17
View File
@@ -25,6 +25,23 @@ pub mod ethernet;
#[cfg(unix)]
pub(crate) mod watcher;
/// Presence lifecycle for a transport bound to a local resource that can
/// disappear and come back: the phase machine, the absence policy, and the
/// damping that keeps a flapping resource from flapping node health with it.
///
/// Transport-agnostic on purpose. Only the *probe* — "is my thing there, and
/// is it still the same one?" — is specific to what is bound, and that stays
/// with the transport that knows how to ask.
///
/// Crate-internal on purpose, for the same reason as `watcher` above: it is a
/// mechanism the crate's own transports share, not a surface an embedder
/// builds against. `ethernet` re-exports the two types it used to own, so the
/// published path stays `transport::ethernet::{AbsencePolicy, Presence}`.
/// Gated with the one transport that binds through it today; widen the gate
/// when a second binder arrives.
#[cfg(any(target_os = "linux", target_os = "macos"))]
pub(crate) mod presence;
#[cfg(ble_available)]
pub mod ble;
@@ -27,7 +27,7 @@
//! enabled it — and deliberately not `IFF_RUNNING`: binding needs no carrier,
//! and a socket outlives a carrier flap. Carrier is reported alongside it
//! rather than steering it. See
//! [`interface_present`](super::io::interface_present) for why.
//! the ethernet transport's `interface_present` for why.
use std::sync::atomic::{AtomicU8, AtomicU32, AtomicU64, Ordering};
use std::sync::{PoisonError, RwLock, RwLockReadGuard, RwLockWriteGuard};