Merge master into next

Carries the per-peer medium-change probe, the work written on top of it,
and the heartbeat-send accounting fix that came up from the maintenance
line. Every source file merged without a conflict, because this line had
not touched any of them.

The changelog was the only conflict, in two hunks, and one of them was
not the union it looked like. The first was: this line's own Fixed
entries and the other line's Node lifecycle subsection are both
additions, and both are kept. The second was not. The medium-change
detection entry already existed here, inherited when this line last took
a merge, and the other side was offering the same entry with an upgrade
note appended rather than a new one. Taking both would have described the
feature twice, once with the note and once without, so the older copy is
replaced in place and stays under Added where it belongs.

That note is the part with a consequence, and it now reaches this line
too: with detection on by default, a node that shortened its liveness
timeout to one or two seconds for fast failover is refused at startup
after the upgrade, citing a key its operator never set.
This commit is contained in:
Johnathan Corgan
2026-09-09 19:38:53 +00:00
18 changed files with 2672 additions and 410 deletions
+69 -16
View File
@@ -121,13 +121,32 @@ with v0.5.x or earlier peers.
#### Node lifecycle
- Transport-medium change detection, controlled by the new `node.netmon.*`
block (on by default). The node samples a coarse fingerprint of its network
attachment — the source addresses the routing table would pick for an off-link
destination, plus the set of up, non-loopback interface addresses — and
reports a change once the picture settles. A handover is not atomic (the old
address goes, briefly nothing has a route, the new one arrives), so a short
debounce coalesces the burst into one event and a fingerprint that settles
back where it started reports nothing. Linux subscribes to `NETLINK_ROUTE`
block (on by default). For each peer whose transport address is a numeric IP
endpoint, the node asks the kernel which local address it would use to reach
*that peer* — a `connect(2)` on a UDP socket, which resolves the route and
sends nothing — and reports a change once some peer held across two
consecutive samples is reached from a different local address, or has stopped
being reachable at all. Asking the question per peer rather than about the
host is what keeps it quiet: a container bridge, a VPN, a `veth` pair or a
tunnel appearing is not the route to any peer and cannot move the
fingerprint, while a peer on the same LAN — reached by its subnet route, not
the default route — is covered, as is a more specific route moving under a
single peer. Peers joining and leaving are ignored on their own, being
ordinary node behaviour rather than a statement about the medium — except
that a peer seen for the first time is checked against its own
`connect()`-ed socket, and reported if that socket is pinned to a source the
routing table would no longer choose, so a medium change in the window
between a peer authenticating and the next sample is not adopted silently
while that peer sits stranded on the old path. A peer
addressed by MAC, by `.onion` or Nym recipient, by a scoped IPv6 literal, or
by a hostname it has not yet been heard from on, has no route to ask about
and contributes nothing; a node with no peers detects nothing, having nothing
bound to the old path to repair. The peer table is read through the node's
existing lock-free entity snapshot, so the detector stays a detached task
holding no node state. A handover is not atomic (the route goes, briefly
there is none, the new one arrives), so a short debounce coalesces the burst
into one event and a fingerprint that settles back where it started reports
nothing. Linux subscribes to `NETLINK_ROUTE`
multicast (the groups `ip monitor` uses) and macOS and FreeBSD to a
`PF_ROUTE` socket, both reacting to the kernel event in milliseconds; every
other platform samples on a timer at `node.netmon.poll_interval_secs`, which
@@ -142,6 +161,18 @@ with v0.5.x or earlier peers.
rebind described under Fixed above.
Bluetooth is not covered: an adapter's state is not an IP attachment and is
invisible to this detector.
**Upgrade note: this couples a new key to one that has already shipped.** A
handover is ridden out for up to eight settling rounds of
`node.netmon.debounce_ms` before a change is reported, and if that worst case
reaches `node.link_dead_timeout_secs` the reaper tears the peering down
before the change is ever acted on, so the node refuses to start rather than
run in that shape. At the shipped defaults the margin is wide (8 × 250ms = 2s
against 30s), but detection is on by default, so **a node that shortened
`node.link_dead_timeout_secs` to 1 or 2 seconds for fast failover will be
refused at startup after the upgrade**, naming a `node.netmon.*` key its
operator never set. Raise the timeout, lower `debounce_ms` so eight rounds
stay under it, or set `node.netmon.enabled: false`. A
`node.link_dead_timeout_secs` of 0 is exempt from the check.
### Changed
@@ -327,6 +358,27 @@ with v0.5.x or earlier peers.
`fipstop` as "Own Loopback". `req_duplicate` returns to meaning only what it
says.
#### Node lifecycle
- A heartbeat whose send failed no longer counts as one that was delivered.
The peer's "last heartbeat" timestamp was stamped before the send and left
alone whatever came back, so a failure suppressed the next attempt for a
full `node.heartbeat_interval_secs` even though the peer had heard nothing —
on a 10s interval against a 30s `link_dead_timeout_secs`, three failures in
a row were the whole budget. The timestamp now moves only on a send that
returned cleanly, and a separate record of the *attempt* spaces the retries
so a peer that keeps failing is retried in seconds rather than either
hammered every tick or left for a full interval. That retry spacing
applies to the failure path only: gating a healthy peer on it as well
would have floored `node.heartbeat_interval_secs` at two seconds, so a
configured value below that would silently not have been honoured.
- A peer that rotates its address no longer keeps sending from a socket
aimed where it used to be. The authenticated-frame path updated the
peer's address and discarded the flag saying it had changed, so the
per-peer `connect()`-ed UDP socket stayed pinned to the old 5-tuple;
the sibling path already cleared it.
#### Data plane
- A peer that stops reading can no longer stall the node. TCP, Tor, Nym and
@@ -366,15 +418,16 @@ with v0.5.x or earlier peers.
medium-change detection added below, which is exactly that missing signal.
Dropping the sockets is self-healing rather than disruptive: the wildcard
listen socket resolves a route per packet, so sends keep working immediately,
and a correctly-bound connected socket is reinstalled on a later tick. Every
peer on a connectionless transport is also heartbeated at once, so the far
side re-pins to the new source address rather than waiting out its own
heartbeat interval. A peer on a connection-oriented transport keeps the
periodic heartbeat instead. That was because such a send awaited an unbounded
`write_all` on a stream the medium change had very likely just stranded, and
this reaction runs on the rx loop; the writer-task change below removes that
hazard, so widening the fan-out to those transports is now open work rather
than something the design forbids. Measured on a live
and a correctly-bound connected socket is reinstalled on a later tick. The
reaction is scoped to the peers the change names: only their sockets are
dropped, and each of those on a connectionless transport is heartbeated at
once, so the far side re-pins to the new source address rather than waiting
out its own heartbeat interval. A peer on a connection-oriented transport
keeps the periodic heartbeat instead. That was because such a send awaited an
unbounded `write_all` on a stream the medium change had very likely just
stranded, and this reaction runs on the rx loop; the writer-task change below
removes that hazard, so widening the fan-out to those transports is now open
work rather than something the design forbids. Measured on a live
node, a WLAN/LAN switch in either direction now costs no reconnection at all —
the Noise session, tree position and routes survive it. Linux and macOS (the
platforms with the connected-socket fast path); elsewhere the heartbeat alone
+71 -11
View File
@@ -210,8 +210,21 @@ an interface arriving or leaving — and rebinds the send path immediately.
| Parameter | Type | Default | Description |
|-----------|------|---------|-------------|
| `node.netmon.enabled` | bool | `true` | Whether medium-change detection runs |
| `node.netmon.poll_interval_secs` | u64 | `5` | How often the host's network attachment is sampled (backstop period where an event-driven backend exists) |
| `node.netmon.debounce_ms` | u64 | `250` | How long to wait for the picture to settle before acting (`0` disables) |
| `node.netmon.poll_interval_secs` | u64 | `5` | How often the path to each peer is sampled (backstop period where an event-driven backend exists). **Must be at least 1 while `enabled`; `0` is refused at startup.** |
| `node.netmon.debounce_ms` | u64 | `250` | How long to wait for the picture to settle before acting (`0` disables). **Refused at startup when `debounce_ms × 8` reaches `node.link_dead_timeout_secs`** — see below. |
Both refusals stop the node rather than degrade it, so they are worth knowing
before they are met.
A handover is ridden out for up to **8** settling rounds (`MAX_DEBOUNCE_ROUNDS`)
of `debounce_ms` each before a change is reported. If that worst case reaches
`node.link_dead_timeout_secs`, the liveness reaper tears the peering down before
the detector ever reports, so the machinery runs and cannot help — the node
refuses to start rather than run in that shape. At the shipped defaults the
margin is wide (8 × 250 ms = 2 s against 30 s), but the constraint couples two
keys in different blocks: **shortening `link_dead_timeout_secs` for fast
failover can make an untouched `debounce_ms` illegal.** The refusal names both
values and the multiplier.
Established UDP peers use a per-peer `connect()`-ed socket for the send fast
path. `connect(2)` makes the kernel resolve the route once and pin the local
@@ -221,15 +234,62 @@ from an abandoned address while the peer answers where it last heard the node
the peering reports itself connected and carries nothing until
`link_dead_timeout_secs` tears it down, typically 60–90s per switch.
On a detected change the node drops those sockets (the wildcard listen socket
resolves a route per packet, so sends keep working, and a correctly bound
connected socket is reinstalled on a later tick) and heartbeats every peer on a
connectionless transport at once so the far side re-pins to the new source
address. A peer reached over TCP, Tor, Nym or BLE is left to its periodic
heartbeat, since sending to it here would block the node's receive loop on a
stream the medium change has very likely just stranded; those transports
re-dial on send. No peering is torn down: sessions, tree positions and routes
survive the switch.
On a detected change the node drops the sockets of the peers the change names
(the wildcard listen socket resolves a route per packet, so sends keep working,
and a correctly bound connected socket is reinstalled on a later tick) and
heartbeats those of them on a connectionless transport at once so the far side
re-pins to the new source address. A peer reached over TCP, Tor, Nym or BLE is
left to its periodic heartbeat, since sending to it here would block the node's
receive loop on a stream the medium change has very likely just stranded. A peer
the change does not name is left alone entirely. No peering is torn down:
sessions, tree positions and routes survive the switch.
**What counts as a change.** For each peer whose transport address is a numeric
IP endpoint, the node asks the kernel which local address it would use to reach
*that peer* — a `connect(2)` on a UDP socket, which resolves the route and sends
nothing. A change is reported when a peer present in two consecutive samples is
now reached from a different local address, or has stopped being reachable at
all.
Because the question is asked per peer, an interface the node does not peer over
cannot trigger anything: a container bridge, a VPN, a `veth` pair or a tunnel
appearing is not the route to any peer, so it does not enter the sample. A peer
on the same LAN, reached by its subnet route rather than the default route, is
covered as well as one across the internet, and so is a more specific route
moving under a single peer.
The converse is the residual. The probe answers for the peer's *current*
address, and that address is the source of the last authentic packet it sent,
so a peer that roams between two of its own addresses which leave this host by
different interfaces is indistinguishable from a local path move. The reaction
is scoped to the peers named in the change, so such a peer moves nothing but
its own send path — but it is the peer, not this host, that decided the
fingerprint changed.
Peers appearing and leaving are ignored on their own — that is ordinary node
behaviour and says nothing about the medium. A peer seen for the first time is
the one exception, and it is not judged against history but against its own
send path: if its `connect()`-ed socket is pinned to a source the routing table
would no longer choose, it is reported. Without that, a medium change in the
window between a peer authenticating and the detector's next sample would be
the detector's first sight of that peer, and would be adopted silently while
the peer's socket stayed pinned to the path the host had just left. A peer
joining onto a path that has not moved has its socket pinned exactly where its
traffic goes, so it still reports nothing. A peer whose address is not a
probeable IP endpoint contributes nothing: a MAC on Ethernet or BLE, a `.onion`
or Nym recipient reached through a local proxy, an IPv6 literal with a scope
suffix, or a peer still carrying the hostname it was configured with (resolving
one would put a DNS lookup on the sample path; the address becomes numeric as
soon as an authenticated packet arrives from the peer). A node holding no peers
detects nothing, which is correct — it has nothing bound to the old path.
The cost is five non-blocking syscalls per peer per sample, read from the
probe's own code rather than measured: `socket(2)` and `bind(2)`, a `connect(2)`
that sends no packet, a `getsockname(2)`, and the `close(2)` the socket takes on
drop. Nothing goes on the wire and no name is resolved.
`node.limits.max_peers` bounds the per-sample total only where it is set: at
`max_peers: 0`, which means unlimited, there is no bound and the cost tracks the
live peer count instead.
Detection uses the best backend the platform has:
+97
View File
@@ -1115,6 +1115,44 @@ impl Config {
}
}
// Medium-change detection. Both checks are about the detector being
// able to do its job at all, not about taste in numbers.
let netmon = &self.node.netmon;
if netmon.enabled {
if netmon.poll_interval_secs == 0 {
return Err(ConfigError::Validation(
"`node.netmon.poll_interval_secs` must be at least 1; it is the backstop \
period behind the kernel event source, and the only detection signal at \
all on a platform without one"
.to_string(),
));
}
// A handover is ridden out for up to `MAX_DEBOUNCE_ROUNDS` rounds
// of `debounce_ms` before the change is reported. If that can
// outlast the liveness timeout, the reaper tears the peering down
// first and the detector never gets to rebind anything — the
// machinery runs and cannot help.
// Saturating: both operands are operator-supplied `u64`s, and an
// overflow here would wrap to a small number and silently accept
// the very configuration this refuses.
let worst_case_debounce_ms = netmon
.debounce_ms
.saturating_mul(u64::from(crate::node::netmon::MAX_DEBOUNCE_ROUNDS));
let dead_timeout_ms = self.node.link_dead_timeout_secs.saturating_mul(1000);
if dead_timeout_ms > 0 && worst_case_debounce_ms >= dead_timeout_ms {
return Err(ConfigError::Validation(format!(
"`node.netmon.debounce_ms` = {} can hold a report for up to {}ms across \
{} settling rounds, which meets or exceeds \
`node.link_dead_timeout_secs` = {}s: the peering would be reaped \
before the medium change was ever acted on",
netmon.debounce_ms,
worst_case_debounce_ms,
crate::node::netmon::MAX_DEBOUNCE_ROUNDS,
self.node.link_dead_timeout_secs,
)));
}
}
let native = &self.node.native_api;
// Both floors refuse a node that would start, answer every setup call
// and then drop every datagram a peer sent. A zero `backlog` makes the
@@ -2347,6 +2385,65 @@ node:
assert!(config.node.discovery.is_none());
}
#[test]
fn test_a_zero_netmon_poll_interval_is_refused() {
// It was silently clamped to 1s, so a typo produced a node that polled
// twenty times more often than asked and said nothing about it.
let mut config = Config::default();
config.node.netmon.poll_interval_secs = 0;
let err = config.validate().expect_err("validation should fail");
assert!(err.to_string().contains("poll_interval_secs"), "{}", err);
}
#[test]
fn test_a_zero_netmon_poll_interval_is_allowed_when_detection_is_off() {
// Nothing reads it, so refusing the node over it would be pedantry.
let mut config = Config::default();
config.node.netmon.enabled = false;
config.node.netmon.poll_interval_secs = 0;
config
.validate()
.expect("a disabled detector imposes no constraint on its own knobs");
}
#[test]
fn test_a_debounce_that_outlasts_the_dead_timeout_is_refused() {
// The detector rides out a handover for up to MAX_DEBOUNCE_ROUNDS
// rounds before reporting. If that can exceed the liveness timeout the
// peering is reaped first and the detector cannot help — the node runs
// the machinery and still takes the outage it was meant to prevent.
let mut config = Config::default();
config.node.link_dead_timeout_secs = 30;
// 8 rounds x 4000ms = 32s > 30s.
config.node.netmon.debounce_ms = 4000;
let err = config.validate().expect_err("validation should fail");
let msg = err.to_string();
assert!(msg.contains("debounce_ms"), "{}", msg);
assert!(msg.contains("link_dead_timeout_secs"), "{}", msg);
// This message is the whole diagnostic for the only refusal an operator
// reaches by editing `node.netmon.debounce_ms`, and substring
// assertions cannot see how it reads. A run of spaces mid-sentence is
// what a continuation join leaves behind, and rustfmt does not touch
// string literals, so nothing else would catch it.
assert!(
!msg.contains(" "),
"the refusal message has a run of literal spaces in it: {:?}",
msg
);
}
#[test]
fn test_the_default_netmon_block_validates() {
// The shipped defaults must not be a config the node refuses to start
// on, which is the failure mode a cross-field check invites.
Config::default()
.validate()
.expect("the default configuration must validate");
}
#[test]
fn test_validate_transport_advert_requires_nostr_enabled() {
let mut config = Config::default();
+38
View File
@@ -18,6 +18,7 @@
//! also advances only on the tick.
use std::collections::HashMap;
use std::net::{IpAddr, SocketAddr};
use std::sync::Arc;
use crate::identity::NodeAddr;
@@ -629,6 +630,43 @@ pub(crate) struct PeerRow {
pub is_parent: bool,
pub is_child: bool,
pub transport_addr: Option<String>,
/// The peer's current transport address as a numeric IP endpoint, when it
/// is one. Not rendered anywhere: this is the medium-change detector's
/// read of the peer table (see [`crate::node::netmon`]), carried here
/// because the detector is a detached task and this snapshot is the
/// node's existing lock-free read side.
///
/// `None` covers everything that is not a probeable IP destination — a
/// MAC on Ethernet or BLE, a `.onion` or Nym recipient, a peer still
/// carrying the hostname it was configured with, an IPv6 literal with a
/// scope suffix. Typed rather than re-parsed from `transport_addr` above
/// so a change to that string's rendering cannot silently leave the
/// detector with nothing to probe.
pub probe_target: Option<SocketAddr>,
/// Source address this peer's per-peer `connect()`-ed UDP socket was
/// pinned to by `connect(2)`, when it has one. Also not rendered, and read by the same detector:
/// it is what the send path is *actually* using, as against the
/// `probe_target` lookup's answer for what the kernel would choose now.
///
/// `None` where there is no such socket — every platform but Linux and
/// macOS, a peer on another transport, and a peer whose socket has not
/// been installed yet or was just released — and also where the kernel
/// declined to name a source, which is not an address and must not be
/// compared as one.
pub bound_source: Option<IpAddr>,
/// Address this peer's transport is bound to, when that bind is not the
/// wildcard. Read by the same detector, which has to put its probe the
/// same constrained question the send path answers.
///
/// `open_connected_fd` binds the transport's configured address verbatim
/// and only then connects, so a non-wildcard `bind_addr` pins the source
/// whatever the routing table says, while an unconstrained probe takes the
/// kernel's choice. Left unequal, those two answers differ permanently and
/// every first-seen peer reports a move that never happened.
///
/// `None` for the wildcard bind, which is the default and the case where
/// the kernel chooses on both sides.
pub probe_bind: Option<IpAddr>,
pub link_info: Option<PeerLinkInfo>,
pub tree_depth: Option<usize>,
/// `effective_depth = tree_depth + link_cost` — the same quantity
+18 -1
View File
@@ -259,6 +259,7 @@ impl Node {
let now_ms = crate::time::mono_ms();
let ce_flag = header.flags & FLAG_CE != 0;
let mut address_changed = false;
if let Some(peer) = self.peers.get_mut(&node_addr) {
// Initiator-side msg3 confirm (see process_authentic_fmp_plaintext):
// a frame authenticated against post-cutover `current` (no pending)
@@ -276,12 +277,28 @@ impl Node {
now_ms,
);
}
peer.set_current_addr(packet.transport_id, packet.remote_addr.clone());
address_changed =
peer.set_current_addr(packet.transport_id, packet.remote_addr.clone());
peer.link_stats_mut()
.record_recv(packet.data.len(), packet.timestamp_ms);
peer.touch(packet.timestamp_ms);
}
// Address rotation invalidates the per-peer connect()-ed UDP socket,
// which is still pinned to the old 5-tuple. The decrypt-worker
// completion path already does this; this one discarded the flag, so a
// peer that roamed kept sending from a socket aimed where it used to
// be. `netmon`'s first-sight rule now rests on this too: it compares
// that socket's pinned source against a probe to the peer's *current*
// address, and a socket left behind makes those two disagree for as
// long as it survives.
#[cfg(any(target_os = "linux", target_os = "macos"))]
if address_changed {
self.clear_connected_udp_for_peer(&node_addr);
}
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
let _ = address_changed;
// Dispatch to link message handler
self.dispatch_link_message(&node_addr, link_message, ce_flag)
.await;
+172 -8
View File
@@ -20,6 +20,57 @@ use crate::transport::{TransportAddr, TransportId};
use std::time::{Duration, Instant};
use tracing::{debug, info, trace, warn};
/// How long a peer whose heartbeat send *failed* waits before the next attempt.
///
/// Applies to the failure path only. Gating a healthy peer on it too would
/// floor `node.heartbeat_interval_secs` at this value without validating or
/// reporting it, which is a configured knob quietly not doing what it says.
///
/// Short against `heartbeat_interval_secs`, because a failed heartbeat means
/// the peer has heard nothing and the point is to recover well inside
/// `link_dead_timeout_secs` rather than after another full interval. Not
/// shorter still, because the send behind it awaits an unbounded `write_all`
/// on a connection-oriented transport, on the rx loop; retrying that every
/// tick would make a stranded stream a stalled node. Once that write is
/// bounded this can come down to the tick.
const HEARTBEAT_RETRY_INTERVAL: Duration = Duration::from_secs(2);
/// Decide whether a peer is due a heartbeat, from the two timestamps it keeps.
///
/// Two gates rather than one. `sent` is when a heartbeat last *landed*, and it
/// alone paces a healthy peer. `attempt` is when one was last *tried*, and it
/// gates only a peer whose last try failed, holding the retry off for
/// [`HEARTBEAT_RETRY_INTERVAL`] so a peer whose send keeps failing is not
/// retried on every tick.
///
/// **The retry gate is deliberately not consulted on the healthy path.** On a
/// peer whose last send succeeded the two timestamps are equal, so gating there
/// would clamp a configured `heartbeat_interval_secs` up to the retry interval,
/// and that setting has no validation floor.
fn heartbeat_due(
sent: Option<Instant>,
attempt: Option<Instant>,
now: Instant,
interval: Duration,
) -> bool {
let landed_due = match sent {
None => true,
Some(last) => now.duration_since(last) >= interval,
};
// An attempt later than the last success is one that failed, and an attempt
// with no success behind it is the same thing on a peer never reached.
let retry_due = match (attempt, sent) {
(Some(last), Some(landed)) if last > landed => {
now.duration_since(last) >= HEARTBEAT_RETRY_INTERVAL
}
(Some(last), None) => now.duration_since(last) >= HEARTBEAT_RETRY_INTERVAL,
_ => true,
};
landed_due && retry_due
}
/// Emit the operator `trace!` point for a processed ReceiverReport outcome.
///
/// These log points used to live inside `MmpMetrics::process_receiver_report`;
@@ -473,11 +524,19 @@ impl Node {
|| (peer.rekey_msg3_payload().is_some()
&& peer.rekey_msg3_resend_count() < max_resends);
// Check if heartbeat is due.
let heartbeat_due = match peer.last_heartbeat_sent() {
None => true,
Some(last) => now.duration_since(last) >= heartbeat_interval,
};
// Check if heartbeat is due. Two gates, not one: a send that
// failed does not satisfy the interval, so a peer that has
// heard nothing stays due instead of being suppressed by an
// attempt that went nowhere, and the retry gap keeps a peer
// whose send keeps failing from being tried on every tick.
// Both are decided by `heartbeat_due`, which is a pure
// function so it can be tested without driving a send.
let heartbeat_due = heartbeat_due(
peer.last_heartbeat_sent(),
peer.last_heartbeat_attempt(),
now,
heartbeat_interval,
);
PeerLivenessSnapshot {
peer: *node_addr,
@@ -513,14 +572,25 @@ impl Node {
self.route_link_dead(peer, now_ms).await;
}
MmpAction::Heartbeat { peer } => {
// Attempt first, success after: the attempt is recorded
// even if the send below fails or never returns, so the
// retry stays spaced; only a send that came back clean
// moves the interval that says the peer has heard from us.
if let Some(p) = self.peers.get_mut(&peer) {
p.mark_heartbeat_sent(now);
p.mark_heartbeat_attempt(now);
}
if let Err(e) = self
match self
.send_encrypted_link_message(&peer, &heartbeat_msg)
.await
{
trace!(peer = %self.peer_display_name(&peer), error = %e, "Failed to send heartbeat");
Ok(()) => {
if let Some(p) = self.peers.get_mut(&peer) {
p.mark_heartbeat_sent(now);
}
}
Err(e) => {
trace!(peer = %self.peer_display_name(&peer), error = %e, "Failed to send heartbeat");
}
}
}
MmpAction::SendLinkReport { .. }
@@ -588,3 +658,97 @@ impl Node {
}
}
}
#[cfg(test)]
mod tests {
use super::{HEARTBEAT_RETRY_INTERVAL, heartbeat_due};
use std::time::{Duration, Instant};
const INTERVAL: Duration = Duration::from_secs(10);
#[test]
fn a_peer_never_heartbeated_is_due_immediately() {
let now = Instant::now();
assert!(heartbeat_due(None, None, now, INTERVAL));
}
#[test]
fn a_peer_whose_heartbeat_landed_waits_the_configured_interval() {
let landed = Instant::now();
assert!(!heartbeat_due(
Some(landed),
Some(landed),
landed + INTERVAL - Duration::from_millis(1),
INTERVAL
));
assert!(heartbeat_due(
Some(landed),
Some(landed),
landed + INTERVAL,
INTERVAL
));
}
#[test]
fn a_healthy_peer_is_paced_by_the_configured_interval_and_not_by_the_retry_floor() {
// The interval a peer configures can be shorter than the retry floor.
// Consulting the retry gate on the healthy path would clamp it, and
// `heartbeat_interval_secs` has no validation floor to prevent that.
let short = Duration::from_secs(1);
assert!(short < HEARTBEAT_RETRY_INTERVAL);
let landed = Instant::now();
assert!(heartbeat_due(
Some(landed),
Some(landed),
landed + short,
short
));
}
#[test]
fn a_failed_attempt_does_not_suppress_the_next_heartbeat_for_a_full_interval() {
// A heartbeat landed at t0 and the next attempt, at t0 + INTERVAL,
// failed. Once the retry interval has passed the peer is due again,
// rather than waiting another whole interval on a send that never
// reached it.
let landed = Instant::now();
let failed = landed + INTERVAL;
assert!(heartbeat_due(
Some(landed),
Some(failed),
failed + HEARTBEAT_RETRY_INTERVAL,
INTERVAL
));
}
#[test]
fn a_failed_attempt_is_not_retried_before_the_retry_interval() {
let landed = Instant::now();
let failed = landed + INTERVAL;
assert!(!heartbeat_due(
Some(landed),
Some(failed),
failed + HEARTBEAT_RETRY_INTERVAL - Duration::from_millis(1),
INTERVAL
));
}
#[test]
fn a_peer_never_reached_is_retried_on_the_retry_interval_not_the_heartbeat_interval() {
// No heartbeat has ever landed, so there is no interval to pace by.
// The attempt alone spaces the retries.
let failed = Instant::now();
assert!(!heartbeat_due(
None,
Some(failed),
failed + Duration::from_millis(1),
INTERVAL
));
assert!(heartbeat_due(
None,
Some(failed),
failed + HEARTBEAT_RETRY_INTERVAL,
INTERVAL
));
}
}
+87 -37
View File
@@ -36,14 +36,14 @@
//! wildcard listen socket resolves a route per packet, so sends keep working
//! immediately, and `activate_connected_udp_sessions` reinstalls a
//! correctly-bound connected socket on a later tick.
//! 2. **Heartbeat every peer whose send path cannot block.** The frame leaves
//! over the new path and carries the node's new source address, so the far
//! side re-pins on receipt instead of waiting out its own
//! 2. **Heartbeat each moved peer whose send path cannot block.** The frame
//! leaves over the new path and carries the node's new source address, so
//! the far side re-pins on receipt instead of waiting out its own
//! `heartbeat_interval_secs`. Without it the forward direction is fixed but
//! the reverse still points at the old address until the node next happens
//! to send. This runs on the rx loop, and it covers the connectionless
//! transports only — see
//! [`Node::heartbeat_all_peers_after_net_change`] for what a peer on a
//! [`Node::heartbeat_moved_peers_after_net_change`] for what a peer on a
//! connection-oriented transport gets instead, and for why that filter has
//! outlived the reason it was written for.
//!
@@ -52,6 +52,32 @@
//! position and the routes all survive the switch. `link_dead_timeout_secs`
//! remains the backstop for a peer that genuinely cannot be reached on the new
//! medium.
//!
//! # Why the reaction is scoped to the peers that moved
//!
//! [`NetChange`] names them, and the reaction acts on exactly that set. It is
//! not an optimisation: keying the sample on peers put the trigger within
//! reach of a remote party for the first time. `probe_target` is the observed
//! source address of every authentic packet, updated with no throttle, so a
//! peer alternating between two addresses that resolve to different local
//! sources can move the fingerprint at will. Node-wide, that peer could drive
//! every other peering's socket teardown, bounded only by the poll interval.
//! Scoped, the only peer in the set is the roamer itself — whose connected
//! socket `dataplane::encrypted` has already cleared on the address change.
//!
//! A peer is absent from that set for one of two reasons. The first is the
//! one the narrowing rests on: its local source address still resolves to the
//! same place, which is the whole content of the fingerprint — a peer that did
//! not move is a peer whose socket is not stale. The second is that it never
//! reached the sample. `PeerRow::probe_target` parses the peer's current
//! address, so a peer still carrying the hostname it was configured with is
//! `None` there and is skipped while any `connect()`-ed socket it holds stays
//! pinned; the node-wide reaction repaired that peer as collateral and this one
//! does not. Where the node dialled out, that is transient —
//! `set_current_addr` replaces the configured string with the observed numeric
//! source on the first authentic frame. Which side supplies `current_addr` on
//! an inbound peering is not established here, so the second group is not
//! claimed to be empty in general.
use std::time::Instant;
@@ -63,38 +89,43 @@ use crate::node::netmon::NetChange;
use crate::proto::link::LinkMessageType;
impl Node {
/// React to a settled transport-medium change.
/// React to a settled transport-medium change, on the peers it names.
pub(in crate::node) async fn handle_net_change(&mut self, change: NetChange) {
let moved: Vec<NodeAddr> = change.summary.moved.iter().map(|m| m.peer).collect();
let peers = self.peers.len();
// Before the heartbeats: they must go out over a socket that resolves
// the route now, not one still pinned to the interface just left.
let sockets_rebound = self.drop_connected_sockets_after_net_change();
let sockets_rebound = self.drop_connected_sockets_after_net_change(&moved);
let heartbeated = self.heartbeat_all_peers_after_net_change().await;
let heartbeated = self.heartbeat_moved_peers_after_net_change(&moved).await;
info!(
generation = change.generation,
change = %change.summary,
peers,
moved = moved.len(),
sockets_rebound,
heartbeated,
"Transport medium changed; rebinding sends and re-pinning peers"
);
}
/// Drop every per-peer `connect()`-ed UDP socket, returning how many were
/// released.
/// Drop the per-peer `connect()`-ed UDP socket of each peer that moved,
/// returning how many were released.
///
/// See the module docs for why they are stale: `connect(2)` pins the local
/// source address to the interface that carried the route at connect time,
/// and never re-evaluates it.
#[cfg(any(target_os = "linux", target_os = "macos"))]
fn drop_connected_sockets_after_net_change(&mut self) -> usize {
let pinned: Vec<NodeAddr> = self
.peers
fn drop_connected_sockets_after_net_change(&mut self, moved: &[NodeAddr]) -> usize {
let pinned: Vec<NodeAddr> = moved
.iter()
.filter(|(_, peer)| peer.connected_udp().is_some())
.map(|(addr, _)| *addr)
.filter(|addr| {
self.peers
.get(*addr)
.is_some_and(|peer| peer.connected_udp().is_some())
})
.copied()
.collect();
for addr in &pinned {
self.clear_connected_udp_for_peer(addr);
@@ -104,13 +135,16 @@ impl Node {
/// No per-peer connected sockets on this platform, so nothing to rebind.
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
fn drop_connected_sockets_after_net_change(&mut self) -> usize {
fn drop_connected_sockets_after_net_change(&mut self, _moved: &[NodeAddr]) -> usize {
0
}
/// Send one heartbeat to every peer whose send path cannot block, so each
/// learns the node's new source address in one RTT rather than at the next
/// due interval. Returns how many went out.
/// Send one heartbeat to each moved peer whose send path cannot block, so
/// each learns the node's new source address in one RTT rather than at the
/// next due interval. Returns how many sends actually succeeded, which is
/// what the operator log reports — a count of peers *selected* would read
/// the same whether every frame left or none did, and a medium change is
/// exactly when sends start failing.
///
/// The filter was written for a hazard that no longer exists, and it is
/// kept deliberately rather than by oversight. It was this: a
@@ -132,36 +166,52 @@ impl Node {
/// the session counter and the MMP sender record.
///
/// So a peer on TCP, Tor, Nym or BLE keeps the periodic heartbeat it had
/// before this detector existed. It is not stranded by the omission: those
/// transports re-dial on send, and `link_dead_timeout_secs` remains the
/// backstop. Doing better for them means dropping the stale connection
/// before this detector existed, and `link_dead_timeout_secs` remains the
/// backstop. Note that it does *not* recover by redialling: `send_async`
/// only dials when the pool holds no connection for the address, and a
/// connection stranded by a medium change is still in the pool. It is
/// evicted after a write to it fails, so the redial happens on the send
/// after the failure, not on the first one. Doing better for them means
/// dropping the stale connection
/// rather than writing into it, which is a different change with a real
/// cost behind it — a Tor peer pays a fresh circuit — and is not this one.
async fn heartbeat_all_peers_after_net_change(&mut self) -> usize {
pub(in crate::node) async fn heartbeat_moved_peers_after_net_change(
&mut self,
moved: &[NodeAddr],
) -> usize {
let now = Instant::now();
let heartbeat = [LinkMessageType::Heartbeat.to_byte()];
let targets: Vec<NodeAddr> = self
.peers
let targets: Vec<NodeAddr> = moved
.iter()
.filter(|(_, peer)| {
peer.transport_id()
.and_then(|id| self.transports.get(&id))
.is_some_and(|t| !t.transport_type().connection_oriented)
.filter(|addr| {
self.peers.get(*addr).is_some_and(|peer| {
peer.transport_id()
.and_then(|id| self.transports.get(&id))
.is_some_and(|t| !t.transport_type().connection_oriented)
})
})
.map(|(addr, _)| *addr)
.copied()
.collect();
let sent = targets.len();
let mut sent = 0usize;
for addr in targets {
if let Some(peer) = self.peers.get_mut(&addr) {
peer.mark_heartbeat_sent(now);
peer.mark_heartbeat_attempt(now);
}
if let Err(e) = self.send_encrypted_link_message(&addr, &heartbeat).await {
debug!(
peer = %self.peer_display_name(&addr),
error = %e,
"Failed to send post-medium-change heartbeat"
);
match self.send_encrypted_link_message(&addr, &heartbeat).await {
Ok(()) => {
if let Some(peer) = self.peers.get_mut(&addr) {
peer.mark_heartbeat_sent(now);
}
sent += 1;
}
Err(e) => {
debug!(
peer = %self.peer_display_name(&addr),
error = %e,
"Failed to send post-medium-change heartbeat"
);
}
}
}
sent
+2 -1
View File
@@ -2180,7 +2180,8 @@ impl Node {
// not health.
let netmon_cfg = self.config().node.netmon.clone();
if netmon_cfg.enabled {
let (rx, task) = crate::node::netmon::spawn_detector(netmon_cfg);
let (rx, task) =
crate::node::netmon::spawn_detector(netmon_cfg, self.entities_snapshot.clone());
self.supervisor.netmon_rx = Some(rx);
self.supervisor.netmon_task = Some(task);
} else {
+17
View File
@@ -2242,6 +2242,23 @@ impl Node {
is_parent,
is_child,
transport_addr: peer.current_addr().map(|a| format!("{}", a)),
probe_target: peer
.current_addr()
.and_then(|a| a.as_str())
.and_then(|s| s.parse::<std::net::SocketAddr>().ok()),
#[cfg(any(target_os = "linux", target_os = "macos"))]
bound_source: peer.connected_udp().and_then(|s| s.pinned_source()),
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
bound_source: None,
probe_bind: peer
.transport_id()
.and_then(|id| self.transports.get(&id))
.and_then(|t| match t {
TransportHandle::Udp(u) => u.local_addr(),
_ => None,
})
.map(|sa| sa.ip())
.filter(|ip| !ip.is_unspecified()),
link_info,
tree_depth: peer.coords().map(|c| c.depth()),
effective_depth,
+496 -244
View File
@@ -27,13 +27,14 @@
//! | macOS, FreeBSD | `PF_ROUTE` socket | kernel event, ~ms |
//! | everything else | timer, `node.netmon.poll_interval_secs` | up to one period |
//!
//! Both kernel sources are [`crate::transport::watcher::LinkWatcher`], shared
//! with the interface binder. It is asked for a wider set of netlink groups
//! here than the binder asks for: presence is a link question, but a default
//! route moving between two interfaces that both stay up emits nothing in the
//! link group, so a presence subscription would never fire for the change this
//! detector exists to catch. `PF_ROUTE` has no group selection and delivers
//! everything regardless.
//! Both kernel sources are [`crate::transport::watcher::LinkWatcher`], which
//! this module is currently the only consumer of. It is asked for
//! [`groups::EGRESS_PATH`](crate::transport::watcher::groups::EGRESS_PATH)
//! rather than the watcher's default link-presence mask: a route moving
//! between two interfaces that both stay up emits nothing in the link group, so
//! a presence subscription would never fire for the change this detector exists
//! to catch. `PF_ROUTE` has no group selection and delivers everything
//! regardless.
//!
//! Still to come, behind the same seam and without touching the handler:
//! `NotifyIpInterfaceChange` on Windows, and an embedder push on iOS. Android
@@ -47,42 +48,105 @@
//!
//! # What the fingerprint captures
//!
//! Two independent signals, because neither alone is sufficient:
//! One local source address per peer: for every peer whose transport address is
//! a numeric IP endpoint, the address the kernel would pick to reach *that
//! peer*. A connected-but-never-sending UDP socket makes the kernel run its
//! route lookup and bind the source address it would use; five syscalls, read
//! off [`NetFingerprint::sample`] rather than measured, no packets, no name
//! resolution, and it works identically on every platform std supports.
//!
//! - **The preferred source addresses.** A connected-but-never-sending UDP
//! socket makes the kernel run its route lookup and pick the source address
//! it *would* use to reach an off-link destination. That address changes
//! exactly when the default route moves between media, which is the
//! WLAN → 5G case. It costs two syscalls and no packets, and works
//! identically on every platform std supports.
//! - **The set of up, non-loopback interface addresses**, on unix, where
//! `getifaddrs(3)` is available through the `libc` dependency the crate
//! already carries. This catches a medium arriving or leaving without
//! displacing the default route. A peer on the same LAN is reached by the
//! on-link subnet route, and the probe above follows the *default* route by
//! construction, so it looks straight past that path: unplug a LAN cable on
//! a host whose default route is cellular and the source addresses do not
//! move, while every connected socket to that peer is stranded.
//! Keying on peers bounds the *reaction* — only the peers a change names are
//! acted on — and does not bound the *sampling*. One roaming peer still makes
//! the detector re-probe every target, up to [`MAX_DEBOUNCE_ROUNDS`] extra
//! times per wake, and the only limit on how often that can start is
//! `node.netmon.poll_interval_secs`.
//!
//! It is a proxy for "could the local end of one of our peerings have
//! moved", and a broad one — it enumerates every address the host has,
//! rather than the local end of the peerings the node actually holds, so a
//! container bridge or a tunnel interface appearing moves it too.
//! That set is exactly the quantity the reaction cares about. The stale
//! `connect(2)` this whole subsystem exists to repair pinned a *local source
//! address chosen for one destination*, so measuring the same thing for the
//! same destinations asks the kernel the question the bug is about, rather
//! than a proxy for it.
//!
//! On platforms without `getifaddrs` (Windows) only the source-address probe
//! contributes, which still catches every default-route move. The native
//! backend is the answer there, not a richer sample.
//! Two consequences fall out of aiming the probe at peers rather than at the
//! host:
//!
//! - **Interfaces the node does not peer over cannot move it.** A container
//! bridge, a VPN coming up, a `veth` pair, a tunnel — none of them is the
//! route to any peer, so none of them enters the fingerprint. This is by
//! construction rather than by a filter that has to keep guessing which
//! interface names are infrastructure, which is what the host-wide address
//! set this replaced could never get right: on a host running containers it
//! moved, and the node dropped every connected socket and heartbeated every
//! peer, for a `docker compose up`.
//! - **On-link peers are visible.** A peer on the same LAN is reached by its
//! subnet route, not the default route, so a host-wide probe at an off-link
//! destination looked straight past it: unplug the LAN cable on a host whose
//! default route is cellular and nothing host-wide moved while every socket
//! to that peer was stranded. Its own probe follows its own route and moves.
//!
//! It is also the right granularity for a *more specific* route changing under
//! one peer while the rest of the host is untouched, which no single host-wide
//! sample can represent at all.
//!
//! # What it deliberately does not capture
//!
//! A BLE adapter's state is invisible to both signals — it is not an IP
//! attachment at all. That signal comes from the radio (BlueZ properties, the
//! Android callback) and belongs on this same channel, pushed by the BLE
//! **A peer whose address is not a probeable IP destination** contributes
//! nothing — an Ethernet or BLE peer addressed by MAC, a `.onion` or a Nym
//! recipient reached through a local proxy, an IPv6 literal carrying a scope
//! suffix. None of them is IP-attached in the way this detector reasons about,
//! and the connected-UDP pinning it repairs cannot happen to them.
//!
//! **A peer still carrying the hostname it was configured with**, because
//! resolving one on the sample path would put a DNS lookup, with its timeouts,
//! inside the detector's tick. In practice the window is small: the address is
//! replaced by the observed numeric source the first time an authenticated
//! packet arrives (`dataplane::encrypted`), and a peer that has never been
//! heard from has no established peering to strand.
//!
//! **A node with no peers** has an empty fingerprint and detects nothing, which
//! is correct — there is nothing bound to the old path to repair.
//!
//! **A BLE adapter's state** is invisible here, as it was before: it is not an
//! IP attachment at all. That signal comes from the radio (BlueZ properties,
//! the Android callback) and belongs on this same channel, pushed by the BLE
//! transport rather than sampled here.
use std::collections::BTreeSet;
//!
//! # Where the peer list comes from
//!
//! [`crate::control::snapshot::EntitySnapshot`], the node's existing lock-free
//! read side, republished from the tick. The detector is deliberately a
//! detached task holding no node state and taking no node lock, so it reads
//! the peer table the same way the off-loop `show_peers` renderer does. The
//! view is at most one `node.tick_interval_secs` stale, which does not matter:
//! a peer that has just appeared is absorbed on the next sample (see below),
//! and one that has just left is dropped from the comparison rather than
//! reported.
//!
//! # Comparing two samples
//!
//! Over the **intersection** of the two peer sets, never their union: a change
//! is reported when some peer present in both samples is now reached from a
//! different local address. Peers joining and leaving is ordinary node
//! behaviour and says nothing about the medium, so on its own it must not fire
//! a reaction that tears down a working send path. A peer whose probe stops
//! answering entirely — the route to it is gone — is a move to "no source
//! address" and does count, because that peer is exactly the one now stranded.
//!
//! A peer seen for the *first* time is the one case the intersection cannot
//! decide, and it cannot simply be skipped: the sample in which a peer first
//! appears may already be the post-change one, and adopting it silently would
//! swallow the event while that peer's socket stayed pinned to the path the
//! host has just left. Such a peer is judged against its own send path
//! instead — the source `connect(2)` pinned its socket to, which needs no
//! history and asks directly whether that socket is already stale. Churn still
//! fires nothing on its own, because a peer joining onto a path that has not
//! moved is pinned exactly where its traffic goes. See
//! [`NetFingerprint::moved`] for the residual this leaves.
//!
use std::collections::BTreeMap;
use std::fmt;
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr, UdpSocket};
use std::sync::Arc;
use std::time::Duration;
#[cfg(unix)]
@@ -92,18 +156,8 @@ use tokio::task::JoinHandle;
use tracing::{debug, trace, warn};
use crate::config::NetmonConfig;
/// Off-link IPv4 probe destination — RFC 5737 TEST-NET-1, which is guaranteed
/// not to be routed anywhere. Nothing is ever sent to it; `connect(2)` on a UDP
/// socket only resolves the route and binds a source address.
const PROBE_V4: SocketAddr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(192, 0, 2, 1)), 9);
/// Off-link IPv6 probe destination — RFC 3849 documentation prefix. Same
/// no-packets-sent contract as [`PROBE_V4`].
const PROBE_V6: SocketAddr = SocketAddr::new(
IpAddr::V6(Ipv6Addr::new(0x2001, 0x0db8, 0, 0, 0, 0, 0, 1)),
9,
);
use crate::control::snapshot::EntitySnapshot;
use crate::identity::NodeAddr;
/// How many resample rounds the debounce will ride out before reporting
/// anyway. A handover emits a burst (address gone, address added, route
@@ -111,15 +165,18 @@ const PROBE_V6: SocketAddr = SocketAddr::new(
/// change again; riding it out coalesces the burst into one event. Bounded so
/// an interface that flaps continuously still produces events rather than
/// starving the handler forever.
const MAX_DEBOUNCE_ROUNDS: u32 = 8;
pub(crate) const MAX_DEBOUNCE_ROUNDS: u32 = 8;
/// Minimum spacing between two reported changes.
///
/// The reaction is not free: it drops every peer's connected UDP socket (each
/// carrying a drain thread) and sends a heartbeat per peer. An interface that
/// flaps cleanly — settling between each transition, so the debounce reports
/// each one — could otherwise drive that several times a second across up to
/// `node.limits.max_peers` peers, which is thread churn rather than recovery.
/// The reaction is not free: it drops the connected UDP socket of each peer
/// the change names (each carrying a drain thread) and sends that peer a
/// heartbeat. An interface flapping cleanly is the worst case for that, because
/// a medium change moves the whole table at once, so the set the reaction is
/// scoped to is every peer: settling between each transition, so the debounce
/// reports each one, could otherwise drive it several times a second across up
/// to `node.limits.max_peers` peers, which is thread churn rather than
/// recovery.
///
/// A genuine change is delayed by at most this long, against a
/// `link_dead_timeout_secs` measured in tens of seconds, so the trade is
@@ -131,117 +188,288 @@ pub(crate) type NetChangeRx = mpsc::Receiver<NetChange>;
/// Sender held by a detection backend.
pub(crate) type NetChangeTx = mpsc::Sender<NetChange>;
/// A coarse fingerprint of how this host is attached to the network.
/// Where this host sits relative to the peers it holds: one local source
/// address per peer, as the routing table would choose it right now.
///
/// Equality is the whole point: the poller reports a change iff two
/// consecutive samples differ. The contents are only ever used for the
/// operator-facing description of what moved.
/// Each peer carries both the answer to that question and the source its
/// connected socket is already pinned to, because a peer seen for the first
/// time has no earlier sample to be compared against and is judged against its
/// own socket instead. [`NetFingerprint::moved`] is the whole definition of
/// "the medium changed" and the only thing the detector asks of a sample.
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct NetFingerprint {
/// Source address the routing table would pick for an off-link IPv4
/// destination. `None` when there is no IPv4 route at all — itself a
/// meaningful state, and distinct from any address.
v4_source: Option<IpAddr>,
/// IPv6 counterpart of `v4_source`.
v6_source: Option<IpAddr>,
/// Every up, non-loopback unicast address on the host. Always empty on
/// platforms with no `getifaddrs`, which makes the fingerprint degrade to
/// the source-address probe rather than to nothing.
local_addrs: BTreeSet<IpAddr>,
/// Peer → where its traffic leaves from. A peer with no probeable address
/// never appears at all.
sources: BTreeMap<NodeAddr, PeerPath>,
}
/// One peer's local end, from two directions.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
struct PeerPath {
/// The local address the kernel would choose to reach this peer right now.
/// `None` when the route lookup fails, which is a real value rather than a
/// missing one: the peer is still ours, and having no route to it is
/// precisely the state worth reacting to.
current: Option<IpAddr>,
/// The local address this peer's connected UDP socket is bound to, if it
/// has one. Unlike `current` this is not a question put to the kernel — it
/// is what the send path is already doing, and it is the only thing that
/// gives a peer the detector has not seen before a baseline to be judged
/// against. See [`NetFingerprint::moved`].
bound: Option<IpAddr>,
}
impl NetFingerprint {
/// Sample the host's current attachment.
/// Probe every target and record the local address the kernel picks.
///
/// Every call is a handful of non-blocking syscalls — a bind, a connect
/// that sends no packet, a `getifaddrs` walk — with no I/O wait, no name
/// resolution, and no allocation beyond the address set. It is called from
/// a dedicated task on a multi-second timer, so it runs inline rather than
/// through `spawn_blocking`.
pub(crate) fn sample() -> Self {
/// Five non-blocking syscalls per target: `socket(2)` and `bind(2)` behind
/// `UdpSocket::bind`, a `connect(2)` that sends no packet, a
/// `getsockname(2)`, and the `close(2)` the socket takes on drop. No I/O
/// wait, no name resolution, and no allocation beyond the map.
///
/// The count matters because a debounced handover resamples: up to
/// `MAX_DEBOUNCE_ROUNDS` rounds plus the settled sample, times the peers
/// held. `node.limits.max_peers` bounds that only where it is set —
/// the value 0 means unlimited, and there the cost tracks the live peer
/// count instead. It runs inline in the detector's own task rather than
/// through `spawn_blocking`, which is what keeps it off every other task
/// regardless.
pub(in crate::node) fn sample(targets: &[ProbeTarget]) -> Self {
Self {
v4_source: preferred_source(PROBE_V4),
v6_source: preferred_source(PROBE_V6),
local_addrs: interface_addrs(),
sources: targets
.iter()
.map(|t| {
(
t.peer,
PeerPath {
current: preferred_source(t.dest, t.bind),
bound: t.bound,
},
)
})
.collect(),
}
}
/// Build a fingerprint directly, so a test can script a sequence of
/// samples instead of reading the host's real attachment.
/// samples instead of probing real peers. No peer has a connected socket;
/// [`Self::for_test_bound`] is the variant that gives one.
#[cfg(test)]
pub(crate) fn for_test(v4_source: Option<IpAddr>, local_addrs: &[IpAddr]) -> Self {
pub(crate) fn for_test(sources: &[(NodeAddr, Option<IpAddr>)]) -> Self {
Self {
v4_source,
v6_source: None,
local_addrs: local_addrs.iter().copied().collect(),
sources: sources
.iter()
.map(|(peer, current)| {
(
*peer,
PeerPath {
current: *current,
bound: None,
},
)
})
.collect(),
}
}
/// Describe the transition from `self` to `next` for the operator log.
fn diff(&self, next: &Self) -> NetChangeSummary {
NetChangeSummary {
added: next
.local_addrs
.difference(&self.local_addrs)
.copied()
/// As [`Self::for_test`], with each peer's connected socket bound where the
/// third element says.
#[cfg(test)]
pub(crate) fn for_test_bound(sources: &[(NodeAddr, Option<IpAddr>, Option<IpAddr>)]) -> Self {
Self {
sources: sources
.iter()
.map(|(peer, current, bound)| {
(
*peer,
PeerPath {
current: *current,
bound: *bound,
},
)
})
.collect(),
removed: self
.local_addrs
.difference(&next.local_addrs)
.copied()
.collect(),
v4_source_moved: self.v4_source != next.v4_source,
v6_source_moved: self.v6_source != next.v6_source,
v4_source: next.v4_source,
v6_source: next.v6_source,
}
}
/// Which peers are now leaving from somewhere other than where their
/// traffic is actually going out.
///
/// Two rules, because there are two ways to know:
///
/// **A peer in both samples** is judged on whether its probe answer moved.
/// The comparison is over the *intersection* of the two peer sets, never
/// the union: a peer that has only just been authenticated, or one that has
/// just been reaped, differs between the samples for reasons that have
/// nothing to do with the host's attachment, and the reaction — drop the
/// peer's connected socket and heartbeat it — is pure churn on a peer whose
/// send path was never stale.
///
/// **A peer only in the newer sample** has no previous probe answer to be
/// compared against, and skipping it outright leaves a hole this detector
/// cannot afford. `last` gains a peer only at the first wake *after* it
/// appears, so a medium change in that window is the detector's first
/// sight of that peer, and adopting it silently would swallow the very
/// event being adopted — while the peer's connected socket stays pinned to
/// the path the host has just left. The window is up to one
/// `poll_interval_secs` after every peer that authenticates, and a medium
/// change wakes the detector, so the two coincide readily rather than
/// rarely.
///
/// So such a peer is judged against `bound` instead: the address its
/// connected socket is *actually* using. That needs no history — it asks
/// whether the send path is already stale, which is the question the whole
/// subsystem exists to answer, and it is exactly the peer that would
/// otherwise be left stranded. Churn still cannot fire anything on its own:
/// a peer joining onto a path that has not moved has `bound == current` and
/// reports nothing.
///
/// A first-seen peer with no `bound` is still skipped, and that is the
/// residual. It covers three groups, and they are not equally harmless.
///
/// Where there is genuinely no connected socket — every platform but Linux
/// and macOS, and every peer on a stream or proxied transport on those two
/// — nothing is pinned to repair, because the wildcard socket resolves a
/// route per packet. Such a peer is not stranded; it loses only the
/// immediate heartbeat that would have told the far side to re-pin, and
/// notices at its next `heartbeat_interval_secs`.
///
/// The third group is a real hole rather than a harmless one, and it is one
/// tick wide. The tick publishes the entity snapshot before it installs
/// connected sockets (`record_stats_history` then
/// `activate_connected_udp_sessions`, in that order), so a socket installed
/// on tick N is first visible to this detector in the snapshot published on
/// tick N+1. A peer that joins and has its socket installed, and whose path
/// then moves before that next publish, is first seen with `bound` still
/// `None` and is skipped — and it *does* hold a pinned socket. The window
/// in which it can be missed is about one `tick_interval_secs` per join,
/// against a `poll_interval_secs` five times longer, but the *consequence*
/// is not so short: the ordinary diff does not recover the peer. Once it is
/// in both samples it is judged on its probe answer alone, and `bound` is
/// consulted only on the first-sight arm, so the socket stays pinned where
/// it was. While the reaction was node-wide such a peer was repaired as
/// collateral the next time any other peer moved; scoping the reaction to
/// the peers a change names removed that. What bounds it now is
/// `node.link_dead_timeout_secs` reaping the peering.
///
/// A non-wildcard `transports.udp.bind_addr` is not part of the residual;
/// see [`ProbeTarget::bind`], which keeps both sides answering the same
/// question rather than skipping the peer.
///
/// An empty result means nothing moved; it is the detector's entire
/// definition of "no change".
pub(in crate::node) fn moved(&self, next: &Self) -> Vec<PeerSourceMove> {
next.sources
.iter()
.filter_map(|(peer, now)| {
let before = match self.sources.get(peer) {
// Seen before: its own previous probe answer.
Some(then) => then.current,
// First sight: what its socket is bound to, if it has one.
None => match now.bound {
Some(bound) => Some(bound),
None => return None,
},
};
(before != now.current).then_some(PeerSourceMove {
peer: *peer,
before,
after: now.current,
})
})
.collect()
}
}
/// What moved between two fingerprints. Operator-facing only — the handler
/// re-evaluates everything regardless of which field changed, because it
/// cannot map an address back to the peers that were reaching over it.
/// One peer to probe: where to aim, and what its send path is already using.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(in crate::node) struct ProbeTarget {
/// The peer this is about.
pub peer: NodeAddr,
/// Its transport address — the destination the route lookup is run for.
pub dest: SocketAddr,
/// The local address its connected UDP socket is bound to, if it has one.
pub bound: Option<IpAddr>,
/// The address to bind the probe to, when the peer's transport binds a
/// specific one rather than the wildcard. `None` means bind unspecified
/// and let the kernel choose, which is the default posture.
///
/// This exists so the probe asks the same question the send path answers.
/// `open_connected_fd` binds the transport's configured address verbatim
/// and only then connects, so under a non-wildcard `transports.udp.bind_addr`
/// the socket's source is that address whatever the routing table says. An
/// unconstrained probe would answer with the kernel's choice instead, and
/// the two would disagree permanently — reporting a first-sight move, on
/// every peer, forever, with nothing having moved.
///
/// It does not preserve per-peer detection under such a bind, and should
/// not be read as if it did. [`preferred_source`] returns the configured
/// address for every destination, which is the same value the pinned socket
/// already holds, so a route moving under one peer while that address stays
/// configured on the host is invisible here. What is still reported is the
/// bind address itself going away, which moves every peer at once. That is
/// a trade rather than a loss: the send path is pinned to the configured
/// address whatever the routing table says, so there is no per-peer pinning
/// left to repair.
pub bind: Option<IpAddr>,
}
/// One peer whose local source address changed between two samples.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) struct PeerSourceMove {
/// The peer that is now reached from somewhere else.
pub peer: NodeAddr,
/// The local address it was reached from, or `None` if there was no route.
pub before: Option<IpAddr>,
/// The local address it is reached from now, or `None` if the route is gone.
pub after: Option<IpAddr>,
}
/// What moved between two fingerprints.
///
/// Operator-facing, and — unlike the host-wide summary this replaced — it now
/// names the peers affected, because the fingerprint is keyed on them. It is
/// also the reaction's whole input: the handler acts on exactly the peers in
/// [`Self::moved`] and touches no others.
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct NetChangeSummary {
/// Local addresses present now but not before.
pub added: Vec<IpAddr>,
/// Local addresses present before but not now.
pub removed: Vec<IpAddr>,
/// The preferred IPv4 source address changed (a default-route move).
pub v4_source_moved: bool,
/// The preferred IPv6 source address changed.
pub v6_source_moved: bool,
/// The preferred IPv4 source address as of this sample.
pub v4_source: Option<IpAddr>,
/// The preferred IPv6 source address as of this sample.
pub v6_source: Option<IpAddr>,
/// Every peer whose local source address changed, in `NodeAddr` order.
pub moved: Vec<PeerSourceMove>,
/// How many peers were probed in the newer of the two samples, so a log
/// line shows what fraction of the table moved.
pub probed: usize,
}
impl fmt::Display for NetChangeSummary {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let mut parts: Vec<String> = Vec::new();
if self.v4_source_moved {
parts.push(match self.v4_source {
Some(ip) => format!("v4 source -> {}", ip),
None => "v4 source lost".to_string(),
});
}
if self.v6_source_moved {
parts.push(match self.v6_source {
Some(ip) => format!("v6 source -> {}", ip),
None => "v6 source lost".to_string(),
});
}
if !self.added.is_empty() {
parts.push(format!("+{} addr", self.added.len()));
}
if !self.removed.is_empty() {
parts.push(format!("-{} addr", self.removed.len()));
}
if parts.is_empty() {
if self.moved.is_empty() {
return write!(f, "no visible difference");
}
write!(f, "{}", parts.join(", "))
write!(f, "{}/{} peers: ", self.moved.len(), self.probed)?;
// Bounded: an operator needs the shape of the change, and a full table
// moving at once is the common case rather than the interesting one.
const NAMED: usize = 3;
for (i, m) in self.moved.iter().take(NAMED).enumerate() {
if i > 0 {
write!(f, ", ")?;
}
// Both ends: the address the stale `connect(2)` had pinned is what
// an operator correlates against route history, so a line naming
// only the destination leaves out the half being diagnosed.
let before = match m.before {
Some(ip) => ip.to_string(),
None => "no route".to_string(),
};
let after = match m.after {
Some(ip) => ip.to_string(),
None => "no route".to_string(),
};
write!(f, "{} {} -> {}", m.peer.short_hex(), before, after)?;
}
if self.moved.len() > NAMED {
write!(f, ", +{} more", self.moved.len() - NAMED)?;
}
Ok(())
}
}
@@ -257,17 +485,41 @@ pub(crate) struct NetChange {
}
impl NetChange {
/// A synthetic change, for tests that exercise the node's *reaction* to a
/// medium change rather than its detection. The summary is empty because
/// the handler never reads it — it re-evaluates every peer regardless of
/// which address moved, having no way to map an address back to the peers
/// that were reaching over it.
/// A synthetic change naming no peers, for tests that assert the node does
/// *nothing* — the reaction is scoped to the peers the summary names, so an
/// empty summary must move nothing.
#[cfg(test)]
pub(crate) fn for_test(generation: u64) -> Self {
let empty = NetFingerprint::default();
Self {
generation,
summary: empty.diff(&empty),
summary: NetChangeSummary {
moved: Vec::new(),
probed: 0,
},
}
}
/// A synthetic change naming `peers` as having moved, for tests that
/// exercise the node's *reaction* rather than its detection.
///
/// The addresses are placeholders: the handler reads only which peers
/// moved, not where from or to.
#[cfg(test)]
pub(crate) fn for_test_moved(generation: u64, peers: &[NodeAddr]) -> Self {
let moved: Vec<PeerSourceMove> = peers
.iter()
.map(|peer| PeerSourceMove {
peer: *peer,
before: Some(IpAddr::V4(Ipv4Addr::new(192, 168, 1, 10))),
after: Some(IpAddr::V4(Ipv4Addr::new(10, 40, 0, 7))),
})
.collect();
Self {
generation,
summary: NetChangeSummary {
probed: moved.len(),
moved,
},
}
}
}
@@ -357,6 +609,17 @@ impl WakeSource {
timer
}
/// The netlink groups this wake source is subscribed to, or `None` if it
/// is not a live netlink source. For the group-mask assertion in the
/// tests — see `the_detector_subscribes_to_the_route_groups_not_just_link`.
#[cfg(all(test, any(target_os = "linux", target_os = "android")))]
fn subscribed_groups(&self) -> Option<u32> {
match &self.source {
Wake::Kernel(watcher) => watcher.subscribed_groups(),
_ => None,
}
}
/// Wait until it is worth sampling again.
async fn wait(&mut self) {
let WakeSource { source, timer } = self;
@@ -402,20 +665,54 @@ impl WakeSource {
/// Spawn the medium-change detector, using the best backend this platform has.
///
/// Returns the receiver the rx loop drains and the task handle the supervisor
/// aborts at teardown. The channel holds a single slot: a change already queued
/// and not yet handled makes a newer one redundant, because the handler's
/// reaction is "re-evaluate every peer and every backoff", which subsumes any
/// number of coalesced changes. A full channel therefore drops rather than
/// queues, and never applies backpressure to the detector.
pub(crate) fn spawn_detector(cfg: NetmonConfig) -> (NetChangeRx, JoinHandle<()>) {
/// aborts at teardown. The channel holds a single slot, and a full channel
/// drops rather than queues: the detector must never apply backpressure to
/// itself, and the node must never work through a backlog of network states
/// the host has already left.
///
/// Dropping is only safe because the dropped change is not the last word on
/// the peers it named. The handler acts on exactly those peers, so discarding
/// one would strand them if the detector had already adopted the sample it was
/// derived from. It has not: the baseline in [`run_detector`] advances only on
/// a successful send, so the next sample re-derives the move against the same
/// baseline and reports it again once the queue drains. Repairing a peer is
/// idempotent, which is what makes re-reporting cheap rather than a loop.
pub(crate) fn spawn_detector(
cfg: NetmonConfig,
peers: Arc<arc_swap::ArcSwap<EntitySnapshot>>,
) -> (NetChangeRx, JoinHandle<()>) {
let (tx, rx) = mpsc::channel(1);
let handle = tokio::spawn(async move {
let wake = build_wake_source(&cfg);
run_detector(tx, cfg, NetFingerprint::sample, wake).await;
let sample = move || NetFingerprint::sample(&probe_targets(&peers.load()));
run_detector(tx, cfg, sample, wake).await;
});
(rx, handle)
}
/// The peers worth probing, read off the node's published entity snapshot.
///
/// Every peer carrying a numeric IP endpoint, paired with it. The filter is
/// [`crate::control::snapshot::PeerRow::probe_target`] being `Some`, which is
/// already exactly "this peer is an IP destination we could `connect(2)` to":
/// a MAC, a `.onion`, a Nym recipient and an unresolved hostname all arrive
/// here as `None` and are skipped, with no per-transport special-casing in
/// this module.
pub(in crate::node) fn probe_targets(snapshot: &EntitySnapshot) -> Vec<ProbeTarget> {
snapshot
.peers
.iter()
.filter_map(|row| {
row.probe_target.map(|dest| ProbeTarget {
peer: row.node_addr,
dest,
bound: row.bound_source,
bind: row.probe_bind,
})
})
.collect()
}
/// Pick the wake source: the event-driven backend where one exists and starts,
/// the timer otherwise.
///
@@ -484,29 +781,42 @@ where
wake.wait().await;
let mut candidate = sample();
if candidate == last {
if last.moved(&candidate).is_empty() {
// Nothing the node is peering over moved. Adopt the sample anyway:
// it is how a peer that has just joined enters the comparison, and
// one that has left leaves it. Skipping this would freeze `last` on
// the peer set the detector started with, and a peer authenticated
// later would never be compared against anything.
last = candidate;
continue;
}
// The picture is moving. Ride out the burst: resample after the
// debounce window until two consecutive samples agree, so the reported
// change is against a settled state rather than a mid-handover one.
// Settling is judged the same way — no peer moved since the previous
// round — so peers joining or leaving mid-handover cannot extend the
// debounce on their own either.
for _ in 0..MAX_DEBOUNCE_ROUNDS {
if debounce.is_zero() {
break;
}
tokio::time::sleep(debounce).await;
let resampled = sample();
if resampled == candidate {
let settled = candidate.moved(&resampled).is_empty();
candidate = resampled;
if settled {
break;
}
candidate = resampled;
}
// The burst may have settled back to where it started (an address that
// flapped away and returned). Nothing changed, so nothing is reported.
if candidate == last {
// The burst may have settled back to where it started (a route that
// flapped away and returned). Nothing moved, so nothing is reported —
// but the sample is still adopted, for the reason above.
let moved = last.moved(&candidate);
if moved.is_empty() {
trace!("Network fingerprint settled back unchanged; no event");
last = candidate;
continue;
}
@@ -524,12 +834,19 @@ where
generation += 1;
let change = NetChange {
generation,
summary: last.diff(&candidate),
summary: NetChangeSummary {
moved,
probed: candidate.sources.len(),
},
};
last = candidate;
match tx.try_send(change) {
Ok(()) => {}
// The baseline advances only here. A dropped change is a peer set
// nobody will act on, and the reaction is scoped to the peers a
// change names, so adopting `candidate` regardless would leave
// those peers unrepaired for good: the next diff would compare the
// post-change sample against itself and report nothing. Holding
// `last` where it was makes the next sample re-derive the move.
Ok(()) => last = candidate,
Err(mpsc::error::TrySendError::Full(dropped)) => {
debug!(
generation = dropped.generation,
@@ -544,16 +861,30 @@ where
}
}
/// The source address the kernel would use to reach `probe`.
/// The source address the kernel would use to reach `probe`, from `bind_to` if
/// the transport constrains it.
///
/// `connect(2)` on a UDP socket is a pure routing-table operation: it resolves
/// the route, binds a source address, and sends nothing. A failure — most often
/// `ENETUNREACH` with no route of that family — is itself a fingerprint value,
/// reported as `None` rather than swallowed.
fn preferred_source(probe: SocketAddr) -> Option<IpAddr> {
let bind: SocketAddr = match probe {
SocketAddr::V4(_) => SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0),
SocketAddr::V6(_) => SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 0),
/// the route, binds a source address, and sends nothing. The socket is dropped
/// here and never written to, so a peer is probed without a single packet
/// reaching it.
///
/// A failure — most often `ENETUNREACH`, no route to that peer at all — is
/// itself a fingerprint value, reported as `None` rather than swallowed. It is
/// the state a stranded peer is in, so losing it would blind the detector to
/// the case it most needs to see.
fn preferred_source(probe: SocketAddr, bind_to: Option<IpAddr>) -> Option<IpAddr> {
// Port 0 always: the probe wants the transport's *address* constraint, not
// its port, and binding the live port would collide with the socket the
// transport already holds there. A family mismatch between the configured
// bind and this peer is not an error to report — the transport could not
// have reached the peer from it either — so fall back to unspecified and
// let the connect below fail on its own terms.
let bind: SocketAddr = match (probe, bind_to) {
(SocketAddr::V4(_), Some(ip @ IpAddr::V4(_))) => SocketAddr::new(ip, 0),
(SocketAddr::V6(_), Some(ip @ IpAddr::V6(_))) => SocketAddr::new(ip, 0),
(SocketAddr::V4(_), _) => SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0),
(SocketAddr::V6(_), _) => SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 0),
};
let socket = UdpSocket::bind(bind).ok()?;
socket.connect(probe).ok()?;
@@ -566,84 +897,5 @@ fn preferred_source(probe: SocketAddr) -> Option<IpAddr> {
Some(local)
}
/// Every up, non-loopback unicast address on the host.
#[cfg(unix)]
fn interface_addrs() -> BTreeSet<IpAddr> {
let mut out = BTreeSet::new();
let mut head: *mut libc::ifaddrs = std::ptr::null_mut();
// SAFETY: `getifaddrs` either returns 0 and writes an owned linked list
// into `head`, or returns non-zero and leaves `head` untouched — so the
// list is only walked on success. Every node is read behind a null check,
// and `freeifaddrs` releases the list exactly once, after the walk.
if unsafe { libc::getifaddrs(&mut head) } != 0 {
return out;
}
let mut cursor = head;
while !cursor.is_null() {
// SAFETY: `cursor` is non-null here and points at a node of the list
// `getifaddrs` allocated, which stays valid until `freeifaddrs` below.
let entry = unsafe { &*cursor };
cursor = entry.ifa_next;
if entry.ifa_addr.is_null() {
continue;
}
let flags = entry.ifa_flags as i32;
let up = flags & libc::IFF_UP != 0 && flags & libc::IFF_RUNNING != 0;
if !up || flags & libc::IFF_LOOPBACK != 0 {
continue;
}
if let Some(ip) = sockaddr_ip(entry.ifa_addr) {
out.insert(ip);
}
}
// SAFETY: `head` is the list `getifaddrs` allocated above, freed once, and
// not read after this point (`cursor` is null by loop exit).
unsafe { libc::freeifaddrs(head) };
out
}
/// No portable interface enumeration without `getifaddrs`. The fingerprint
/// degrades to the source-address probe, which still catches every
/// default-route move; the native `NotifyIpInterfaceChange` backend is the
/// answer here rather than a richer poll.
#[cfg(not(unix))]
fn interface_addrs() -> BTreeSet<IpAddr> {
BTreeSet::new()
}
/// Read an `IpAddr` out of a kernel-supplied `sockaddr`, if it is one of the
/// two families we fingerprint.
#[cfg(unix)]
fn sockaddr_ip(sa: *const libc::sockaddr) -> Option<IpAddr> {
// SAFETY: `sa` is non-null (checked by the caller) and points at a
// kernel-supplied `sockaddr` whose `sa_family` selects the concrete layout
// that follows. Both branches copy out through `read_unaligned`, so nothing
// here assumes the pointer is aligned for the larger type.
let family = unsafe { std::ptr::addr_of!((*sa).sa_family).read_unaligned() } as i32;
match family {
libc::AF_INET => {
// SAFETY: family is AF_INET, so the allocation is at least a
// `sockaddr_in`.
let raw: libc::sockaddr_in = unsafe { std::ptr::read_unaligned(sa.cast()) };
// `s_addr` holds the octets in network order, so its native-endian
// bytes are the address octets in order.
Some(IpAddr::V4(Ipv4Addr::from(
raw.sin_addr.s_addr.to_ne_bytes(),
)))
}
libc::AF_INET6 => {
// SAFETY: family is AF_INET6, so the allocation is at least a
// `sockaddr_in6`.
let raw: libc::sockaddr_in6 = unsafe { std::ptr::read_unaligned(sa.cast()) };
Some(IpAddr::V6(Ipv6Addr::from(raw.sin6_addr.s6_addr)))
}
_ => None,
}
}
#[cfg(test)]
mod tests;
+830 -70
View File
File diff suppressed because it is too large Load Diff
+215
View File
@@ -34,6 +34,70 @@ fn set_link_dead_timeout(node: &mut crate::node::Node, secs: u64) {
});
}
/// Set `node.heartbeat_interval_secs` on an already-constructed node, the same
/// way `set_link_dead_timeout` does. This is the knob the retry gate must not
/// floor.
fn set_heartbeat_interval(node: &mut crate::node::Node, secs: u64) {
node.replace_context(|ctx| {
let mut cfg = (*ctx.config).clone();
cfg.node.heartbeat_interval_secs = secs;
ctx.config = std::sync::Arc::new(cfg);
});
}
/// A heartbeat whose send failed is not recorded as having landed, and the
/// failed attempt is not retried on the very next tick.
///
/// The failure is forced by taking the node's transport handles away, so the
/// encrypted send fails before any I/O with `TransportNotFound`. Marking the
/// send before it happens, which is what this replaced, would record the peer
/// as heartbeated and suppress the next attempt for a whole interval although
/// the peer heard nothing.
#[tokio::test]
async fn a_failed_heartbeat_send_is_not_recorded_as_landed() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
assert!(nodes[0].node.get_peer(&addr_1).is_some());
// Whatever landed during convergence is the baseline this asserts against.
let landed_before = nodes[0]
.node
.get_peer(&addr_1)
.unwrap()
.last_heartbeat_sent();
// Due on every tick, so the only variable is what the send does.
set_heartbeat_interval(&mut nodes[0].node, 0);
nodes[0].node.transports.clear();
nodes[0].node.check_link_heartbeats().await;
let peer = nodes[0].node.get_peer(&addr_1).expect("peer present");
let failed_at = peer
.last_heartbeat_attempt()
.expect("the attempt is recorded even though the send failed");
assert_eq!(
peer.last_heartbeat_sent(),
landed_before,
"a heartbeat whose send failed was recorded as having landed"
);
// The retry gate spaces the next attempt out rather than letting a failing
// peer be retried on every tick.
nodes[0].node.check_link_heartbeats().await;
let peer = nodes[0].node.get_peer(&addr_1).expect("peer present");
assert_eq!(
peer.last_heartbeat_attempt(),
Some(failed_at),
"a peer whose send failed was retried inside the retry interval"
);
cleanup_nodes(&mut nodes).await;
}
/// A peer past the link-dead timeout is NOT reaped while an FMP rekey is in
/// progress with its msg1 budget unexhausted.
#[tokio::test]
@@ -160,3 +224,154 @@ async fn heartbeat_unaffected_without_rekey() {
cleanup_nodes(&mut nodes).await;
}
/// Rewind a peer's heartbeat bookkeeping by `age`, as if that long had passed
/// since its last successful send.
///
/// The sweep reads `std::time::Instant`, which tokio's paused clock does not
/// move, so elapsed time is staged on the peer rather than waited out. Sets
/// both timestamps, which is the state a *healthy* peer is in.
fn age_heartbeat(node: &mut crate::node::Node, addr: &NodeAddr, age: Duration) {
let then = std::time::Instant::now() - age;
node.peers
.get_mut(addr)
.expect("peer present")
.mark_heartbeat_sent(then);
}
/// **The retry gate must not floor a healthy peer's configured interval.**
///
/// A successful send stamps `last_heartbeat_sent` and `last_heartbeat_attempt`
/// with the same instant. Gating every peer on the attempt timestamp therefore
/// gates the healthy path too, and the effective interval becomes the larger of
/// the configured value and `HEARTBEAT_RETRY_INTERVAL` — so a configured 1s
/// becomes 2s, silently, with nothing validating the value and nothing saying
/// why. `src/node/tests/tcp.rs` already configures 1s against a 3s dead
/// timeout, which is the margin that would quietly halve.
#[tokio::test]
async fn a_healthy_peer_is_heartbeated_on_its_configured_interval() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
set_heartbeat_interval(&mut nodes[0].node, 1);
// Past the configured interval, short of the failure-retry interval. That
// window is the whole defect: healthy, due, and gated anyway.
age_heartbeat(&mut nodes[0].node, &addr_1, Duration::from_millis(1_200));
let before = nodes[0]
.node
.get_peer(&addr_1)
.expect("peer 1 is established")
.last_heartbeat_sent()
.expect("staged above");
nodes[0].node.check_link_heartbeats().await;
let after = nodes[0]
.node
.get_peer(&addr_1)
.expect("peer 1 is still established")
.last_heartbeat_sent()
.expect("still sent");
assert!(
after > before,
"a healthy peer must be heartbeated on its configured interval, not \
floored at the failure-retry interval"
);
cleanup_nodes(&mut nodes).await;
}
/// The other half: after a send that *failed*, the retry is spaced out rather
/// than reattempted on the very next tick.
///
/// Without that spacing a peer whose send keeps failing is retried every tick,
/// and the send behind it can await an unbounded stream write on the rx loop.
/// The peer is re-pinned onto a UDP transport that was never started, so its
/// send fails with `NotStarted` before touching a socket.
#[tokio::test]
async fn a_failing_peer_is_retried_after_the_gap_and_not_before() {
use crate::transport::{TransportAddr, TransportHandle, TransportId, packet_channel};
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
set_heartbeat_interval(&mut nodes[0].node, 1);
let dead_id = TransportId::new(91);
let (tx, _rx) = packet_channel(64);
nodes[0].node.transports.insert(
dead_id,
TransportHandle::Udp(crate::transport::udp::UdpTransport::new(
dead_id,
None,
crate::config::UdpConfig::default(),
tx,
)),
);
nodes[0]
.node
.peers
.get_mut(&addr_1)
.expect("peer 1 is established")
.set_current_addr(dead_id, TransportAddr::from_string("10.0.0.2:2121"));
// Long overdue and healthy-looking, so the sweep will try.
age_heartbeat(&mut nodes[0].node, &addr_1, Duration::from_secs(10));
nodes[0].node.check_link_heartbeats().await;
let attempt_1 = nodes[0]
.node
.get_peer(&addr_1)
.expect("peer 1 is established")
.last_heartbeat_attempt()
.expect("a failed send is still an attempt");
assert!(
nodes[0]
.node
.get_peer(&addr_1)
.unwrap()
.last_heartbeat_sent()
.expect("staged")
< attempt_1,
"the failed send must not have stamped a success"
);
// Immediately after: still inside the gap, so no second attempt.
nodes[0].node.check_link_heartbeats().await;
assert_eq!(
nodes[0]
.node
.get_peer(&addr_1)
.expect("peer 1 is established")
.last_heartbeat_attempt(),
Some(attempt_1),
"a failing peer must not be retried on the very next tick"
);
// Stage the gap as elapsed, keeping the attempt newer than the success so
// the peer still reads as "last one failed".
let past = std::time::Instant::now() - Duration::from_secs(3);
nodes[0]
.node
.peers
.get_mut(&addr_1)
.expect("peer present")
.mark_heartbeat_attempt(past);
nodes[0].node.check_link_heartbeats().await;
let attempt_2 = nodes[0]
.node
.get_peer(&addr_1)
.expect("peer 1 is established")
.last_heartbeat_attempt()
.expect("still attempted");
assert!(
attempt_2 > past,
"a failing peer must be retried once the gap has passed"
);
cleanup_nodes(&mut nodes).await;
}
+419 -12
View File
@@ -9,7 +9,7 @@ use super::spanning_tree::*;
use super::*;
use crate::config::PeerConfig;
use crate::config::TcpConfig;
use crate::node::netmon::NetChange;
use crate::node::netmon::{NetChange, NetFingerprint, ProbeTarget};
use crate::transport::tcp::TcpTransport;
use crate::transport::{TransportAddr, TransportHandle, TransportId, packet_channel};
@@ -37,7 +37,7 @@ fn identity_of(nodes: &[TestNode], j: usize) -> PeerIdentity {
/// The socket is opened against a discard port on loopback: nothing is ever
/// sent through it, and the test only cares whether the handle survives a
/// medium change.
#[cfg(target_os = "linux")]
#[cfg(any(target_os = "linux", target_os = "macos"))]
fn install_connected_udp(node: &mut Node, addr: &NodeAddr, transport_id: TransportId) {
let local: std::net::SocketAddr = "0.0.0.0:0".parse().unwrap();
let peer_sa: std::net::SocketAddr = "127.0.0.1:9".parse().unwrap();
@@ -78,7 +78,7 @@ fn install_connected_udp(node: &mut Node, addr: &NodeAddr, transport_id: Transpo
/// Observed in the field as a peering that carried exactly one packet after a
/// route change and then stalled until the 30s liveness timeout, reporting
/// itself connected the whole time.
#[cfg(target_os = "linux")]
#[cfg(any(target_os = "linux", target_os = "macos"))]
#[tokio::test]
async fn a_medium_change_drops_connected_sockets_pinned_to_the_old_path() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
@@ -102,7 +102,7 @@ async fn a_medium_change_drops_connected_sockets_pinned_to_the_old_path() {
nodes[0]
.node
.handle_net_change(NetChange::for_test(1))
.handle_net_change(NetChange::for_test_moved(1, &[addr_1]))
.await;
assert!(
@@ -116,6 +116,69 @@ async fn a_medium_change_drops_connected_sockets_pinned_to_the_old_path() {
);
}
/// **A peer the change did not name keeps its socket.**
///
/// The reaction is scoped to `change.summary.moved`, and that is not an
/// optimisation. Keying the sample on peers put the trigger within reach of a
/// remote party: `probe_target` is the observed source of every authentic
/// packet, updated with no throttle, so a peer alternating between two
/// addresses can move the fingerprint at will. Node-wide, that peer could tear
/// down every other peering's send path on repeat. If this test starts failing
/// because the untouched peer lost its socket, that lever is back.
#[cfg(any(target_os = "linux", target_os = "macos"))]
#[tokio::test]
async fn a_peer_the_change_did_not_name_keeps_its_socket() {
let mut nodes = run_tree_test(3, &[(0, 1), (0, 2)], false).await;
verify_tree_convergence(&nodes);
let moved = *nodes[1].node.node_addr();
let untouched = *nodes[2].node.node_addr();
let transport_id = nodes[0].transport_id;
install_connected_udp(&mut nodes[0].node, &moved, transport_id);
install_connected_udp(&mut nodes[0].node, &untouched, transport_id);
let heartbeat_before = nodes[0]
.node
.get_peer(&untouched)
.unwrap()
.last_heartbeat_sent();
nodes[0]
.node
.handle_net_change(NetChange::for_test_moved(1, &[moved]))
.await;
assert!(
nodes[0]
.node
.get_peer(&moved)
.unwrap()
.connected_udp()
.is_none(),
"the peer that moved must lose its pinned socket"
);
assert!(
nodes[0]
.node
.get_peer(&untouched)
.unwrap()
.connected_udp()
.is_some(),
"a peer whose source address did not move must keep its socket"
);
assert_eq!(
nodes[0]
.node
.get_peer(&untouched)
.unwrap()
.last_heartbeat_sent(),
heartbeat_before,
"and must not be heartbeated for another peer's move"
);
cleanup_nodes(&mut nodes).await;
}
/// The rebind must not cost the peering. Everything above the socket — the
/// Noise session, the tree position, the routes — is unaffected by which local
/// address the node sends from, so a medium change that tore peers down would
@@ -132,7 +195,7 @@ async fn a_medium_change_keeps_every_peering_intact() {
nodes[0]
.node
.handle_net_change(NetChange::for_test(1))
.handle_net_change(NetChange::for_test_moved(1, &[addr_1]))
.await;
let peer = nodes[0]
@@ -169,7 +232,7 @@ async fn every_peer_is_heartbeated_so_the_far_side_re_pins() {
nodes[0]
.node
.handle_net_change(NetChange::for_test(1))
.handle_net_change(NetChange::for_test_moved(1, &[addr_1]))
.await;
let after = nodes[0]
@@ -192,11 +255,23 @@ async fn every_peer_is_heartbeated_so_the_far_side_re_pins() {
/// A peer on a connection-oriented transport is deliberately left out of the
/// immediate fan-out.
///
/// Its send would await `write_all` on a stream the medium change has very
/// likely just stranded — unbounded, and on the rx loop, where it would hold
/// every other arm of the select behind it. Such a peer keeps the periodic
/// heartbeat it had before this detector existed. If this ever starts passing
/// because the peer *was* heartbeated, the rx loop has a new way to stall.
/// The hazard the filter was written for is gone: every connection-oriented
/// send now enqueues onto its connection's bounded queue and returns, so none
/// of them can await the wire from the rx loop any more. The exclusion is kept
/// anyway, so that widening the fan-out is its own change with its own
/// evidence rather than a side effect of the one that bounded the write. Such
/// a peer keeps the periodic heartbeat it had before this detector existed.
///
/// **So this test guards a deliberate boundary, not a stall.** If the fan-out
/// is widened on purpose, this test is the thing to change, and changing it is
/// how that decision gets recorded.
///
/// The attempt stamp is the observation that sees the exclusion. The fan-out
/// records it for every peer it picks, before the send, and records the sent
/// stamp only for a send that returned. A connection-oriented send fails at
/// the readiness gate, so the sent stamp would sit still either way — whether
/// the peer was excluded or picked and failed — and on its own it cannot tell
/// the two apart.
#[tokio::test]
async fn a_peer_on_a_connection_oriented_transport_is_left_to_the_periodic_heartbeat() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
@@ -234,7 +309,7 @@ async fn a_peer_on_a_connection_oriented_transport_is_left_to_the_periodic_heart
nodes[0]
.node
.handle_net_change(NetChange::for_test(1))
.handle_net_change(NetChange::for_test_moved(1, &[addr_1]))
.await;
let after = nodes[0]
@@ -246,6 +321,15 @@ async fn a_peer_on_a_connection_oriented_transport_is_left_to_the_periodic_heart
before, after,
"a connection-oriented peer must not be heartbeated from the rx loop"
);
assert!(
nodes[0]
.node
.get_peer(&addr_1)
.unwrap()
.last_heartbeat_attempt()
.is_none(),
"a connection-oriented peer must not even be attempted from the rx loop"
);
cleanup_nodes(&mut nodes).await;
}
@@ -257,3 +341,326 @@ async fn a_change_with_no_peers_is_harmless() {
node.handle_net_change(NetChange::for_test(1)).await;
assert!(node.peers.is_empty());
}
/// The detector reads the peer table through the published entity snapshot,
/// and this is the seam: what a peer's transport address is determines whether
/// it arrives on the other side as something to probe. Nothing else in the
/// tree exercises `PeerRow::probe_target`, because nothing renders it — so if
/// the publish site stopped populating it, every other test here would still
/// pass while the detector silently probed an empty table and never reported
/// anything again.
///
/// Each case re-pins the same established peer, because the address is the
/// only variable that matters: the projection is a property of the address,
/// not of the transport it was learned on. (The harness's own peers sit on a
/// synthetic `loopback:1` transport, which is itself correctly unprobeable.)
#[tokio::test]
async fn only_a_peer_with_an_ip_endpoint_reaches_the_probe() {
// (address as the peer carries it, the destination the detector should
// probe, why)
let cases: [(TransportAddr, Option<&str>, &str); 6] = [
(
TransportAddr::from_string("10.0.0.2:2121"),
Some("10.0.0.2:2121"),
"an ordinary IPv4 peer is the whole point",
),
(
TransportAddr::from_string("[2001:db8::1]:2121"),
Some("[2001:db8::1]:2121"),
"IPv6 literals round-trip through the row",
),
(
TransportAddr::from_bytes(&[0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff]),
None,
"a MAC has no IP destination to ask the routing table about",
),
(
TransportAddr::from_string("example.com:2121"),
None,
"resolving a hostname would put DNS on the detector's sample path",
),
(
TransportAddr::from_string("abcdefghij234567.onion:2121"),
None,
"a .onion is reached through a local proxy, not a route",
),
(
TransportAddr::from_string("[fe80::1%eth0]:2121"),
None,
"a scoped link-local literal is not a parseable SocketAddr",
),
];
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
let transport_id = nodes[0]
.node
.peers
.get(&addr_1)
.and_then(|p| p.transport_id())
.expect("peer 1 has a transport");
for (addr, expected, why) in cases {
nodes[0]
.node
.peers
.get_mut(&addr_1)
.expect("peer 1 is established")
.set_current_addr(transport_id, addr.clone());
// The snapshot is published from the tick, which is its only writer.
nodes[0].node.record_stats_history();
let snapshot = nodes[0].node.entities_snapshot.load_full();
let got = crate::node::netmon::probe_targets(&snapshot)
.into_iter()
.find(|t| t.peer == addr_1)
.map(|t| t.dest);
let want = expected.map(|s| s.parse::<std::net::SocketAddr>().unwrap());
assert_eq!(got, want, "{}: {}", addr, why);
}
cleanup_nodes(&mut nodes).await;
}
/// The seed the join-window fix rests on, wired end to end.
///
/// A peer the detector has not seen before is judged against the source its
/// connected socket was pinned to, and that value has to be the address
/// `connect(2)` actually chose — not the wildcard the bind was requested with.
/// `ConnectedPeerSocket::local_addr()` is the wildcard (`0.0.0.0:port`), and
/// reading *that* would compare an unspecified address against a real one for
/// every peer, so every peer joining would report a medium change: precisely
/// the "peer churn fires the fan-out" behaviour the intersection rule exists to
/// prevent. Nothing renders `bound_source`, so no other test would notice.
#[cfg(any(target_os = "linux", target_os = "macos"))]
#[tokio::test]
async fn a_peers_connected_socket_publishes_the_source_it_was_pinned_to() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
let transport_id = nodes[0]
.node
.peers
.get(&addr_1)
.and_then(|p| p.transport_id())
.expect("peer 1 has a transport");
// No socket yet: nothing to seed from, and the peer must say so rather
// than offering the wildcard.
nodes[0].node.record_stats_history();
let row = |n: &Node| {
n.entities_snapshot
.load_full()
.peers
.iter()
.find(|r| r.node_addr == addr_1)
.expect("peer 1 has a row")
.clone()
};
assert_eq!(
row(&nodes[0].node).bound_source,
None,
"a peer with no connected socket has no pinned source to be judged against"
);
// The helper connects to 127.0.0.1:9, so the kernel pins the loopback
// source — a real address, and demonstrably not the `0.0.0.0` the bind was
// requested with.
install_connected_udp(&mut nodes[0].node, &addr_1, transport_id);
// The harness peers sit on a synthetic `loopback:1` address, which is
// correctly not probeable. Re-pin to a numeric endpoint on the same
// transport so the row reaches the probe at all — the pinned source is a
// property of the socket, not of the address, and survives this.
nodes[0]
.node
.peers
.get_mut(&addr_1)
.expect("peer 1 is established")
.set_current_addr(transport_id, TransportAddr::from_string("10.0.0.2:2121"));
nodes[0].node.record_stats_history();
assert_eq!(
row(&nodes[0].node).bound_source,
Some(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)),
"the published source must be what connect(2) pinned, not the wildcard bind"
);
// Publishing it is only half the wiring. Nothing else asserts that
// `probe_targets` carries `bound_source` through to the target, so
// substituting `None` there leaves the whole suite green while silently
// restoring the bug the first-sight rule exists to fix — the same failure
// class as reading the wildcard `local_addr()`, one layer further on.
let snapshot = nodes[0].node.entities_snapshot.load_full();
let target = crate::node::netmon::probe_targets(&snapshot)
.into_iter()
.find(|t| t.peer == addr_1)
.expect("an established UDP peer must be probeable");
assert_eq!(
target.bound,
Some(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)),
"the pinned source must reach the probe target, not stop at the row"
);
// And the last link: `sample()` has to carry it into the fingerprint, or a
// first-seen peer is judged against nothing again. An empty previous
// fingerprint is exactly the first-sight case, and the peer's socket is
// pinned to loopback while the probe answers for a routable destination,
// so the two disagree and a move must be reported.
let sampled = NetFingerprint::sample(&[ProbeTarget {
peer: addr_1,
dest: "192.0.2.1:9".parse().unwrap(),
bound: Some(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)),
bind: None,
}]);
assert!(
!NetFingerprint::default().moved(&sampled).is_empty(),
"sample() must carry the pinned source into the fingerprint, or first \
sight has nothing to judge against"
);
cleanup_nodes(&mut nodes).await;
}
/// A heartbeat that did not go out must not be counted as one that did — in
/// the operator log, or in the peer's own idea of when it was last heard from.
///
/// A medium change is exactly the condition under which sends start failing,
/// so a count of peers *selected* would read identically whether every frame
/// left or none did, and the peer would then be suppressed for a full
/// `heartbeat_interval_secs` on the strength of a send that never landed.
///
/// The peer is re-pinned onto a UDP transport that was never started, which is
/// connectionless — so the fan-out selects it — and fails its send with
/// `NotStarted` before touching a socket.
#[tokio::test]
async fn a_heartbeat_that_failed_is_not_counted_and_does_not_suppress_the_next() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
let peer_1 = identity_of(&nodes, 1);
configure_auto_peer(&mut nodes[0].node, &peer_1);
let dead_id = TransportId::new(88);
let (tx, _rx) = packet_channel(64);
nodes[0].node.transports.insert(
dead_id,
TransportHandle::Udp(crate::transport::udp::UdpTransport::new(
dead_id,
None,
crate::config::UdpConfig::default(),
tx,
)),
);
nodes[0]
.node
.peers
.get_mut(&addr_1)
.expect("peer 1 is established")
.set_current_addr(dead_id, TransportAddr::from_string("10.0.0.2:2121"));
let before = nodes[0]
.node
.get_peer(&addr_1)
.unwrap()
.last_heartbeat_sent();
let sent = nodes[0]
.node
.heartbeat_moved_peers_after_net_change(&[addr_1])
.await;
assert_eq!(
sent, 0,
"the count reports sends that succeeded, not peers picked out"
);
let peer = nodes[0].node.get_peer(&addr_1).unwrap();
assert_eq!(
peer.last_heartbeat_sent(),
before,
"a failed heartbeat must not move the interval that says the peer has heard from us"
);
assert!(
peer.last_heartbeat_attempt().is_some(),
"the attempt is still recorded, or a failing peer would be retried every tick"
);
cleanup_nodes(&mut nodes).await;
}
/// The transport's bind address has to reach the probe, or the detector asks
/// the routing table a different question than the send path answers.
///
/// `open_connected_fd` binds `transports.udp.bind_addr` verbatim before it
/// connects, so under a non-wildcard bind the source is pinned to that address
/// whatever the route says. A probe left unconstrained takes the kernel's
/// choice instead, the two answers differ permanently, and every first-seen
/// peer reports a move that never happened. Nothing renders the field, so
/// substituting `None` at either the publish or the read leaves the rest of
/// the suite green.
///
/// The transport is started on `127.0.0.1:0`: `start_async` fills `local_addr`
/// from the socket the kernel actually bound, which is what the publish reads
/// and what the `!is_unspecified()` filter admits. No privileges are needed.
#[tokio::test]
async fn a_transports_bind_address_reaches_the_probe_target() {
let mut nodes = run_tree_test(2, &[(0, 1)], false).await;
verify_tree_convergence(&nodes);
let addr_1 = *nodes[1].node.node_addr();
let bound_id = TransportId::new(99);
let (tx, _rx) = packet_channel(64);
let mut udp = crate::transport::udp::UdpTransport::new(
bound_id,
None,
crate::config::UdpConfig {
bind_addr: Some("127.0.0.1:0".to_string()),
..Default::default()
},
tx,
);
udp.start_async()
.await
.expect("bind a UDP socket on loopback");
assert_eq!(
udp.local_addr().map(|sa| sa.ip()),
Some(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)),
"precondition: the transport is bound to a real address, not the wildcard"
);
nodes[0]
.node
.transports
.insert(bound_id, TransportHandle::Udp(udp));
// The harness peers sit on a synthetic `loopback:1` address, which is
// correctly not probeable. Re-pin onto the bound transport with a numeric
// endpoint so the row reaches the probe at all.
nodes[0]
.node
.peers
.get_mut(&addr_1)
.expect("peer 1 is established")
.set_current_addr(bound_id, TransportAddr::from_string("10.0.0.2:2121"));
// The snapshot is published from the tick, which is its only writer.
nodes[0].node.record_stats_history();
let snapshot = nodes[0].node.entities_snapshot.load_full();
let target = crate::node::netmon::probe_targets(&snapshot)
.into_iter()
.find(|t| t.peer == addr_1)
.expect("a peer with a numeric endpoint must be probeable");
assert_eq!(
target.bind,
Some(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)),
"the transport's bind address must reach the probe target, not stop at the row"
);
cleanup_nodes(&mut nodes).await;
}
+32 -3
View File
@@ -270,8 +270,16 @@ pub struct ActivePeer {
send_rr: bool,
// === Heartbeat ===
/// When we last sent a heartbeat to this peer.
/// When a heartbeat to this peer last *succeeded*. A send that failed does
/// not move this: it told the peer nothing, and treating it as if it had
/// would leave the peer un-heartbeated for a full interval on the strength
/// of a send that never landed.
last_heartbeat_sent: Option<Instant>,
/// When a heartbeat to this peer was last *attempted*, whatever came of it.
/// Paired with the above so a peer whose send failed is retried sooner than
/// the heartbeat interval without being retried on every tick — see
/// `HEARTBEAT_RETRY_INTERVAL`.
last_heartbeat_attempt: Option<Instant>,
// === Handshake Resend ===
/// Wire-format msg2 for resend on duplicate msg1 (responder only).
@@ -354,6 +362,7 @@ impl ActivePeer {
send_sr: true,
send_rr: true,
last_heartbeat_sent: None,
last_heartbeat_attempt: None,
handshake_msg2: None,
session_established_at: now,
rekey_jitter_secs: draw_rekey_jitter(),
@@ -448,6 +457,7 @@ impl ActivePeer {
send_sr,
send_rr,
last_heartbeat_sent: None,
last_heartbeat_attempt: None,
handshake_msg2: None,
session_established_at: now,
rekey_jitter_secs: draw_rekey_jitter(),
@@ -876,14 +886,33 @@ impl ActivePeer {
// === Heartbeat ===
/// When we last sent a heartbeat to this peer.
/// When a heartbeat to this peer last succeeded.
pub fn last_heartbeat_sent(&self) -> Option<Instant> {
self.last_heartbeat_sent
}
/// Record that we sent a heartbeat.
/// Record that a heartbeat reached the transport without error.
///
/// Call this *after* the send, and only when it returned cleanly. Marking
/// before the send makes a failed heartbeat indistinguishable from a
/// delivered one, which then suppresses the next attempt for a full
/// `heartbeat_interval_secs` although the peer has heard nothing.
pub fn mark_heartbeat_sent(&mut self, now: Instant) {
self.last_heartbeat_sent = Some(now);
self.last_heartbeat_attempt = Some(now);
}
/// When a heartbeat to this peer was last attempted, whatever came of it.
pub(crate) fn last_heartbeat_attempt(&self) -> Option<Instant> {
self.last_heartbeat_attempt
}
/// Record that a heartbeat send was attempted.
///
/// Call this *before* the send, so an attempt that fails, or one that never
/// returns, still spaces the next one out.
pub(crate) fn mark_heartbeat_attempt(&mut self, now: Instant) {
self.last_heartbeat_attempt = Some(now);
}
// === State Updates ===
+14 -4
View File
@@ -39,8 +39,11 @@ pub(crate) struct PeerLivenessSnapshot {
/// An FMP rekey handshake is genuinely in flight with retransmission budget
/// left; suppresses teardown of an otherwise-silent rekey link.
pub rekey_active: bool,
/// A heartbeat is due (`last_heartbeat_sent` is none, or elapsed since it is
/// >= the heartbeat interval).
/// A heartbeat is due. Two conditions, both resolved shell-side: elapsed
/// since the last heartbeat that *landed* is >= the heartbeat interval (or
/// none has), and — only when the last attempt failed — the failure-retry
/// gap has passed since that attempt. The retry gap deliberately does not
/// apply to a healthy peer, or it would floor the configured interval.
pub heartbeat_due: bool,
}
@@ -169,8 +172,15 @@ pub(crate) enum MmpAction {
/// Reap a dead peer: the shell runs `remove_active_peer` +
/// `schedule_reconnect` (with its wall-clock `now_ms`).
ReapPeer { peer: NodeAddr },
/// Send a heartbeat to `peer`: the shell runs `mark_heartbeat_sent` and the
/// encrypted link send.
/// Send a heartbeat to `peer`: the shell runs `mark_heartbeat_attempt`,
/// then the encrypted link send, and `mark_heartbeat_sent` **only if that
/// send returned cleanly**.
///
/// The order and the condition are the contract, not an implementation
/// detail. Stamping the success before the send makes a failed heartbeat
/// indistinguishable from a delivered one, which suppresses the next
/// attempt for a full interval even though the peer heard nothing; the
/// separate attempt stamp is what spaces the retries without doing that.
Heartbeat { peer: NodeAddr },
/// Build (shell: `proto/mmp/` `build_report` + `encode`) and send the given
/// link report over the encrypted link. The interval-advancing
+52 -1
View File
@@ -4,7 +4,7 @@
//! and closes it on drop. See that function's docs for why established
//! peers get their own connected socket.
use std::net::SocketAddr;
use std::net::{IpAddr, SocketAddr};
use std::os::unix::io::{AsRawFd, OwnedFd, RawFd};
/// A `connect()`-ed UDP socket for one established peer.
@@ -28,6 +28,9 @@ pub(crate) struct ConnectedPeerSocket {
fd: OwnedFd,
peer_addr: SocketAddr,
local_addr: SocketAddr,
/// The source address `connect(2)` actually pinned, read back once at
/// construction. See [`ConnectedPeerSocket::pinned_source`].
pinned_source: Option<IpAddr>,
}
impl ConnectedPeerSocket {
@@ -35,10 +38,12 @@ impl ConnectedPeerSocket {
/// `crate::transport::udp::open_connected_fd`) into an owning
/// handle. Takes ownership of the fd; the `OwnedFd` closes it on drop.
pub(crate) fn from_fd(fd: OwnedFd, peer_addr: SocketAddr, local_addr: SocketAddr) -> Self {
let pinned_source = pinned_source_of(fd.as_raw_fd());
Self {
fd,
peer_addr,
local_addr,
pinned_source,
}
}
@@ -50,6 +55,52 @@ impl ConnectedPeerSocket {
pub fn local_addr(&self) -> SocketAddr {
self.local_addr
}
/// The source address the kernel bound when this socket was
/// `connect(2)`-ed, as opposed to [`Self::local_addr`], which is the
/// wildcard the bind was *requested* with and carries no interface
/// information at all.
///
/// This is the quantity the whole connected-socket fast path turns on: the
/// kernel resolves the route once at connect time and pins the source
/// address to whichever interface was carrying it then, and never
/// re-evaluates. Read back once here rather than per call, because it
/// cannot change for the life of the socket — that being exactly the
/// problem. `crate::node::netmon` compares it against the address the
/// routing table would choose now, which is how a peer whose socket is
/// already stale is recognised without any earlier sample to compare
/// against.
///
/// `None` if `getsockname` fails or reports a family this does not decode,
/// which is treated as "no answer" rather than guessed at.
pub(crate) fn pinned_source(&self) -> Option<IpAddr> {
self.pinned_source
}
}
/// `getsockname` on a connected UDP socket, reduced to the local IP.
///
/// Returns `None` on any failure: the caller's contract is that an unknown
/// pinned source is indistinguishable from not having one, and both mean "do
/// not draw a conclusion from this socket".
fn pinned_source_of(fd: RawFd) -> Option<IpAddr> {
let mut storage: libc::sockaddr_storage = unsafe { std::mem::zeroed() };
let mut len = std::mem::size_of::<libc::sockaddr_storage>() as libc::socklen_t;
// SAFETY: `fd` is the socket this handle owns, and `storage` / `len` are a
// correctly sized and initialised out-parameter pair for `getsockname`,
// which writes at most `len` bytes and updates `len` to what it wrote.
let rc =
unsafe { libc::getsockname(fd, &mut storage as *mut _ as *mut libc::sockaddr, &mut len) };
if rc < 0 {
return None;
}
let addr = super::super::unix::sockaddr_to_socket_addr(&storage).ok()?;
// An unspecified source means the kernel declined to choose — no route of
// that family — which is not an address and must not be compared as one.
if addr.ip().is_unspecified() {
return None;
}
Some(addr.ip())
}
impl AsRawFd for ConnectedPeerSocket {
+40
View File
@@ -57,6 +57,30 @@ impl Drop for LinkEventSocket {
}
impl LinkEventSocket {
/// The multicast group mask this socket is actually subscribed to, read
/// back from the kernel rather than remembered from the bind.
///
/// `getsockname` on a netlink socket fills `sockaddr_nl.nl_groups` with the
/// legacy 32-bit subscription mask, which covers every group in
/// [`groups`]. Reading it back is the only way to tell a watcher that
/// *asked* for the right groups from one that got them: a bind with a
/// wrong mask succeeds just as happily as a bind with the right one, and
/// then silently never delivers the messages the caller subscribed for.
#[cfg(all(test, any(target_os = "linux", target_os = "android")))]
fn bound_groups(&self) -> Option<u32> {
let mut sa: libc::sockaddr_nl = unsafe { std::mem::zeroed() };
let mut len = std::mem::size_of::<libc::sockaddr_nl>() as libc::socklen_t;
// SAFETY: `self.fd` is the netlink socket this struct owns, and `sa` /
// `len` are a correctly sized out-parameter pair for `getsockname`.
let rc = unsafe {
libc::getsockname(self.fd, &mut sa as *mut _ as *mut libc::sockaddr, &mut len)
};
if rc < 0 {
return None;
}
Some(sa.nl_groups)
}
fn recv(&self, buf: &mut [u8]) -> std::io::Result<usize> {
let n = unsafe { libc::recv(self.fd, buf.as_mut_ptr() as *mut libc::c_void, buf.len(), 0) };
if n < 0 {
@@ -218,6 +242,22 @@ impl LinkWatcher {
}
}
/// The netlink multicast groups this watcher is actually subscribed to, as
/// the kernel reports them.
///
/// `None` when there is no live source, and on every platform whose backend
/// is `PF_ROUTE`, which has no group selection to report.
///
/// This exists to be asserted on. A bind with the wrong group mask succeeds
/// exactly like a bind with the right one and then silently never delivers
/// what the caller subscribed for, so nothing short of reading the
/// subscription back can tell the two apart without provoking a real
/// kernel event — which needs privileges CI does not have.
#[cfg(all(test, any(target_os = "linux", target_os = "android")))]
pub(crate) fn subscribed_groups(&self) -> Option<u32> {
self.inner.as_ref()?.get_ref().bound_groups()
}
/// Whether an event source is actually backing this watcher.
pub fn is_event_driven(&self) -> bool {
self.inner.is_some()
+3 -2
View File
@@ -4,8 +4,9 @@ fipsctl
fipstop
fips-gateway
# Generated test configs
generated-configs/
# Generated test configs. The suffixed form is what a run with
# FIPS_CI_NAME_SUFFIX set writes, so the glob has to cover both.
generated-configs*/
# Simulation results
sim-results/