Files
fips/src/node/dataplane/rx_loop.rs
T
ArjenandJohnathan Corgan 7c8cf01905 refactor(transport): classify send failures centrally
InterfaceUnavailable existed to make absence branchable, and then stopped
being branchable at the transport boundary: every non-MTU error was flattened
into NodeError::SendFailed { reason: format!(...) }, so no caller downstream
could tell a two-second interface flap from a permanent fault. Both got the
same treatment, which for a half-built handshake means being torn down and
filed as peer misbehaviour.

The classification belongs on the error rather than at each call site, and the
question worth asking is not what went wrong but whether waiting fixes it: a
transient failure was refused by a condition the daemon is already working to
resolve, so the state built around it — a half-finished handshake, a route, a
queued packet — is worth keeping. TransportError::is_transient answers that
once, and NodeError::SendUnavailable carries the answer across the node
boundary instead of discarding it.

Deliberately narrow: only InterfaceUnavailable. Timeout and ConnectionRefused
describe a remote that did not answer, which is a statement about the peer
rather than about this node's ability to transmit, and their retry paths sit
at a different layer. The test pins that narrowness in both directions,
because the failure mode of this abstraction is someone adding a variant to
the transient list and quietly making callers hold state open for a fault that
will never clear.

No behaviour change yet. This is the plumbing half; the callers that should
act on it — route withdrawal on detach, and not counting a local interface
flap as a handshake reject — are recorded in reference/ and deferred, because
both are routing changes that want their own test story.

feat(node): withdraw a transport's peers when its interface goes away

Losing an interface withdrew nothing. The peers stayed in the registry, the
routes through them stayed selectable, and this node kept advertising
reachability it no longer had — so transit traffic was dropped in silence and
other nodes kept routing toward us for those destinations, until the liveness
reaper noticed up to link_dead_timeout_secs (30 s) later.

Measured on real hardware: a dongle detached at 07:37:19 took the node's parent
with it, and no new parent was chosen until 07:37:46. Twenty-seven seconds
routing through a link that had already gone, with four alternative peers
available the whole time. The alternatives are the point — a mesh that can
route around a dead link should not be the last to hear the link is dead.

The detach edge is both earlier and more certain than inactivity, so it is the
better trigger. reap_peers_on_transport routes through the same
route_link_dead the liveness reaper uses rather than open-coding a second
teardown: every consequence of losing a peer — sessions, path MTU release,
session indices, decrypt-worker unregistration, the link, the control machine,
tree cleanup and re-announce, bloom withdrawal — already hangs off that one
path, and a parallel one would drift from it.

Not policy-filtered. Whether an interface's absence is normal is a statement
about node *health*; it says nothing about whether the routes over it still
work. An optional interface's peers are exactly as unreachable.

PathBroken needs no new wiring. Once the peers are gone resolve_next_hop
returns None, which takes the NoRoute path — and that one already synthesises
the routing error, rate limiting included. The cure for the silent drop was to
stop having a route, not to add a second error path.

Deliberately undamped. A flapping interface cannot drive a reap storm through
here: ChurnGuard suppresses `announce` after three short-lived bindings, which
leaves `announced` false, which makes `detached()` return `retract: false` —
so no presence edge is published at all during churn. The edges this reacts to
are already rate-limited at the source, reaping an already-reaped transport is
a no-op, and a second damper would only add a way for the two to disagree.

The trade taken: immediate reaping costs a re-peer for an absence shorter than
the dead timeout that then recovers — a `wifi reload` returns in ~5 s and today
costs nothing, where this costs ~15 s of re-peering. Accepted, because
black-holing is silent, poisons other nodes' routing and needs the full timeout
to clear, where a re-peer is bounded, visible and self-healing. A grace period
remains available if that proves wrong; reference/ records its shape.

The integration assertion is the one that proves the wiring rather than the
unit: link_dead_timeout_secs is left at its 30 s default, so a withdrawal
inside 15 s can only have come from the detach edge. Verified against the
defect — with the reap disabled, that assertion fails and every other case in
the suite still passes.

fix(node): keep a half-built link when msg2 hits a transient transport

