mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-06 03:28:24 +00:00
The medium-change detector sampled two host-wide signals: the source address the routing table would pick for an off-link destination, and the set of up, non-loopback interface addresses. The second one was the problem. It enumerated every address the host had, so a docker bridge coming up, a VPN connecting or a container network appearing moved the fingerprint with no peering affected at all — and the reaction to a moved fingerprint is to drop every connected UDP socket and heartbeat every peer. Self-healing, so it cost work rather than connectivity, but on any host running containers it could fire repeatedly for nothing. Ask the question per peer instead. For each peer whose transport address is a numeric IP endpoint, `connect(2)` a UDP socket to it and read back the local address — the same no-packets operation, aimed at the peers we actually hold rather than at a documentation prefix. The fingerprint becomes the set of local addresses the kernel would use to reach our peers. That is the quantity the reaction cares about. The stale `connect(2)` this subsystem exists to repair pinned a local source address chosen for one destination, so measuring the same thing for the same destinations asks the kernel the question the bug is about rather than a proxy for it. Three things follow: - Interfaces the node does not peer over cannot move it, by construction rather than by a filter guessing which interface names are infrastructure. `docker compose up` moves nothing. - On-link peers become visible. A peer on the same LAN is reached by its subnet route, and the old probe followed the *default* route by construction, so it looked straight past that path. - A more specific route moving under one peer is representable at all, which no single host-wide sample could be. Samples are compared over the *intersection* of their peer sets, never the union, so peers joining and leaving cannot fire the fan-out on their own. The sample is still adopted on the no-change path, or `last` would freeze on the peer set the detector started with. A peer whose probe stops answering is a move to "no route" and does count: that peer is exactly the one now stranded. **A peer's first sample is judged against its socket, not against nothing.** The intersection rule skips a peer present in only one sample, which is right for peer churn and wrong for the sample in which a peer first appears, because that sample may already be the post-change one. `last` gains a peer only at the first wake after it shows up in the entity snapshot, so a medium change inside that window is consumed rather than delayed: the peer's connected socket stays pinned to the path the host has just left, and the peering black-holes until the liveness timeout tears it down. A first-seen peer is therefore compared against the source its own connected socket is bound to, where it has one. One residual on that rule, stated exactly rather than understated: the tick publishes the entity snapshot *before* it installs connected sockets, so a socket installed on tick N is first visible on tick N+1, and a peer whose path moves inside that window is first seen with `bound` still `None` while genuinely holding a pinned socket. About one `tick_interval_secs` per join. A hole, not a harmless skip. **The probe binds the way the send path binds.** `open_connected_fd` binds `local_addr` verbatim before connecting, so the socket keeps a configured address whatever the route says, while the probe took the kernel's choice. Under a non-wildcard `bind_addr` the two answered different questions and every first-seen peer reported a phantom move. The probe now binds what the transport binds, address only and port 0; under the default wildcard bind nothing changes. Operator-facing corrections in the same surface. `PeerSourceMove::before` was recorded on every move and then wildcarded away by the only thing that read it, so the log said where a peer moved to but not where from — and the address it moved *from* is the one the stale `connect(2)` had pinned. Both ends are rendered now. The peer id used a private four-byte hex helper that duplicated `NodeAddr::short_hex` and dropped the `...` suffix every other operator surface prints; the duplicate is gone. The cost note said three syscalls per target: `UdpSocket::bind` is a `socket(2)` and a `bind(2)`, and the socket takes a `close(2)` on drop, so it is five. In the same place, "bounded by `node.limits.max_peers`" does not hold when that value is 0, which the configuration defines as unlimited. This deletes `interface_addrs()`'s only call site, and with it the `getifaddrs` walk and its `sockaddr` decoding. It therefore absorbs the Android `getifaddrs` issue rather than leaving it to be fixed separately. The peer table is reached through `entities_snapshot` rather than a new sharing primitive, so the detector stays a detached task holding no node state and taking no node lock. `PeerRow` gains a typed `probe_target` rather than having the detector re-parse the display string next to it, so a change to that string's rendering cannot silently leave the detector with an empty table and no way to notice. Peers that are not probeable IP destinations contribute nothing and need no per-transport special-casing here: a MAC on Ethernet or BLE, a .onion or Nym recipient behind a local proxy, a scoped IPv6 literal, and a peer still carrying its configured hostname all arrive as `None`. The last is deliberate — resolving one would put a DNS lookup with its timeouts on the sample path — and the window is small, since the address is replaced by the observed numeric source the first time an authenticated packet arrives. A node with no peers detects nothing, which is right: nothing is bound to the old path. Also corrects two doc claims that did not match the code, both in the text being rewritten: `transport::watcher` has no consumer besides this module, so it is not "shared with the interface binder"; and a connection-oriented transport does not "re-dial on send" in the case that matters, because `send_async` only dials when the pool holds no connection for the address and a connection stranded by a medium change is still in the pool — it is evicted after a write to it fails. Tests, each run against the defect it guards rather than only against the fix: - The churn rules are mutation-checked: iterating the union instead of the intersection fails four tests, and dropping the sample-adoption fails the one that pins a newly joined peer entering the comparison. - The snapshot seam is table-driven over six address shapes and fails if the publish site stops populating `probe_target` — nothing renders that field, so nothing else would have caught it. - A namespace test brings up a dummy interface with its own subnet and asserts the fingerprint does not move, then puts a more specific route to the peer out of that same interface and asserts it does, so the negative half cannot pass because sampling quietly stopped working. The same test now asserts that an unconstrained probe answers with the carrier and a constrained one with the other address, and that the two differ — the disagreement that would otherwise report a phantom move on every first-seen peer. Ignoring the constraint in `preferred_source` fails it. - The two links carrying the pinned source were asserted by nothing. Substituting `None` where `probe_targets` reads the row, or where `sample` writes the fingerprint, left the whole suite green while silently restoring the bug the first-sight rule exists to fix. Both are asserted now, and both mutations fail. - Every live-probe test passed if `preferred_source` returned `None` for everything: two all-`None` samples are self-consistent, the recorded-keys test never inspected a value, and the loopback assertion skipped through its `if let`. A probe to loopback must now answer with loopback, which holds on any host that can run the suite, including one started with `--network none`. - `reports_are_spaced_out_under_clean_flapping` polled at one second against a one-second pacing floor, so the two were indistinguishable and deleting the pacing block still passed. The wake is 100ms now. - The pinned-source publish test was Linux-gated though `open_connected_fd` and the field it asserts are available on macOS too. Widened, along with the two sibling tests on the same helper. - The detector's netlink subscription asserts the route groups rather than logging a decline, so a wrong group mask reds instead of passing quietly.
494 lines
20 KiB
Rust
494 lines
20 KiB
Rust
//! Kernel link-event sources.
|
|
//!
|
|
//! A watcher that resolves when the kernel reports that something about the
|
|
//! host's network links changed. Callers use it to react to interface state in
|
|
//! sub-second time instead of polling for it.
|
|
//!
|
|
//! | Platform | Source |
|
|
//! | -------- | ------ |
|
|
//! | Linux, Android | netlink `RTNLGRP_LINK` (`RTM_NEWLINK` / `RTM_DELLINK`) |
|
|
//! | macOS, FreeBSD | `PF_ROUTE` socket, `RTM_IFINFO` |
|
|
//! | Fallback | none — the watcher never fires, and callers poll |
|
|
//!
|
|
//! The messages themselves are deliberately **not parsed**. An event is a hint
|
|
//! to re-run whatever question the caller actually cares about, which is
|
|
//! cheap and authoritative; decoding `nlmsghdr`/`ifinfomsg` payloads to reach
|
|
//! the same answer would add a parser whose bugs would become the caller's
|
|
//! bugs. Any event on the socket wakes the caller, which then asks its own
|
|
//! question directly.
|
|
//!
|
|
//! Construction is best-effort, and that is the contract: a kernel or sandbox
|
|
//! that refuses the socket yields a watcher that never fires. `changed()` then
|
|
//! parks forever, which is what makes it safe to `select!` against a poll
|
|
//! ticker — the ticker simply always wins, and the caller degrades to polling
|
|
//! without a special case.
|
|
//!
|
|
//! Callers that need a *different* event group should extend
|
|
//! `open_link_socket` rather than opening a second socket beside this one:
|
|
//! `RTNLGRP_LINK` carries link state only, so a route change with both
|
|
//! interfaces up produces no event here.
|
|
|
|
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 {
|
|
/// The multicast group mask this socket is actually subscribed to, read
|
|
/// back from the kernel rather than remembered from the bind.
|
|
///
|
|
/// `getsockname` on a netlink socket fills `sockaddr_nl.nl_groups` with the
|
|
/// legacy 32-bit subscription mask, which covers every group in
|
|
/// [`groups`]. Reading it back is the only way to tell a watcher that
|
|
/// *asked* for the right groups from one that got them: a bind with a
|
|
/// wrong mask succeeds just as happily as a bind with the right one, and
|
|
/// then silently never delivers the messages the caller subscribed for.
|
|
#[cfg(all(test, any(target_os = "linux", target_os = "android")))]
|
|
fn bound_groups(&self) -> Option<u32> {
|
|
let mut sa: libc::sockaddr_nl = unsafe { std::mem::zeroed() };
|
|
let mut len = std::mem::size_of::<libc::sockaddr_nl>() as libc::socklen_t;
|
|
// SAFETY: `self.fd` is the netlink socket this struct owns, and `sa` /
|
|
// `len` are a correctly sized out-parameter pair for `getsockname`.
|
|
let rc = unsafe {
|
|
libc::getsockname(self.fd, &mut sa as *mut _ as *mut libc::sockaddr, &mut len)
|
|
};
|
|
if rc < 0 {
|
|
return None;
|
|
}
|
|
Some(sa.nl_groups)
|
|
}
|
|
|
|
fn recv(&self, buf: &mut [u8]) -> std::io::Result<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)
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Netlink multicast groups a watcher can subscribe to where the backend is
|
|
/// netlink.
|
|
///
|
|
/// Spelled as literals because the constants' names and availability differ
|
|
/// across libc versions; the values are ABI. These are the `RTMGRP_*` bitmask
|
|
/// form taken by `sockaddr_nl.nl_groups`, not the `RTNLGRP_*` ordinals.
|
|
#[cfg(any(target_os = "linux", target_os = "android"))]
|
|
pub mod groups {
|
|
/// Interfaces appearing, disappearing, or changing state.
|
|
pub const LINK: u32 = 0x1;
|
|
/// IPv4 addresses added to or removed from an interface.
|
|
pub const IPV4_IFADDR: u32 = 0x10;
|
|
/// IPv4 route table changes, including the default route moving.
|
|
pub const IPV4_ROUTE: u32 = 0x40;
|
|
/// IPv6 addresses added to or removed from an interface.
|
|
pub const IPV6_IFADDR: u32 = 0x100;
|
|
/// IPv6 route table changes.
|
|
pub const IPV6_ROUTE: u32 = 0x400;
|
|
|
|
/// Everything that can change which local address the host would use to
|
|
/// reach a given destination.
|
|
///
|
|
/// [`LINK`] alone does not cover it. A default route moving between two
|
|
/// interfaces that both stay up emits no link message at all — verified
|
|
/// with `ip monitor`, which reports zero events in the link group for
|
|
/// that change and two in the route group. A watcher that wants to hear
|
|
/// about egress-path changes rather than interface presence needs this.
|
|
pub const EGRESS_PATH: u32 = LINK | IPV4_IFADDR | IPV4_ROUTE | IPV6_IFADDR | IPV6_ROUTE;
|
|
}
|
|
|
|
/// Open the platform's link-event socket, non-blocking.
|
|
///
|
|
/// `groups` is ignored where the backend is `PF_ROUTE`: it has no group
|
|
/// selection and delivers every routing message to every reader, so a caller
|
|
/// that wants more than link events already has them there.
|
|
#[cfg(any(target_os = "linux", target_os = "android"))]
|
|
fn open_link_socket(groups: u32) -> std::io::Result<LinkEventSocket> {
|
|
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 = groups;
|
|
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(any(target_os = "linux", target_os = "android")))]
|
|
fn open_link_socket(_groups: u32) -> 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 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 Default for LinkWatcher {
|
|
fn default() -> Self {
|
|
Self::new()
|
|
}
|
|
}
|
|
|
|
impl LinkWatcher {
|
|
/// Open a watcher for interface presence.
|
|
///
|
|
/// Subscribes to link events only, which is what a caller asking "is this
|
|
/// interface here?" needs.
|
|
pub fn new() -> Self {
|
|
#[cfg(any(target_os = "linux", target_os = "android"))]
|
|
let groups = groups::LINK;
|
|
#[cfg(not(any(target_os = "linux", target_os = "android")))]
|
|
let groups = 0;
|
|
Self::with_groups(groups)
|
|
}
|
|
|
|
/// Open a watcher over an explicit set of netlink multicast groups.
|
|
///
|
|
/// Only meaningful on Linux, where the group mask decides what the kernel
|
|
/// sends; elsewhere `PF_ROUTE` delivers everything regardless and the mask
|
|
/// is ignored. See [`groups`] for the values, and `groups::EGRESS_PATH`
|
|
/// for the set that covers a change of egress path rather than of
|
|
/// interface presence.
|
|
pub fn with_groups(groups: u32) -> Self {
|
|
let inner = match open_link_socket(groups) {
|
|
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),
|
|
}
|
|
}
|
|
|
|
/// The netlink multicast groups this watcher is actually subscribed to, as
|
|
/// the kernel reports them.
|
|
///
|
|
/// `None` when there is no live source, and on every platform whose backend
|
|
/// is `PF_ROUTE`, which has no group selection to report.
|
|
///
|
|
/// This exists to be asserted on. A bind with the wrong group mask succeeds
|
|
/// exactly like a bind with the right one and then silently never delivers
|
|
/// what the caller subscribed for, so nothing short of reading the
|
|
/// subscription back can tell the two apart without provoking a real
|
|
/// kernel event — which needs privileges CI does not have.
|
|
#[cfg(all(test, any(target_os = "linux", target_os = "android")))]
|
|
pub(crate) fn subscribed_groups(&self) -> Option<u32> {
|
|
self.inner.as_ref()?.get_ref().bound_groups()
|
|
}
|
|
|
|
/// Whether an event source is actually backing this watcher.
|
|
pub fn is_event_driven(&self) -> bool {
|
|
self.inner.is_some()
|
|
}
|
|
|
|
/// 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 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"
|
|
);
|
|
}
|
|
}
|
|
|
|
/// The wider egress-path mask must open too.
|
|
///
|
|
/// Same no-privilege argument as the link group above: these are all
|
|
/// read-only `NETLINK_ROUTE` multicast groups. If this one cannot bind
|
|
/// while `new()` can, the mask is wrong rather than the environment
|
|
/// restricted — and the caller that needs it would silently fall back to
|
|
/// polling.
|
|
#[cfg(target_os = "linux")]
|
|
#[tokio::test]
|
|
async fn the_egress_path_mask_opens_a_source() {
|
|
let w = LinkWatcher::with_groups(groups::EGRESS_PATH);
|
|
assert!(
|
|
w.is_event_driven(),
|
|
"the egress-path group mask must bind on Linux"
|
|
);
|
|
}
|
|
|
|
/// The mask actually reaches the socket.
|
|
///
|
|
/// `EGRESS_PATH` is a superset of `LINK`, so a watcher built on it must
|
|
/// still be a watcher — this pins that widening the mask does not make
|
|
/// the bind fail in a way `is_event_driven` would report as a missing
|
|
/// source, which is how a wrong constant would present.
|
|
#[cfg(target_os = "linux")]
|
|
#[test]
|
|
fn the_egress_path_mask_is_a_superset_of_link() {
|
|
assert_eq!(groups::EGRESS_PATH & groups::LINK, groups::LINK);
|
|
assert_ne!(groups::EGRESS_PATH, groups::LINK);
|
|
}
|
|
|
|
/// 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");
|
|
}
|
|
}
|