mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-05 19:18:25 +00:00
feat(transport): bind and rebind network interfaces dynamically
An interface-bound transport is a long-lived object that is *sometimes bound*. The interface it names may not exist when the daemon starts, may appear minutes later, may vanish and return mid-operation, and may never appear at all. Until now the first observation was final: an interface missing at start was logged once and skipped for the life of the process, and one that disappeared at runtime published a health change but was never rebound. On OpenWrt that is a live bug. procd starts fips before wifi has created fips-mesh0 / fips-ap0, both transports are skipped and never retried, and the 802.11s peer link forms anyway (that is mac80211, not the daemon) — so the node looks healthy and reaches nothing. Give every interface-bound transport presence state and a binder task. start_async now returns Ok with the transport ABSENT rather than failing; the binder binds when the interface appears, tears the socket down when it goes away, and rebinds when it returns. Start-time absence and runtime detach are one transition and one code path. TransportError::InterfaceUnavailable makes absence branchable: a missing interface and a typo'd interface name were the same flat StartFailed (String), so nothing downstream could tell a state from a fault. A bind failure that is *not* absence — no CAP_NET_RAW, no readable /dev/bpf* — still fails the start, because retrying a socket that can never open behind a Degraded nobody is watching is worse than dying loudly at boot. **Presence is IFF_UP, not IFF_UP|IFF_RUNNING.** Carrier and bindability are different questions and only the second belongs in a bind gate. An AF_PACKET socket on a carrier-less bridge is valid and starts carrying traffic the instant a member port comes up, with no rebind. Gating on carrier would report a wifi-only router Degraded forever for an empty br-lan, turn every carrier flap the socket would have survived into an unbind/rebind cycle, and deadlock an 802.11s interface that reports RUNNING only once it has peered — peering needs beacons, beacons need a bound socket. Carrier is reported beside presence in show_transports instead. **A name is not a device.** Both backends bind by device: AF_PACKET stores sll_ifindex, a BPF descriptor follows the interface it was attached to. A netdev deleted and recreated under the same name leaves the socket attached to something gone while the name resolves perfectly well, and nothing notices — a stale AF_PACKET socket never becomes readable, so the receive loop neither errors nor exits, and send failures go to the caller rather than the binder. A listen-only node (announce: false, no beacon sender to fail) sat present and deaf indefinitely after a wifi reload. The bound index is captured at bind and compared on every poll. A name reappearing with a different MAC is different hardware: cached neighbors are dropped rather than resumed onto. **Detection is event-driven** where the kernel offers a source — netlink RTNLGRP_LINK on Linux, PF_ROUTE on macOS and FreeBSD — with a 1 s getifaddrs poll underneath as a backstop. Link-event payloads are not parsed: an event is a hint to re-run the probe, which is cheap and authoritative, and a parser's bugs would be presence bugs. Probes are coalesced to ten a second because PF_ROUTE has no group filter and delivers every routing message on the host; a persistently failing event source is logged once, backed off, and abandoned for the poll after five errors rather than spinning a core silently on ENOBUFS. **Degraded becomes a level, not a latch.** The supervisor's reason set was monotonic, which was correct while no child could recover; with recovery it would have come to mean "something broke at some point since boot" rather than "something is broken now". Absence lives in its own reversible set, health is recomputed on every transition in both directions, and an absent transport still counts as up so a single-interface node that boots before its wifi degrades rather than exiting on NoTransports. There is deliberately no restart action in the FSM: rebinding belongs next to the file descriptor, and a supervisor-authored retry would be a second mechanism racing the first for the same socket. **The node's egress MTU follows the bound set.** transport_mtu filters on is_bound(), not is_operational() — an interface-bound transport is operational from the moment it starts, so the weaker predicate let hardware that had never appeared clamp the whole node's IPv6 MTU. Because that minimum now moves at runtime, the TUN reader and writer read the TCP MSS ceiling from a shared atomic instead of a u16 captured at spawn; every other consumer (show_status, the snapshot, the session-layer fragmentation check) already read it live, so the clamp was the one place the daemon could report one effective MTU and enforce another. It moves in both directions: a narrow interface appearing tightens it, its departure releases it. MSS is negotiated per connection, so a change binds connections opened after it. **Logging is per edge, never per attempt**, with one deadline after it. The edge itself is not an error — info at boot, warn on a runtime detach, since an interface bound a moment later is the ordinary case this mechanism exists to absorb and crying error at t=0 then "recovered" at t=0.2s is the failure the rule exists to prevent. Ten seconds is the whole grace: past it, absence is no longer a race against a radio or a container, so a required interface still missing is reported once at error. Start-time absence and a runtime detach share that one deadline, as they share everything else here. Said once, not repeated: how long an absence has lasted is state, published as interface.since_secs and as Degraded, and a monitor can threshold it per deployment rather than the daemon compiling a schedule in. Successful rebinds are damped. Backoff covers failed binds; the nastier case is binds that keep succeeding into a socket that dies moments later, which a receive loop giving up on a persistent error while the interface stays UP produces once per second forever. Consecutive bindings dying inside ten seconds back off on the 1 s → 30 s curve, and past three the binder stops announcing each bind as a recovery until one lasts. The receive loop backs off and exits on a dead socket instead of spinning on Err with a warn per iteration, and the ad-hoc ENXIO socket reopen in the beacon sender is gone: both hand recovery to the presence machine, one mechanism for every cause rather than one hack per symptom. Beacons pause while absent because the task simply does not exist then. transports.ethernet.*.optional (default false) decides how loudly absence is reported. Naming an interface is a statement that you expect it, so the default complains: Degraded, and the error at the deadline. optional: true is silent — info on the edge, no health impact, no error — for hardware legitimately not always there. It describes the interface's presence, not the transport's importance, and no value of it makes a missing interface fatal at startup. show_transports grows an `interface` block — presence, carrier, policy, how long the phase has been held, bind count, failed attempts. The original boot race was expensive precisely because nothing an operator could see said the node was deaf.
This commit is contained in:
@@ -302,6 +302,25 @@ pub struct EthernetConfig {
|
||||
/// Announcement beacon interval in seconds. Default: 30.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub beacon_interval_secs: Option<u64>,
|
||||
|
||||
/// Whether the absence of this interface is normal. Default: false.
|
||||
///
|
||||
/// Naming an interface in configuration is a statement that you expect it,
|
||||
/// so the default is to complain: while the interface is missing the node
|
||||
/// reports `Degraded`, the edge is logged (`info` at startup, `warn` on a
|
||||
/// runtime detach), and an absence that outlasts the bring-up window — ten
|
||||
/// seconds, past which it is no longer a race against a radio or a
|
||||
/// container — is reported once at `error`. Set `optional: true` for
|
||||
/// hardware that is legitimately not always there — a dock adapter, a
|
||||
/// radio that only some boards carry — and its absence becomes silent
|
||||
/// (`info` on the edge, no health impact, no error).
|
||||
///
|
||||
/// This describes **the interface's presence, not the transport's
|
||||
/// importance**. An optional interface that is present is used exactly as
|
||||
/// hard as any other; setting it does not deprioritize the transport, and
|
||||
/// no value of this field makes a missing interface fatal at startup.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub optional: Option<bool>,
|
||||
}
|
||||
|
||||
impl EthernetConfig {
|
||||
@@ -340,6 +359,11 @@ impl EthernetConfig {
|
||||
self.accept_connections.unwrap_or(false)
|
||||
}
|
||||
|
||||
/// Whether absence of the interface is normal. Default: false.
|
||||
pub fn optional(&self) -> bool {
|
||||
self.optional.unwrap_or(false)
|
||||
}
|
||||
|
||||
/// Get the beacon interval, clamped to minimum. Default: 30s.
|
||||
pub fn beacon_interval_secs(&self) -> u64 {
|
||||
self.beacon_interval_secs
|
||||
@@ -1071,4 +1095,40 @@ mod tests {
|
||||
serde_yaml::from_str("interface: eth0\nbogus: true\n");
|
||||
assert!(bogus.is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ethernet_absence_is_an_error_unless_opted_out() {
|
||||
// Naming an interface is a statement that you expect it, so the
|
||||
// default has to be the loud one. A default of `true` here would make
|
||||
// every missing interface silent, which is the failure mode the
|
||||
// whole presence mechanism exists to stop hiding.
|
||||
let bare: EthernetConfig = serde_yaml::from_str("interface: eth0\n").unwrap();
|
||||
assert_eq!(bare.optional, None);
|
||||
assert!(!bare.optional(), "absence must default to an error");
|
||||
|
||||
let opted: EthernetConfig =
|
||||
serde_yaml::from_str("interface: enx00e04c680001\noptional: true\n").unwrap();
|
||||
assert!(opted.optional());
|
||||
|
||||
let explicit: EthernetConfig =
|
||||
serde_yaml::from_str("interface: eth0\noptional: false\n").unwrap();
|
||||
assert!(!explicit.optional());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ethernet_optional_survives_a_round_trip() {
|
||||
// The packaged OpenWrt config ships `optional: true` on the mesh and
|
||||
// access blocks; a serializer that dropped it would silently turn
|
||||
// every stock router Degraded.
|
||||
let opted: EthernetConfig =
|
||||
serde_yaml::from_str("interface: fips-mesh0\noptional: true\n").unwrap();
|
||||
let round: EthernetConfig =
|
||||
serde_yaml::from_str(&serde_yaml::to_string(&opted).unwrap()).unwrap();
|
||||
assert!(round.optional());
|
||||
|
||||
// ... and the default stays absent from the output rather than being
|
||||
// written back as an explicit `false`.
|
||||
let bare: EthernetConfig = serde_yaml::from_str("interface: eth0\n").unwrap();
|
||||
assert!(!serde_yaml::to_string(&bare).unwrap().contains("optional"));
|
||||
}
|
||||
}
|
||||
|
||||
+38
-2
@@ -1361,8 +1361,17 @@ pub(crate) fn show_connections_from_handle(
|
||||
|
||||
/// `show_transports` — Transport instances.
|
||||
pub fn show_transports(node: &Node) -> Value {
|
||||
let transports: Vec<Value> = node
|
||||
.transport_ids()
|
||||
// Ascending id, which is creation order: UDP, then Ethernet, then TCP,
|
||||
// Tor, Nym, BLE, with each type's instances in config order. The map
|
||||
// behind `transport_ids` is a `HashMap`, so without this the array order
|
||||
// is whatever the hash seed produced — arbitrary, and different on every
|
||||
// daemon restart. Anything scripting against this output, and every view
|
||||
// rendering it, inherits that. Sorting by id groups the list by transport
|
||||
// type for free, because the ids were handed out that way.
|
||||
let mut ids: Vec<_> = node.transport_ids().copied().collect();
|
||||
ids.sort_by_key(|id| id.as_u32());
|
||||
let transports: Vec<Value> = ids
|
||||
.iter()
|
||||
.map(|id| {
|
||||
let handle = node.get_transport(id).unwrap();
|
||||
let mut t_json = json!({
|
||||
@@ -1390,6 +1399,21 @@ pub fn show_transports(node: &Node) -> Value {
|
||||
t_json["tor_monitoring"] = serde_json::to_value(&monitoring).unwrap_or_default();
|
||||
}
|
||||
|
||||
// Interface presence, for the transports that have an interface.
|
||||
// Absent from the payload entirely for the ones that do not, rather
|
||||
// than reported as a permanently-`present` interface named "".
|
||||
if let Some(p) = handle.interface_presence() {
|
||||
t_json["interface"] = json!({
|
||||
"name": handle.interface_name().unwrap_or_default(),
|
||||
"presence": p.presence,
|
||||
"carrier": p.carrier,
|
||||
"policy": p.policy,
|
||||
"since_secs": p.since_secs,
|
||||
"binds": p.binds,
|
||||
"failed_attempts": p.failed_attempts,
|
||||
});
|
||||
}
|
||||
|
||||
t_json["stats"] = handle.transport_stats();
|
||||
|
||||
t_json
|
||||
@@ -1434,6 +1458,18 @@ pub(crate) fn show_transports_from_handle(handle: &super::read_handle::ControlRe
|
||||
t_json["tor_monitoring"] = monitoring.clone();
|
||||
}
|
||||
|
||||
if let Some(iface) = &t.interface {
|
||||
t_json["interface"] = json!({
|
||||
"name": iface.name,
|
||||
"presence": iface.presence,
|
||||
"carrier": iface.carrier,
|
||||
"policy": iface.policy,
|
||||
"since_secs": iface.since_secs,
|
||||
"binds": iface.binds,
|
||||
"failed_attempts": iface.failed_attempts,
|
||||
});
|
||||
}
|
||||
|
||||
t_json["stats"] = t.stats.clone();
|
||||
|
||||
t_json
|
||||
|
||||
@@ -740,6 +740,21 @@ pub(crate) struct TransportRow {
|
||||
pub onion_address: Option<String>,
|
||||
pub tor_monitoring: Option<serde_json::Value>,
|
||||
pub stats: serde_json::Value,
|
||||
/// Interface presence for interface-bound transports; `None` for the rest.
|
||||
pub interface: Option<InterfaceRow>,
|
||||
}
|
||||
|
||||
/// Interface name, presence and policy for an interface-bound transport, as
|
||||
/// `show_transports` renders it.
|
||||
#[derive(Clone, PartialEq)]
|
||||
pub(crate) struct InterfaceRow {
|
||||
pub name: String,
|
||||
pub presence: &'static str,
|
||||
pub carrier: bool,
|
||||
pub policy: &'static str,
|
||||
pub since_secs: u64,
|
||||
pub binds: u64,
|
||||
pub failed_attempts: u32,
|
||||
}
|
||||
|
||||
/// MMP trend labels for a peer's link-layer block in `show_mmp` (each present
|
||||
|
||||
@@ -109,6 +109,16 @@ impl Node {
|
||||
}
|
||||
};
|
||||
|
||||
// Interface-presence receiver, or a dummy channel — same pattern and
|
||||
// same reason as the child-liveness receiver above.
|
||||
let (mut presence_rx, _presence_guard) = match self.transport_presence_rx.take() {
|
||||
Some(rx) => (rx, None),
|
||||
None => {
|
||||
let (tx, rx) = tokio::sync::mpsc::channel(1);
|
||||
(rx, Some(tx))
|
||||
}
|
||||
};
|
||||
|
||||
let tick_period = Duration::from_secs(self.config().node.tick_interval_secs);
|
||||
let mut tick = tokio::time::interval(tick_period);
|
||||
|
||||
@@ -334,6 +344,42 @@ impl Node {
|
||||
self.supervisor.state = ns;
|
||||
}
|
||||
}
|
||||
// A transport child exiting leaves the bound set, so
|
||||
// it can be the one that was holding the node's egress
|
||||
// MTU down. `is_bound()` is `is_operational()` plus the
|
||||
// presence refinement, and this moves the first half.
|
||||
self.refresh_tun_mss_ceiling();
|
||||
}
|
||||
}
|
||||
// Interface presence. An interface-bound transport's binder
|
||||
// reports attach and detach; the FSM folds it into health.
|
||||
// Unlike `ChildExited` this is reversible in both directions —
|
||||
// the interface coming back republishes `Running` — which is
|
||||
// the whole point of `Degraded` being a level rather than a
|
||||
// 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;
|
||||
}
|
||||
}
|
||||
// The bound set just changed, so the node's egress MTU
|
||||
// floor may have. Both directions: an interface that
|
||||
// binds can be the narrow one, and one that detaches
|
||||
// can be the reason the clamp was tight.
|
||||
self.refresh_tun_mss_ceiling();
|
||||
}
|
||||
}
|
||||
Some(ipv6_packet) = tun_outbound_rx.recv() => {
|
||||
|
||||
+85
-10
@@ -1381,6 +1381,17 @@ impl Node {
|
||||
self.child_exit_tx = Some(child_exit_tx);
|
||||
self.child_exit_rx = Some(child_exit_rx);
|
||||
|
||||
// Interface-presence channel. Created before `create_transports` so
|
||||
// every interface-bound transport gets the sender at construction and
|
||||
// its very first bind attempt — the one `start_async` makes inline —
|
||||
// is already reportable. A boot race therefore reaches the FSM while
|
||||
// it is still `Starting`, and start-completion health resolves to
|
||||
// `Degraded` on the first publish rather than publishing `Full` and
|
||||
// correcting it a moment later.
|
||||
let (presence_tx, presence_rx) = tokio::sync::mpsc::channel(16);
|
||||
self.transport_presence_tx = Some(presence_tx);
|
||||
self.transport_presence_rx = Some(presence_rx);
|
||||
|
||||
// Initialize transports first (before TUN, before Nostr discovery).
|
||||
// Creation allocates each transport's id; the supervisor FSM authors
|
||||
// the start order over those ids.
|
||||
@@ -1642,12 +1653,20 @@ impl Node {
|
||||
info!(" address: {}", device.address());
|
||||
info!(" mtu: {}", mtu);
|
||||
|
||||
// Calculate max MSS for TCP clamping
|
||||
// Seed the shared MSS ceiling from whatever is bound
|
||||
// right now. Both TUN threads read it live from here
|
||||
// on, so a transport binding or unbinding later moves
|
||||
// the clamp instead of leaving it at this instant's
|
||||
// value — see `crate::upper::tun::MssCeiling`.
|
||||
self.refresh_tun_mss_ceiling();
|
||||
let max_mss = self.tun_mss_ceiling.clone();
|
||||
let effective_mtu = self.effective_ipv6_mtu();
|
||||
let max_mss = effective_mtu.saturating_sub(40).saturating_sub(20); // IPv6 + TCP headers
|
||||
|
||||
info!("effective MTU: {} bytes", effective_mtu);
|
||||
debug!(" max TCP MSS: {} bytes", max_mss);
|
||||
debug!(
|
||||
" max TCP MSS: {} bytes",
|
||||
max_mss.load(std::sync::atomic::Ordering::Relaxed)
|
||||
);
|
||||
|
||||
// On macOS and FreeBSD, create a shutdown pipe. Writing to it
|
||||
// unblocks the reader thread's select() loop without closing
|
||||
@@ -1671,8 +1690,8 @@ impl Node {
|
||||
// Create writer (dups the fd for independent write access).
|
||||
// Pass path_mtu_lookup so inbound SYN-ACK clamp can read
|
||||
// per-destination path MTU learned via discovery.
|
||||
let (writer, tun_tx) =
|
||||
device.create_writer(max_mss, self.path_mtu_lookup.clone())?;
|
||||
let (writer, tun_tx) = device
|
||||
.create_writer(max_mss.clone(), self.path_mtu_lookup.clone())?;
|
||||
|
||||
// Spawn writer thread. On exit it self-reports
|
||||
// `Child::Tun` (sync context → `blocking_send`); TUN
|
||||
@@ -1698,7 +1717,6 @@ impl Node {
|
||||
// self-reports `Child::Tun` on exit (sync context →
|
||||
// `blocking_send`). Exactly one cfg variant compiles,
|
||||
// so the single clone is moved into that closure.
|
||||
let transport_mtu = self.transport_mtu();
|
||||
let path_mtu_lookup = self.path_mtu_lookup.clone();
|
||||
let reader_child_tx = self.child_exit_tx.clone();
|
||||
#[cfg(any(target_os = "macos", target_os = "freebsd"))]
|
||||
@@ -1709,7 +1727,7 @@ impl Node {
|
||||
our_addr,
|
||||
reader_tun_tx,
|
||||
outbound_tx,
|
||||
transport_mtu,
|
||||
max_mss,
|
||||
path_mtu_lookup,
|
||||
shutdown_read_fd,
|
||||
);
|
||||
@@ -1725,7 +1743,7 @@ impl Node {
|
||||
our_addr,
|
||||
reader_tun_tx,
|
||||
outbound_tx,
|
||||
transport_mtu,
|
||||
max_mss,
|
||||
path_mtu_lookup,
|
||||
);
|
||||
if let Some(tx) = &reader_child_tx {
|
||||
@@ -1856,6 +1874,16 @@ impl Node {
|
||||
}
|
||||
};
|
||||
|
||||
// Drain any presence edges this child's start produced *before*
|
||||
// reporting the child itself. An interface-bound transport whose
|
||||
// interface is missing reports absence from inside `start_async`
|
||||
// and then reports `SubstrateUp` (absence is a state, not a start
|
||||
// failure), so ordering the drain first means start-completion
|
||||
// health already knows about the absence when `pending` empties.
|
||||
// Otherwise a boot race publishes `Full` and corrects itself a
|
||||
// moment later, and every consumer sees a spurious transition.
|
||||
let _ = self.drain_transport_presence();
|
||||
|
||||
let feedback_actions = self.supervisor.fsm.step(feedback);
|
||||
for action in &feedback_actions {
|
||||
if let Action::PublishState(ns) = action {
|
||||
@@ -1864,6 +1892,12 @@ impl Node {
|
||||
}
|
||||
}
|
||||
|
||||
// Late edges: a transport that bound after its `SubstrateUp` was
|
||||
// reported, or one that detached during a later child's bring-up.
|
||||
if let Some(ns) = self.drain_transport_presence() {
|
||||
start_outcome = Some(ns);
|
||||
}
|
||||
|
||||
// Seams that never triggered inside the loop: the "Transports
|
||||
// initialized" info! when there was no non-transport child, and the
|
||||
// peer-connect when there was no Tun/Dns child (today it still runs,
|
||||
@@ -1902,8 +1936,10 @@ impl Node {
|
||||
// children. Enumerate them for the operator, then proceed —
|
||||
// a degraded node serves traffic.
|
||||
warn!(
|
||||
degraded_children = ?self.supervisor.fsm.failed(),
|
||||
"Node started DEGRADED: one or more configured optional children failed to start"
|
||||
degraded_children = ?self.supervisor.fsm.degraded_children(),
|
||||
absent_interfaces = ?self.supervisor.fsm.absent(),
|
||||
"Node started DEGRADED: one or more configured optional children failed to \
|
||||
start, or a configured interface is absent"
|
||||
);
|
||||
}
|
||||
_ => {}
|
||||
@@ -2237,6 +2273,45 @@ impl Node {
|
||||
}
|
||||
}
|
||||
|
||||
/// Feed every queued interface-presence edge to the supervisor FSM,
|
||||
/// returning the last [`NodeState`] it asked to publish (if any).
|
||||
///
|
||||
/// Non-blocking: it drains what is already queued and returns. Used during
|
||||
/// bring-up, where the rx_loop's presence arm is not running yet — from
|
||||
/// then on that arm owns the same translation.
|
||||
pub(in crate::node) fn drain_transport_presence(&mut self) -> Option<NodeState> {
|
||||
let mut edges = Vec::new();
|
||||
if let Some(rx) = self.transport_presence_rx.as_mut() {
|
||||
while let Ok(edge) = rx.try_recv() {
|
||||
edges.push(edge);
|
||||
}
|
||||
}
|
||||
|
||||
let mut published = None;
|
||||
let saw_edge = !edges.is_empty();
|
||||
for edge in edges {
|
||||
let child = Child::Transport(edge.transport_id);
|
||||
let event = if edge.present {
|
||||
Event::ChildPresent { child }
|
||||
} else {
|
||||
Event::ChildAbsent { child }
|
||||
};
|
||||
for action in self.supervisor.fsm.step(event) {
|
||||
if let Action::PublishState(ns) = action {
|
||||
published = Some(ns);
|
||||
}
|
||||
}
|
||||
}
|
||||
if saw_edge {
|
||||
// A bind or unbind changes which transports are bound, and so the
|
||||
// node's egress MTU floor. During bring-up this runs before the
|
||||
// TUN threads exist, which is exactly when it must: they read the
|
||||
// ceiling this leaves behind.
|
||||
self.refresh_tun_mss_ceiling();
|
||||
}
|
||||
published
|
||||
}
|
||||
|
||||
/// Reconstruct the supervised up-set from observed runtime presence, so the
|
||||
/// FSM authors the teardown order regardless of how the node reached
|
||||
/// `Running`. Worker pools are deliberately excluded: today's teardown never
|
||||
|
||||
@@ -67,6 +67,39 @@
|
||||
//! when a task/thread dies at runtime) is **deferred**: start-completion health
|
||||
//! resolution is start-framed, and liveness monitoring is a substantial unbuilt
|
||||
//! mechanism. This commit is start-time health only.
|
||||
//!
|
||||
//! ## Scope: interface presence, and `Degraded` as a level (this commit)
|
||||
//!
|
||||
//! Interface-bound transports are now *sometimes bound*: a transport whose
|
||||
//! interface is missing at start comes up [`Absent`] and binds later, and one
|
||||
//! whose interface goes away at runtime unbinds and rebinds when it returns.
|
||||
//! Two things follow for this machine.
|
||||
//!
|
||||
//! - **A second reason set.** [`Event::ChildAbsent`] / [`Event::ChildPresent`]
|
||||
//! move a child in and out of `absent`, which feeds `Degraded` exactly like
|
||||
//! `failed` does. It is kept separate because it is *reversible* and `failed`
|
||||
//! is not: a child that failed to start stays failed for the bring-up, while
|
||||
//! an absent interface is expected to come back.
|
||||
//! - **`Degraded` is a level, not a latch.** Nothing ever removed from `failed`,
|
||||
//! which was correct while no child could recover — a monotonic set accurately
|
||||
//! described a one-way door. Once recovery exists the assumption inverts: plug
|
||||
//! the WAN back in and the node would stay `Degraded` until the process
|
||||
//! restarted, and `Degraded` would come to mean "something broke at some point
|
||||
//! since boot" rather than "something is broken now".
|
||||
//! [`Self::classify_health`](SupervisorFsm::classify_health) is therefore
|
||||
//! recomputed on every transition **in both directions**.
|
||||
//!
|
||||
//! An absent transport still counts as *up*. It came up — `start_async` returns
|
||||
//! `Ok` with the transport absent — so it does not push a single-transport node
|
||||
//! into the fatal [`FailReason::NoTransports`], which would make a node that
|
||||
//! merely booted before its wifi exit instead of waiting. Absence degrades; it
|
||||
//! never kills.
|
||||
//!
|
||||
//! There is deliberately **no restart action**. Rebinding is owned by the
|
||||
//! transport's own binder task, which is where the file descriptor and the
|
||||
//! presence watcher live; the FSM is told what happened and republishes health.
|
||||
//!
|
||||
//! [`Absent`]: crate::transport::ethernet::Presence::Absent
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::sync::Arc;
|
||||
@@ -174,6 +207,22 @@ pub(crate) enum Event {
|
||||
/// The child whose task/thread exited.
|
||||
child: Child,
|
||||
},
|
||||
/// An interface-bound child lost its interface — it was never there at
|
||||
/// start, or it went away at runtime. The child stays *up* (the transport
|
||||
/// object survives detach, only its socket goes) but contributes
|
||||
/// `Degraded`. Valid while `Running`; while `Starting` the edge is recorded
|
||||
/// so start-completion health already reflects it.
|
||||
ChildAbsent {
|
||||
/// The child whose interface is absent.
|
||||
child: Child,
|
||||
},
|
||||
/// An interface-bound child's interface came back and it rebound. Clears
|
||||
/// the absence and republishes health, which is how `Degraded` becomes
|
||||
/// reversible.
|
||||
ChildPresent {
|
||||
/// The child whose interface is present again.
|
||||
child: Child,
|
||||
},
|
||||
}
|
||||
|
||||
/// A driver-scheduled timer the supervisor can arm. Only the
|
||||
@@ -320,7 +369,16 @@ pub(crate) struct SupervisorFsm {
|
||||
up: HashSet<Child>,
|
||||
/// Configured children that failed to start during the current bring-up.
|
||||
/// Feeds the `Degraded` health determination when `pending` empties.
|
||||
///
|
||||
/// One-way within a bring-up: a start failure is not recoverable, so
|
||||
/// nothing removes from this set until the next `Start`.
|
||||
failed: HashSet<Child>,
|
||||
/// Children that are up but whose network interface is currently absent.
|
||||
///
|
||||
/// Reversible, unlike [`Self::failed`] — that is the whole reason it is a
|
||||
/// second set rather than more entries in the first. Feeds `Degraded` the
|
||||
/// same way, and empties as interfaces come back.
|
||||
absent: HashSet<Child>,
|
||||
}
|
||||
|
||||
impl SupervisorFsm {
|
||||
@@ -330,6 +388,7 @@ impl SupervisorFsm {
|
||||
state: SupState::Created,
|
||||
up: HashSet::new(),
|
||||
failed: HashSet::new(),
|
||||
absent: HashSet::new(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -348,6 +407,7 @@ impl SupervisorFsm {
|
||||
},
|
||||
up: up.into_iter().collect(),
|
||||
failed: HashSet::new(),
|
||||
absent: HashSet::new(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -357,13 +417,26 @@ impl SupervisorFsm {
|
||||
&self.state
|
||||
}
|
||||
|
||||
/// The configured children that failed to start during bring-up. The driver
|
||||
/// reads this on the `Degraded` start outcome to enumerate the degraded
|
||||
/// children in an operator-visible `warn!`.
|
||||
/// The configured children that failed to start during bring-up. Kept for
|
||||
/// the tests that pin the failure-vs-absence split; the driver reports
|
||||
/// [`Self::degraded_children`], which is the union of the two.
|
||||
#[cfg(test)]
|
||||
pub(in crate::node) fn failed(&self) -> &HashSet<Child> {
|
||||
&self.failed
|
||||
}
|
||||
|
||||
/// Children whose interface is currently absent.
|
||||
pub(in crate::node) fn absent(&self) -> &HashSet<Child> {
|
||||
&self.absent
|
||||
}
|
||||
|
||||
/// Every child currently contributing `Degraded` — the ones that failed to
|
||||
/// start plus the ones whose interface is away. This is what an operator
|
||||
/// wants named when the node reports `Degraded`.
|
||||
pub(in crate::node) fn degraded_children(&self) -> HashSet<Child> {
|
||||
self.failed.union(&self.absent).copied().collect()
|
||||
}
|
||||
|
||||
/// Whether the machine is in the bounded-drain window. The driver uses this
|
||||
/// after the rx loop returns to decide between the drain-teardown path and
|
||||
/// the immediate-`stop()` fallback.
|
||||
@@ -398,6 +471,8 @@ impl SupervisorFsm {
|
||||
Event::DrainDeadlineElapsed => self.on_drain_deadline_elapsed(),
|
||||
Event::ChildStopped { child } => self.on_child_stopped(child),
|
||||
Event::ChildExited { child } => self.on_child_exited(child),
|
||||
Event::ChildAbsent { child } => self.on_child_absent(child),
|
||||
Event::ChildPresent { child } => self.on_child_present(child),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -444,6 +519,7 @@ impl SupervisorFsm {
|
||||
|
||||
self.up.clear();
|
||||
self.failed.clear();
|
||||
self.absent.clear();
|
||||
|
||||
// A node with no children at all resolves health immediately. Zero
|
||||
// transports up → `Failed` (this is the behavioral
|
||||
@@ -497,19 +573,30 @@ impl SupervisorFsm {
|
||||
self.classify_health()
|
||||
}
|
||||
|
||||
/// Classify health from the current `up` / `failed` sets and set the
|
||||
/// resulting state, returning the [`NodeState`] the driver should publish.
|
||||
/// Shared by start-completion ([`Self::resolve_start_health`]) and runtime
|
||||
/// child-exit ([`Self::on_child_exited`]):
|
||||
/// Classify health from the current `up` / `failed` / `absent` sets and set
|
||||
/// the resulting state, returning the [`NodeState`] the driver should
|
||||
/// publish. Shared by start-completion ([`Self::resolve_start_health`]),
|
||||
/// runtime child-exit ([`Self::on_child_exited`]) and the presence edges
|
||||
/// ([`Self::on_child_absent`] / [`Self::on_child_present`]):
|
||||
///
|
||||
/// - zero transports up → [`SupState::Failed`] / [`NodeState::Failed`];
|
||||
/// - ≥1 transport up but some child in `failed` → [`Health::Degraded`] /
|
||||
/// [`NodeState::Degraded`];
|
||||
/// - everything up and nothing failed → [`Health::Full`] / [`NodeState::Running`].
|
||||
/// - ≥1 transport up but some child in `failed` or `absent` →
|
||||
/// [`Health::Degraded`] / [`NodeState::Degraded`];
|
||||
/// - everything up, nothing failed, nothing absent → [`Health::Full`] /
|
||||
/// [`NodeState::Running`].
|
||||
///
|
||||
/// Worker-pool failures are captured in `failed` like any other optional
|
||||
/// child, so they contribute `Degraded` (never `Failed`) — the inline crypto
|
||||
/// fallback keeps the node correct without the pools.
|
||||
///
|
||||
/// A transport whose interface is absent is still counted among
|
||||
/// `transports_up`: it came up, it is simply not bound. Excluding it would
|
||||
/// make a single-ethernet node that booted before its wifi resolve to the
|
||||
/// fatal [`FailReason::NoTransports`] and exit — which is the failure this
|
||||
/// whole mechanism exists to remove. Absence degrades; it never kills.
|
||||
///
|
||||
/// Recomputed in **both** directions. This function is the reason
|
||||
/// `Degraded` is a level rather than a latch.
|
||||
fn classify_health(&mut self) -> NodeState {
|
||||
let transports_up = self
|
||||
.up
|
||||
@@ -521,10 +608,10 @@ impl SupervisorFsm {
|
||||
reason: FailReason::NoTransports,
|
||||
};
|
||||
NodeState::Failed
|
||||
} else if !self.failed.is_empty() {
|
||||
} else if !self.failed.is_empty() || !self.absent.is_empty() {
|
||||
self.state = SupState::Running {
|
||||
health: Health::Degraded {
|
||||
reasons: self.failed.clone(),
|
||||
reasons: self.degraded_children(),
|
||||
},
|
||||
};
|
||||
NodeState::Degraded
|
||||
@@ -536,6 +623,53 @@ impl SupervisorFsm {
|
||||
}
|
||||
}
|
||||
|
||||
/// An interface-bound child lost its interface.
|
||||
///
|
||||
/// The child stays in the up-set: the transport object survives detach —
|
||||
/// config, id, statistics and neighbor buffer persist, only the descriptor
|
||||
/// and its loops go. Only the reason set changes.
|
||||
///
|
||||
/// While `Starting` the edge is recorded silently; start-completion health
|
||||
/// picks it up when `pending` empties, so a boot race resolves to
|
||||
/// `Degraded` on the first publish rather than publishing `Full` and
|
||||
/// immediately correcting it. Inert outside `Starting` / `Running`: a
|
||||
/// teardown in flight owns its own bookkeeping.
|
||||
fn on_child_absent(&mut self, child: Child) -> Vec<Action> {
|
||||
match self.state {
|
||||
SupState::Starting { .. } => {
|
||||
self.absent.insert(child);
|
||||
Vec::new()
|
||||
}
|
||||
SupState::Running { .. } => {
|
||||
if !self.absent.insert(child) {
|
||||
return Vec::new();
|
||||
}
|
||||
vec![Action::PublishState(self.classify_health())]
|
||||
}
|
||||
_ => Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
/// An interface-bound child's interface came back.
|
||||
///
|
||||
/// Clears the absence and republishes. A child that is not currently
|
||||
/// recorded absent produces nothing — a duplicate edge is not an event.
|
||||
fn on_child_present(&mut self, child: Child) -> Vec<Action> {
|
||||
match self.state {
|
||||
SupState::Starting { .. } => {
|
||||
self.absent.remove(&child);
|
||||
Vec::new()
|
||||
}
|
||||
SupState::Running { .. } => {
|
||||
if !self.absent.remove(&child) {
|
||||
return Vec::new();
|
||||
}
|
||||
vec![Action::PublishState(self.classify_health())]
|
||||
}
|
||||
_ => Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
fn on_stop(&mut self) -> Vec<Action> {
|
||||
if !matches!(self.state, SupState::Running { .. }) {
|
||||
return Vec::new();
|
||||
@@ -588,6 +722,7 @@ impl SupervisorFsm {
|
||||
if let SupState::Stopping { pending } = &mut self.state {
|
||||
pending.remove(&child);
|
||||
self.up.remove(&child);
|
||||
self.absent.remove(&child);
|
||||
if pending.is_empty() {
|
||||
self.state = SupState::Stopped;
|
||||
}
|
||||
@@ -617,6 +752,8 @@ impl SupervisorFsm {
|
||||
if !self.up.remove(&child) {
|
||||
return Vec::new();
|
||||
}
|
||||
// An exit supersedes an absence: the child is gone, not waiting.
|
||||
self.absent.remove(&child);
|
||||
self.failed.insert(child);
|
||||
vec![Action::PublishState(self.classify_health())]
|
||||
}
|
||||
@@ -1378,4 +1515,289 @@ mod tests {
|
||||
);
|
||||
assert!(s.failed().is_empty());
|
||||
}
|
||||
|
||||
// ── Interface presence ────────────────────────────────────────────────
|
||||
//
|
||||
// The absence set is separate from `failed` because it is reversible, and
|
||||
// reversibility is what turns `Degraded` from a latch into a level.
|
||||
|
||||
/// Helper: bring a full node up cleanly and leave it `Running{Full}`.
|
||||
fn running_node() -> SupervisorFsm {
|
||||
let mut s = SupervisorFsm::new();
|
||||
s.step(start_full());
|
||||
for child in [
|
||||
Child::Transport(tid(1)),
|
||||
Child::Transport(tid(2)),
|
||||
Child::EncryptWorkers,
|
||||
Child::DecryptWorkers,
|
||||
Child::Nostr,
|
||||
Child::Mdns,
|
||||
Child::Tun,
|
||||
Child::Dns,
|
||||
] {
|
||||
s.step(Event::SubstrateUp { child });
|
||||
}
|
||||
assert_eq!(
|
||||
s.state(),
|
||||
&SupState::Running {
|
||||
health: Health::Full
|
||||
}
|
||||
);
|
||||
s
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_absent_interface_degrades_a_running_node() {
|
||||
let mut s = running_node();
|
||||
assert_eq!(
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1))
|
||||
}),
|
||||
vec![Action::PublishState(NodeState::Degraded)]
|
||||
);
|
||||
assert!(s.absent().contains(&Child::Transport(tid(1))));
|
||||
// Absence is not failure: the two sets stay distinct.
|
||||
assert!(s.failed().is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_returning_interface_clears_degraded() {
|
||||
// The regression this whole design turns on: plug the WAN back in and
|
||||
// the node must leave `Degraded`, not carry it until the process
|
||||
// restarts. `Degraded` is a level, not a latch.
|
||||
let mut s = running_node();
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1)),
|
||||
});
|
||||
assert_eq!(
|
||||
s.step(Event::ChildPresent {
|
||||
child: Child::Transport(tid(1))
|
||||
}),
|
||||
vec![Action::PublishState(NodeState::Running)]
|
||||
);
|
||||
assert_eq!(
|
||||
s.state(),
|
||||
&SupState::Running {
|
||||
health: Health::Full
|
||||
}
|
||||
);
|
||||
assert!(s.absent().is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn health_reflects_the_last_interface_still_away() {
|
||||
// Two interfaces away, one returns: still Degraded. Only the empty
|
||||
// absence set publishes Full.
|
||||
let mut s = running_node();
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1)),
|
||||
});
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(2)),
|
||||
});
|
||||
assert_eq!(
|
||||
s.step(Event::ChildPresent {
|
||||
child: Child::Transport(tid(1))
|
||||
}),
|
||||
vec![Action::PublishState(NodeState::Degraded)]
|
||||
);
|
||||
assert_eq!(
|
||||
s.step(Event::ChildPresent {
|
||||
child: Child::Transport(tid(2))
|
||||
}),
|
||||
vec![Action::PublishState(NodeState::Running)]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_duplicate_presence_edge_is_not_an_event() {
|
||||
let mut s = running_node();
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1)),
|
||||
});
|
||||
assert_eq!(
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1))
|
||||
}),
|
||||
vec![],
|
||||
"re-reporting the same absence must not republish"
|
||||
);
|
||||
s.step(Event::ChildPresent {
|
||||
child: Child::Transport(tid(1)),
|
||||
});
|
||||
assert_eq!(
|
||||
s.step(Event::ChildPresent {
|
||||
child: Child::Transport(tid(1))
|
||||
}),
|
||||
vec![],
|
||||
"re-reporting the same return must not republish"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn absence_during_bringup_resolves_degraded_on_the_first_publish() {
|
||||
// The boot race. The transport reports absence from inside its own
|
||||
// start, then reports up. Start-completion health must already know,
|
||||
// so the node never publishes `Full` and corrects itself.
|
||||
let mut s = SupervisorFsm::new();
|
||||
s.step(start_full());
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1)),
|
||||
});
|
||||
for child in [
|
||||
Child::Transport(tid(1)),
|
||||
Child::Transport(tid(2)),
|
||||
Child::EncryptWorkers,
|
||||
Child::DecryptWorkers,
|
||||
Child::Nostr,
|
||||
Child::Mdns,
|
||||
Child::Tun,
|
||||
] {
|
||||
assert_eq!(s.step(Event::SubstrateUp { child }), vec![]);
|
||||
}
|
||||
assert_eq!(
|
||||
s.step(Event::SubstrateUp { child: Child::Dns }),
|
||||
vec![Action::PublishState(NodeState::Degraded)]
|
||||
);
|
||||
let mut reasons = HashSet::new();
|
||||
reasons.insert(Child::Transport(tid(1)));
|
||||
assert_eq!(
|
||||
s.state(),
|
||||
&SupState::Running {
|
||||
health: Health::Degraded { reasons }
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_lone_absent_transport_degrades_rather_than_fails() {
|
||||
// A single-ethernet node that boots before its wifi. The transport
|
||||
// came up — absence is a state, not a start failure — so it counts
|
||||
// among the transports up and the node serves. Resolving to `Failed`
|
||||
// here would make the daemon exit on the very race this mechanism
|
||||
// exists to survive.
|
||||
let mut s = SupervisorFsm::new();
|
||||
s.step(Event::Start {
|
||||
transports: vec![tid(1)],
|
||||
encrypt_workers: false,
|
||||
decrypt_workers: false,
|
||||
nostr: false,
|
||||
mdns: false,
|
||||
tun: false,
|
||||
dns: false,
|
||||
});
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1)),
|
||||
});
|
||||
assert_eq!(
|
||||
s.step(Event::SubstrateUp {
|
||||
child: Child::Transport(tid(1))
|
||||
}),
|
||||
vec![Action::PublishState(NodeState::Degraded)]
|
||||
);
|
||||
assert!(
|
||||
!matches!(s.state(), SupState::Failed { .. }),
|
||||
"an absent interface must never be the fatal no-transports case"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_exit_supersedes_an_absence() {
|
||||
// A transport that is away and then exits is gone, not waiting: it
|
||||
// leaves the up-set and moves from `absent` to `failed`, so a later
|
||||
// spurious `ChildPresent` cannot resurrect it into `Full`.
|
||||
let mut s = running_node();
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1)),
|
||||
});
|
||||
s.step(Event::ChildExited {
|
||||
child: Child::Transport(tid(1)),
|
||||
});
|
||||
assert!(s.absent().is_empty());
|
||||
assert!(s.failed().contains(&Child::Transport(tid(1))));
|
||||
assert_eq!(
|
||||
s.step(Event::ChildPresent {
|
||||
child: Child::Transport(tid(1))
|
||||
}),
|
||||
vec![]
|
||||
);
|
||||
assert!(matches!(
|
||||
s.state(),
|
||||
SupState::Running {
|
||||
health: Health::Degraded { .. }
|
||||
}
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn presence_edges_are_inert_outside_starting_and_running() {
|
||||
// Teardown and drain own their own bookkeeping; a late edge from a
|
||||
// binder that has not noticed the stop must not author actions.
|
||||
let mut created = SupervisorFsm::new();
|
||||
assert_eq!(
|
||||
created.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1))
|
||||
}),
|
||||
vec![]
|
||||
);
|
||||
|
||||
let mut draining = running_node();
|
||||
draining.step(Event::Drain { deadline_ms: 1 });
|
||||
assert_eq!(
|
||||
draining.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1))
|
||||
}),
|
||||
vec![]
|
||||
);
|
||||
assert!(
|
||||
draining.is_draining(),
|
||||
"a presence edge must not end a drain"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_new_start_clears_the_previous_absences() {
|
||||
let mut s = running_node();
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1)),
|
||||
});
|
||||
s.step(Event::Stop);
|
||||
for child in [
|
||||
Child::Dns,
|
||||
Child::Nostr,
|
||||
Child::Mdns,
|
||||
Child::Transport(tid(1)),
|
||||
Child::Transport(tid(2)),
|
||||
Child::Tun,
|
||||
] {
|
||||
s.step(Event::ChildStopped { child });
|
||||
}
|
||||
s.step(start_full());
|
||||
assert!(s.absent().is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn degraded_children_is_the_union_of_both_reason_sets() {
|
||||
let mut s = SupervisorFsm::new();
|
||||
s.step(start_full());
|
||||
s.step(Event::ChildAbsent {
|
||||
child: Child::Transport(tid(1)),
|
||||
});
|
||||
s.step(Event::SubstrateFailed { child: Child::Mdns });
|
||||
for child in [
|
||||
Child::Transport(tid(1)),
|
||||
Child::Transport(tid(2)),
|
||||
Child::EncryptWorkers,
|
||||
Child::DecryptWorkers,
|
||||
Child::Nostr,
|
||||
Child::Tun,
|
||||
Child::Dns,
|
||||
] {
|
||||
s.step(Event::SubstrateUp { child });
|
||||
}
|
||||
let degraded = s.degraded_children();
|
||||
assert!(degraded.contains(&Child::Mdns));
|
||||
assert!(degraded.contains(&Child::Transport(tid(1))));
|
||||
assert_eq!(degraded.len(), 2);
|
||||
}
|
||||
}
|
||||
|
||||
+109
-3
@@ -374,6 +374,16 @@ pub struct Node {
|
||||
/// SYN/SYN-ACK clamp can use the smaller of the local-egress floor
|
||||
/// and the learned per-destination path MTU.
|
||||
path_mtu_lookup: crate::upper::tun::PathMtuLookup,
|
||||
/// Node-global TCP MSS ceiling, shared live with the TUN reader and writer
|
||||
/// threads and recomputed whenever the set of *bound* transports changes.
|
||||
///
|
||||
/// Sits beside `path_mtu_lookup` because it answers the other half of the
|
||||
/// same question at the same moment: that map supplies the per-destination
|
||||
/// ceiling, this supplies the local-egress one, and the clamp takes the
|
||||
/// smaller. Both have to be read live — a transport that binds after start
|
||||
/// can be the narrow one, and one that unbinds can be the reason the node
|
||||
/// was clamped at all.
|
||||
tun_mss_ceiling: crate::upper::tun::MssCeiling,
|
||||
/// Which transport last supplied a *link seed* into `path_mtu_lookup`,
|
||||
/// per destination.
|
||||
///
|
||||
@@ -418,6 +428,20 @@ pub struct Node {
|
||||
/// rx_loop select arm that feeds `Event::ChildExited` to the supervisor FSM.
|
||||
child_exit_rx: Option<tokio::sync::mpsc::Receiver<crate::node::lifecycle::supervisor::Child>>,
|
||||
|
||||
// === Interface Presence Channel ===
|
||||
/// Sender half of the interface-presence channel, cloned into every
|
||||
/// interface-bound transport so its binder task can report attach and
|
||||
/// detach. Held on `self` for the rx_loop's lifetime as the keep-alive
|
||||
/// sender, exactly like [`Self::child_exit_tx`].
|
||||
///
|
||||
/// Separate from the child-exit channel because presence is *reversible*:
|
||||
/// an exit is one-way, an interface comes back.
|
||||
transport_presence_tx: Option<crate::transport::PresenceTx>,
|
||||
/// Receiver half of the interface-presence channel, `take()`-en by the
|
||||
/// rx_loop select arm that feeds `Event::ChildAbsent` / `Event::ChildPresent`
|
||||
/// to the supervisor FSM.
|
||||
transport_presence_rx: Option<crate::transport::PresenceRx>,
|
||||
|
||||
// === Per-Peer Control Machines ===
|
||||
/// Per-peer lifecycle control FSMs, keyed by the stable `LinkId` that spans
|
||||
/// the handshake→active lifetime. Each machine owns its handshake crypto
|
||||
@@ -813,6 +837,8 @@ impl Node {
|
||||
packet_rx: None,
|
||||
child_exit_tx: None,
|
||||
child_exit_rx: None,
|
||||
transport_presence_tx: None,
|
||||
transport_presence_rx: None,
|
||||
peer_machines: HashMap::new(),
|
||||
peer_timers: HashMap::new(),
|
||||
peers: HashMap::new(),
|
||||
@@ -880,6 +906,12 @@ impl Node {
|
||||
peer_acl,
|
||||
host_map,
|
||||
path_mtu_lookup: Arc::new(std::sync::RwLock::new(HashMap::new())),
|
||||
// Seeded at the IPv6 minimum, which is what `transport_mtu()`
|
||||
// itself falls back to when nothing is bound. Refreshed before
|
||||
// the TUN threads start and on every change to the bound set.
|
||||
tun_mss_ceiling: Arc::new(std::sync::atomic::AtomicU16::new(
|
||||
crate::upper::icmp::mss_ceiling(crate::upper::tun::IPV6_MIN_MTU),
|
||||
)),
|
||||
path_mtu_seeded_by: Arc::new(std::sync::RwLock::new(HashMap::new())),
|
||||
#[cfg(unix)]
|
||||
decrypt_registered_sessions: std::collections::HashSet::new(),
|
||||
@@ -978,6 +1010,8 @@ impl Node {
|
||||
packet_rx: None,
|
||||
child_exit_tx: None,
|
||||
child_exit_rx: None,
|
||||
transport_presence_tx: None,
|
||||
transport_presence_rx: None,
|
||||
peer_machines: HashMap::new(),
|
||||
peer_timers: HashMap::new(),
|
||||
peers: HashMap::new(),
|
||||
@@ -1042,6 +1076,12 @@ impl Node {
|
||||
peer_acl,
|
||||
host_map,
|
||||
path_mtu_lookup: Arc::new(std::sync::RwLock::new(HashMap::new())),
|
||||
// Seeded at the IPv6 minimum, which is what `transport_mtu()`
|
||||
// itself falls back to when nothing is bound. Refreshed before
|
||||
// the TUN threads start and on every change to the bound set.
|
||||
tun_mss_ceiling: Arc::new(std::sync::atomic::AtomicU16::new(
|
||||
crate::upper::icmp::mss_ceiling(crate::upper::tun::IPV6_MIN_MTU),
|
||||
)),
|
||||
path_mtu_seeded_by: Arc::new(std::sync::RwLock::new(HashMap::new())),
|
||||
#[cfg(unix)]
|
||||
decrypt_registered_sessions: std::collections::HashSet::new(),
|
||||
@@ -1098,6 +1138,11 @@ impl Node {
|
||||
let mut eth =
|
||||
EthernetTransport::new(transport_id, name, eth_config, packet_tx.clone());
|
||||
eth.set_local_pubkey(xonly);
|
||||
// The binder task reports attach and detach here, so node
|
||||
// health tracks the interface in both directions.
|
||||
if let Some(tx) = self.transport_presence_tx.clone() {
|
||||
eth.set_presence_tx(tx);
|
||||
}
|
||||
transports.push(TransportHandle::Ethernet(eth));
|
||||
}
|
||||
}
|
||||
@@ -1411,6 +1456,46 @@ impl Node {
|
||||
crate::upper::icmp::effective_ipv6_mtu(self.transport_mtu())
|
||||
}
|
||||
|
||||
/// The TCP MSS ceiling the TUN threads are currently clamping to.
|
||||
#[cfg(test)]
|
||||
pub(crate) fn tun_mss_ceiling(&self) -> u16 {
|
||||
self.tun_mss_ceiling
|
||||
.load(std::sync::atomic::Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Recompute the shared TUN MSS ceiling from the currently bound
|
||||
/// transports, and log it if it moved.
|
||||
///
|
||||
/// Called wherever the bound set can change — the presence edges that
|
||||
/// bind and unbind an interface-bound transport, and a child exiting —
|
||||
/// so the clamp the TUN threads apply keeps agreeing with the
|
||||
/// `effective_ipv6_mtu` this node reports in `show_status`.
|
||||
///
|
||||
/// Moves in **both** directions, deliberately. A narrow interface
|
||||
/// appearing has to tighten the ceiling or the clamp is wrong for
|
||||
/// traffic that will egress over it; that same interface going away has
|
||||
/// to release it, or unplugging a low-MTU adapter leaves the node
|
||||
/// over-clamped until it restarts. It is the same argument that makes
|
||||
/// `Degraded` a level rather than a latch: nothing here is one-way once
|
||||
/// a transport can come back.
|
||||
///
|
||||
/// Existing flows are not re-clamped — MSS is negotiated per connection
|
||||
/// at SYN time, so a change applies to connections opened after it.
|
||||
pub(crate) fn refresh_tun_mss_ceiling(&self) {
|
||||
use std::sync::atomic::Ordering;
|
||||
|
||||
let ceiling = crate::upper::icmp::mss_ceiling(self.transport_mtu());
|
||||
let previous = self.tun_mss_ceiling.swap(ceiling, Ordering::Relaxed);
|
||||
if previous != ceiling {
|
||||
tracing::info!(
|
||||
previous_max_mss = previous,
|
||||
max_mss = ceiling,
|
||||
effective_ipv6_mtu = self.effective_ipv6_mtu(),
|
||||
"Node egress MTU changed; TCP MSS ceiling updated for new connections"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// Get the transport MTU governing the global TUN-boundary MSS clamp.
|
||||
///
|
||||
/// Returns the **minimum** MTU across all operational transports, or
|
||||
@@ -1426,10 +1511,17 @@ impl Node {
|
||||
/// to vary across HashMap iteration order + async-startup race) makes
|
||||
/// the clamp deterministic across daemon restarts.
|
||||
pub fn transport_mtu(&self) -> u16 {
|
||||
// `is_bound`, not `is_operational`. An interface-bound transport is
|
||||
// "operational" from the moment it starts, whether or not its
|
||||
// interface exists — so filtering on that let a transport whose
|
||||
// interface has never appeared clamp the whole node's IPv6 MTU to a
|
||||
// number derived from hardware that is not present. Before dynamic
|
||||
// binding a transport that could not bind was never inserted here at
|
||||
// all, so the distinction did not exist to get wrong.
|
||||
let min_operational = self
|
||||
.transports
|
||||
.values()
|
||||
.filter(|h| h.is_operational())
|
||||
.filter(|h| h.is_bound())
|
||||
.map(|h| h.mtu())
|
||||
.min();
|
||||
if let Some(mtu) = min_operational {
|
||||
@@ -2318,8 +2410,13 @@ impl Node {
|
||||
.collect();
|
||||
|
||||
// --- transports (show_transports) ---
|
||||
let transport_rows: Vec<snap::TransportRow> = self
|
||||
.transport_ids()
|
||||
// Ascending id, matching `show_transports`; see the note there. The
|
||||
// off-loop renderer reads this table verbatim, so the two paths would
|
||||
// otherwise disagree about ordering as well as being arbitrary.
|
||||
let mut transport_ids: Vec<_> = self.transport_ids().copied().collect();
|
||||
transport_ids.sort_by_key(|id| id.as_u32());
|
||||
let transport_rows: Vec<snap::TransportRow> = transport_ids
|
||||
.iter()
|
||||
.map(|id| {
|
||||
let handle = self.get_transport(id).unwrap();
|
||||
snap::TransportRow {
|
||||
@@ -2335,6 +2432,15 @@ impl Node {
|
||||
.tor_monitoring()
|
||||
.map(|m| serde_json::to_value(&m).unwrap_or_default()),
|
||||
stats: handle.transport_stats(),
|
||||
interface: handle.interface_presence().map(|p| snap::InterfaceRow {
|
||||
name: handle.interface_name().unwrap_or_default().to_string(),
|
||||
presence: p.presence,
|
||||
carrier: p.carrier,
|
||||
policy: p.policy,
|
||||
since_secs: p.since_secs,
|
||||
binds: p.binds,
|
||||
failed_attempts: p.failed_attempts,
|
||||
}),
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
|
||||
@@ -1842,6 +1842,114 @@ async fn test_transport_mtu_returns_min_across_operational() {
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn the_tun_mss_ceiling_follows_a_transport_arriving_and_leaving() {
|
||||
// The TUN reader and writer used to be handed a `u16` computed once at
|
||||
// spawn. Every other consumer of `transport_mtu()` reads it live, so once
|
||||
// a transport could bind minutes after start the daemon reported one
|
||||
// effective MTU in `show_status` and clamped to another.
|
||||
//
|
||||
// Both directions. A narrow transport arriving has to tighten the ceiling
|
||||
// or traffic egressing over it is clamped too loose; the same transport
|
||||
// leaving has to release it, or unplugging a low-MTU adapter leaves the
|
||||
// node over-clamped until it restarts.
|
||||
let mut node = make_node();
|
||||
let (packet_tx, packet_rx) = packet_channel(64);
|
||||
node.supervisor.packet_tx = Some(packet_tx);
|
||||
node.packet_rx = Some(packet_rx);
|
||||
|
||||
let wide = make_udp_transport_with_mtu(1, 1452).await;
|
||||
node.transports.insert(TransportId::new(1), wide);
|
||||
node.refresh_tun_mss_ceiling();
|
||||
let wide_ceiling = node.tun_mss_ceiling();
|
||||
assert_eq!(
|
||||
wide_ceiling,
|
||||
crate::upper::icmp::mss_ceiling(1452),
|
||||
"the seeded ceiling must match the only bound transport"
|
||||
);
|
||||
|
||||
// A narrower transport arrives after the TUN threads would already be
|
||||
// running. The shared ceiling has to tighten.
|
||||
let narrow = make_udp_transport_with_mtu(2, 1280).await;
|
||||
node.transports.insert(TransportId::new(2), narrow);
|
||||
node.refresh_tun_mss_ceiling();
|
||||
let narrow_ceiling = node.tun_mss_ceiling();
|
||||
assert_eq!(narrow_ceiling, crate::upper::icmp::mss_ceiling(1280));
|
||||
assert!(
|
||||
narrow_ceiling < wide_ceiling,
|
||||
"a narrower transport must tighten the clamp, not be ignored"
|
||||
);
|
||||
assert_eq!(
|
||||
narrow_ceiling,
|
||||
crate::upper::icmp::mss_ceiling(node.transport_mtu()),
|
||||
"the clamp and the reported MTU must not disagree"
|
||||
);
|
||||
|
||||
// ...and leaving has to release it again.
|
||||
if let Some(mut gone) = node.transports.remove(&TransportId::new(2)) {
|
||||
gone.stop().await.ok();
|
||||
}
|
||||
node.refresh_tun_mss_ceiling();
|
||||
assert_eq!(
|
||||
node.tun_mss_ceiling(),
|
||||
wide_ceiling,
|
||||
"the ceiling must rise again when the narrow transport goes away"
|
||||
);
|
||||
|
||||
for transport in node.transports.values_mut() {
|
||||
transport.stop().await.ok();
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_presence_edge_refreshes_the_tun_mss_ceiling_without_being_asked() {
|
||||
// The one above pins the arithmetic; this pins the wiring. A presence
|
||||
// edge arriving has to refresh the ceiling on its own — if the refresh is
|
||||
// dropped from the edge handlers the value silently stops tracking, which
|
||||
// is the defect in its original form.
|
||||
let mut node = make_node();
|
||||
let (packet_tx, packet_rx) = packet_channel(64);
|
||||
node.supervisor.packet_tx = Some(packet_tx);
|
||||
node.packet_rx = Some(packet_rx);
|
||||
|
||||
let (presence_tx, presence_rx) = tokio::sync::mpsc::channel(4);
|
||||
node.transport_presence_tx = Some(presence_tx.clone());
|
||||
node.transport_presence_rx = Some(presence_rx);
|
||||
|
||||
// Nothing bound: the conservative seed.
|
||||
node.refresh_tun_mss_ceiling();
|
||||
let seeded = node.tun_mss_ceiling();
|
||||
assert_eq!(seeded, crate::upper::icmp::mss_ceiling(1280));
|
||||
|
||||
// A wide transport appears, and an edge announces it. No explicit
|
||||
// refresh call here — draining the edge is the whole trigger.
|
||||
let wide = make_udp_transport_with_mtu(1, 1452).await;
|
||||
node.transports.insert(TransportId::new(1), wide);
|
||||
presence_tx
|
||||
.send(crate::transport::TransportPresence {
|
||||
transport_id: TransportId::new(1),
|
||||
present: true,
|
||||
})
|
||||
.await
|
||||
.expect("presence edge queued");
|
||||
node.drain_transport_presence();
|
||||
|
||||
assert_eq!(
|
||||
node.tun_mss_ceiling(),
|
||||
crate::upper::icmp::mss_ceiling(1452),
|
||||
"draining a presence edge must refresh the ceiling on its own"
|
||||
);
|
||||
assert_ne!(
|
||||
node.tun_mss_ceiling(),
|
||||
seeded,
|
||||
"the ceiling stayed at its seed, so the edge did not refresh it"
|
||||
);
|
||||
|
||||
for transport in node.transports.values_mut() {
|
||||
transport.stop().await.ok();
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_transport_mtu_fallback_when_no_operational_transports() {
|
||||
// No transports configured at all → falls back to 1280 (IPv6 minimum).
|
||||
|
||||
@@ -9,6 +9,97 @@ use crate::transport::TransportError;
|
||||
/// Broadcast MAC address.
|
||||
pub const ETHERNET_BROADCAST: [u8; 6] = [0xff; 6];
|
||||
|
||||
/// Whether the named interface exists and is administratively up.
|
||||
///
|
||||
/// Presence is `IFF_UP` — the interface exists and the operator has enabled
|
||||
/// it — and deliberately **not** `IFF_RUNNING`.
|
||||
///
|
||||
/// Carrier is a different question from bindability, and only the second one
|
||||
/// belongs in a bind gate. An `AF_PACKET` socket on a carrier-less bridge is
|
||||
/// perfectly valid and starts carrying traffic the instant a member port comes
|
||||
/// up, with no rebind: the socket outlives the carrier. Gating on `IFF_RUNNING`
|
||||
/// bought nothing and cost three things —
|
||||
///
|
||||
/// - `br-lan` on a router with nothing plugged into its LAN ports is `UP` with
|
||||
/// `NO-CARRIER`, so a perfectly healthy wifi-only router reported `Degraded`
|
||||
/// forever;
|
||||
/// - every carrier flap the socket would have survived became an unbind /
|
||||
/// rebind cycle, which is churn the presence machine then has to damp;
|
||||
/// - an 802.11s mesh interface that reports `RUNNING` only once it has peered
|
||||
/// cannot peer, because peering needs beacons, which need a bound socket,
|
||||
/// which the gate refuses. A deadlock reachable on shipped hardware.
|
||||
///
|
||||
/// The signal `IFF_RUNNING` does carry — "is anything plugged in" — is not
|
||||
/// lost; it is reported alongside presence by [`interface_carrier`] and
|
||||
/// surfaced in `show_transports`, where an operator can read it without it
|
||||
/// steering the daemon.
|
||||
///
|
||||
/// `getifaddrs` rather than an `SIOCGIFFLAGS` ioctl: it needs no socket, so
|
||||
/// the presence watcher can poll before any file descriptor exists, and it is
|
||||
/// spelled the same on Linux and the BSDs.
|
||||
#[cfg(unix)]
|
||||
pub fn interface_present(interface: &str) -> bool {
|
||||
interface_has_flags(interface, libc::IFF_UP as u32)
|
||||
}
|
||||
|
||||
/// The kernel's index for the named interface, or `None` if it does not exist.
|
||||
///
|
||||
/// A name is not a device, and neither is a name that is still there. Both
|
||||
/// backends bind by index — `AF_PACKET` stores `sll_ifindex`, and a BPF
|
||||
/// descriptor follows the device it was attached to — so an interface deleted
|
||||
/// and recreated under the same name leaves the socket attached to a device
|
||||
/// that no longer exists while the *name* resolves perfectly well. Comparing
|
||||
/// the live index against the one captured at bind is what tells those apart.
|
||||
#[cfg(unix)]
|
||||
pub fn interface_index(interface: &str) -> Option<u32> {
|
||||
let c_name = std::ffi::CString::new(interface).ok()?;
|
||||
// Cheaper than `getifaddrs`: one syscall, no allocation, no walk.
|
||||
match unsafe { libc::if_nametoindex(c_name.as_ptr()) } {
|
||||
0 => None,
|
||||
idx => Some(idx),
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether the named interface currently has carrier (`IFF_RUNNING`).
|
||||
///
|
||||
/// Reported, never acted on — see [`interface_present`]. `false` for an
|
||||
/// interface that does not exist, which keeps "no carrier" and "no interface"
|
||||
/// 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)
|
||||
}
|
||||
|
||||
/// Whether the named interface exists and has every flag in `wanted` set.
|
||||
#[cfg(unix)]
|
||||
fn interface_has_flags(interface: &str, wanted: u32) -> bool {
|
||||
let Ok(c_name) = std::ffi::CString::new(interface) else {
|
||||
return false;
|
||||
};
|
||||
|
||||
let mut addrs: *mut libc::ifaddrs = std::ptr::null_mut();
|
||||
if unsafe { libc::getifaddrs(&mut addrs) } != 0 {
|
||||
return false;
|
||||
}
|
||||
|
||||
let mut matched = false;
|
||||
let mut cur = addrs;
|
||||
while !cur.is_null() {
|
||||
let entry = unsafe { &*cur };
|
||||
if !entry.ifa_name.is_null()
|
||||
&& unsafe { libc::strcmp(entry.ifa_name, c_name.as_ptr()) } == 0
|
||||
&& entry.ifa_flags & wanted == wanted
|
||||
{
|
||||
matched = true;
|
||||
break;
|
||||
}
|
||||
cur = entry.ifa_next;
|
||||
}
|
||||
|
||||
unsafe { libc::freeifaddrs(addrs) };
|
||||
matched
|
||||
}
|
||||
|
||||
// Platform-specific PacketSocket implementation.
|
||||
#[cfg(target_os = "linux")]
|
||||
#[path = "io_linux.rs"]
|
||||
|
||||
@@ -59,6 +59,14 @@ impl PacketSocket {
|
||||
if ret < 0 {
|
||||
let err = std::io::Error::last_os_error();
|
||||
unsafe { libc::close(fd) };
|
||||
// The interface can disappear between the index lookup and the
|
||||
// bind. That is absence arriving a few microseconds late, not a
|
||||
// configuration fault, so it reports as absence.
|
||||
if matches!(err.raw_os_error(), Some(libc::ENODEV) | Some(libc::ENXIO)) {
|
||||
return Err(TransportError::InterfaceUnavailable {
|
||||
interface: interface.to_string(),
|
||||
});
|
||||
}
|
||||
return Err(TransportError::StartFailed(format!(
|
||||
"bind(AF_PACKET, {}) failed: {}",
|
||||
interface, err
|
||||
@@ -231,11 +239,11 @@ fn get_if_index(_fd: RawFd, interface: &str) -> Result<i32, TransportError> {
|
||||
|
||||
let idx = unsafe { libc::if_nametoindex(c_name.as_ptr()) };
|
||||
if idx == 0 {
|
||||
return Err(TransportError::StartFailed(format!(
|
||||
"interface not found: {} ({})",
|
||||
interface,
|
||||
std::io::Error::last_os_error()
|
||||
)));
|
||||
// Absence, not a fault: the caller's presence watcher rebinds when the
|
||||
// interface shows up. See `TransportError::InterfaceUnavailable`.
|
||||
return Err(TransportError::InterfaceUnavailable {
|
||||
interface: interface.to_string(),
|
||||
});
|
||||
}
|
||||
Ok(idx as i32)
|
||||
}
|
||||
|
||||
@@ -393,10 +393,18 @@ fn bind_to_interface(fd: RawFd, interface: &str) -> Result<(), TransportError> {
|
||||
|
||||
let ret = unsafe { libc::ioctl(fd, BIOCSETIF, ifreq.as_ptr()) };
|
||||
if ret < 0 {
|
||||
let err = std::io::Error::last_os_error();
|
||||
// BIOCSETIF answers ENXIO for an interface that is not there. The
|
||||
// interface can also vanish between the index lookup and this ioctl,
|
||||
// so absence is reported as absence rather than as a bind fault.
|
||||
if matches!(err.raw_os_error(), Some(libc::ENXIO) | Some(libc::ENODEV)) {
|
||||
return Err(TransportError::InterfaceUnavailable {
|
||||
interface: interface.to_string(),
|
||||
});
|
||||
}
|
||||
return Err(TransportError::StartFailed(format!(
|
||||
"BIOCSETIF({}) failed: {}",
|
||||
interface,
|
||||
std::io::Error::last_os_error()
|
||||
interface, err
|
||||
)));
|
||||
}
|
||||
Ok(())
|
||||
@@ -498,11 +506,11 @@ fn get_if_index(interface: &str) -> Result<i32, TransportError> {
|
||||
|
||||
let idx = unsafe { libc::if_nametoindex(c_name.as_ptr()) };
|
||||
if idx == 0 {
|
||||
return Err(TransportError::StartFailed(format!(
|
||||
"interface not found: {} ({})",
|
||||
interface,
|
||||
std::io::Error::last_os_error()
|
||||
)));
|
||||
// Absence, not a fault — see the Linux twin and
|
||||
// `TransportError::InterfaceUnavailable`.
|
||||
return Err(TransportError::InterfaceUnavailable {
|
||||
interface: interface.to_string(),
|
||||
});
|
||||
}
|
||||
Ok(idx as i32)
|
||||
}
|
||||
|
||||
+1357
-197
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,718 @@
|
||||
//! Interface presence state for interface-bound transports.
|
||||
//!
|
||||
//! A transport bound to a network interface is a long-lived object that is
|
||||
//! *sometimes bound*. The interface it names may not exist when the daemon
|
||||
//! starts, may appear minutes later, may vanish and return mid-operation, and
|
||||
//! may never appear at all. This module holds the state that makes that
|
||||
//! observable and drives the rebind loop in [`super`].
|
||||
//!
|
||||
//! ```text
|
||||
//! Absent ──attach──> Binding ──ok──> Present
|
||||
//! ^ │ │
|
||||
//! └──── fail/backoff ─┘ │
|
||||
//! └──────────── detach ──────────────┘
|
||||
//! ```
|
||||
//!
|
||||
//! Two invariants do the work:
|
||||
//!
|
||||
//! - **The transport object survives detach.** Config, `TransportId`,
|
||||
//! statistics, and the neighbor buffer persist; only the file descriptor and
|
||||
//! its loops go. A transport is never destroyed because its interface went
|
||||
//! away.
|
||||
//! - **Start-time absence and runtime detach are the same transition.** A node
|
||||
//! that boots before wifi and a node whose wifi reloads at 03:00 take one
|
||||
//! code path.
|
||||
//!
|
||||
//! Presence tracks `IFF_UP` — the interface exists and the operator has
|
||||
//! 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.
|
||||
|
||||
use std::sync::atomic::{AtomicU8, AtomicU32, AtomicU64, Ordering};
|
||||
use std::sync::{PoisonError, RwLock, RwLockReadGuard, RwLockWriteGuard};
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
/// Read a lock, ignoring poisoning.
|
||||
///
|
||||
/// Every value guarded in this module is plain data — an `Instant`, an
|
||||
/// `Option<[u8; 6]>` — that a panic mid-write cannot leave logically
|
||||
/// inconsistent, so poisoning carries no information worth propagating.
|
||||
/// Treating it as a failure is what would hurt: the callers here are on the
|
||||
/// presence path, and "assume the worst" there means a transport that reports
|
||||
/// itself bound while every send fails, or a binder that tears down and
|
||||
/// rebinds every second forever. A stuck state is a worse outcome than
|
||||
/// reading a byte written by a thread that later panicked.
|
||||
fn read<T>(lock: &RwLock<T>) -> RwLockReadGuard<'_, T> {
|
||||
lock.read().unwrap_or_else(PoisonError::into_inner)
|
||||
}
|
||||
|
||||
/// Write a lock, ignoring poisoning. See [`read`].
|
||||
fn write<T>(lock: &RwLock<T>) -> RwLockWriteGuard<'_, T> {
|
||||
lock.write().unwrap_or_else(PoisonError::into_inner)
|
||||
}
|
||||
|
||||
/// How absence of the configured interface is reported.
|
||||
///
|
||||
/// Describes *the interface's presence*, not the transport's importance: an
|
||||
/// optional interface that is present is used exactly as hard as any other.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum AbsencePolicy {
|
||||
/// Naming an interface in configuration is a statement that you expect it,
|
||||
/// so the default is to complain: absence degrades node health from the
|
||||
/// first edge, and if it outlasts [`ABSENCE_ERROR_AFTER`] — the window in
|
||||
/// which it could still have been an ordinary bring-up race — it is
|
||||
/// reported once at `error`.
|
||||
Required,
|
||||
/// Absence is normal for this interface (a dock adapter, a radio that only
|
||||
/// exists on some hardware): no health impact, `info` on the edge.
|
||||
Optional,
|
||||
}
|
||||
|
||||
impl AbsencePolicy {
|
||||
/// `optional: true` in configuration selects [`AbsencePolicy::Optional`].
|
||||
pub fn from_optional(optional: bool) -> Self {
|
||||
if optional {
|
||||
Self::Optional
|
||||
} else {
|
||||
Self::Required
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether absence should be hidden from node health.
|
||||
pub fn is_optional(self) -> bool {
|
||||
matches!(self, Self::Optional)
|
||||
}
|
||||
|
||||
/// Operator-facing label, used by `show_transports`.
|
||||
pub fn as_str(self) -> &'static str {
|
||||
match self {
|
||||
Self::Required => "required",
|
||||
Self::Optional => "optional",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Where a transport sits in the presence cycle.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum Presence {
|
||||
/// The interface is not there (or is there without carrier). No socket, no
|
||||
/// loops; the watcher is waiting.
|
||||
Absent,
|
||||
/// The interface appeared and a bind is in flight, or a non-absence bind
|
||||
/// failure is backing off.
|
||||
Binding,
|
||||
/// Bound, with a live socket and running loops.
|
||||
Present,
|
||||
}
|
||||
|
||||
impl Presence {
|
||||
fn from_u8(v: u8) -> Self {
|
||||
match v {
|
||||
1 => Self::Binding,
|
||||
2 => Self::Present,
|
||||
_ => Self::Absent,
|
||||
}
|
||||
}
|
||||
|
||||
fn as_u8(self) -> u8 {
|
||||
match self {
|
||||
Self::Absent => 0,
|
||||
Self::Binding => 1,
|
||||
Self::Present => 2,
|
||||
}
|
||||
}
|
||||
|
||||
/// Operator-facing label, used by `show_transports`.
|
||||
pub fn as_str(self) -> &'static str {
|
||||
match self {
|
||||
Self::Absent => "absent",
|
||||
Self::Binding => "binding",
|
||||
Self::Present => "present",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl std::fmt::Display for Presence {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
f.write_str(self.as_str())
|
||||
}
|
||||
}
|
||||
|
||||
/// Shared, lock-light presence state.
|
||||
///
|
||||
/// Written by the transport's binder task, read by `send`, by the control
|
||||
/// plane, and by tests. Held behind an `Arc` so the binder task can outlive
|
||||
/// any particular borrow of the transport.
|
||||
#[derive(Debug)]
|
||||
pub struct PresenceState {
|
||||
phase: AtomicU8,
|
||||
/// Successful binds since the transport was created. `1` after a clean
|
||||
/// start; every increment past that is a rebind.
|
||||
binds: AtomicU64,
|
||||
/// Failed bind attempts since the last successful bind. Reset on bind.
|
||||
attempts: AtomicU32,
|
||||
/// When the current phase was entered.
|
||||
since: RwLock<Instant>,
|
||||
/// MAC observed at the last successful bind. A name reappearing with a
|
||||
/// different MAC is different hardware, not the same device returning.
|
||||
last_mac: RwLock<Option<[u8; 6]>>,
|
||||
}
|
||||
|
||||
impl Default for PresenceState {
|
||||
fn default() -> Self {
|
||||
Self::new()
|
||||
}
|
||||
}
|
||||
|
||||
impl PresenceState {
|
||||
/// A fresh tracker in [`Presence::Absent`].
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
phase: AtomicU8::new(Presence::Absent.as_u8()),
|
||||
binds: AtomicU64::new(0),
|
||||
attempts: AtomicU32::new(0),
|
||||
since: RwLock::new(Instant::now()),
|
||||
last_mac: RwLock::new(None),
|
||||
}
|
||||
}
|
||||
|
||||
/// Current phase.
|
||||
pub fn presence(&self) -> Presence {
|
||||
Presence::from_u8(self.phase.load(Ordering::Acquire))
|
||||
}
|
||||
|
||||
/// Whether the transport currently holds a bound socket.
|
||||
pub fn is_present(&self) -> bool {
|
||||
self.presence() == Presence::Present
|
||||
}
|
||||
|
||||
/// How long the current presence *episode* has been held.
|
||||
///
|
||||
/// Episode, not phase: `Binding` is part of the absence episode until it
|
||||
/// succeeds. An interface that has been gone for a week while a bind is
|
||||
/// retried and refused every second must report a week, not one second —
|
||||
/// otherwise [`ABSENCE_ERROR_AFTER`] is never reached and the
|
||||
/// operator-facing `since_secs` reads as a healthy young absence forever.
|
||||
pub fn since(&self) -> Duration {
|
||||
read(&self.since).elapsed()
|
||||
}
|
||||
|
||||
/// Successful binds since creation (`1` after a clean start).
|
||||
pub fn binds(&self) -> u64 {
|
||||
self.binds.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Failed bind attempts since the last successful bind.
|
||||
pub fn attempts(&self) -> u32 {
|
||||
self.attempts.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// MAC observed at the last successful bind, if any.
|
||||
pub fn last_mac(&self) -> Option<[u8; 6]> {
|
||||
*read(&self.last_mac)
|
||||
}
|
||||
|
||||
/// Move to `phase`, returning `true` if this was an actual edge.
|
||||
///
|
||||
/// Edge-vs-level is what the logging policy keys on: logged once on
|
||||
/// entering absence and once on recovery, never per retry attempt.
|
||||
///
|
||||
/// The episode clock ([`Self::since`]) restarts only when *boundness*
|
||||
/// changes — Present↔not-Present. A failed bind walks
|
||||
/// `Absent → Binding → Absent`, and resetting the clock on those would
|
||||
/// hide a permanent absence behind a timer that never gets past one
|
||||
/// second.
|
||||
pub fn transition(&self, phase: Presence) -> bool {
|
||||
let prev = Presence::from_u8(self.phase.swap(phase.as_u8(), Ordering::AcqRel));
|
||||
if prev == phase {
|
||||
return false;
|
||||
}
|
||||
if (prev == Presence::Present) != (phase == Presence::Present) {
|
||||
*write(&self.since) = Instant::now();
|
||||
}
|
||||
true
|
||||
}
|
||||
|
||||
/// Record a successful bind at `mac`.
|
||||
///
|
||||
/// Returns `true` when the interface came back as *different hardware* —
|
||||
/// the name reappeared with a MAC other than the one last bound. The
|
||||
/// caller drops cached neighbor state rather than silently resuming onto
|
||||
/// a different adapter.
|
||||
pub fn record_bind(&self, mac: [u8; 6]) -> bool {
|
||||
let changed = match *read(&self.last_mac) {
|
||||
Some(prev) => prev != mac,
|
||||
None => false,
|
||||
};
|
||||
*write(&self.last_mac) = Some(mac);
|
||||
self.binds.fetch_add(1, Ordering::Relaxed);
|
||||
self.attempts.store(0, Ordering::Relaxed);
|
||||
self.transition(Presence::Present);
|
||||
changed
|
||||
}
|
||||
|
||||
/// Record a failed bind attempt, returning the new attempt count.
|
||||
pub fn record_attempt(&self) -> u32 {
|
||||
self.attempts.fetch_add(1, Ordering::Relaxed) + 1
|
||||
}
|
||||
}
|
||||
|
||||
/// Backoff for bind failures that are *not* absence — permission denied,
|
||||
/// buffer sizing, a BPF device shortage. Absence itself does not back off
|
||||
/// where an event source is available: there is nothing to poll.
|
||||
///
|
||||
/// 1 s doubling to a 30 s ceiling.
|
||||
pub fn bind_backoff(attempts: u32) -> Duration {
|
||||
const BASE_SECS: u64 = 1;
|
||||
const CEILING_SECS: u64 = 30;
|
||||
let shift = attempts.saturating_sub(1).min(5);
|
||||
Duration::from_secs((BASE_SECS << shift).min(CEILING_SECS))
|
||||
}
|
||||
|
||||
/// How long a *required* interface may be absent before it is an error.
|
||||
///
|
||||
/// Absence is a state the presence machine handles, so it is not an error for
|
||||
/// happening — a daemon that wins the race against its own radio, or a cable
|
||||
/// out for two seconds, is the ordinary case this mechanism exists to absorb,
|
||||
/// and calling that an error at t=0 and "recovered" at t=0.2 s is cry-wolf.
|
||||
/// Past this window it is no longer a race: something an operator has to fix
|
||||
/// is wrong, and the log should say so once.
|
||||
///
|
||||
/// One window for both shapes of absence. A node that boots before its wifi
|
||||
/// and a node whose wifi reloads at 03:00 take one code path everywhere else
|
||||
/// in this module; giving them different deadlines would reintroduce exactly
|
||||
/// the start-versus-runtime asymmetry the presence machine removed.
|
||||
///
|
||||
/// Tuned against the platforms this exists for: comfortably past a veth or a
|
||||
/// container coming up, short enough that a mesh radio which never appears is
|
||||
/// named while somebody is still watching the boot. Raising it hides a real
|
||||
/// fault for longer; lowering it starts reporting ordinary bring-up races.
|
||||
pub const ABSENCE_ERROR_AFTER: Duration = Duration::from_secs(10);
|
||||
|
||||
/// Minimum lifetime for a binding to count as a real recovery.
|
||||
///
|
||||
/// A socket that dies sooner than this never really came back.
|
||||
pub const MIN_STABLE_BINDING: Duration = Duration::from_secs(10);
|
||||
|
||||
/// Consecutive short-lived bindings before the binder stops treating a
|
||||
/// successful bind as a recovery.
|
||||
pub const CHURN_THRESHOLD: u32 = 3;
|
||||
|
||||
/// Damping for the rebind loop.
|
||||
///
|
||||
/// Backoff covers *failed* binds; this covers the opposite and nastier case —
|
||||
/// binds that keep **succeeding** into a socket that dies moments later. A
|
||||
/// receive loop that gives up on a persistent error while the interface stays
|
||||
/// `UP` produces exactly that: tear down, rebind, succeed, fail again, once
|
||||
/// per second, forever. Undamped it is an `error!`/`info!` pair and a
|
||||
/// `Degraded`→`Running` health flap every cycle, which defeats both the
|
||||
/// "log edges, not attempts" rule and the meaning of `Degraded`.
|
||||
///
|
||||
/// So: count consecutive bindings that die young, back off between them on
|
||||
/// the same 1 s → 30 s curve, and once the streak reaches
|
||||
/// [`CHURN_THRESHOLD`] stop announcing each bind as a recovery — hold the
|
||||
/// node at its degraded reading until a binding actually survives
|
||||
/// [`MIN_STABLE_BINDING`]. A binding that holds ends the streak.
|
||||
///
|
||||
/// Pure state, driven by an injected clock, so the policy is testable without
|
||||
/// a network interface.
|
||||
#[derive(Debug, Default)]
|
||||
pub struct ChurnGuard {
|
||||
/// Consecutive bindings that died younger than [`MIN_STABLE_BINDING`].
|
||||
streak: u32,
|
||||
/// When the current binding was established.
|
||||
bound_at: Option<Instant>,
|
||||
/// Whether the current binding was announced as a recovery.
|
||||
announced: bool,
|
||||
}
|
||||
|
||||
/// What the caller should do about a successful bind.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub struct BindOutcome {
|
||||
/// Log the recovery and publish presence now. `false` while churning:
|
||||
/// the bind is held back until it proves it will last.
|
||||
pub announce: bool,
|
||||
}
|
||||
|
||||
/// What the caller should do about a detach.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub struct DetachOutcome {
|
||||
/// Publish absence. `false` when the binding that just died was never
|
||||
/// announced, so there is nothing to retract.
|
||||
pub retract: bool,
|
||||
/// Log the detach edge. `false` once churning — the fact has been said.
|
||||
pub log_edge: bool,
|
||||
/// This detach is the one that crossed [`CHURN_THRESHOLD`]; say so once.
|
||||
pub entered_churn: bool,
|
||||
/// Wait this long before trying to bind again.
|
||||
pub backoff: Option<Duration>,
|
||||
}
|
||||
|
||||
impl ChurnGuard {
|
||||
/// A guard with no history.
|
||||
pub fn new() -> Self {
|
||||
Self::default()
|
||||
}
|
||||
|
||||
/// Consecutive short-lived bindings, for logging and tests.
|
||||
pub fn streak(&self) -> u32 {
|
||||
self.streak
|
||||
}
|
||||
|
||||
/// Record a successful bind.
|
||||
pub fn bound(&mut self, now: Instant) -> BindOutcome {
|
||||
self.bound_at = Some(now);
|
||||
let announce = self.streak < CHURN_THRESHOLD;
|
||||
if announce {
|
||||
self.announced = true;
|
||||
}
|
||||
BindOutcome { announce }
|
||||
}
|
||||
|
||||
/// Called on every tick while bound. Returns `true` exactly once, at the
|
||||
/// moment a held-back binding has proved stable and should be announced.
|
||||
pub fn stabilized(&mut self, now: Instant) -> bool {
|
||||
let Some(bound_at) = self.bound_at else {
|
||||
return false;
|
||||
};
|
||||
if now.duration_since(bound_at) < MIN_STABLE_BINDING {
|
||||
return false;
|
||||
}
|
||||
let newly_announced = !self.announced;
|
||||
self.streak = 0;
|
||||
self.announced = true;
|
||||
newly_announced
|
||||
}
|
||||
|
||||
/// Record a detach.
|
||||
pub fn detached(&mut self, now: Instant) -> DetachOutcome {
|
||||
let young = self
|
||||
.bound_at
|
||||
.is_some_and(|t| now.duration_since(t) < MIN_STABLE_BINDING);
|
||||
self.bound_at = None;
|
||||
|
||||
if young {
|
||||
self.streak += 1;
|
||||
} else {
|
||||
self.streak = 0;
|
||||
}
|
||||
|
||||
DetachOutcome {
|
||||
retract: std::mem::take(&mut self.announced),
|
||||
log_edge: self.streak < CHURN_THRESHOLD,
|
||||
entered_churn: self.streak == CHURN_THRESHOLD,
|
||||
backoff: (self.streak > 0).then(|| bind_backoff(self.streak)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use std::sync::Arc;
|
||||
|
||||
#[test]
|
||||
fn presence_starts_absent() {
|
||||
let p = PresenceState::new();
|
||||
assert_eq!(p.presence(), Presence::Absent);
|
||||
assert_eq!(p.binds(), 0);
|
||||
assert_eq!(p.attempts(), 0);
|
||||
assert!(p.last_mac().is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transition_reports_edges_only() {
|
||||
let p = PresenceState::new();
|
||||
assert!(p.transition(Presence::Binding));
|
||||
assert!(!p.transition(Presence::Binding));
|
||||
assert!(p.transition(Presence::Present));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn bind_records_mac_and_clears_attempts() {
|
||||
let p = PresenceState::new();
|
||||
p.record_attempt();
|
||||
p.record_attempt();
|
||||
assert_eq!(p.attempts(), 2);
|
||||
|
||||
let mac = [0x02, 0, 0, 0, 0, 1];
|
||||
assert!(!p.record_bind(mac), "first bind is not a hardware change");
|
||||
assert_eq!(p.attempts(), 0);
|
||||
assert_eq!(p.binds(), 1);
|
||||
assert_eq!(p.presence(), Presence::Present);
|
||||
assert_eq!(p.last_mac(), Some(mac));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rebind_on_same_mac_is_not_a_hardware_change() {
|
||||
let p = PresenceState::new();
|
||||
let mac = [0x02, 0, 0, 0, 0, 1];
|
||||
p.record_bind(mac);
|
||||
p.transition(Presence::Absent);
|
||||
assert!(!p.record_bind(mac));
|
||||
assert_eq!(p.binds(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rebind_on_different_mac_is_a_hardware_change() {
|
||||
let p = PresenceState::new();
|
||||
p.record_bind([0x02, 0, 0, 0, 0, 1]);
|
||||
p.transition(Presence::Absent);
|
||||
assert!(p.record_bind([0x02, 0, 0, 0, 0, 2]));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn backoff_climbs_to_a_ceiling() {
|
||||
assert_eq!(bind_backoff(0), Duration::from_secs(1));
|
||||
assert_eq!(bind_backoff(1), Duration::from_secs(1));
|
||||
assert_eq!(bind_backoff(2), Duration::from_secs(2));
|
||||
assert_eq!(bind_backoff(3), Duration::from_secs(4));
|
||||
assert_eq!(bind_backoff(6), Duration::from_secs(30));
|
||||
assert_eq!(bind_backoff(u32::MAX), Duration::from_secs(30));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_error_deadline_outlasts_an_ordinary_bring_up_race() {
|
||||
// The window has to clear the races it exists to absorb — a veth
|
||||
// arriving a fraction of a second late, a container starting — while
|
||||
// staying short enough that a radio which never appears is named
|
||||
// during the boot somebody is watching. It is also the one deadline:
|
||||
// start-time absence and a runtime detach share it.
|
||||
assert!(
|
||||
ABSENCE_ERROR_AFTER >= Duration::from_secs(5),
|
||||
"shorter than a bring-up race would report the ordinary case"
|
||||
);
|
||||
assert!(
|
||||
ABSENCE_ERROR_AFTER <= Duration::from_secs(60),
|
||||
"longer and a required interface that never appears goes unsaid \
|
||||
for the whole boot"
|
||||
);
|
||||
}
|
||||
|
||||
// ── The absence clock measures an episode, not a phase ────────────────
|
||||
|
||||
#[test]
|
||||
fn a_failed_bind_does_not_restart_the_absence_clock() {
|
||||
// An interface that is present but refuses to bind walks
|
||||
// Absent → Binding → Absent on every retry. If those edges reset the
|
||||
// clock, `since_secs` reads as a one-second-old absence forever and
|
||||
// the error deadline is never reached — so a permission error would
|
||||
// sit silently behind a healthy-looking counter.
|
||||
let p = PresenceState::new();
|
||||
std::thread::sleep(Duration::from_millis(30));
|
||||
let before = p.since();
|
||||
|
||||
assert!(p.transition(Presence::Binding));
|
||||
assert!(p.transition(Presence::Absent));
|
||||
assert!(p.transition(Presence::Binding));
|
||||
assert!(p.transition(Presence::Absent));
|
||||
|
||||
assert!(
|
||||
p.since() >= before,
|
||||
"the absence clock ran backwards across failed binds"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_clock_restarts_only_when_boundness_changes() {
|
||||
let p = PresenceState::new();
|
||||
std::thread::sleep(Duration::from_millis(30));
|
||||
|
||||
// Absent → Present restarts it: a new episode began.
|
||||
p.record_bind([0x02, 0, 0, 0, 0, 1]);
|
||||
assert!(p.since() < Duration::from_millis(30));
|
||||
|
||||
std::thread::sleep(Duration::from_millis(30));
|
||||
let bound_for = p.since();
|
||||
|
||||
// Present → Present is not an edge at all.
|
||||
assert!(!p.transition(Presence::Present));
|
||||
assert!(p.since() >= bound_for);
|
||||
|
||||
// Present → Absent restarts it: the episode ended.
|
||||
assert!(p.transition(Presence::Absent));
|
||||
assert!(p.since() < Duration::from_millis(30));
|
||||
}
|
||||
|
||||
// ── Poisoning must not be a stuck state ───────────────────────────────
|
||||
|
||||
#[test]
|
||||
fn a_poisoned_lock_still_reports_presence() {
|
||||
// `.ok()`-style handling would make a poisoned lock read as "no MAC,
|
||||
// no socket, tasks dead" — a transport reporting itself present while
|
||||
// every send fails, and a binder rebinding once a second forever.
|
||||
// Poisoning carries no information about plain data, so it is ignored.
|
||||
let p = Arc::new(PresenceState::new());
|
||||
p.record_bind([0x02, 0, 0, 0, 0, 7]);
|
||||
|
||||
let poisoner = Arc::clone(&p);
|
||||
let panicked = std::thread::spawn(move || {
|
||||
let _guard = poisoner.last_mac.write().unwrap();
|
||||
panic!("poison the lock while holding it");
|
||||
})
|
||||
.join();
|
||||
assert!(panicked.is_err(), "the helper thread was supposed to panic");
|
||||
assert!(
|
||||
p.last_mac.is_poisoned(),
|
||||
"the lock was supposed to be poisoned"
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
p.last_mac(),
|
||||
Some([0x02, 0, 0, 0, 0, 7]),
|
||||
"a poisoned lock must not erase the binding"
|
||||
);
|
||||
// And the clock still answers rather than collapsing to zero.
|
||||
let _ = p.since();
|
||||
}
|
||||
|
||||
// ── Rebind churn ──────────────────────────────────────────────────────
|
||||
|
||||
#[test]
|
||||
fn a_healthy_bind_and_detach_is_not_churn() {
|
||||
let mut g = ChurnGuard::new();
|
||||
let t0 = Instant::now();
|
||||
assert!(g.bound(t0).announce, "a first bind is a recovery");
|
||||
|
||||
// Held well past the stability floor, then lost.
|
||||
let out = g.detached(t0 + MIN_STABLE_BINDING + Duration::from_secs(60));
|
||||
assert!(out.retract, "an announced binding must be retracted");
|
||||
assert!(out.log_edge, "an isolated detach is worth a line");
|
||||
assert!(!out.entered_churn);
|
||||
assert_eq!(out.backoff, None, "one clean outage must not back off");
|
||||
assert_eq!(g.streak(), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn short_lived_bindings_back_off() {
|
||||
// The failure this guards: a receive loop that gives up on a
|
||||
// persistent error while the interface stays UP. Bind succeeds, dies,
|
||||
// rebinds, dies — once per second, forever, undamped.
|
||||
let mut g = ChurnGuard::new();
|
||||
let mut t = Instant::now();
|
||||
|
||||
for expected in [1u64, 2, 4] {
|
||||
g.bound(t);
|
||||
t += Duration::from_secs(1);
|
||||
let out = g.detached(t);
|
||||
assert_eq!(
|
||||
out.backoff,
|
||||
Some(Duration::from_secs(expected)),
|
||||
"streak {} should back off {expected}s",
|
||||
g.streak()
|
||||
);
|
||||
}
|
||||
assert_eq!(g.streak(), 3);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn churn_stops_announcing_and_stops_logging() {
|
||||
let mut g = ChurnGuard::new();
|
||||
let mut t = Instant::now();
|
||||
|
||||
// The detaches below the threshold are still news and still logged.
|
||||
for _ in 0..CHURN_THRESHOLD - 1 {
|
||||
assert!(g.bound(t).announce);
|
||||
t += Duration::from_secs(1);
|
||||
let out = g.detached(t);
|
||||
assert!(out.log_edge, "the first few detaches are still news");
|
||||
assert!(!out.entered_churn);
|
||||
}
|
||||
|
||||
// The detach that crosses the threshold reports the churn instead of
|
||||
// the edge: one line saying "this keeps happening", not two saying
|
||||
// "it happened" and "it keeps happening".
|
||||
assert!(g.bound(t).announce);
|
||||
t += Duration::from_secs(1);
|
||||
let crossing = g.detached(t);
|
||||
assert!(crossing.entered_churn, "crossing must be announced once");
|
||||
assert!(!crossing.log_edge, "the churn line replaces the edge line");
|
||||
|
||||
// Past the threshold: bindings are no longer announced as recoveries,
|
||||
// so node health stays put instead of flapping every second, and the
|
||||
// edges stop being logged.
|
||||
for _ in 0..5 {
|
||||
assert!(!g.bound(t).announce, "a churning bind is not a recovery");
|
||||
t += Duration::from_secs(1);
|
||||
let out = g.detached(t);
|
||||
assert!(!out.log_edge, "churn must not log per cycle");
|
||||
assert!(!out.retract, "nothing was announced, so nothing to retract");
|
||||
assert!(
|
||||
!out.entered_churn,
|
||||
"the threshold is crossed once, not repeatedly"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn backoff_during_churn_is_capped() {
|
||||
let mut g = ChurnGuard::new();
|
||||
let mut t = Instant::now();
|
||||
let mut last = None;
|
||||
for _ in 0..12 {
|
||||
g.bound(t);
|
||||
t += Duration::from_secs(1);
|
||||
last = g.detached(t).backoff;
|
||||
}
|
||||
assert_eq!(
|
||||
last,
|
||||
Some(Duration::from_secs(30)),
|
||||
"churn backoff must climb to the ceiling and stop"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_binding_that_lasts_ends_the_streak_and_announces_once() {
|
||||
let mut g = ChurnGuard::new();
|
||||
let mut t = Instant::now();
|
||||
|
||||
// Churn into the held-back state.
|
||||
for _ in 0..CHURN_THRESHOLD + 1 {
|
||||
g.bound(t);
|
||||
t += Duration::from_secs(1);
|
||||
g.detached(t);
|
||||
}
|
||||
assert!(!g.bound(t).announce);
|
||||
|
||||
// Not yet stable: still nothing to say.
|
||||
assert!(!g.stabilized(t + Duration::from_secs(1)));
|
||||
|
||||
// Survived the floor: announce exactly once, and the streak is over.
|
||||
let stable_at = t + MIN_STABLE_BINDING;
|
||||
assert!(
|
||||
g.stabilized(stable_at),
|
||||
"a binding that lasts is a recovery"
|
||||
);
|
||||
assert!(
|
||||
!g.stabilized(stable_at + Duration::from_secs(60)),
|
||||
"recovery is announced once, not on every tick"
|
||||
);
|
||||
assert_eq!(g.streak(), 0);
|
||||
|
||||
// And the next detach behaves like an ordinary one again.
|
||||
let out = g.detached(stable_at + Duration::from_secs(60));
|
||||
assert!(out.retract);
|
||||
assert!(out.log_edge);
|
||||
assert_eq!(out.backoff, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stabilized_is_silent_for_an_ordinary_binding() {
|
||||
// A bind that was announced immediately must not be announced again
|
||||
// when it passes the stability floor.
|
||||
let mut g = ChurnGuard::new();
|
||||
let t = Instant::now();
|
||||
assert!(g.bound(t).announce);
|
||||
assert!(!g.stabilized(t + MIN_STABLE_BINDING + Duration::from_secs(1)));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn policy_labels() {
|
||||
assert_eq!(AbsencePolicy::from_optional(true), AbsencePolicy::Optional);
|
||||
assert_eq!(AbsencePolicy::from_optional(false), AbsencePolicy::Required);
|
||||
assert!(AbsencePolicy::Optional.is_optional());
|
||||
assert!(!AbsencePolicy::Required.is_optional());
|
||||
assert_eq!(AbsencePolicy::Required.as_str(), "required");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,305 @@
|
||||
//! 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::{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,
|
||||
}
|
||||
|
||||
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),
|
||||
}
|
||||
}
|
||||
|
||||
/// 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")
|
||||
};
|
||||
|
||||
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: nothing more to take this round.
|
||||
Ok(Ok(_)) => 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 {
|
||||
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();
|
||||
// Both answers are legitimate — a sandbox may refuse the socket — so
|
||||
// this pins that asking is safe, not which answer comes back.
|
||||
let _ = w.is_event_driven();
|
||||
}
|
||||
|
||||
/// 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),
|
||||
};
|
||||
|
||||
// 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),
|
||||
};
|
||||
let fired = tokio::time::timeout(std::time::Duration::from_millis(50), w.changed()).await;
|
||||
assert!(fired.is_err(), "sourceless watcher resolved");
|
||||
}
|
||||
}
|
||||
@@ -102,6 +102,54 @@ pub fn packet_channel(buffer: usize) -> (PacketTx, PacketRx) {
|
||||
tokio::sync::mpsc::channel(buffer)
|
||||
}
|
||||
|
||||
/// Operator-visible interface presence, rendered by `show_transports`.
|
||||
///
|
||||
/// Worth as much as the retry itself. The original boot-race bug was expensive
|
||||
/// precisely because the 802.11s peer link formed regardless of the daemon, so
|
||||
/// nothing an operator could see said the node was deaf.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub struct InterfacePresence {
|
||||
/// `absent`, `binding`, or `present`.
|
||||
pub presence: &'static str,
|
||||
/// Whether the interface currently has carrier (`IFF_RUNNING`).
|
||||
///
|
||||
/// Reported, never acted on. Presence is `IFF_UP`, because binding does
|
||||
/// not need carrier and a socket outlives a carrier flap — but "is
|
||||
/// anything plugged in" is still what an operator wants to know when a
|
||||
/// bound transport is carrying nothing, so it is reported here instead of
|
||||
/// steering the daemon.
|
||||
pub carrier: bool,
|
||||
/// `required` or `optional`.
|
||||
pub policy: &'static str,
|
||||
/// How long the current phase has been held.
|
||||
pub since_secs: u64,
|
||||
/// Successful binds since the transport was created (`1` after a clean
|
||||
/// start; more means it has rebound).
|
||||
pub binds: u64,
|
||||
/// Failed bind attempts since the last successful bind.
|
||||
pub failed_attempts: u32,
|
||||
}
|
||||
|
||||
/// A presence edge published by an interface-bound transport.
|
||||
///
|
||||
/// Absence and return are the same transition seen from two sides, so one
|
||||
/// event type carries both: `present: false` on detach (including a start
|
||||
/// where the interface was never there), `present: true` on every successful
|
||||
/// bind after the first observation.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub struct TransportPresence {
|
||||
/// The transport whose interface changed presence.
|
||||
pub transport_id: TransportId,
|
||||
/// Whether the interface is now bound.
|
||||
pub present: bool,
|
||||
}
|
||||
|
||||
/// Channel sender for transport presence edges.
|
||||
pub type PresenceTx = tokio::sync::mpsc::Sender<TransportPresence>;
|
||||
|
||||
/// Channel receiver for transport presence edges.
|
||||
pub type PresenceRx = tokio::sync::mpsc::Receiver<TransportPresence>;
|
||||
|
||||
// ============================================================================
|
||||
// Errors
|
||||
// ============================================================================
|
||||
@@ -118,6 +166,20 @@ pub enum TransportError {
|
||||
#[error("transport failed to start: {0}")]
|
||||
StartFailed(String),
|
||||
|
||||
/// The named network interface is not usable right now: it does not exist,
|
||||
/// or it exists but is administratively down (no `IFF_UP`).
|
||||
///
|
||||
/// Distinct from [`TransportError::StartFailed`] because absence is a
|
||||
/// *state*, not a fault. Interface-bound transports treat it as "not bound
|
||||
/// yet" and keep a presence watcher running; a `StartFailed` carrying the
|
||||
/// same text could not be told apart from a typo'd interface name or a
|
||||
/// missing capability.
|
||||
#[error("interface unavailable: {interface}")]
|
||||
InterfaceUnavailable {
|
||||
/// The configured interface name.
|
||||
interface: String,
|
||||
},
|
||||
|
||||
#[error("transport shutdown failed: {0}")]
|
||||
ShutdownFailed(String),
|
||||
|
||||
@@ -864,6 +926,52 @@ impl TransportHandle {
|
||||
}
|
||||
}
|
||||
|
||||
/// Interface presence for interface-bound transports: the phase label, the
|
||||
/// absence policy, how long the phase has been held, and the failed-bind
|
||||
/// count since the last successful bind.
|
||||
///
|
||||
/// `None` for transports that are not bound to a named interface — for
|
||||
/// those, presence is not a concept and an operator should not be shown an
|
||||
/// always-`present` column.
|
||||
pub fn interface_presence(&self) -> Option<InterfacePresence> {
|
||||
match self {
|
||||
#[cfg(any(target_os = "linux", target_os = "macos"))]
|
||||
TransportHandle::Ethernet(t) => {
|
||||
let state = t.presence_state();
|
||||
Some(InterfacePresence {
|
||||
presence: t.presence().as_str(),
|
||||
carrier: t.has_carrier(),
|
||||
policy: t.absence_policy().as_str(),
|
||||
since_secs: state.since().as_secs(),
|
||||
binds: state.binds(),
|
||||
failed_attempts: state.attempts(),
|
||||
})
|
||||
}
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether this transport can actually put a frame on the wire *now*.
|
||||
///
|
||||
/// [`Self::is_operational`] answers a different question: it means the
|
||||
/// transport was started, which for an interface-bound transport no longer
|
||||
/// implies a live socket — that is the whole point of presence. Callers
|
||||
/// that are choosing a transport to use, or deriving a value from one,
|
||||
/// want this; callers reasoning about lifecycle want `is_operational`.
|
||||
///
|
||||
/// `true` for every transport that is not interface-bound, so this is
|
||||
/// `is_operational` with the presence refinement applied where it exists.
|
||||
pub fn is_bound(&self) -> bool {
|
||||
if !self.is_operational() {
|
||||
return false;
|
||||
}
|
||||
match self {
|
||||
#[cfg(any(target_os = "linux", target_os = "macos"))]
|
||||
TransportHandle::Ethernet(t) => t.presence() == ethernet::Presence::Present,
|
||||
_ => true,
|
||||
}
|
||||
}
|
||||
|
||||
/// Get the interface name (Ethernet only, returns None for other transports).
|
||||
pub fn interface_name(&self) -> Option<&str> {
|
||||
match self {
|
||||
|
||||
+81
-42
@@ -80,6 +80,35 @@ impl PathMtuEntry {
|
||||
/// address).
|
||||
pub type PathMtuLookup = Arc<RwLock<HashMap<FipsAddress, PathMtuEntry>>>;
|
||||
|
||||
/// The node-global TCP MSS ceiling, shared live with the TUN reader and
|
||||
/// writer threads.
|
||||
///
|
||||
/// Shared rather than copied because the value it is derived from moves at
|
||||
/// runtime. `Node::transport_mtu()` is the minimum across *bound* transports,
|
||||
/// and since dynamic interface binding a transport can bind minutes after
|
||||
/// start or unbind mid-operation — so a narrow interface appearing must
|
||||
/// tighten the clamp, and its departure must release it. Every other consumer
|
||||
/// of `transport_mtu()` already reads it live (`show_status`, the snapshot,
|
||||
/// the session-layer fragmentation check); these two threads captured a `u16`
|
||||
/// at spawn and were the only place left where the daemon could report one
|
||||
/// effective MTU and clamp to another.
|
||||
///
|
||||
/// A relaxed load per packet, beside the `PathMtuLookup` read that already
|
||||
/// happens on the same packet — strictly the cheaper of the two. Ordering is
|
||||
/// irrelevant: this is a clamp, and a packet that reads the previous value
|
||||
/// during the store is clamped by the ceiling that was correct a microsecond
|
||||
/// earlier. The per-flow ceiling has always had that property.
|
||||
/// The IPv6 minimum link MTU (RFC 8200): every compliant path accepts a packet
|
||||
/// this large, so an MSS derived from it fits anywhere.
|
||||
///
|
||||
/// Two callers, and they must not disagree: the cold-flow fallback in
|
||||
/// [`per_flow_max_mss`], and the seed for the node's [`MssCeiling`] before any
|
||||
/// transport has bound — which is the same value `Node::transport_mtu()` falls
|
||||
/// back to when nothing is bound.
|
||||
pub const IPV6_MIN_MTU: u16 = 1280;
|
||||
|
||||
pub type MssCeiling = Arc<std::sync::atomic::AtomicU16>;
|
||||
|
||||
/// Compute the effective TCP MSS ceiling for a packet given its peer
|
||||
/// address bytes (a 16-byte IPv6 destination on outbound, source on
|
||||
/// inbound). Returns `min(global_max_mss, learned_path_max_mss)` when
|
||||
@@ -115,7 +144,6 @@ pub(crate) fn per_flow_max_mss(
|
||||
// RFC 8200 IPv6-minimum MTU (1280) → effective FIPS-encapsulated
|
||||
// payload (1203) → TCP segment after IPv6+TCP headers (1143).
|
||||
// Used as the conservative ceiling for empty-lookup destinations.
|
||||
const IPV6_MIN_MTU: u16 = 1280;
|
||||
let conservative_max_mss = mss_ceiling(IPV6_MIN_MTU);
|
||||
let empty_lookup_ceiling = std::cmp::min(global_max_mss, conservative_max_mss);
|
||||
|
||||
@@ -435,13 +463,14 @@ impl TunDevice {
|
||||
/// a channel sender for submitting packets to be written.
|
||||
///
|
||||
/// `max_mss` is the global TCP MSS ceiling derived from the local
|
||||
/// `transport_mtu()` floor. `path_mtu_lookup` is a read-only handle to
|
||||
/// the per-destination path MTU map populated by discovery; the writer
|
||||
/// reads it on each inbound SYN-ACK to compute a per-flow ceiling that
|
||||
/// honors learned narrow paths through the mesh.
|
||||
/// `transport_mtu()` floor, shared live so a transport that binds or
|
||||
/// unbinds after start moves it (see [`MssCeiling`]). `path_mtu_lookup`
|
||||
/// is a read-only handle to the per-destination path MTU map populated by
|
||||
/// discovery; the writer reads both on each inbound SYN-ACK to compute a
|
||||
/// per-flow ceiling that honors learned narrow paths through the mesh.
|
||||
pub fn create_writer(
|
||||
&self,
|
||||
max_mss: u16,
|
||||
max_mss: MssCeiling,
|
||||
path_mtu_lookup: PathMtuLookup,
|
||||
) -> Result<(TunWriter, TunTx), TunError> {
|
||||
let fd = self.device.as_raw_fd();
|
||||
@@ -517,7 +546,7 @@ pub struct TunWriter {
|
||||
file: File,
|
||||
rx: mpsc::Receiver<Vec<u8>>,
|
||||
name: String,
|
||||
max_mss: u16,
|
||||
max_mss: MssCeiling,
|
||||
path_mtu_lookup: PathMtuLookup,
|
||||
}
|
||||
|
||||
@@ -531,16 +560,23 @@ impl TunWriter {
|
||||
pub fn run(mut self) {
|
||||
use super::tcp_mss::clamp_tcp_mss;
|
||||
|
||||
debug!(name = %self.name, max_mss = self.max_mss, "TUN writer starting");
|
||||
debug!(
|
||||
name = %self.name,
|
||||
max_mss = self.max_mss.load(std::sync::atomic::Ordering::Relaxed),
|
||||
"TUN writer starting"
|
||||
);
|
||||
|
||||
for mut packet in self.rx {
|
||||
// Read per packet, not once: a transport binding or unbinding
|
||||
// moves the node's egress floor at runtime. See `MssCeiling`.
|
||||
let global_max_mss = self.max_mss.load(std::sync::atomic::Ordering::Relaxed);
|
||||
// Per-destination clamp: peer IPv6 source address (bytes 8..24)
|
||||
// identifies the flow's remote end. If discovery has learned a
|
||||
// smaller path MTU for that peer, tighten the ceiling.
|
||||
let effective_max_mss = if packet.len() >= 24 {
|
||||
per_flow_max_mss(&self.path_mtu_lookup, &packet[8..24], self.max_mss)
|
||||
per_flow_max_mss(&self.path_mtu_lookup, &packet[8..24], global_max_mss)
|
||||
} else {
|
||||
self.max_mss
|
||||
global_max_mss
|
||||
};
|
||||
// Clamp TCP MSS on inbound SYN-ACK packets
|
||||
if clamp_tcp_mss(&mut packet, effective_max_mss) {
|
||||
@@ -622,17 +658,17 @@ pub fn run_tun_reader(
|
||||
our_addr: FipsAddress,
|
||||
tun_tx: TunTx,
|
||||
outbound_tx: TunOutboundTx,
|
||||
transport_mtu: u16,
|
||||
max_mss: MssCeiling,
|
||||
path_mtu_lookup: PathMtuLookup,
|
||||
) {
|
||||
let (name, mut buf, max_mss) = tun_reader_setup(device.name(), mtu, transport_mtu);
|
||||
let (name, mut buf) = tun_reader_setup(device.name(), mtu, &max_mss);
|
||||
|
||||
loop {
|
||||
match device.read_packet(&mut buf) {
|
||||
Ok(n) if n > 0 => {
|
||||
if !handle_tun_packet(
|
||||
&mut buf[..n],
|
||||
max_mss,
|
||||
max_mss.load(std::sync::atomic::Ordering::Relaxed),
|
||||
&name,
|
||||
our_addr,
|
||||
&tun_tx,
|
||||
@@ -685,13 +721,13 @@ pub fn run_tun_reader(
|
||||
our_addr: FipsAddress,
|
||||
tun_tx: TunTx,
|
||||
outbound_tx: TunOutboundTx,
|
||||
transport_mtu: u16,
|
||||
max_mss: MssCeiling,
|
||||
path_mtu_lookup: PathMtuLookup,
|
||||
shutdown_fd: std::os::unix::io::RawFd,
|
||||
) {
|
||||
let _shutdown_fd = ShutdownFd(shutdown_fd);
|
||||
let tun_fd = device.device().as_raw_fd();
|
||||
let (name, mut buf, max_mss) = tun_reader_setup(device.name(), mtu, transport_mtu);
|
||||
let (name, mut buf) = tun_reader_setup(device.name(), mtu, &max_mss);
|
||||
|
||||
// Set TUN fd to non-blocking so we can use select + read without blocking
|
||||
// past the point where select returns readable.
|
||||
@@ -741,7 +777,7 @@ pub fn run_tun_reader(
|
||||
Ok(n) if n > 0 => {
|
||||
if !handle_tun_packet(
|
||||
&mut buf[..n],
|
||||
max_mss,
|
||||
max_mss.load(std::sync::atomic::Ordering::Relaxed),
|
||||
&name,
|
||||
our_addr,
|
||||
&tun_tx,
|
||||
@@ -774,30 +810,25 @@ pub fn run_tun_reader(
|
||||
// _shutdown_fd closes on drop
|
||||
}
|
||||
|
||||
/// Common setup for TUN reader: allocates buffer, computes max MSS.
|
||||
fn tun_reader_setup(device_name: &str, mtu: u16, transport_mtu: u16) -> (String, Vec<u8>, u16) {
|
||||
use super::icmp::effective_ipv6_mtu;
|
||||
|
||||
/// Common setup for TUN reader: allocates the buffer and names the device.
|
||||
///
|
||||
/// The MSS ceiling is deliberately *not* returned. It is read from the shared
|
||||
/// [`MssCeiling`] on every packet, because a transport binding or unbinding
|
||||
/// moves it after this function has run; returning it here is what let the
|
||||
/// reader clamp to a floor derived from the transports that happened to be
|
||||
/// bound at spawn.
|
||||
fn tun_reader_setup(device_name: &str, mtu: u16, max_mss: &MssCeiling) -> (String, Vec<u8>) {
|
||||
let name = device_name.to_string();
|
||||
let buf = vec![0u8; mtu as usize + 100];
|
||||
|
||||
const IPV6_HEADER: u16 = 40;
|
||||
const TCP_HEADER: u16 = 20;
|
||||
let effective_mtu = effective_ipv6_mtu(transport_mtu);
|
||||
let max_mss = effective_mtu
|
||||
.saturating_sub(IPV6_HEADER)
|
||||
.saturating_sub(TCP_HEADER);
|
||||
|
||||
debug!(
|
||||
name = %name,
|
||||
tun_mtu = mtu,
|
||||
transport_mtu = transport_mtu,
|
||||
effective_mtu = effective_mtu,
|
||||
max_mss = max_mss,
|
||||
max_mss = max_mss.load(std::sync::atomic::Ordering::Relaxed),
|
||||
"TUN reader starting"
|
||||
);
|
||||
|
||||
(name, buf, max_mss)
|
||||
(name, buf)
|
||||
}
|
||||
|
||||
/// Process a single TUN packet. Returns `false` if the reader should exit.
|
||||
@@ -1076,12 +1107,13 @@ mod windows_tun {
|
||||
/// packets independently. Returns the writer and a channel sender for
|
||||
/// submitting packets to be written.
|
||||
///
|
||||
/// `max_mss` is the global TCP MSS ceiling. `path_mtu_lookup` is a
|
||||
/// read-only handle to per-destination path MTU learned via
|
||||
/// discovery.
|
||||
/// `max_mss` is the global TCP MSS ceiling, shared live so a transport
|
||||
/// that binds or unbinds after start moves it (see [`MssCeiling`]).
|
||||
/// `path_mtu_lookup` is a read-only handle to per-destination path MTU
|
||||
/// learned via discovery.
|
||||
pub fn create_writer(
|
||||
&self,
|
||||
max_mss: u16,
|
||||
max_mss: MssCeiling,
|
||||
path_mtu_lookup: PathMtuLookup,
|
||||
) -> Result<(TunWriter, TunTx), TunError> {
|
||||
let (tx, rx) = mpsc::channel();
|
||||
@@ -1119,7 +1151,7 @@ mod windows_tun {
|
||||
session: Arc<wintun::Session>,
|
||||
rx: mpsc::Receiver<Vec<u8>>,
|
||||
name: String,
|
||||
max_mss: u16,
|
||||
max_mss: MssCeiling,
|
||||
path_mtu_lookup: PathMtuLookup,
|
||||
}
|
||||
|
||||
@@ -1132,14 +1164,21 @@ mod windows_tun {
|
||||
use super::per_flow_max_mss;
|
||||
use crate::upper::tcp_mss::clamp_tcp_mss;
|
||||
|
||||
debug!(name = %self.name, max_mss = self.max_mss, "TUN writer starting");
|
||||
debug!(
|
||||
name = %self.name,
|
||||
max_mss = self.max_mss.load(std::sync::atomic::Ordering::Relaxed),
|
||||
"TUN writer starting"
|
||||
);
|
||||
|
||||
for mut packet in self.rx {
|
||||
// Read per packet, not once: a transport binding or unbinding
|
||||
// moves the node's egress floor. See `MssCeiling`.
|
||||
let global_max_mss = self.max_mss.load(std::sync::atomic::Ordering::Relaxed);
|
||||
// Per-destination clamp (peer source IPv6 = bytes 8..24)
|
||||
let effective_max_mss = if packet.len() >= 24 {
|
||||
per_flow_max_mss(&self.path_mtu_lookup, &packet[8..24], self.max_mss)
|
||||
per_flow_max_mss(&self.path_mtu_lookup, &packet[8..24], global_max_mss)
|
||||
} else {
|
||||
self.max_mss
|
||||
global_max_mss
|
||||
};
|
||||
// Clamp TCP MSS on inbound SYN-ACK packets
|
||||
if clamp_tcp_mss(&mut packet, effective_max_mss) {
|
||||
@@ -1188,17 +1227,17 @@ mod windows_tun {
|
||||
our_addr: FipsAddress,
|
||||
tun_tx: TunTx,
|
||||
outbound_tx: TunOutboundTx,
|
||||
transport_mtu: u16,
|
||||
max_mss: MssCeiling,
|
||||
path_mtu_lookup: PathMtuLookup,
|
||||
) {
|
||||
let (name, mut buf, max_mss) = super::tun_reader_setup(device.name(), mtu, transport_mtu);
|
||||
let (name, mut buf) = super::tun_reader_setup(device.name(), mtu, &max_mss);
|
||||
|
||||
loop {
|
||||
match device.read_packet(&mut buf) {
|
||||
Ok(n) if n > 0 => {
|
||||
if !super::handle_tun_packet(
|
||||
&mut buf[..n],
|
||||
max_mss,
|
||||
max_mss.load(std::sync::atomic::Ordering::Relaxed),
|
||||
&name,
|
||||
our_addr,
|
||||
&tun_tx,
|
||||
|
||||
Reference in New Issue
Block a user