A send refused because the interface is absent or mid-rebind was treated as a
failed handshake: the link was removed, the reverse-address entry dropped, the
session index freed, the control machine torn down, the queued PromoteToActive
aborted — and the whole thing recorded as
RejectReason::Handshake(HandshakeReject::BadState).

That counter means "the remote sent something invalid". A local interface flap
is not the remote's fault, and an operator reading the rejects would conclude
it was. The initiator, meanwhile, resends msg1 into a link that no longer
exists and has to rebuild from nothing.

The binder is already working to bring the interface back, so the half-built
link is now left exactly where it is for that resend to land on. Nothing leaks
by staying: an initiator that never resends leaves a stale connection, which
`check_timeouts` reaps at `handshake_timeout_secs` like every other abandoned
handshake. Only a genuinely terminal error still tears down.

This is the first consumer of `TransportError::is_transient`, which is what
the central classification was for — before it, the distinction did not
survive as far as this call site.

The rekey msg1 send site gets the severity half only. Its teardown was already
benign: it returns before `set_rekey_state`, so the cycle simply does not
start and is retried when rekey next comes due, with nothing torn down and
nothing charged to the peer. Only the `warn!` was wrong for a local,
self-clearing condition the presence machine has already reported.

The test drives a real absent Ethernet transport rather than a stub, so the
error under test is the one production raises, from the code path that raises
it. Verified against the defect: with the transient branch disabled, the link
is destroyed and the assertion fails.

The deferral gives the epoch-mismatch restart arm in handle_msg1 a third
outcome, and its post-promote debug_assert! did not admit it. That arm's
assertion required the machine to be Established or absent; a transient msg2
failure returns before PromoteToActive and leaves it registered at
Handshaking{ReceivedMsg1}, so a debug build panics there. The assertion is
widened to name that phase exactly, which keeps it red for any other state,
and a_transient_msg2_failure_on_the_restart_path_leaves_the_fresh_leg_pending
covers the arm the existing transient test does not reach. Three comments
around those two arms claimed a send failure always removes the machine; each
now names the transient case as well. Release builds were never affected: the
tail is gated on Established, so a deferred machine simply skips it.

The route_link_dead doc comment is put back on route_link_dead. Inserting
reap_peers_on_transport between the comment and the function it described left
both blocks running together, so rustdoc attached the whole thing to the new
function and route_link_dead lost its documentation. The restored text also
names the second caller this commit adds, and generalises the sentence about
where now_ms comes from, since both callers now hoist it once per batch.

The transport-layer design's grace-period paragraph is rewritten. It argued
that no linger timer was needed because peers survive a detach untouched and
the liveness reaper is the effective bound. The detach reap makes both halves
false, so the paragraph now states the trade it actually makes: peers go at
the edge, a half-built link is still held, and a short absence that recovers
costs a re-peer, which is preferred to silent black-holing.

Changelog entries for both user-visible halves.

The insert_transport_for_test helper is no longer added here. Nothing used it
until two commits later, so cargo clippy --all-targets -- -D warnings failed
on dead code at this point in the history; it now lands with its caller.
2026-09-10 19:18:09 +00:00

734 lines
39 KiB
Rust

//! RX event loop and packet dispatch.
use crate::control::{ControlSocket, commands};
use crate::node::reject::{RejectReason, TransportReject};
use crate::node::{Node, NodeError};
use crate::proto::fmp::wire::{
COMMON_PREFIX_SIZE, CommonPrefix, FMP_VERSION, PHASE_ESTABLISHED, PHASE_MSG1, PHASE_MSG2,
expected_payload_len,
};
use crate::transport::ReceivedPacket;
use std::time::Duration;
use tracing::{debug, info, warn};
/// Inside the packet_rx burst drain, run a fallback drain every
/// N packets so bounced FMP plaintexts can't sit behind a full
/// 256-packet UDP burst. Used on unix; on Windows the decrypt-worker
/// pool isn't spawned so the fallback channel is always empty —
/// hold the constants at module scope anyway so the burst-loop
/// dispatch in `run_rx_loop` doesn't need a `#[cfg]` on every site.
const FALLBACK_INTERLEAVE_EVERY: usize = 32;
/// How many fallback events to drain per interleave step. Bounded so
/// the inner loop can keep making forward progress on packet_rx.
#[cfg(unix)]
const FALLBACK_INTERLEAVE_BUDGET: usize = 32;
impl Node {
/// Run the receive event loop.
///
/// Processes packets from all transports, dispatching based on
/// the phase field in the 4-byte common prefix:
/// - Phase 0x0: Encrypted frame (session data)
/// - Phase 0x1: Handshake message 1 (initiator -> responder)
/// - Phase 0x2: Handshake message 2 (responder -> initiator)
///
/// Also processes outbound IPv6 packets from the TUN reader for session
/// encapsulation and routing through the mesh.
///
/// Also processes DNS-resolved identities for identity cache population.
///
/// Also runs a periodic tick (1s) to clean up stale handshake connections
/// that never received a response. This prevents resource leaks when peers
/// are unreachable.
///
/// This method takes ownership of the packet_rx channel and runs
/// until the channel is closed (typically when stop() is called).
pub async fn run_rx_loop(&mut self) -> Result<(), NodeError> {
// No shutdown observer → today's infinite loop, byte-identical. All
// existing callers/tests use this; `pending()` never fires, so the
// shutdown/deadline arms below stay permanently disabled.
self.run_rx_loop_with_shutdown(std::future::pending()).await
}
/// The rx event loop, which serves until `shutdown` fires and then drains
/// **in place** before returning.
///
/// The channel receivers are moved into this frame's locals and live across
/// both serve and drain, so — unlike a `select!`-cancelled loop — they are
/// never destructively dropped mid-flight; they are released only on clean
/// exit, after which teardown does not need them.
///
/// - While serving (`drain_deadline == None`) the loop is behaviorally
/// identical to before: the shutdown arm, the deadline arm, and the
/// peers-empty early-exit are all guarded off, so the hot per-packet path
/// and the `biased` order of the real arms are unchanged.
/// - When `shutdown` fires, the loop calls [`Node::enter_drain`] once
/// (broadcast Disconnect, gate the reconciler off) and arms the bounded
/// deadline, then keeps servicing inbound/tick/peer-removal until all
/// peers clear or the deadline elapses, then returns. The caller
/// ([`Node::finish_shutdown`]) closes the window and tears down.
pub async fn run_rx_loop_with_shutdown(
&mut self,
shutdown: impl std::future::Future<Output = ()>,
) -> Result<(), NodeError> {
tokio::pin!(shutdown);
// `None` = serving; `Some(deadline)` = draining (bounded window).
let mut drain_deadline: Option<tokio::time::Instant> = None;
let mut packet_rx = self.packet_rx.take().ok_or(NodeError::NotStarted)?;
// Take the TUN outbound receiver, or create a dummy channel that never
// produces messages (when TUN is disabled). Holding the sender prevents
// the channel from closing.
let (mut tun_outbound_rx, _tun_guard) = match self.supervisor.tun_outbound_rx.take() {
Some(rx) => (rx, None),
None => {
let (tx, rx) = tokio::sync::mpsc::channel(1);
(rx, Some(tx))
}
};
// Take the DNS identity receiver, or create a dummy channel (when DNS
// is disabled). Same pattern as TUN outbound.
let (mut dns_identity_rx, _dns_guard) = match self.supervisor.dns_identity_rx.take() {
Some(rx) => (rx, None),
None => {
let (tx, rx) = tokio::sync::mpsc::channel(1);
(rx, Some(tx))
}
};
// Take the runtime child-liveness receiver, or a dummy channel (when the
// node was seeded straight into Running without a start()). Holding the
// dummy sender in the guard keeps the channel open. Same pattern as TUN
// outbound / DNS identity.
let (mut child_exit_rx, _child_exit_guard) = match self.child_exit_rx.take() {
Some(rx) => (rx, None),
None => {
let (tx, rx) = tokio::sync::mpsc::channel(1);
(rx, Some(tx))
}
};
// Interface-presence receiver, or a dummy channel — same pattern and
// same reason as the child-liveness receiver above.
let (mut presence_rx, _presence_guard) = match self.transport_presence_rx.take() {
Some(rx) => (rx, None),
None => {
let (tx, rx) = tokio::sync::mpsc::channel(1);
(rx, Some(tx))
}
};
let tick_period = Duration::from_secs(self.config().node.tick_interval_secs);
let mut tick = tokio::time::interval(tick_period);
// Set up control socket channel
let (control_tx, mut control_rx) =
tokio::sync::mpsc::channel::<crate::control::ControlMessage>(32);
if self.config().node.control.enabled {
let config = self.config().node.control.clone();
let tx = control_tx.clone();
let read_handle = self.control_read_handle();
tokio::spawn(async move {
match ControlSocket::bind(&config) {
Ok(socket) => {
socket.accept_loop(tx, read_handle).await;
}
Err(e) => {
warn!(error = %e, "Failed to bind control socket");
}
}
});
}
// Drop unused sender to avoid keeping channel open if control is disabled
drop(control_tx);
// Native datagram API socket. Experimental, off by default, and built
// on Linux, FreeBSD and macOS: Windows has no way to pass a descriptor
// between processes at all, which is the mechanism itself. Bound
// synchronously so a bad path or a socket already in use is reported
// here, before the node starts serving, rather than at whatever later
// moment a spawned bind happened to run.
// Two channels, as the TUN plane has: registrations go one way and are
// rare, datagrams go the other and arrive in bursts. Keeping them apart
// lets the data arm drain in batches without a registration waiting
// behind a burst of traffic.
let (native_out_tx, mut native_outbound_rx) =
tokio::sync::mpsc::channel::<crate::native::link::Outbound>(1024);
let _native_out_guard = native_out_tx.clone();
let (mut native_rx, _native_guard) = {
let (tx, rx) = tokio::sync::mpsc::channel::<crate::native::link::NativeMessage>(64);
// The `cfg` sits on the binding rather than on the arm that drains
// the receiver, because `tokio::select!` does not accept one. Where
// there is no listener the channel exists and nothing ever sends,
// and the guard keeps it open so the arm never sees a closed
// receiver.
#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
let guard = {
let mut guard = Some(tx.clone());
if self.config().node.native_api.enabled {
match crate::native::NativeApi::bind(&self.config().node.native_api) {
Ok(socket) => {
// The accept loop owns a sender, so the dummy guard
// is dropped: the channel closes when the last
// client task ends, not while one is still serving.
guard = None;
let npub = self.npub();
tokio::spawn(socket.accept_loop(tx, native_out_tx.clone(), npub));
}
Err(e) => {
warn!(error = %e, "Failed to bind native API socket");
}
}
}
guard
};
#[cfg(not(any(target_os = "linux", target_os = "freebsd", target_os = "macos")))]
let guard = Some(tx.clone());
(rx, guard)
};
// Transport-medium change receiver, or a dummy channel when detection
// is disabled (or the node was seeded straight into Running without a
// start()). Same guard pattern as TUN outbound and DNS identity: the
// held sender keeps the channel open so the arm never sees it closed.
let (mut netmon_rx, _netmon_guard) = match self.supervisor.netmon_rx.take() {
Some(rx) => (rx, None),
None => {
let (tx, rx) = tokio::sync::mpsc::channel(1);
(rx, Some(tx))
}
};
// Decrypt-worker fallback receiver. The worker pushes each
// authenticated FMP plaintext here so rx_loop can finish the
// per-peer side-effects (stats, MMP, ECN, link dispatch).
// Always declared so the `tokio::select!` arm doesn't need
// a `cfg` (which the macro doesn't support); on Windows the
// channel just never sees events.
let (mut decrypt_fallback_rx, _decrypt_fallback_guard) = {
#[cfg(unix)]
{
match self.decrypt_fallback_rx.take() {
Some(rx) => (rx, None),
None => {
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
(rx, Some(tx))
}
}
}
#[cfg(not(unix))]
{
// On non-unix nothing ever sends, but the macro arm
// still needs an existing rx. Keep the sender alive to
// avoid the channel closing into an Err loop.
let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<()>();
(rx, Some(tx))
}
};
info!("RX event loop started");
// Optional per-stage perf profiler (FIPS_PERF=1). No-op otherwise.
crate::perf_profile::maybe_spawn_reporter();
loop {
// Bounded drain mode: break as soon as all peers have cleared. In
// normal mode (`None`) this short-circuits before touching
// `self.peers`, so the loop is byte-identical.
if drain_deadline.is_some() && self.peers.is_empty() {
info!("Drain complete: all peers cleared, ending drain loop");
break;
}
tokio::select! {
biased;
// Decrypt-worker fallback drains FIRST. Under sustained
// inbound bursts the packet_rx drain (up to 256 packets)
// can starve fallback work for tens of ms — TCP doesn't
// tolerate that (late ACKs → dup-ACK fast retransmits →
// cwnd collapse). Promoting fallback gives the kernel's
// TCP machinery a fair chance to ACK in time.
Some(event) = decrypt_fallback_rx.recv() => {
#[cfg(unix)]
{
self.process_decrypt_worker_event(event).await;
let mut drained = 0;
while drained < 255 {
match decrypt_fallback_rx.try_recv() {
Ok(ev) => {
self.process_decrypt_worker_event(ev).await;
drained += 1;
}
Err(_) => break,
}
}
}
#[cfg(not(unix))]
let _ = event;
}
packet = packet_rx.recv() => {
match packet {
Some(p) => self.process_packet(p).await,
None => break, // channel closed
}
// Drain remaining ready inbound packets in a tight loop
// before yielding back to select! — every yield is a
// futex hop on tokio's multi-thread scheduler, and at
// line rate the kernel UDP queue typically has several
// datagrams available per wake. Caps at a batch
// boundary so other branches (tick, control) eventually
// get a turn even under sustained load.
//
// **Interleave fallback drain** every N packets so
// bounced FMP plaintexts (heartbeats, post-FMP-
// decrypt forwarding payloads, control frames) don't
// sit in the fallback queue for a full 256-packet
// burst. Even with the priority-first ordering of
// the outer select!, once we're inside this inner
// loop only this interleave can free queued
// fallbacks. On multihop forwarding paths this is
// the difference between back-to-back encrypt-
// worker dispatches happening promptly vs piling up
// behind the rx burst.
let mut drained: usize = 1; // count the packet processed above
while drained < 256 {
if drained.is_multiple_of(FALLBACK_INTERLEAVE_EVERY) {
#[cfg(unix)]
{
let mut fb_drained = 0;
while fb_drained < FALLBACK_INTERLEAVE_BUDGET {
match decrypt_fallback_rx.try_recv() {
Ok(ev) => {
self.process_decrypt_worker_event(ev).await;
fb_drained += 1;
}
Err(_) => break,
}
}
}
}
match packet_rx.try_recv() {
Ok(p) => {
self.process_packet(p).await;
drained += 1;
}
Err(_) => break,
}
}
// Trailing fallback drain so the last bounced
// packets of the burst aren't held up by the
// next select! iteration.
#[cfg(unix)]
{
let mut fb_drained = 0;
while fb_drained < 256 {
match decrypt_fallback_rx.try_recv() {
Ok(ev) => {
self.process_decrypt_worker_event(ev).await;
fb_drained += 1;
}
Err(_) => break,
}
}
}
}
// Runtime child-liveness. Placed AFTER `packet_rx` so the hot
// inbound path keeps its `biased` priority. A directly-observable
// child (TUN threads, DNS/mDNS/Nostr) exited on its own; feed the
// FSM, which republishes health (Degraded here — a Running node
// always has ≥1 transport up). `on_child_exited` only ever emits
// `PublishState`; other variants are ignored defensively.
maybe_child = child_exit_rx.recv() => {
if let Some(child) = maybe_child {
// Drop anything the dead child published for embedders
// (e.g. the DNS responder's bound address) before
// republishing health, so nothing outside the node can
// observe an address the listener no longer answers on.
self.retract_child_publications(child);
let actions = self
.supervisor
.fsm
.step(crate::node::lifecycle::supervisor::Event::ChildExited { child });
for action in actions {
if let crate::node::lifecycle::supervisor::Action::PublishState(ns) =
action
{
self.supervisor.state = ns;
}
}
// A transport child exiting leaves the bound set, so
// it can be the one that was holding the node's egress
// MTU down. `is_bound()` is `is_operational()` plus the
// presence refinement, and this moves the first half.
self.refresh_tun_mss_ceiling();
}
}
// Interface presence. An interface-bound transport's binder
// reports attach and detach; the FSM folds it into health.
// Unlike `ChildExited` this is reversible in both directions —
// the interface coming back republishes `Running` — which is
// the whole point of `Degraded` being a level rather than a
// latch.
maybe_presence = presence_rx.recv() => {
if let Some(edge) = maybe_presence {
// Health is policy-filtered; the MTU floor below is
// not. An `optional` interface's absence is normal and
// must not move the node off `Full`, but it changes
// the bound set all the same.
if edge.health_relevant {
let child = crate::node::lifecycle::supervisor::Child::Transport(
edge.transport_id,
);
let event = if edge.present {
crate::node::lifecycle::supervisor::Event::ChildPresent { child }
} else {
crate::node::lifecycle::supervisor::Event::ChildAbsent { child }
};
let actions = self.supervisor.fsm.step(event);
for action in actions {
if let crate::node::lifecycle::supervisor::Action::PublishState(
ns,
) = action
{
self.supervisor.state = ns;
}
}
}
// The bound set just changed, so the node's egress MTU
// floor may have. Both directions: an interface that
// binds can be the narrow one, and one that detaches
// can be the reason the clamp was tight.
self.refresh_tun_mss_ceiling();
// A peer reachable only through an interface that has
// gone is not reachable. Withdraw it now rather than
// leaving the liveness reaper to notice up to
// `link_dead_timeout_secs` later, during which this
// node both drops transit traffic in silence and keeps
// advertising reachability it does not have.
//
// Not policy-filtered: whether an interface's absence
// is normal is a statement about node *health*, not
// about whether the routes over it still work.
if !edge.present {
let reaped =
self.reap_peers_on_transport(edge.transport_id).await;
if reaped > 0 {
info!(
transport_id = %edge.transport_id,
peers = reaped,
"Withdrew peers whose interface went away"
);
}
}
}
}
Some(ipv6_packet) = tun_outbound_rx.recv() => {
self.handle_tun_outbound(ipv6_packet).await;
let mut drained = 0;
while drained < 256 {
match tun_outbound_rx.try_recv() {
Ok(p) => {
self.handle_tun_outbound(p).await;
drained += 1;
}
Err(_) => break,
}
}
}
Some(identity) = dns_identity_rx.recv() => {
debug!(
node_addr = %identity.node_addr,
"Registering identity from DNS resolution"
);
self.register_identity(identity.node_addr, identity.pubkey);
}
// A transport medium change (WLAN -> LAN -> 5G, a BLE adapter
// arriving or leaving). Placed after the hot inbound path so
// the `biased` priority of packet processing is unchanged; the
// detector coalesces into a single-slot channel, so this arm
// never sees a burst and needs no drain loop.
Some(change) = netmon_rx.recv() => {
self.handle_net_change(change).await;
}
// Native API datagrams a client wrote to its descriptor. Drained
// in a burst like the TUN arm, for the same reason: one wake-up
// should clear what a client handed over, not one datagram.
Some(out) = native_outbound_rx.recv() => {
self.handle_native_outbound(out.key, out.peer, out.payload).await;
let mut drained = 0;
while drained < 256 {
match native_outbound_rx.try_recv() {
Ok(next) => {
self.handle_native_outbound(next.key, next.peer, next.payload).await;
drained += 1;
}
Err(_) => break,
}
}
}
// Native API registry requests. Placed after the hot inbound
// path so a burst of client registrations cannot delay packet
// processing. No `cfg` here: `tokio::select!` does not accept
// one, so the channel exists on every platform and only the
// listener that feeds it is gated.
Some(message) = native_rx.recv() => {
self.handle_native(message);
}
Some((request, response_tx)) = control_rx.recv() => {
// Only mutating COMMAND requests (`connect` / `disconnect`)
// reach the rx_loop now. Every pure-read `show_*` query is
// served off-loop from the read handle in the control accept
// task (`snapshot_dispatch`), so it never round-trips here —
// the data-plane dispatch path carries no `show_*` arm. A
// `show_*` that somehow arrives (none does) falls through to
// `commands::dispatch`, which returns "unknown command".
let response = commands::dispatch(
self,
&request.command,
request.params.as_ref(),
).await;
let _ = response_tx.send(response);
}
deadline = tick.tick() => {
// Tick-body instrumentation. The gate is read ONCE per tick
// into `instr_on`, which is then passed explicitly to every
// `instr_step!` invocation — macro hygiene makes a call-site
// local invisible inside the macro body. With the
// `profiling` feature off, `gate()` is a `const fn`
// returning false and the macro is a pure pass-through, so
// the whole arm compiles to the uninstrumented sequence.
//
// `tick_entry` records how late this entry is against the
// deadline the interval scheduled it for. That is the
// measurement this instrumentation exists for: the arm is
// polled LAST under `biased;`, so the
// lateness IS the time it spent waiting behind the packet,
// TUN and control arms. `tick()` hands back its scheduled
// deadline, so this is a subtraction rather than a model.
// The whole-tick span below measures the body alone.
let instr_on = crate::instr::gate();
crate::instr::tick_entry(instr_on, deadline.into_std(), std::time::Instant::now());
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::WholeTick, {
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckTimeouts,
self.check_timeouts().await);
// Discard flows the rx_loop announced and a listener's
// own task never took. Cheap: it walks only the pending
// map, which the backlog bounds.
self.native_expire();
let now_ms = Self::now_ms();
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ReloadPeerAcl,
self.reload_peer_acl().await);
// The host map hot-reloads on the same tick as the ACL. It
// is polled separately from `reload_peer_acl` because the
// ACL's embedded alias reloader and this snapshot are
// distinct resources; the `path_mtu_lookup` cache and the
// `nostr_rendezvous` subsystem are deliberately excluded
// from `Reloadable` since neither reloads from a backing
// file (see `node::reloadable`). The `path_mtu_lookup`
// cache is nevertheless swept on this tick, by
// `purge_expired_path_mtu` below: that is expiry, not
// reload.
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ReloadHostMap,
self.reload_host_map().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PollPendingConnects,
self.poll_pending_connects().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PollNostrRendezvous,
self.poll_nostr_rendezvous().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PollLanRendezvous,
self.poll_lan_rendezvous().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::DrivePeerTimers,
self.drive_peer_timers(now_ms).await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ResendPendingRekeys,
self.resend_pending_rekeys(now_ms).await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ResendPendingSessionHandshakes,
self.resend_pending_session_handshakes(now_ms).await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ResendPendingSessionMsg3,
self.resend_pending_session_msg3(now_ms).await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PurgeIdleSessions,
self.purge_idle_sessions(now_ms));
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PurgeExpiredPathMtu,
self.purge_expired_path_mtu(now_ms));
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ProcessPendingRetries,
self.process_pending_retries(now_ms).await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckTreeState,
self.check_tree_state().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckBloomState,
self.check_bloom_state().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ComputeMeshSize,
self.compute_mesh_size());
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::RecordStatsHistory,
self.record_stats_history());
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckMmpReports,
self.check_mmp_reports().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckSessionMmpReports,
self.check_session_mmp_reports().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckLinkHeartbeats,
self.check_link_heartbeats().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckRekey,
self.check_rekey().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckSessionRekey,
self.check_session_rekey().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckPendingLookups,
self.check_pending_lookups(now_ms).await);
// After CheckPendingLookups so a probe's resolve stage
// observes this tick's lookup progress rather than the
// previous tick's.
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PollProbes,
self.poll_probes().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PollTransportDiscovery,
self.poll_transport_discovery().await);
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::SampleTransportCongestion,
self.sample_transport_congestion());
#[cfg(any(target_os = "linux", target_os = "macos"))]
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ActivateConnectedUdpSessions,
self.activate_connected_udp_sessions().await);
// Debug-build sweep of the peer-lifecycle map invariant
// (leaked machines / machine-less legs); two map scans,
// compiled out of release builds.
#[cfg(debug_assertions)]
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::DebugAssertPeerMapsCoherent,
self.debug_assert_peer_maps_coherent());
});
crate::instr::tick_gauges(instr_on, self.peers.len() as u64);
}
// Shutdown signal → enter the bounded drain in place, ONCE.
// Gated on `is_none()` so it only fires while serving; after
// entering drain the arm is disabled (the completed signal is
// never polled again) and the deadline arm below bounds the
// window. Placed after the real arms so their `biased` priority
// is unchanged, and inert while serving with `pending()`.
_ = &mut shutdown, if drain_deadline.is_none() => {
self.enter_drain().await;
drain_deadline =
Some(tokio::time::Instant::now() + self.config().node.drain_timeout());
}
// Bounded drain deadline (drain mode only). Placed LAST so the
// `biased` priority of the normal arms is unchanged, and gated
// on `is_some()` so in normal mode the branch is disabled — the
// future is created but never polled and never fires.
_ = tokio::time::sleep_until(
drain_deadline.unwrap_or_else(tokio::time::Instant::now)
), if drain_deadline.is_some() => {
info!("Drain deadline elapsed, ending drain loop");
break;
}
}
}
info!("RX event loop stopped (channel closed)");
Ok(())
}
/// Process a single received packet.
///
/// Dispatches based on the phase field in the 4-byte common prefix.
///
/// Visible to the rest of `crate::node` so tests can drive a single
/// packet through the dispatch, the same reach `handle_msg1` and
/// `handle_msg2` already have.
pub(in crate::node) async fn process_packet(&mut self, packet: ReceivedPacket) {
if packet.data.len() < COMMON_PREFIX_SIZE {
return; // Drop packets too short for common prefix
}
let prefix = match CommonPrefix::parse(&packet.data) {
Some(p) => p,
None => return, // Malformed prefix
};
if prefix.version != FMP_VERSION {
debug!(
version = prefix.version,
transport_id = %packet.transport_id,
"Unknown FMP version, dropping"
);
// If the packet arrived on an adopted Nostr-NAT bootstrap
// transport, the originating peer is necessarily on a
// different FMP-protocol version than us — the discovery
// sweep would otherwise re-traverse them every cycle even
// though no msg1/msg2 exchange can ever succeed. Bump the
// discovery-layer cooldown to the long protocol-mismatch
// window and emit a single WARN per fresh observation.
if self
.supervisor
.nostr_rendezvous
.is_bootstrap_transport(&packet.transport_id)
&& let Some(npub) = self
.supervisor
.nostr_rendezvous
.bootstrap_transport_npub(&packet.transport_id)
.cloned()
&& let Some(handle) = self.nostr_rendezvous_handle()
{
let now_ms = Self::now_ms();
let cooldown_secs = handle.protocol_mismatch_cooldown_secs();
if handle.record_protocol_mismatch(&npub, now_ms) {
warn!(
peer_npub = %npub,
transport_id = %packet.transport_id,
peer_version = prefix.version,
our_version = FMP_VERSION,
cooldown_secs,
"Nostr-discovered peer speaks a different FMP version; suppressing retraversal"
);
}
}
return;
}
// Drop a frame whose declared payload length disagrees with the
// frame that arrived, before that field can be used as a parsing
// input.
//
// Every transport's packets converge here, but the two families
// reach this line differently. TCP, Tor and Nym read their frame
// boundary out of this same field, so for them the comparison holds
// by construction and never fires. UDP, Ethernet and BLE deliver one
// whole frame per packet, where the arrived length is known exactly
// and nothing compares the two today. A short read on those
// transports is a truncated frame, which fails the AEAD tag or the
// exact-size handshake parse already; this changes which reason it
// is dropped for, not whether it is dropped.
//
// A `None` means the phase carries no fixed relationship and the
// frame is left alone rather than rejected, so an unrecognised phase
// still reaches the dispatch below and is handled there.
if let Some(expected) = expected_payload_len(prefix.phase, packet.data.len())
&& prefix.payload_len != expected
{
debug!(
phase = prefix.phase,
declared = prefix.payload_len,
expected,
len = packet.data.len(),
transport_id = %packet.transport_id,
"FMP payload_len disagrees with frame length, dropping"
);
self.stats_mut()
.record_reject(RejectReason::Transport(TransportReject::PayloadLenMismatch));
return;
}
match prefix.phase {
PHASE_ESTABLISHED => {
self.handle_encrypted_frame(packet).await;
}
PHASE_MSG1 => {
self.handle_msg1(packet).await;
}
PHASE_MSG2 => {
self.handle_msg2(packet).await;
}
_ => {
debug!(
phase = prefix.phase,
transport_id = %packet.transport_id,
"Unknown FMP phase, dropping"
);
}
}
}
}