From 5ab7cd3e89d3f3f70c47272f94c8a01204e392ea Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 2 Oct 2026 15:14:25 +0000 Subject: [PATCH 1/8] fix(testing): hold the chaos final snapshot until every node answers The settle before the final tree snapshot compared only the nodes that answered. A node restored at teardown that had not opened its control socket yet was simply left out of each read, so three reads missing the same node agreed with each other and the snapshot was taken without it. A scenario whose node-count floor equals its node count, such as churn-mixed with ten nodes, then failed its baseline with nine nodes answering although the tenth had come back. A read that any node in the topology did not answer now never counts toward agreement, so the snapshot waits for the restored node and then for three agreeing reads with every node in them. The ninety-second bound is unchanged and running out is still not a failure in itself: the snapshot is taken and the assertions judge it, so a node that never comes back is still reported absent, now after the bound rather than after ten seconds. Running out names the nodes still not answering. The docstring and comments now describe what the settle does, including that the bound is checked after each read and so can be overrun by one interval and one read. --- testing/chaos/sim/runner.py | 57 +++++++++++++++++++++++++------------ 1 file changed, 39 insertions(+), 18 deletions(-) diff --git a/testing/chaos/sim/runner.py b/testing/chaos/sim/runner.py index 7af8660a..957bd828 100644 --- a/testing/chaos/sim/runner.py +++ b/testing/chaos/sim/runner.py @@ -48,8 +48,9 @@ from .veth import VethManager log = logging.getLogger(__name__) -# The final snapshot waits for this many identical consecutive tree reads, -# taken this far apart, so the tree must hold still for two intervals. +# The final snapshot waits until every node answers and this many identical +# consecutive tree reads, taken this far apart, agree, so the tree must hold +# still, with every node in it, for two intervals. SETTLE_READS = 3 SETTLE_INTERVAL_SECS = 5 SETTLE_TIMEOUT_SECS = 90 @@ -744,7 +745,8 @@ class SimRunner: # Take final tree snapshot while nodes are still running, once the # tree has stopped moving. A node restored a moment ago is its own # root until it re-parents, so a snapshot taken straight after the - # restore reads a mesh still converging. + # restore reads a mesh still converging. It may not answer at all + # yet either, and the settle waits for it. self._settle_tree() self._take_snapshot("final") @@ -893,29 +895,39 @@ class SimRunner: return result def _settle_tree(self): - """Wait until consecutive tree reads agree, or the settle time runs out. + """Wait until every node answers and consecutive tree reads agree. - Compares each answering node's root and parent. A fixed delay would - either waste time on a mesh that settled at once or cut off one that - had not. Running out is logged and is not a failure in itself: the - final snapshot is taken anyway, and the assertions judge what it - shows. + Returns once SETTLE_READS reads in a row, SETTLE_INTERVAL_SECS apart, + have each been answered by every node in the topology and show the + same root and parent for each, or once SETTLE_TIMEOUT_SECS has passed. + The bound is checked after each read, so the settle can return up to + one interval and one read past it. A fixed delay would either waste + time on a mesh that settled at once or cut off one that had not. - A read that no node answered never counts toward agreement. It does - not catch a node that stays its own root for longer than the reads - span, which is a tree that is stable and wrong, and is left to the - assertions. + A read that any node did not answer never counts toward agreement. A + node restored at teardown may not have opened its control socket yet, + and a tree that holds still without it is not the tree the final + snapshot is meant to record. + + Running out is logged, naming any node that still does not answer, + and is not a failure in itself: the final snapshot is taken anyway + and the assertions judge what it shows. A node that never comes back + is therefore reported as absent by the assertions that count nodes. + This does not catch a node that answers but stays its own root for + longer than the reads span, which is a tree that is stable and + wrong, and is left to the assertions. """ started = time.time() previous = None agreeing = 0 while not self._interrupted: trees = snapshot_all_trees(self.topology) + missing = sorted(set(self.topology.nodes) - trees.keys()) shape = { nid: (data.get("root"), data.get("parent")) for nid, data in trees.items() } - if not shape: + if missing: agreeing = 0 else: agreeing = agreeing + 1 if shape == previous else 1 @@ -925,10 +937,19 @@ class SimRunner: log.info("Tree settled after %.0fs", waited) return if waited >= SETTLE_TIMEOUT_SECS: - log.warning( - "Tree still changing after %.0fs; taking the final snapshot anyway", - waited, - ) + if missing: + log.warning( + "%s still not answering after %.0fs; " + "taking the final snapshot anyway", + ", ".join(missing), + waited, + ) + else: + log.warning( + "Tree still changing after %.0fs; " + "taking the final snapshot anyway", + waited, + ) return self._sleep(SETTLE_INTERVAL_SECS) From 3c2dfa9d418742220ef08a31673f3358199d385f Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 2 Oct 2026 15:34:52 +0000 Subject: [PATCH 2/8] Drive the first-RTT parent re-evaluation from the spanning-tree core's periodic classification The first-RTT branch of the ReceiverReport handler carried its own copy of the parent-switch / self-root ladder, which already lives in Stp::classify_periodic and drives the periodic re-evaluation. The handler now matches on that classification, so the two paths cannot drift apart. The arm bodies, log lines and metrics are unchanged. The first-RTT path ignores the periodic rebroadcast outcome: the periodic tick is what rebroadcasts an unchanged declaration. Two characterization tests pin the paths no test reached before: a root whose only peer is larger sends no TreeAnnounce on the first RTT sample, and a node whose visible roots are all larger than itself promotes itself to root. --- src/node/handlers/mmp.rs | 155 +++++++++++++++++--------------- src/node/tests/mmp_chartests.rs | 129 +++++++++++++++++++++++++- src/proto/stp/core.rs | 2 + src/proto/stp/mod.rs | 6 +- 4 files changed, 220 insertions(+), 72 deletions(-) diff --git a/src/node/handlers/mmp.rs b/src/node/handlers/mmp.rs index 43c0af8a..f6f6c723 100644 --- a/src/node/handlers/mmp.rs +++ b/src/node/handlers/mmp.rs @@ -15,7 +15,7 @@ use crate::proto::mmp::{ LinkReportKind, LinkReportSnapshot, MmpAction, PeerLivenessSnapshot, ReceiverReport, RrLog, SenderReport, }; -use crate::proto::stp::ParentEval; +use crate::proto::stp::{Stp, TreeDecision}; use crate::transport::{TransportAddr, TransportId}; use std::time::{Duration, Instant}; use tracing::{debug, info, trace, warn}; @@ -256,78 +256,93 @@ impl Node { // Compute the flap-dampening / hold-down veto at the edge; a mandatory // switch bypasses it, a discretionary one is taken only if not suppressed. let switch_suppressed = self.tree_state.is_switch_suppressed(mono_now_ms); - let new_parent = match self - .tree_state - .evaluate_parent(&peer_costs, &std::collections::BTreeSet::new()) - { - ParentEval::Mandatory(p) => Some(p), - ParentEval::Discretionary(p) if !switch_suppressed => Some(p), - ParentEval::Discretionary(_) | ParentEval::None => None, - }; - if let Some(new_parent) = new_parent { - let new_seq = self.tree_state.my_declaration().sequence() + 1; - let flap_dampened = - self.tree_state - .set_parent(new_parent, new_seq, now_secs, mono_now_ms); - self.tree_state.recompute_coords(); - // Clone identity once: sign_declaration borrows &mut tree_state while - // the identity() accessor borrows all of &self, so an owned copy avoids - // the split-borrow conflict on this infrequent parent-switch path. - let our_identity = self.identity().clone(); - if let Err(e) = - sign_declaration(self.tree_state.my_declaration_mut(), &our_identity) - { - warn!(error = %e, "Failed to sign declaration after first-RTT parent eval"); - self.metrics() - .tree - .record_reject(TreeReject::OutboundSignFailed); - return; + match Stp::classify_periodic( + &self.tree_state, + &peer_costs, + &std::collections::BTreeSet::new(), + switch_suppressed, + ) { + TreeDecision::Switch { + new_parent, + new_seq, + } => { + let flap_dampened = + self.tree_state + .set_parent(new_parent, new_seq, now_secs, mono_now_ms); + self.tree_state.recompute_coords(); + // Clone identity once: sign_declaration borrows &mut tree_state while + // the identity() accessor borrows all of &self, so an owned copy avoids + // the split-borrow conflict on this infrequent parent-switch path. + let our_identity = self.identity().clone(); + if let Err(e) = + sign_declaration(self.tree_state.my_declaration_mut(), &our_identity) + { + warn!(error = %e, "Failed to sign declaration after first-RTT parent eval"); + self.metrics() + .tree + .record_reject(TreeReject::OutboundSignFailed); + return; + } + // Surgical invalidation — see CoordCache::invalidate_via_node doc. + self.coord_cache + .invalidate_via_node(our_identity.node_addr()); + self.reset_lookup_backoff(); + self.metrics().tree.parent_switches.inc(); + info!( + new_parent = %self.peer_display_name(&new_parent), + new_seq = new_seq, + new_root = %self.tree_state.root(), + depth = self.tree_state.my_coords().depth(), + trigger = "first-rtt", + "Parent switched after first RTT measurement" + ); + if flap_dampened { + self.note_flap("first-rtt"); + } + self.send_tree_announce_to_all().await; + let all_peers: Vec = self.peers.keys().copied().collect(); + self.bloom_state.mark_all_updates_needed(all_peers); } - // Surgical invalidation — see CoordCache::invalidate_via_node doc. - self.coord_cache - .invalidate_via_node(our_identity.node_addr()); - self.reset_lookup_backoff(); - self.metrics().tree.parent_switches.inc(); - info!( - new_parent = %self.peer_display_name(&new_parent), - new_seq = new_seq, - new_root = %self.tree_state.root(), - depth = self.tree_state.my_coords().depth(), - trigger = "first-rtt", - "Parent switched after first RTT measurement" - ); - if flap_dampened { - self.note_flap("first-rtt"); + TreeDecision::SelfRoot => { + self.tree_state.become_root(now_secs); + // Clone identity once (see the parent-switch branch above for why). + let our_identity = self.identity().clone(); + if let Err(e) = + sign_declaration(self.tree_state.my_declaration_mut(), &our_identity) + { + warn!(error = %e, "Failed to sign self-root declaration after first-RTT"); + self.metrics() + .tree + .record_reject(TreeReject::OutboundSignFailed); + return; + } + // Surgical invalidation — see CoordCache::invalidate_other_roots doc. + self.coord_cache + .invalidate_other_roots(our_identity.node_addr()); + self.reset_lookup_backoff(); + self.metrics().tree.parent_switches.inc(); + info!( + new_root = %self.tree_state.root(), + trigger = "first-rtt", + "Self-promoted to root after first RTT: smallest visible NodeAddr" + ); + self.send_tree_announce_to_all().await; + let all_peers: Vec = self.peers.keys().copied().collect(); + self.bloom_state.mark_all_updates_needed(all_peers); } - self.send_tree_announce_to_all().await; - let all_peers: Vec = self.peers.keys().copied().collect(); - self.bloom_state.mark_all_updates_needed(all_peers); - } else if !self.tree_state.is_root() && self.tree_state.should_be_root() { - self.tree_state.become_root(now_secs); - // Clone identity once (see the parent-switch branch above for why). - let our_identity = self.identity().clone(); - if let Err(e) = - sign_declaration(self.tree_state.my_declaration_mut(), &our_identity) - { - warn!(error = %e, "Failed to sign self-root declaration after first-RTT"); - self.metrics() - .tree - .record_reject(TreeReject::OutboundSignFailed); - return; + // Nothing changed. The periodic tick rebroadcasts; this path does not. + TreeDecision::PeriodicRebroadcast => {} + // classify_periodic never yields these: there is no announcing + // peer, so the loop-drop / ancestry-update arms cannot arise, and + // ParentLost is the removal drive's outcome. + TreeDecision::LoopDrop + | TreeDecision::AncestryUpdate { .. } + | TreeDecision::ParentLost + | TreeDecision::NoChange => { + unreachable!( + "classify_periodic yields only Switch / SelfRoot / PeriodicRebroadcast" + ) } - // Surgical invalidation — see CoordCache::invalidate_other_roots doc. - self.coord_cache - .invalidate_other_roots(our_identity.node_addr()); - self.reset_lookup_backoff(); - self.metrics().tree.parent_switches.inc(); - info!( - new_root = %self.tree_state.root(), - trigger = "first-rtt", - "Self-promoted to root after first RTT: smallest visible NodeAddr" - ); - self.send_tree_announce_to_all().await; - let all_peers: Vec = self.peers.keys().copied().collect(); - self.bloom_state.mark_all_updates_needed(all_peers); } } } diff --git a/src/node/tests/mmp_chartests.rs b/src/node/tests/mmp_chartests.rs index efd8db80..99beeb68 100644 --- a/src/node/tests/mmp_chartests.rs +++ b/src/node/tests/mmp_chartests.rs @@ -33,7 +33,7 @@ use crate::node::session::{EndToEndState, SessionEntry}; use crate::noise::HandshakeState; use crate::peer::ActivePeer; use crate::proto::mmp::{MmpMode, ReceiverReport}; -use crate::proto::stp::{ParentDeclaration, TreeCoordinate}; +use crate::proto::stp::{CoordEntry, ParentDeclaration, TreeCoordinate}; // =========================================================================== // Helpers @@ -437,3 +437,130 @@ async fn non_first_receiver_report_does_not_retrigger_tree() { "a non-first ReceiverReport does not re-enter the first-RTT tree branch" ); } + +/// Insert a peer whose NodeAddr is strictly larger than the node's own, with +/// link MMP but no RTT yet, aged so a crafted ReceiverReport yields a first +/// RTT sample. The caller registers the peer's tree position. +fn insert_larger_unmeasured_peer(node: &mut Node) -> NodeAddr { + let my_addr = *node.node_addr(); + let (identity, addr) = loop { + let id = make_peer_identity(); + let a = *id.node_addr(); + if a > my_addr { + break (id, a); + } + }; + let mut peer = ActivePeer::new(identity, LinkId::new(1), 0); + peer.test_init_mmp(MmpMode::Full); + peer.test_backdate_session_start(std::time::Duration::from_secs(10)); + node.peers.insert(addr, peer); + addr +} + +/// Sum of the tree-announce fan-out counters. Every attempt to send an +/// announce to a peer moves exactly one of them. +fn tree_announce_attempts(node: &Node) -> u64 { + let tree = &node.metrics().tree; + tree.sent.get() + tree.send_failed.get() + tree.rate_limited.get() +} + +/// A root node whose only peer has a larger address has nothing to change on +/// the first RTT sample: it stays root, records no switch, and sends no +/// TreeAnnounce. The periodic tick rebroadcasts; the first-RTT path does not. +#[tokio::test] +async fn first_rtt_on_a_root_with_only_larger_peers_sends_no_tree_announce() { + let mut node = make_node(); + let addr = insert_larger_unmeasured_peer(&mut node); + node.tree_state_mut().update_peer( + ParentDeclaration::self_root(addr, 1, 0), + TreeCoordinate::root(addr), + ); + + assert!( + node.tree_state().is_root(), + "precondition: node starts as its own root" + ); + let switches_before = node.metrics().tree.parent_switches.get(); + let attempts_before = tree_announce_attempts(&node); + + node.handle_receiver_report(&addr, &craft_rr_payload(10, 5, 500)) + .await; + + assert!( + node.get_peer(&addr).unwrap().has_srtt(), + "precondition: this report was the peer's first RTT sample" + ); + assert!(node.tree_state().is_root(), "node remains its own root"); + assert_eq!( + node.metrics().tree.parent_switches.get(), + switches_before, + "no parent switch is recorded" + ); + assert_eq!( + tree_announce_attempts(&node), + attempts_before, + "the first-RTT path does not rebroadcast an unchanged declaration" + ); +} + +/// A node holding a parent whose tree has since re-rooted at a larger address +/// than its own promotes itself to root on the first RTT sample. +#[tokio::test] +async fn first_rtt_self_promotes_when_no_visible_root_is_smaller() { + let mut node = make_node(); + let addr = insert_larger_unmeasured_peer(&mut node); + + // The peer first sits under a root smaller than us, and we take it as + // parent, so our root is that smaller node. + let far_root = NodeAddr::from_bytes([0u8; 16]); + assert!( + far_root < *node.node_addr(), + "precondition: the far root is smaller than the node" + ); + node.tree_state_mut().update_peer( + ParentDeclaration::new(addr, far_root, 1, 0), + TreeCoordinate::new(vec![ + CoordEntry::new(addr, 1, 0), + CoordEntry::new(far_root, 1, 0), + ]) + .unwrap(), + ); + let seq = node.tree_state().my_declaration().sequence() + 1; + node.tree_state_mut() + .set_parent(addr, seq, 0, crate::time::mono_ms()); + node.tree_state_mut().recompute_coords(); + assert_eq!( + node.tree_state().root(), + &far_root, + "precondition: the node sits under the far root" + ); + + // The peer then re-roots at itself, larger than us: no visible root is + // smaller than the node any more. + node.tree_state_mut().update_peer( + ParentDeclaration::self_root(addr, 2, 0), + TreeCoordinate::root(addr), + ); + assert!( + !node.tree_state().is_root() && node.tree_state().should_be_root(), + "precondition: not root, but should be" + ); + let switches_before = node.metrics().tree.parent_switches.get(); + + node.handle_receiver_report(&addr, &craft_rr_payload(10, 5, 500)) + .await; + + assert!( + node.get_peer(&addr).unwrap().has_srtt(), + "precondition: this report was the peer's first RTT sample" + ); + assert!( + node.tree_state().is_root(), + "the first-RTT path promoted the node to root" + ); + assert_eq!( + node.metrics().tree.parent_switches.get(), + switches_before + 1, + "the self-promotion is recorded as one parent switch" + ); +} diff --git a/src/proto/stp/core.rs b/src/proto/stp/core.rs index 8181e94f..04e2fffe 100644 --- a/src/proto/stp/core.rs +++ b/src/proto/stp/core.rs @@ -165,6 +165,8 @@ impl Stp { /// `classify_announce`, the periodic path has no same-parent loop-drop / /// ancestry-update arms — a periodic tick has no announcing peer, so those cases /// never arise; the no-change tail is a re-broadcast rather than a true no-op. + /// The first-RTT re-evaluation in `node::handlers::mmp` is driven by it too, + /// and ignores `PeriodicRebroadcast`. pub(crate) fn classify_periodic( tree: &TreeState, peer_costs: &BTreeMap, diff --git a/src/proto/stp/mod.rs b/src/proto/stp/mod.rs index 2c679e24..677bfd96 100644 --- a/src/proto/stp/mod.rs +++ b/src/proto/stp/mod.rs @@ -34,7 +34,11 @@ pub use crate::proto::coord::{CoordEntry, CoordError, TreeCoordinate}; pub(crate) use crate::proto::coord::{ coords_wire_size, decode_coords, decode_optional_coords, encode_coords, encode_empty_coords, }; -pub(crate) use core::{ParentEval, Stp, TreeDecision}; +// Callers outside this module take the decision from `Stp`; only the tests +// name the parent evaluation itself. +#[cfg(test)] +pub(crate) use core::ParentEval; +pub(crate) use core::{Stp, TreeDecision}; pub use declaration::ParentDeclaration; pub use state::TreeState; pub use wire::TreeAnnounce; From e431390381daa413781b047b1213179dd587e818 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 2 Oct 2026 15:41:47 +0000 Subject: [PATCH 3/8] Decide the discovery path-MTU write with the shared keep-tighter rule The LookupResponse handler open-coded the keep-tighter comparison that should_apply_path_mtu already holds for the session carriers. It now calls that function, so the discovery write and the session tightens follow one rule. The equal-value case still keeps the stored entry and its learn time, so a later answer of the same value cannot extend the entry's deadline. No handler test reached that equal-value case: a replayed response is dropped as unsolicited before it gets to the write, and the keep-tighter test uses a strictly looser value. A new test answers two genuine lookups for one target with the same path MTU and checks the entry is left exactly as it was. The replay test's comment and the handler's comment on the keep-tighter arm are corrected to say what that arm bounds and what actually stops a replay, and the shared function's documentation now says it writes only a strictly tighter value, which is what it does. --- src/node/handlers/lookup.rs | 20 ++++++++----- src/node/tests/discovery.rs | 60 ++++++++++++++++++++++++++++++++++--- src/proto/fsp/core.rs | 11 +++---- src/proto/fsp/mod.rs | 2 +- 4 files changed, 76 insertions(+), 17 deletions(-) diff --git a/src/node/handlers/lookup.rs b/src/node/handlers/lookup.rs index 88cae3c0..3f783115 100644 --- a/src/node/handlers/lookup.rs +++ b/src/node/handlers/lookup.rs @@ -7,6 +7,7 @@ use crate::node::Node; use crate::node::reject::DiscoveryReject; +use crate::proto::fsp::should_apply_path_mtu; use crate::proto::lookup::{ LookupAction, LookupRequest, LookupResponse, MAX_RECENT_LOOKUP_REQUESTS, }; @@ -416,7 +417,9 @@ impl Node { let fips_addr = crate::FipsAddress::from_node_addr(&target); match self.path_mtu_lookup.write() { Ok(mut map) => match map.get(&fips_addr).copied() { - Some(existing) if existing.mtu <= path_mtu => { + Some(existing) + if !should_apply_path_mtu(Some(existing.mtu), path_mtu) => + { // Keep the tighter learned value; never loosen // the clamp. A reactive MtuExceeded or // PathMtuNotification tighten takes precedence @@ -424,12 +427,15 @@ impl Node { // (cross-carrier keep-tighter). // // This arm deliberately leaves `learned_ms` - // alone. That is what bounds a replayed - // response: the replay of a value already - // stored takes this arm, so the entry still - // expires at first-write plus the TTL rather - // than being pushed out again on every - // injection. Refreshing the stamp here would + // alone. A later answered lookup that reports + // the value already stored takes this arm, so + // the entry still expires at first-write plus + // the TTL rather than being pushed out again + // by every answer of the same value. (A + // replayed response never gets here: the + // pending lookup is gone once the first answer + // is accepted, so the copy is dropped as + // unsolicited.) Refreshing the stamp here would // read as a tidy-up and would silently restore // indefinite pinning. debug!( diff --git a/src/node/tests/discovery.rs b/src/node/tests/discovery.rs index 2bf4b64c..3b29b4a0 100644 --- a/src/node/tests/discovery.rs +++ b/src/node/tests/discovery.rs @@ -1653,10 +1653,11 @@ async fn test_lookup_response_path_mtu_expires_without_a_session() { #[tokio::test] async fn test_replayed_lookup_response_does_not_extend_the_path_mtu_deadline() { - // The response carries no replay dedupe, so a captured one can be - // re-injected indefinitely. What bounds the damage is that a replay of a - // value already stored takes the keep-tighter arm, which does not touch - // the learn time: each injection buys one TTL, not one per packet. + // The response carries no replay dedupe of its own, so a captured one can + // be re-injected indefinitely. Accepting the first response clears the + // pending lookup, so each replay is dropped as unsolicited before it + // reaches the path-MTU write: each injection buys one TTL, not one per + // packet. The equal-value arm of that write is pinned by the next test. let mut node = make_node(); let from = make_node_addr(0xAA); @@ -1694,6 +1695,57 @@ async fn test_replayed_lookup_response_does_not_extend_the_path_mtu_deadline() { ); } +#[tokio::test] +async fn test_a_later_solicited_response_of_the_same_path_mtu_keeps_the_learn_time() { + // Two genuine lookups for one target, answered with the same path_mtu. + // The second answer is solicited, so it reaches the path-MTU write, and an + // equal value must keep the stored entry, learn time included. Refreshing + // the stamp on equality would let every answer of the same value push the + // deadline out again. + let mut node = make_node(); + let from = make_node_addr(0xAA); + + let target_identity = Identity::generate(); + let target = *target_identity.node_addr(); + let target_fips = crate::FipsAddress::from_node_addr(&target); + let root = make_node_addr(0xF0); + let coords = TreeCoordinate::from_addrs(vec![target, root]).unwrap(); + node.register_identity(target, target_identity.pubkey_full()); + + let answer = |request_id: u64| { + let proof = + target_identity.sign(&LookupResponse::proof_bytes(request_id, &target, &coords)); + let mut response = LookupResponse::new(request_id, target, coords.clone(), proof); + response.path_mtu = 1300; + response.encode()[1..].to_vec() + }; + + seed_pending_lookup(&mut node, target, 805); + node.handle_lookup_response(&from, &answer(805)).await; + let first = node + .path_mtu_lookup_entry(&target_fips) + .expect("precondition: the first response wrote an entry"); + assert!( + first.learned_ms.is_some(), + "precondition: the entry carries a learn time" + ); + + // Real elapsed wall-clock, so a refreshed stamp would differ. + std::thread::sleep(std::time::Duration::from_millis(5)); + seed_pending_lookup(&mut node, target, 806); + node.handle_lookup_response(&from, &answer(806)).await; + + assert!( + !node.lookup.pending_lookups.contains_key(&target), + "precondition: the second response was accepted as solicited" + ); + assert_eq!( + node.path_mtu_lookup_entry(&target_fips), + Some(first), + "an equal path_mtu must leave the entry exactly as it was, learn time included" + ); +} + // ============================================================================ // Open-Discovery Sweep — cache-injection unit test // ============================================================================ diff --git a/src/proto/fsp/core.rs b/src/proto/fsp/core.rs index f45526a0..8a3b25c4 100644 --- a/src/proto/fsp/core.rs +++ b/src/proto/fsp/core.rs @@ -483,10 +483,10 @@ impl Fsp { } /// Decide whether a path-MTU update should tighten the shared lookup: emit - /// `TightenPathMtuLookup` only when `candidate` is at least as tight as the - /// `existing` value (keep-tighter, never loosen). The `existing` read and - /// the applied write are performed shell-side under one `path_mtu_lookup` - /// write guard, so the decision stays atomic. + /// `TightenPathMtuLookup` only when there is no `existing` value or + /// `candidate` is strictly tighter than it (keep-tighter, never loosen). + /// The `existing` read and the applied write are performed shell-side + /// under one `path_mtu_lookup` write guard, so the decision stays atomic. pub(crate) fn plan_path_mtu_tighten( &self, fips_addr: FipsAddress, @@ -524,7 +524,8 @@ pub(crate) fn initiation_winner(our_node_addr: &NodeAddr, their_node_addr: &Node /// Decide whether a path-MTU update should be applied to the shared /// `FipsAddress`-keyed lookup: keep the tighter of existing-or-candidate, never /// loosen. Returns `true` when `candidate` should be written (there is no -/// existing value, or the candidate is at least as tight). +/// existing value, or the candidate is strictly tighter). An equal candidate is +/// not written, so the stored entry, and any learn time it carries, is kept. pub(crate) fn should_apply_path_mtu(existing: Option, candidate: u16) -> bool { !matches!(existing, Some(existing) if existing <= candidate) } diff --git a/src/proto/fsp/mod.rs b/src/proto/fsp/mod.rs index 7c9b37f5..15df2ff5 100644 --- a/src/proto/fsp/mod.rs +++ b/src/proto/fsp/mod.rs @@ -33,7 +33,7 @@ mod tests; pub(crate) use core::{ DecryptSlot, EpochReaction, Fsp, FspAction, InitialMsg3ResendSnapshot, RekeyCfg, RekeyMsg3ResendSnapshot, SessionSnapshot, cutover_timer_elapsed, initiation_winner, - mark_ipv6_ecn_ce, push_bounded_pending, + mark_ipv6_ecn_ce, push_bounded_pending, should_apply_path_mtu, }; pub use wire::{ FspInnerFlags, SessionAck, SessionFlags, SessionMessageType, SessionMsg3, SessionSetup, From 1d8c22a91b24c94d53424066be2feef8e18a9d2a Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 2 Oct 2026 16:44:05 +0000 Subject: [PATCH 4/8] Share one hop-limit rule between the forwarding pre-check and the routing core The forwarding handler resolves a next hop only for datagrams the routing core can forward, which keeps the coordinate-cache touch in that resolution scoped to genuine forwards. Its TTL test (ttl > 1) restated the core's drop (decrement, then drop at zero) in a different form, and no test could see the two disagree: reverting the pre-check to its older ttl != 0 form passed every unit test. ttl_after_hop now holds the rule. The routing core drops on it, SessionDatagram::can_forward uses it, and a new SessionDatagramRef::can_forward is what the handler calls. A unit test pins the function, a routing test checks the core and can_forward agree for every TTL, and a handler test checks that a last-hop transit datagram does not refresh the destination's cached coordinates while a ttl=2 one does. --- src/node/dataplane/forwarding.rs | 13 ++++---- src/node/tests/forwarding.rs | 47 +++++++++++++++++++++++++++++ src/proto/link.rs | 52 +++++++++++++++++++++++++++++++- src/proto/routing/core.rs | 16 +++++----- src/proto/routing/tests/core.rs | 44 ++++++++++++++++++++++++++- 5 files changed, 155 insertions(+), 17 deletions(-) diff --git a/src/node/dataplane/forwarding.rs b/src/node/dataplane/forwarding.rs index 76164e66..2af800f7 100644 --- a/src/node/dataplane/forwarding.rs +++ b/src/node/dataplane/forwarding.rs @@ -55,13 +55,12 @@ impl Node { self.try_warm_coord_cache_ref(&datagram_ref, payload.len()); // Pre-resolve the next hop only for datagrams the core can actually - // forward: not locally destined, and carrying a TTL that survives the - // decrement (`ttl > 1` — the shell-side mirror of the core's - // would-leave-zero drop). This keeps `find_next_hop`'s coord-cache - // LRU-touch side effect scoped to genuine forwards, as it was when the - // TTL test ran inline ahead of it. Warming above has already run, so - // the resolution observes freshly cached coords. - let next_hop = if datagram_ref.dest_addr != my_addr && datagram_ref.ttl > 1 { + // forward: not locally destined, and passing `can_forward`, which is + // the core's own hop-limit rule. This keeps `find_next_hop`'s + // coord-cache LRU-touch side effect scoped to genuine forwards, as it + // was when the TTL test ran inline ahead of it. Warming above has + // already run, so the resolution observes freshly cached coords. + let next_hop = if datagram_ref.dest_addr != my_addr && datagram_ref.can_forward() { self.resolve_next_hop(&datagram_ref.dest_addr) } else { None diff --git a/src/node/tests/forwarding.rs b/src/node/tests/forwarding.rs index c092e2d8..80023fc9 100644 --- a/src/node/tests/forwarding.rs +++ b/src/node/tests/forwarding.rs @@ -155,6 +155,53 @@ async fn test_forwarding_ttl_two_transit_clears_the_gate() { ); } +/// The next hop is resolved only for a datagram the core can forward, and that +/// resolution refreshes the destination's cached coordinates. A last-hop +/// transit datagram (ttl=1) is dropped by the core, so it must not refresh +/// them; a ttl=2 datagram reaches the resolution and does. +#[tokio::test] +async fn test_forwarding_last_hop_transit_does_not_refresh_destination_coords() { + let mut node = make_node(); + let from = make_node_addr(0xAA); + let src = make_node_addr(0x01); + let dest = make_node_addr(0x02); + let root = make_node_addr(0xF0); + let coords = TreeCoordinate::from_addrs(vec![dest, root]).unwrap(); + + let now_ms = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis() as u64; + let stamped_ms = now_ms - node.coord_cache().default_ttl_ms() / 2; + node.coord_cache_mut() + .insert_verified(dest, coords, stamped_ms); + let last_used = |node: &Node| node.coord_cache().get_entry(&dest).unwrap().last_used(); + assert_eq!( + last_used(&node), + stamped_ms, + "precondition: the entry carries the past stamp" + ); + + for ttl in [1u8, 2] { + let dg = SessionDatagram::new(src, dest, vec![0x10, 0x00, 0x00, 0x00]).with_ttl(ttl); + let encoded = dg.encode(); + node.handle_session_datagram(&from, &encoded[1..], false) + .await; + if ttl == 1 { + assert_eq!( + last_used(&node), + stamped_ms, + "a transit ttl=1 datagram is dropped, so it must not refresh the destination's coords" + ); + } else { + assert!( + last_used(&node) > stamped_ms, + "a transit ttl=2 datagram reaches next-hop resolution, which refreshes the coords" + ); + } + } +} + // --- Local delivery --- #[tokio::test] diff --git a/src/proto/link.rs b/src/proto/link.rs index ebc4d26b..7c982fa4 100644 --- a/src/proto/link.rs +++ b/src/proto/link.rs @@ -144,6 +144,20 @@ pub struct SessionDatagramRef<'a> { pub payload: &'a [u8], } +/// The TTL a transit datagram leaves this node with, or `None` when it may not +/// be transmitted because it would leave with zero. +/// +/// Follows IP semantics: the decrement comes first, and `saturating_sub` folds +/// an already-exhausted arrival (TTL 0) into the same outcome as a last-hop +/// arrival (TTL 1). This is the one rule behind the routing core's hop-limit +/// drop and both `can_forward`s. +pub(crate) fn ttl_after_hop(ttl: u8) -> Option { + match ttl.saturating_sub(1) { + 0 => None, + left => Some(left), + } +} + /// SessionDatagram fixed header size: msg_type(1) + ttl(1) + path_mtu(2) + src_addr(16) + dest_addr(16). pub const SESSION_DATAGRAM_HEADER_SIZE: usize = 36; @@ -191,7 +205,7 @@ impl SessionDatagram { /// True only at TTL 2 or more: at TTL 1 the decrement leaves zero, so the /// datagram is dropped rather than forwarded. pub fn can_forward(&self) -> bool { - self.ttl > 1 + ttl_after_hop(self.ttl).is_some() } /// Encode as link-layer message (msg_type + ttl + path_mtu + src_addr + dest_addr + payload). @@ -239,6 +253,12 @@ impl<'a> SessionDatagramRef<'a> { }) } + /// Check whether this datagram would survive a transit hop, by the same + /// rule the routing core drops on (`ttl_after_hop`). + pub fn can_forward(&self) -> bool { + ttl_after_hop(self.ttl).is_some() + } + /// Materialize an owned datagram for forwarding/re-encoding paths. pub fn into_owned(self) -> SessionDatagram { SessionDatagram { @@ -424,6 +444,36 @@ mod tests { assert!(dg.with_ttl(255).can_forward()); } + #[test] + fn ttl_after_hop_drops_only_what_would_leave_at_zero() { + assert_eq!( + ttl_after_hop(0), + None, + "an exhausted arrival leaves at zero" + ); + assert_eq!(ttl_after_hop(1), None, "a last-hop arrival leaves at zero"); + assert_eq!(ttl_after_hop(2), Some(1)); + assert_eq!(ttl_after_hop(255), Some(254)); + + let dg = SessionDatagram::new(make_node_addr(1), make_node_addr(2), vec![0x42]); + for ttl in 0..=u8::MAX { + assert_eq!( + ttl_after_hop(ttl).is_some(), + ttl > 1, + "ttl={ttl}: only a TTL of 2 or more survives the hop" + ); + let owned = dg.clone().with_ttl(ttl); + let encoded = owned.encode(); + let view = SessionDatagramRef::decode(&encoded[1..]).unwrap(); + assert_eq!( + view.can_forward(), + owned.can_forward(), + "ttl={ttl}: the borrowed and owned views must agree" + ); + assert_eq!(view.can_forward(), ttl_after_hop(ttl).is_some()); + } + } + #[test] fn test_session_datagram_decrement_ttl() { let base = SessionDatagram::new(make_node_addr(1), make_node_addr(2), vec![0x42]); diff --git a/src/proto/routing/core.rs b/src/proto/routing/core.rs index e73caa72..44215a25 100644 --- a/src/proto/routing/core.rs +++ b/src/proto/routing/core.rs @@ -17,7 +17,7 @@ use super::limits::LimitVerdict; use super::state::Router; use super::wire::{CoordsRequired, MtuExceeded, PathBroken}; -use crate::proto::link::{SessionDatagram, SessionDatagramRef}; +use crate::proto::link::{SessionDatagram, SessionDatagramRef, ttl_after_hop}; use crate::{NodeAddr, TreeCoordinate}; /// Read-only view of routing state the routing core needs. @@ -112,7 +112,8 @@ impl Router { /// datagram that would leave with a TTL of zero is not transmitted. /// /// The shell pre-resolves `next_hop` only for datagrams this can actually - /// forward (dest not local and TTL surviving the decrement), so + /// forward (dest not local, and `SessionDatagramRef::can_forward`, which + /// applies the same [`ttl_after_hop`] rule this drops on), so /// `find_next_hop`'s LRU-touch side effect stays scoped to genuine /// forwards. `route` still re-checks local delivery and the TTL /// authoritatively. @@ -132,15 +133,14 @@ impl Router { } // TTL enforcement on the transit path: decrement first, then drop if - // the datagram would leave with a TTL of zero. `saturating_sub` folds - // the already-exhausted arrival (ttl=0) into the same test as the - // last-hop arrival (ttl=1); neither is transmitted. - let forwarded_ttl = dg.ttl.saturating_sub(1); - if forwarded_ttl == 0 { + // the datagram would leave with a TTL of zero. The already-exhausted + // arrival (ttl=0) and the last-hop arrival (ttl=1) are both dropped; + // neither is transmitted. + let Some(forwarded_ttl) = ttl_after_hop(dg.ttl) else { return RouteOutcome::Drop { reason: DropReason::TtlExhausted, }; - } + }; let nh = match next_hop { Some(nh) => nh, diff --git a/src/proto/routing/tests/core.rs b/src/proto/routing/tests/core.rs index 75d25bab..7c9f3174 100644 --- a/src/proto/routing/tests/core.rs +++ b/src/proto/routing/tests/core.rs @@ -1,7 +1,7 @@ //! Tests for the sans-IO routing decision core. use super::util::{MockPeer, MockRoutingView, make_coords, make_datagram_ref, make_next_hop}; -use crate::proto::link::SessionDatagramRef; +use crate::proto::link::{SessionDatagramRef, ttl_after_hop}; use crate::proto::routing::RoutingSignalType; use crate::proto::routing::{ DropReason, LimitVerdict, RouteAction, RouteOutcome, Router, RoutingView, select_best_candidate, @@ -514,3 +514,45 @@ fn synth_mtu_exceeded_rate_limit_gate_suppresses_second_call() { assert!(third.action.is_some()); assert_eq!(third.verdict, LimitVerdict::Admit); } + +/// The routing core's hop-limit drop and the shell's `can_forward` pre-check +/// agree for every TTL: a transit datagram with a next hop is dropped as +/// TTL-exhausted exactly when `can_forward` is false, and otherwise leaves +/// with the TTL `ttl_after_hop` gives. +#[test] +fn route_drops_for_hop_limit_exactly_when_can_forward_is_false() { + let my_addr = make_node_addr(0x10); + let nh_addr = make_node_addr(0x30); + let rv = MockRoutingView::new(false); + for ttl in 0..=u8::MAX { + let mut router = Router::new(); + let dg = make_datagram_ref(ttl, make_node_addr(0x20)); + let out = router.route( + &dg, + &my_addr, + false, + Some(make_next_hop(nh_addr, 1400)), + &rv, + ); + match out { + RouteOutcome::Drop { + reason: DropReason::TtlExhausted, + } => assert!( + !dg.can_forward(), + "ttl={ttl}: the core dropped a datagram the pre-check would forward" + ), + RouteOutcome::Forward { bytes, .. } => { + assert!( + dg.can_forward(), + "ttl={ttl}: the core forwarded a datagram the pre-check would not" + ); + assert_eq!( + Some(decode_forward(&bytes).ttl), + ttl_after_hop(ttl), + "ttl={ttl}: the forwarded TTL must be the shared rule's" + ); + } + _ => panic!("ttl={ttl}: expected Drop(TtlExhausted) or Forward"), + } + } +} From 3ca40ac5df02f7e1634dcea03efd0ce6e429b25d Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 2 Oct 2026 15:55:04 +0000 Subject: [PATCH 5/8] Correct the network monitor's description of BLE adapter state The module documentation said a BLE adapter coming or going was among the medium changes this module catches, and that the BLE transport pushes adapter state onto the network-change channel. Neither is true: the detector is the only source of a network change, and nothing in the BLE transport sends on that channel. The text now says so, and names what does exist: on Android the app installs and clears its radio through the slot the node hands it, which only the BLE transport observes. --- src/node/netmon/mod.rs | 21 ++++++++++++--------- 1 file changed, 12 insertions(+), 9 deletions(-) diff --git a/src/node/netmon/mod.rs b/src/node/netmon/mod.rs index b67d8f16..74cec120 100644 --- a/src/node/netmon/mod.rs +++ b/src/node/netmon/mod.rs @@ -1,11 +1,10 @@ //! Transport-medium change detection. //! -//! A node that moves between media (WLAN → LAN, WLAN → 5G, a BLE adapter -//! coming or going) would otherwise learn about it only as *silence*: the peer -//! sits in the table until `node.link_dead_timeout_secs` reaps it, and the -//! reconnect then waits out whatever backoff the old medium had already -//! accumulated. The host kernel knew within milliseconds; the node would find -//! out half a minute later. +//! A node that moves between IP media (WLAN → LAN, WLAN → 5G) would otherwise +//! learn about it only as *silence*: the peer sits in the table until +//! `node.link_dead_timeout_secs` reaps it, and the reconnect then waits out +//! whatever backoff the old medium had already accumulated. The host kernel +//! knew within milliseconds; the node would find out half a minute later. //! //! This module closes that gap. It samples a coarse [`NetFingerprint`] of the //! host's network attachment and publishes a [`NetChange`] on the channel the @@ -110,9 +109,13 @@ //! 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. +//! IP attachment at all, and nothing routes it onto this channel. The detector +//! below is the only source of a [`NetChange`]; the BLE transport publishes +//! none. On Android the embedder-facing half does exist: the app installs and +//! clears its radio through the `BleRadioSlot` returned by +//! `Node::enable_app_owned_ble_radio`, and the BLE transport re-resolves the +//! slot when it changes. That reaches only the BLE transport. The rest of the +//! node sees a radio going away as the loss of the links it carried. //! //! # Where the peer list comes from //! From a9423f801ac6faf8d3ac835a9710a94ece650535 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 2 Oct 2026 15:48:20 +0000 Subject: [PATCH 6/8] Report crypto worker deaths as degraded health instead of hiding them The encrypt and decrypt worker pools discarded their thread handles, so a worker that panicked left the node reporting full health while every packet hashed to that worker was dropped behind a DEBUG line. A worker thread that could not be started panicked start-up instead. The pools now keep each worker's thread handle and report how many are live. A worker that cannot be started is logged and left dead, and the node starts degraded instead of aborting; a pool with no live worker is not installed, so its traffic takes the main-loop path. A sweep on the rx loop tick notices a worker that exits at runtime, logs a warning naming the pool and the live and configured counts, and reports the pool to the supervisor, which publishes Degraded. Losing workers never fails the node: the pools are an offload, and Windows runs without them. A dispatch refused by an exited worker is now counted and logged at WARN, rate-limited, in place of the DEBUG line. --- src/instr/recorder.rs | 15 +- src/node/dataplane/rx_loop.rs | 16 +- src/node/decrypt_worker.rs | 220 ++++++++++++++++---- src/node/encrypt_worker.rs | 338 +++++++++++++++++++++++++------ src/node/lifecycle/mod.rs | 66 +++--- src/node/lifecycle/supervisor.rs | 67 +++++- src/node/lifecycle/workers.rs | 240 ++++++++++++++++++++++ src/node/mod.rs | 2 + src/node/tests/unit.rs | 241 ++++++++++++++++++++++ src/node/worker_set.rs | 245 ++++++++++++++++++++++ 10 files changed, 1304 insertions(+), 146 deletions(-) create mode 100644 src/node/lifecycle/workers.rs create mode 100644 src/node/worker_set.rs diff --git a/src/instr/recorder.rs b/src/instr/recorder.rs index 93bd51b6..98b430d4 100644 --- a/src/instr/recorder.rs +++ b/src/instr/recorder.rs @@ -41,7 +41,7 @@ impl Domain { /// /// `as usize` indexes the counter arrays, so the discriminants are dense and /// `WholeTick` is last (it defines `N_STEPS`). Variants are declared -/// unconditionally — see [`Step::emitted`] for how the two platform- and +/// unconditionally — see [`Step::emitted`] for how the three platform- and /// profile-conditional steps are kept out of the emitted table. #[derive(Copy, Clone, Debug, PartialEq, Eq)] #[repr(usize)] @@ -73,6 +73,7 @@ pub(crate) enum Step { PollTransportDiscovery, SampleTransportCongestion, ActivateConnectedUdpSessions, + PollWorkerLiveness, DebugAssertPeerMapsCoherent, /// The whole tick-arm body, from before `check_timeouts` to after the last /// step. Composes safely with the per-step spans because the macro @@ -112,6 +113,7 @@ pub(crate) const STEPS: [Step; N_STEPS] = [ Step::PollTransportDiscovery, Step::SampleTransportCongestion, Step::ActivateConnectedUdpSessions, + Step::PollWorkerLiveness, Step::DebugAssertPeerMapsCoherent, Step::WholeTick, ]; @@ -146,6 +148,7 @@ impl Step { Step::PollTransportDiscovery => "poll_transport_discovery", Step::SampleTransportCongestion => "sample_transport_congestion", Step::ActivateConnectedUdpSessions => "activate_connected_udp_sessions", + Step::PollWorkerLiveness => "poll_worker_liveness", Step::DebugAssertPeerMapsCoherent => "debug_assert_peer_maps_coherent", Step::WholeTick => "whole_tick", } @@ -153,7 +156,7 @@ impl Step { /// Whether this step gets a row in this build. /// - /// Two steps are conditionally compiled at their call sites. Emitting a row + /// Three steps are conditionally compiled at their call sites. Emitting a row /// for them in a build where the call site does not exist would publish a /// count that is structurally zero forever, which reads as "this step never /// runs" rather than "this step is not in this build". The predicates below @@ -164,6 +167,7 @@ impl Step { Step::ActivateConnectedUdpSessions => { cfg!(any(target_os = "linux", target_os = "macos")) } + Step::PollWorkerLiveness => cfg!(unix), Step::DebugAssertPeerMapsCoherent => cfg!(debug_assertions), _ => true, } @@ -381,12 +385,15 @@ mod tests { #[test] fn emitted_row_count_matches_build() { let emitted = STEPS.iter().filter(|s| s.emitted()).count(); - // 26 unconditional subsystem steps + the whole-tick span, plus the two - // conditionally-compiled steps where this build has them. + // 26 unconditional subsystem steps + the whole-tick span, plus the + // three conditionally-compiled steps where this build has them. let mut expected = 27; if cfg!(any(target_os = "linux", target_os = "macos")) { expected += 1; } + if cfg!(unix) { + expected += 1; + } if cfg!(debug_assertions) { expected += 1; } diff --git a/src/node/dataplane/rx_loop.rs b/src/node/dataplane/rx_loop.rs index bb26d78a..425e7c8d 100644 --- a/src/node/dataplane/rx_loop.rs +++ b/src/node/dataplane/rx_loop.rs @@ -345,17 +345,7 @@ impl Node { // 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; - } - } + self.step_child_exited(child); // 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 @@ -583,6 +573,10 @@ impl Node { #[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); + // Crypto worker threads that exited since the last tick. + #[cfg(unix)] + instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PollWorkerLiveness, + self.poll_worker_liveness()); // Debug-build sweep of the peer-lifecycle map invariant // (leaked machines / machine-less legs); two map scans, // compiled out of release builds. diff --git a/src/node/decrypt_worker.rs b/src/node/decrypt_worker.rs index d4d5173c..63449205 100644 --- a/src/node/decrypt_worker.rs +++ b/src/node/decrypt_worker.rs @@ -36,6 +36,7 @@ #![cfg_attr(not(unix), allow(dead_code))] use crate::NodeAddr; +use crate::node::worker_set::{WorkerLiveness, WorkerSet, worth_logging}; use crate::transport::{TransportAddr, TransportId}; use crossbeam_channel::{Receiver, Sender, TrySendError, bounded}; use portable_atomic::{AtomicU64, Ordering}; @@ -220,26 +221,49 @@ pub(crate) enum WorkerMsg { /// shard. #[derive(Clone)] pub(crate) struct DecryptWorkerPool { - senders: Arc<[Sender]>, + workers: Arc>>, +} + +/// Start the production worker loop on `rx` in a named OS thread. +fn spawn_worker( + idx: usize, + rx: Receiver, +) -> std::io::Result> { + std::thread::Builder::new() + .name(format!("fips-decrypt-{idx}")) + .spawn(move || run_worker(idx, rx)) } impl DecryptWorkerPool { + /// Spawn `n` worker OS threads (at least one). A worker thread that + /// cannot be started is logged and left dead; the caller reads how + /// many started from [`Self::liveness`]. pub fn spawn(n: usize) -> Self { - let n = n.max(1); - let mut senders = Vec::with_capacity(n); - for i in 0..n { - let (tx, rx) = bounded::(WORKER_CHANNEL_CAP); - std::thread::Builder::new() - .name(format!("fips-decrypt-{i}")) - .spawn(move || run_worker(i, rx)) - .expect("failed to spawn fips-decrypt OS thread"); - senders.push(tx); - } + Self::start_with(n, spawn_worker) + } + + /// Build the pool, starting each worker with `spawn`. Production passes + /// [`spawn_worker`]; tests pass workers that fail to start or exit on cue. + fn start_with( + n: usize, + spawn: impl FnMut(usize, Receiver) -> std::io::Result>, + ) -> Self { Self { - senders: senders.into(), + workers: Arc::new(WorkerSet::start( + "decrypt", + n, + || bounded::(WORKER_CHANNEL_CAP), + spawn, + )), } } + /// Whether each worker is still running, and how many dispatches a dead + /// one refused. + pub(crate) fn liveness(&self) -> &dyn WorkerLiveness { + &*self.workers + } + /// Stable hash from session key → worker index. Same hash is used /// for session registration and per-packet dispatch so packets and /// registration arrive at the same shard. @@ -247,24 +271,37 @@ impl DecryptWorkerPool { use std::hash::{Hash, Hasher}; let mut h = std::collections::hash_map::DefaultHasher::new(); cache_key.hash(&mut h); - (h.finish() as usize) % self.senders.len() + (h.finish() as usize) % self.workers.len() + } + + /// Count a message refused by the exited worker `idx`, and log it at a + /// bounded rate. + fn note_refused(&self, idx: usize, what: &'static str) { + let n = self.workers.note_refused(); + if worth_logging(n) { + warn!( + pool = "decrypt", + worker = idx, + refused = n + 1, + what, + "Decrypt worker has exited; message refused" + ); + } } /// Dispatch a per-packet decrypt job. Drops if the per-worker /// channel is full (sustained rate overrun); the rx_loop's drain /// caps inbound at the same scale upstream so the cliff is - /// bounded. + /// bounded. A job for an exited worker is dropped, counted and + /// logged at WARN. pub fn dispatch_job(&self, job: DecryptJob) { - if self.senders.is_empty() { - return; - } let idx = self.worker_idx_for(job.cache_key); - match self.senders[idx].try_send(WorkerMsg::Job(job)) { + match self.workers.sender(idx).try_send(WorkerMsg::Job(job)) { Ok(()) => {} Err(TrySendError::Full(_)) => { static FULL_COUNT: AtomicU64 = AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, Ordering::Relaxed); - if n < 8 || n.is_multiple_of(10000) { + if worth_logging(n) { warn!( worker = idx, drops = n + 1, @@ -272,9 +309,7 @@ impl DecryptWorkerPool { ); } } - Err(TrySendError::Disconnected(_)) => { - debug!(worker = idx, "DecryptWorker thread gone; dropping job"); - } + Err(TrySendError::Disconnected(_)) => self.note_refused(idx, "inbound packet"), } } @@ -303,11 +338,12 @@ impl DecryptWorkerPool { cache_key: (TransportId, u32), state: OwnedSessionState, ) -> bool { - if self.senders.is_empty() { - return false; - } let idx = self.worker_idx_for(cache_key); - match self.senders[idx].try_send(WorkerMsg::RegisterSession { cache_key, state }) { + match self + .workers + .sender(idx) + .try_send(WorkerMsg::RegisterSession { cache_key, state }) + { Ok(()) => true, Err(TrySendError::Full(_)) => { warn!( @@ -317,10 +353,7 @@ impl DecryptWorkerPool { false } Err(TrySendError::Disconnected(_)) => { - debug!( - worker = idx, - "DecryptWorker thread gone; ignoring registration" - ); + self.note_refused(idx, "session registration"); false } } @@ -329,11 +362,11 @@ impl DecryptWorkerPool { /// Drop a session from its worker (rekey, peer removed). Fire and /// forget — if the worker is gone we don't care. pub fn unregister_session(&self, cache_key: (TransportId, u32)) { - if self.senders.is_empty() { - return; - } let idx = self.worker_idx_for(cache_key); - let _ = self.senders[idx].try_send(WorkerMsg::UnregisterSession { cache_key }); + let _ = self + .workers + .sender(idx) + .try_send(WorkerMsg::UnregisterSession { cache_key }); } } @@ -774,3 +807,122 @@ mod tests { } } } + +#[cfg(test)] +impl DecryptWorkerPool { + /// A pool whose workers behave as `plan` says; `Run` workers are the + /// production loop. + pub(crate) fn for_test(plan: Vec) -> Self { + Self::start_with( + plan.len(), + crate::node::worker_set::test_spawner(plan, spawn_worker), + ) + } +} + +/// The pool's view of its workers: which are live and what a dead one does to +/// a dispatch or a registration. Every wait is bounded. +#[cfg(test)] +mod pool_tests { + use super::*; + use crate::node::worker_set::{TestWorker, wait_for}; + use ring::aead::UnboundKey; + use std::sync::mpsc; + + /// A session key the pool hashes to worker `idx`. + fn key_on_worker(pool: &DecryptWorkerPool, idx: usize) -> (TransportId, u32) { + (0u32..) + .map(|n| (TransportId::new(1), n)) + .find(|key| pool.worker_idx_for(*key) == idx) + .expect("some key hashes to every worker") + } + + fn session_state() -> OwnedSessionState { + let key = UnboundKey::new(&ring::aead::CHACHA20_POLY1305, &[0u8; 32]).unwrap(); + OwnedSessionState { + fmp_cipher: LessSafeKey::new(key), + fmp_replay: ReplayWindow::new(), + source_npub: None, + } + } + + fn job(cache_key: (TransportId, u32)) -> DecryptJob { + let (fallback_tx, _) = tokio::sync::mpsc::unbounded_channel::(); + DecryptJob { + packet_data: vec![0u8; 48], + cache_key, + _transport_id: cache_key.0, + _remote_addr: TransportAddr::from_string("127.0.0.1:1234"), + timestamp_ms: 1_000, + source_node_addr: NodeAddr::from_bytes([1u8; 16]), + fmp_counter: 1, + fmp_flags: 0, + fmp_header: [0u8; 16], + fmp_ciphertext_offset: 16, + fallback_tx, + } + } + + #[test] + fn a_decrypt_worker_that_panics_is_counted_dead() { + let (die_tx, die_rx) = mpsc::channel::<()>(); + let mut die_rx = Some(die_rx); + let pool = DecryptWorkerPool::start_with(2, |idx, rx| { + if idx != 1 { + return spawn_worker(idx, rx); + } + let die = die_rx.take().expect("worker 1 starts once"); + std::thread::Builder::new().spawn(move || { + let _rx = rx; + let _ = die.recv(); + panic!("simulated decrypt worker panic"); + }) + }); + assert_eq!(pool.liveness().live_workers(), 2); + die_tx.send(()).expect("worker 1 gone before its signal"); + assert!( + wait_for(|| pool.liveness().live_workers() == 1), + "a worker that panicked is still counted live" + ); + assert_eq!(pool.liveness().dead_workers(), vec![1]); + } + + #[test] + fn a_worker_that_fails_to_spawn_leaves_the_pool_serving_the_rest() { + let pool = DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + assert_eq!(pool.liveness().worker_count(), 2); + assert_eq!(pool.liveness().live_workers(), 1); + let key = key_on_worker(&pool, 0); + assert!( + pool.register_session(key, session_state()), + "the live worker refused a registration" + ); + pool.dispatch_job(job(key)); + assert_eq!(pool.liveness().refused_dispatches(), 0); + } + + #[test] + fn dispatch_to_an_exited_decrypt_worker_warns_and_counts() { + let pool = DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + let key = key_on_worker(&pool, 1); + let ((), logs) = crate::testutil::capture_logs(|| pool.dispatch_job(job(key))); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" worker=1"), "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"decrypt\""), "{warnings:?}"); + assert_eq!(pool.liveness().refused_dispatches(), 1); + } + + #[test] + fn registration_with_an_exited_decrypt_worker_warns_counts_and_is_refused() { + let pool = DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + let key = key_on_worker(&pool, 1); + let (registered, logs) = + crate::testutil::capture_logs(|| pool.register_session(key, session_state())); + assert!(!registered, "a dead worker cannot own a session"); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" worker=1"), "{warnings:?}"); + assert_eq!(pool.liveness().refused_dispatches(), 1); + } +} diff --git a/src/node/encrypt_worker.rs b/src/node/encrypt_worker.rs index d8cd4da7..ba1881bc 100644 --- a/src/node/encrypt_worker.rs +++ b/src/node/encrypt_worker.rs @@ -50,11 +50,14 @@ // warnings rather than gate every function individually. #![cfg_attr(not(unix), allow(dead_code))] +#[cfg(test)] +use crate::node::worker_set::{TestWorker, test_spawner}; +use crate::node::worker_set::{WorkerLiveness, WorkerSet, worth_logging}; use crate::proto::fmp::wire::ESTABLISHED_HEADER_SIZE; use crate::proto::fsp::wire::FSP_HEADER_SIZE; use crate::transport::udp::io::AsyncUdpSocket; #[cfg(not(target_os = "macos"))] -use crossbeam_channel::{Receiver, SendError, Sender, TrySendError, bounded}; +use crossbeam_channel::{Receiver, Sender, TrySendError, bounded}; use ring::aead::{Aad, LessSafeKey, Nonce}; #[cfg(any(target_os = "macos", test))] use std::collections::VecDeque; @@ -406,6 +409,36 @@ type WorkerSender = MacWorkerSender; #[cfg(not(target_os = "macos"))] type WorkerSender = Sender; +#[cfg(target_os = "macos")] +type WorkerReceiver = MacWorkerReceiver; + +#[cfg(not(target_os = "macos"))] +type WorkerReceiver = Receiver; + +fn worker_channel() -> (WorkerSender, WorkerReceiver) { + #[cfg(target_os = "macos")] + { + mac_worker_channel(WORKER_CHANNEL_CAP) + } + #[cfg(not(target_os = "macos"))] + { + bounded::(WORKER_CHANNEL_CAP) + } +} + +/// Start the production worker loop on `rx` in a named OS thread. +fn spawn_worker(idx: usize, rx: WorkerReceiver) -> std::io::Result> { + let builder = std::thread::Builder::new().name(format!("fips-encrypt-{idx}")); + #[cfg(target_os = "macos")] + { + builder.spawn(move || run_worker_macos(idx, rx)) + } + #[cfg(not(target_os = "macos"))] + { + builder.spawn(move || run_worker(idx, rx)) + } +} + /// Handle to the encrypt worker pool. /// /// Workers are **dedicated `std::thread`s** with **`crossbeam_channel`** @@ -425,7 +458,7 @@ type WorkerSender = Sender; /// destinations hash to different workers. #[derive(Clone)] pub(crate) struct EncryptWorkerPool { - senders: Arc<[WorkerSender]>, + workers: Arc>, #[cfg(target_os = "macos")] macos_senders: Arc, #[cfg(target_os = "macos")] @@ -437,31 +470,21 @@ impl EncryptWorkerPool { /// dispatches jobs hash-by-destination to them. The workers exit /// when all senders for their channel are dropped (i.e. when the /// returned `EncryptWorkerPool` and all clones go away). + /// + /// A worker thread that cannot be started is logged and left dead; + /// the caller reads how many started from [`Self::liveness`]. pub fn spawn(n: usize) -> Self { - let n = n.max(1); - let mut senders = Vec::with_capacity(n); - for i in 0..n { - #[cfg(target_os = "macos")] - { - let (tx, rx) = mac_worker_channel(WORKER_CHANNEL_CAP); - std::thread::Builder::new() - .name(format!("fips-encrypt-{i}")) - .spawn(move || run_worker_macos(i, rx)) - .expect("failed to spawn fips-encrypt OS thread"); - senders.push(tx); - } - #[cfg(not(target_os = "macos"))] - { - let (tx, rx) = bounded::(WORKER_CHANNEL_CAP); - std::thread::Builder::new() - .name(format!("fips-encrypt-{i}")) - .spawn(move || run_worker(i, rx)) - .expect("failed to spawn fips-encrypt OS thread"); - senders.push(tx); - } - } + Self::start_with(n, spawn_worker) + } + + /// Build the pool, starting each worker with `spawn`. Production passes + /// [`spawn_worker`]; tests pass workers that fail to start or exit on cue. + fn start_with( + n: usize, + spawn: impl FnMut(usize, WorkerReceiver) -> std::io::Result>, + ) -> Self { Self { - senders: senders.into(), + workers: Arc::new(WorkerSet::start("encrypt", n, worker_channel, spawn)), #[cfg(target_os = "macos")] macos_senders: Arc::new(MacSequencedSendFlows::default()), #[cfg(target_os = "macos")] @@ -469,12 +492,19 @@ impl EncryptWorkerPool { } } + /// Whether each worker is still running, and how many dispatches a dead + /// one refused. + pub(crate) fn liveness(&self) -> &dyn WorkerLiveness { + &*self.workers + } + /// Dispatch a job to the worker that owns its destination flow. /// The hash is over `dest_addr` so every packet for one peer's /// kernel `SocketAddr` lands on the same worker and stays in /// order — required for TCP's fast-retransmit logic above to - /// behave on a single-flow run. Fire-and-forget — the worker - /// handles send errors itself via stats counters. + /// behave on a single-flow run. The worker handles send errors + /// itself via stats counters. A job whose worker has exited is + /// dropped, counted, and logged at WARN. /// /// Uses `try_send` for the common uncontended case, then blocks /// only when the bounded worker channel is full. These jobs carry @@ -483,12 +513,18 @@ impl EncryptWorkerPool { /// retransmits. Blocking here pushes back toward the TUN reader /// and lets the kernel/app TCP stack pace the flow instead. pub fn dispatch(&self, job: FmpSendJob) { - if self.senders.is_empty() { - debug!("EncryptWorkerPool has no workers; dropping job"); - return; - } let (idx, job) = self.prepare_dispatch(job); - self.dispatch_to_worker(idx, job); + if !self.dispatch_to_worker(idx, job) { + let n = self.workers.note_refused(); + if worth_logging(n) { + warn!( + pool = "encrypt", + worker = idx, + refused = n + 1, + "Encrypt worker has exited; dropping packet" + ); + } + } } #[cfg(target_os = "macos")] @@ -503,7 +539,7 @@ impl EncryptWorkerPool { }; let mut h = std::collections::hash_map::DefaultHasher::new(); key.hash(&mut h); - let idx = (h.finish() as usize) % self.senders.len(); + let idx = (h.finish() as usize) % self.workers.len(); return (idx, QueuedFmpSendJob::direct(job)); } @@ -518,68 +554,81 @@ impl EncryptWorkerPool { .next_worker .fetch_add(1, std::sync::atomic::Ordering::Relaxed) / macos_worker_stride(); - let idx = ticket % self.senders.len(); + let idx = ticket % self.workers.len(); (idx, QueuedFmpSendJob::macos_sequenced(job, flow)) } #[cfg(not(target_os = "macos"))] fn prepare_dispatch(&self, job: FmpSendJob) -> (usize, QueuedFmpSendJob) { - use std::hash::{Hash, Hasher}; - let mut h = std::collections::hash_map::DefaultHasher::new(); - job.dest_addr.hash(&mut h); - let idx = (h.finish() as usize) % self.senders.len(); + let idx = self.worker_index_for(job.dest_addr); (idx, QueuedFmpSendJob::direct(job)) } + /// The worker that owns `dest`'s flow. + #[cfg(not(target_os = "macos"))] + fn worker_index_for(&self, dest: SocketAddr) -> usize { + use std::hash::{Hash, Hasher}; + let mut h = std::collections::hash_map::DefaultHasher::new(); + dest.hash(&mut h); + (h.finish() as usize) % self.workers.len() + } + + /// Queue `job` on worker `idx`, returning `false` when that worker has + /// exited and the job was not queued. #[cfg(target_os = "macos")] - fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) { - match self.senders[idx].try_push(job) { - Ok(()) => {} + fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> bool { + let sender = self.workers.sender(idx); + match sender.try_push(job) { + Ok(()) => true, Err(MacWorkerTryPushError::Full(job)) => { static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed); - if n < 8 || n.is_multiple_of(10000) { + if worth_logging(n) { warn!( worker = idx, full_events = n + 1, "EncryptWorker channel full; applying outbound backpressure" ); } - if let Err(MacWorkerPushError) = self.senders[idx].push_blocking(*job) { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); - } - } - Err(MacWorkerTryPushError::Closed) => { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); + sender.push_blocking(*job).is_ok() } + Err(MacWorkerTryPushError::Closed) => false, } } + /// Queue `job` on worker `idx`, returning `false` when that worker has + /// exited and the job was not queued. #[cfg(not(target_os = "macos"))] - fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) { - match self.senders[idx].try_send(job) { - Ok(()) => {} + fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> bool { + let sender = self.workers.sender(idx); + match sender.try_send(job) { + Ok(()) => true, Err(TrySendError::Full(job)) => { static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed); - if n < 8 || n.is_multiple_of(10000) { + if worth_logging(n) { warn!( worker = idx, full_events = n + 1, "EncryptWorker channel full; applying outbound backpressure" ); } - if let Err(SendError(_)) = self.senders[idx].send(job) { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); - } - } - Err(TrySendError::Disconnected(_)) => { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); + sender.send(job).is_ok() } + Err(TrySendError::Disconnected(_)) => false, } } } +#[cfg(test)] +impl EncryptWorkerPool { + /// A pool whose workers behave as `plan` says; `Run` workers are the + /// production loop. + pub(crate) fn for_test(plan: Vec) -> Self { + Self::start_with(plan.len(), test_spawner(plan, spawn_worker)) + } +} + #[cfg(target_os = "macos")] #[derive(Clone, Copy, Debug, Hash, PartialEq, Eq)] struct MacSendFlowKey { @@ -2892,3 +2941,176 @@ mod mac_ordered_tests { drop(gap); } } + +/// The pool's view of its workers: which are live, what a dead one does to a +/// dispatch, and that a live but full queue still blocks. Every wait is +/// bounded so a regression fails instead of hanging. +#[cfg(test)] +mod pool_tests { + use super::*; + use crate::node::worker_set::wait_for; + #[cfg(not(target_os = "macos"))] + use crate::transport::udp::io::UdpRawSocket; + #[cfg(not(target_os = "macos"))] + use ring::aead::UnboundKey; + #[cfg(not(target_os = "macos"))] + use std::net::UdpSocket; + use std::sync::mpsc; + #[cfg(not(target_os = "macos"))] + use std::time::Duration; + + #[cfg(not(target_os = "macos"))] + struct Rig { + _rt: tokio::runtime::Runtime, + socket: AsyncUdpSocket, + } + + #[cfg(not(target_os = "macos"))] + impl Rig { + fn new() -> Self { + let rt = tokio::runtime::Builder::new_current_thread() + .enable_io() + .build() + .expect("tokio rt"); + let enter = rt.enter(); + let socket = UdpRawSocket::open("127.0.0.1:0".parse().unwrap(), 1 << 20, 1 << 20) + .expect("open send socket") + .into_async() + .expect("into_async"); + drop(enter); + Self { _rt: rt, socket } + } + + fn job(&self, dest: SocketAddr, counter: u64) -> FmpSendJob { + let key = UnboundKey::new(&ring::aead::CHACHA20_POLY1305, &[3u8; 32]).expect("key"); + let mut wire_buf = + Vec::with_capacity(ESTABLISHED_HEADER_SIZE + 8 + crate::noise::TAG_SIZE); + wire_buf.extend_from_slice(&[0x5A; ESTABLISHED_HEADER_SIZE]); + wire_buf.extend_from_slice(&counter.to_le_bytes()); + FmpSendJob { + cipher: LessSafeKey::new(key), + counter, + wire_buf, + fsp_seal: None, + socket: self.socket.clone(), + dest_addr: dest, + #[cfg(any(target_os = "linux", target_os = "macos"))] + connected_socket: None, + drop_on_backpressure: false, + queued_at: None, + } + } + } + + /// A bound receiver whose address the pool dispatches to worker `idx`. + #[cfg(not(target_os = "macos"))] + fn receiver_on_worker(pool: &EncryptWorkerPool, idx: usize) -> UdpSocket { + loop { + let sock = UdpSocket::bind("127.0.0.1:0").expect("bind receiver"); + if pool.worker_index_for(sock.local_addr().unwrap()) == idx { + sock.set_read_timeout(Some(Duration::from_secs(5))) + .expect("read timeout"); + return sock; + } + } + } + + #[test] + fn an_encrypt_worker_that_panics_is_counted_dead() { + let (die_tx, die_rx) = mpsc::channel::<()>(); + let mut die_rx = Some(die_rx); + let pool = EncryptWorkerPool::start_with(2, |idx, rx| { + if idx != 1 { + return spawn_worker(idx, rx); + } + let die = die_rx.take().expect("worker 1 starts once"); + std::thread::Builder::new().spawn(move || { + let _rx = rx; + let _ = die.recv(); + panic!("simulated encrypt worker panic"); + }) + }); + assert_eq!(pool.liveness().live_workers(), 2); + die_tx.send(()).expect("worker 1 gone before its signal"); + assert!( + wait_for(|| pool.liveness().live_workers() == 1), + "a worker that panicked is still counted live" + ); + assert_eq!(pool.liveness().dead_workers(), vec![1]); + assert_eq!(pool.liveness().worker_count(), 2); + } + + #[cfg(not(target_os = "macos"))] + #[test] + fn a_worker_that_fails_to_spawn_leaves_the_pool_serving_the_rest() { + let rig = Rig::new(); + let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + assert_eq!(pool.liveness().worker_count(), 2); + assert_eq!(pool.liveness().live_workers(), 1); + + let recv = receiver_on_worker(&pool, 0); + pool.dispatch(rig.job(recv.local_addr().unwrap(), 1)); + let mut buf = [0u8; 128]; + recv.recv_from(&mut buf) + .expect("the live worker did not send the job dispatched to it"); + assert_eq!(pool.liveness().refused_dispatches(), 0); + } + + #[cfg(not(target_os = "macos"))] + #[test] + fn dispatch_to_an_exited_encrypt_worker_warns_and_counts() { + let rig = Rig::new(); + let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + let recv = receiver_on_worker(&pool, 1); + let ((), logs) = crate::testutil::capture_logs(|| { + pool.dispatch(rig.job(recv.local_addr().unwrap(), 1)); + }); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" worker=1"), "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"encrypt\""), "{warnings:?}"); + assert_eq!(pool.liveness().refused_dispatches(), 1); + } + + /// A live worker that has fallen behind must hold the rx loop back, not + /// have its packets dropped: these are tunnelled packets, and a drop here + /// reads as loss to TCP inside the tunnel. + #[cfg(not(target_os = "macos"))] + #[test] + fn a_full_live_encrypt_queue_still_blocks_dispatch() { + let rig = Rig::new(); + let (release_tx, release_rx) = mpsc::channel::<()>(); + let mut release_rx = Some(release_rx); + let pool = EncryptWorkerPool::start_with(1, |idx, rx| { + let release = release_rx.take().expect("one worker"); + std::thread::Builder::new().spawn(move || { + let _ = release.recv(); + run_worker(idx, rx); + }) + }); + let recv = UdpSocket::bind("127.0.0.1:0").expect("bind receiver"); + let dest = recv.local_addr().unwrap(); + for counter in 0..WORKER_CHANNEL_CAP as u64 { + pool.dispatch(rig.job(dest, counter)); + } + + let (done_tx, done_rx) = mpsc::channel::<()>(); + let blocked_pool = pool.clone(); + let last = rig.job(dest, WORKER_CHANNEL_CAP as u64); + std::thread::spawn(move || { + blocked_pool.dispatch(last); + let _ = done_tx.send(()); + }); + std::thread::sleep(Duration::from_millis(200)); + assert!( + matches!(done_rx.try_recv(), Err(mpsc::TryRecvError::Empty)), + "dispatch returned while the live worker's queue was full" + ); + + release_tx.send(()).expect("worker gone before release"); + done_rx + .recv_timeout(Duration::from_secs(5)) + .expect("dispatch still blocked after the worker drained"); + assert_eq!(pool.liveness().refused_dispatches(), 0); + } +} diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index 10b9b930..38e469e8 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -1,6 +1,7 @@ //! Node lifecycle management: start, stop, and peer connection initiation. pub(crate) mod supervisor; +mod workers; use super::{Node, NodeError, NodeState}; use supervisor::{Action, Child, Event, PeeringDesired, SupervisorFsm}; @@ -1609,44 +1610,33 @@ impl Node { } } Child::EncryptWorkers => { - // Hash-by-destination pins a TCP flow to one worker - // (preserves wire ordering); additional workers light up - // under multi-flow load. Infallible → always up. + // A worker that cannot be started degrades the node; it + // never stops start-up. #[cfg(unix)] - { - self.supervisor.encrypt_workers = Some( - super::encrypt_worker::EncryptWorkerPool::spawn(encrypt_worker_count), - ); - info!( - workers = encrypt_worker_count, - "Spawned FMP-encrypt worker pool" - ); + let start = self.start_encrypt_workers(encrypt_worker_count); + #[cfg(not(unix))] + let start = workers::PoolStart::ALL_LIVE; - // `FIPS_DECRYPT_WORKERS=0` disables the pool entirely - // and forces the in-line rx_loop decrypt path. When 0 - // no DecryptWorkers child is emitted, so this info! - // sits here — exactly where the decrypt spawn would be - // in today's sequence (after the encrypt spawn+info, - // before nostr). - if decrypt_worker_count == 0 { - info!("FIPS_DECRYPT_WORKERS=0 → in-line decrypt in rx_loop"); - } + // `FIPS_DECRYPT_WORKERS=0` disables the pool entirely + // and forces the in-line rx_loop decrypt path. When 0 + // no DecryptWorkers child is emitted, so this info! + // sits here — exactly where the decrypt spawn would be + // in today's sequence (after the encrypt spawn+info, + // before nostr). + #[cfg(unix)] + if decrypt_worker_count == 0 { + info!("FIPS_DECRYPT_WORKERS=0 → in-line decrypt in rx_loop"); } - Event::SubstrateUp { child } + start.event(child) } Child::DecryptWorkers => { - // Shard-owned decrypt pool. Infallible → always up. + // Shard-owned decrypt pool. A worker that cannot be + // started degrades the node; it never stops start-up. #[cfg(unix)] - { - self.supervisor.decrypt_workers = Some( - super::decrypt_worker::DecryptWorkerPool::spawn(decrypt_worker_count), - ); - info!( - workers = decrypt_worker_count, - "Spawned FMP-decrypt worker pool" - ); - } - Event::SubstrateUp { child } + let start = self.start_decrypt_workers(decrypt_worker_count); + #[cfg(not(unix))] + let start = workers::PoolStart::ALL_LIVE; + start.event(child) } Child::Nostr => { match NostrRendezvous::start( @@ -2388,6 +2378,18 @@ impl Node { } } + /// Tell the supervisor FSM that `child` exited on its own at runtime, and + /// publish the health it resolves. Shared by the child-exit channel's arm + /// in the rx loop and the tick's worker-liveness sweep, which steps the FSM + /// directly because it runs on the loop that drains that channel. + pub(in crate::node) fn step_child_exited(&mut self, child: Child) { + for action in self.supervisor.fsm.step(Event::ChildExited { child }) { + if let Action::PublishState(ns) = action { + self.supervisor.state = ns; + } + } + } + /// Feed every queued interface-presence edge to the supervisor FSM, /// returning the last [`NodeState`] it asked to publish (if any). /// diff --git a/src/node/lifecycle/supervisor.rs b/src/node/lifecycle/supervisor.rs index 301ad22c..edb4e658 100644 --- a/src/node/lifecycle/supervisor.rs +++ b/src/node/lifecycle/supervisor.rs @@ -63,10 +63,11 @@ //! - the degenerate no-children path now resolves to `Failed` (zero transports), //! **not** the old immediate-`Running`. //! -//! Runtime child-liveness monitoring (a `ChildExited` event re-routing health -//! when a task/thread dies at runtime) is **deferred**: start-completion health -//! resolution is start-framed, and liveness monitoring is a substantial unbuilt -//! mechanism. This commit is start-time health only. +//! Runtime child-liveness (a `ChildExited` event re-routing health when a +//! task/thread dies at runtime) came later. Its producers are the children that +//! report their own exit (TUN threads, the DNS task, the mDNS/Nostr monitor) +//! and, for the crypto worker pools, the rx loop tick's liveness sweep, which +//! reports a pool's child as exited when it loses a worker. //! //! ## Scope: interface presence, and `Degraded` as a level (this commit) //! @@ -298,7 +299,8 @@ pub(crate) enum Health { Full, /// ≥1 transport is up, but one or more configured optional children failed /// to start (a transport beyond the first, Nostr, mDNS, TUN, DNS, or a - /// worker-pool spawn). The node is operational (serving) but degraded. + /// worker pool that did not start every worker). The node is operational + /// (serving) but degraded. Degraded { /// The configured children that failed to start. reasons: HashSet, @@ -860,8 +862,8 @@ pub(crate) struct Supervisor { /// Off-task FMP-encrypt + UDP-send worker pool. Unix-only — /// the worker issues direct sendmmsg(2) / sendmsg+UDP_GSO calls - /// on raw fds via `AsRawFd`. None on Windows or when the worker - /// pool failed to spawn. + /// on raw fds via `AsRawFd`. None on Windows or when no worker + /// thread could be started. #[cfg(unix)] pub(crate) encrypt_workers: Option, @@ -869,9 +871,16 @@ pub(crate) struct Supervisor { /// `encrypt_workers`. Workers are shards: each owns its session /// state directly in a thread-local `HashMap` (no `RwLock`, /// no `Mutex` per packet). Hash-by-cache-key dispatch. + /// None on Windows, with `FIPS_DECRYPT_WORKERS=0`, or when no worker + /// thread could be started. #[cfg(unix)] pub(crate) decrypt_workers: Option, + /// Pools a test has built for start-up to install in place of spawning + /// its own. + #[cfg(all(test, unix))] + pub(in crate::node) staged_pools: StagedPools, + /// Transport-medium change detection: the receiver the rx loop drains and /// the detector task behind it. /// @@ -889,6 +898,15 @@ pub(crate) struct Supervisor { pub(in crate::node) fsm: SupervisorFsm, } +/// Crypto worker pools a test hands to start-up, each taken by the first +/// start of its pool's child. +#[cfg(all(test, unix))] +#[derive(Default)] +pub(in crate::node) struct StagedPools { + pub encrypt: Option, + pub decrypt: Option, +} + impl Supervisor { /// A fresh supervisor with all handles empty and the FSM in `Created`, /// matching the field initializers `Node::new` previously used. @@ -914,6 +932,8 @@ impl Supervisor { encrypt_workers: None, #[cfg(unix)] decrypt_workers: None, + #[cfg(all(test, unix))] + staged_pools: StagedPools::default(), netmon_rx: None, netmon_task: None, fsm: SupervisorFsm::new(), @@ -1796,6 +1816,39 @@ mod tests { assert!(s.absent().is_empty()); } + /// The crypto worker pools are a performance offload with a main-loop path + /// behind them, so losing one degrades the node and can never be what + /// fails it. + #[test] + fn a_worker_pool_exit_is_degraded_and_never_failed() { + let mut s = SupervisorFsm::running_with([ + Child::Transport(tid(1)), + Child::EncryptWorkers, + Child::DecryptWorkers, + ]); + assert_eq!( + s.step(Event::ChildExited { + child: Child::EncryptWorkers + }), + vec![Action::PublishState(NodeState::Degraded)] + ); + assert_eq!( + s.step(Event::ChildExited { + child: Child::DecryptWorkers + }), + vec![Action::PublishState(NodeState::Degraded)] + ); + assert!(matches!( + s.state(), + SupState::Running { + health: Health::Degraded { .. } + } + )); + let degraded = s.degraded_children(); + assert!(degraded.contains(&Child::EncryptWorkers)); + assert!(degraded.contains(&Child::DecryptWorkers)); + } + #[test] fn degraded_children_is_the_union_of_both_reason_sets() { let mut s = SupervisorFsm::new(); diff --git a/src/node/lifecycle/workers.rs b/src/node/lifecycle/workers.rs new file mode 100644 index 00000000..98f6fe8a --- /dev/null +++ b/src/node/lifecycle/workers.rs @@ -0,0 +1,240 @@ +//! Supervision of the crypto worker pools: what a pool that started short +//! means for the node, and the tick's sweep for workers that have since +//! exited. +//! +//! The pools are a performance offload, and Windows never starts them at +//! all, so losing workers is `Degraded` at most and never fatal. Inbound +//! packets for a session already held by a missing decrypt worker are +//! dropped until the session rekeys or the link is re-established; a session +//! that would be registered on the missing worker after the loss is decrypted +//! on the main loop instead. +//! The pools report facts (how many workers, how many live); the decisions +//! below turn them into supervisor events. + +use super::supervisor::{Child, Event}; +#[cfg(unix)] +use crate::node::Node; +#[cfg(unix)] +use crate::node::worker_set::WorkerLiveness; +#[cfg(unix)] +use tracing::{info, warn}; + +/// What a freshly started pool means for the node. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(in crate::node) struct PoolStart { + /// Keep the pool. A pool with no live worker is dropped, so every packet + /// takes the main-loop path directly instead of being refused first. + pub keep: bool, + /// Report the pool's child as up. Anything short of every worker live is + /// reported as failed to start, so start-completion health resolves to + /// `Degraded` on the first publish rather than `Full` corrected a tick + /// later. + pub up: bool, +} + +impl PoolStart { + /// Every worker started. Off Unix the pools are never started, and their + /// children report up as before. + #[cfg(not(unix))] + pub(in crate::node) const ALL_LIVE: Self = Self { + keep: true, + up: true, + }; + + /// The supervisor event that reports this outcome for `child`. + pub(in crate::node) fn event(self, child: Child) -> Event { + if self.up { + Event::SubstrateUp { child } + } else { + Event::SubstrateFailed { child } + } + } +} + +/// Classify a pool that started `live` of `configured` workers. +#[cfg(unix)] +pub(in crate::node) fn pool_start(live: usize, configured: usize) -> PoolStart { + PoolStart { + keep: live > 0, + up: live >= configured, + } +} + +/// Whether a pool now at `live` workers has lost any since it last reported +/// `reported`. Every loss is reported, the first one and each after it, so the +/// count of live workers stays visible as it falls. Nothing restarts a worker, +/// so the count never rises. +#[cfg(unix)] +pub(in crate::node) fn lost_workers(reported: usize, live: usize) -> bool { + live < reported +} + +#[cfg(unix)] +impl Node { + /// Start the encrypt pool and install it, or the main-loop path when no + /// worker started. Under test, a pool staged in + /// [`StagedPools`](super::supervisor::StagedPools) is installed instead of + /// spawning one, so a pool that starts short can be driven through + /// start-up. + pub(in crate::node) fn start_encrypt_workers(&mut self, n: usize) -> PoolStart { + #[cfg(test)] + if let Some(pool) = self.supervisor.staged_pools.encrypt.take() { + return self.install_encrypt_workers(pool); + } + self.install_encrypt_workers(crate::node::encrypt_worker::EncryptWorkerPool::spawn(n)) + } + + /// Install a started encrypt pool, or the main-loop path when none of + /// its workers started. + fn install_encrypt_workers( + &mut self, + pool: crate::node::encrypt_worker::EncryptWorkerPool, + ) -> PoolStart { + let start = report_pool_start("encrypt", pool.liveness()); + if start.up { + info!( + workers = pool.liveness().worker_count(), + "Spawned FMP-encrypt worker pool" + ); + } + self.supervisor.encrypt_workers = start.keep.then_some(pool); + start + } + + /// Start the decrypt pool and install it, or the main-loop path when no + /// worker started. Under test, a pool staged in + /// [`StagedPools`](super::supervisor::StagedPools) is installed instead of + /// spawning one, so a pool that starts short can be driven through + /// start-up. + pub(in crate::node) fn start_decrypt_workers(&mut self, n: usize) -> PoolStart { + #[cfg(test)] + if let Some(pool) = self.supervisor.staged_pools.decrypt.take() { + return self.install_decrypt_workers(pool); + } + self.install_decrypt_workers(crate::node::decrypt_worker::DecryptWorkerPool::spawn(n)) + } + + /// Install a started decrypt pool, or the main-loop path when none of + /// its workers started. + fn install_decrypt_workers( + &mut self, + pool: crate::node::decrypt_worker::DecryptWorkerPool, + ) -> PoolStart { + let start = report_pool_start("decrypt", pool.liveness()); + if start.up { + info!( + workers = pool.liveness().worker_count(), + "Spawned FMP-decrypt worker pool" + ); + } + self.supervisor.decrypt_workers = start.keep.then_some(pool); + start + } + + /// The tick's worker-liveness sweep. For each pool that has lost a worker + /// since the last sweep, log the loss with the live count and report the + /// pool's child as exited. The FSM republishes `Degraded` the first time; + /// a pool already out of its up-set republishes nothing, but each further + /// loss is still logged with the new count. + /// + /// Steps the FSM directly rather than through the child-exit channel: that + /// channel is drained by the same loop that runs this sweep, so a send into + /// a full channel from here would never complete. + pub(in crate::node) fn poll_worker_liveness(&mut self) { + let mut exited = Vec::with_capacity(2); + if let Some(pool) = &self.supervisor.encrypt_workers + && report_worker_loss("encrypt", pool.liveness()) + { + exited.push(Child::EncryptWorkers); + } + if let Some(pool) = &self.supervisor.decrypt_workers + && report_worker_loss("decrypt", pool.liveness()) + { + exited.push(Child::DecryptWorkers); + } + for child in exited { + self.step_child_exited(child); + } + } +} + +/// Classify how a pool started, logging a pool that started short. +#[cfg(unix)] +fn report_pool_start(pool: &'static str, workers: &dyn WorkerLiveness) -> PoolStart { + let live = workers.live_workers(); + let configured = workers.worker_count(); + let start = pool_start(live, configured); + if !start.up { + warn!( + pool, + live, configured, "Crypto worker pool started with fewer workers than configured" + ); + } + start +} + +/// Record a pool's live count, logging and returning `true` when it has lost a +/// worker since the last call. +#[cfg(unix)] +fn report_worker_loss(pool: &'static str, workers: &dyn WorkerLiveness) -> bool { + let live = workers.live_workers(); + let reported = workers.swap_reported_live(live); + if !lost_workers(reported, live) { + return false; + } + warn!( + pool, + live, + configured = workers.worker_count(), + dead = ?workers.dead_workers(), + "Crypto worker thread exited" + ); + true +} + +#[cfg(all(test, unix))] +mod tests { + use super::*; + + #[test] + fn worker_spawn_outcome_maps_k_of_n() { + assert_eq!( + pool_start(4, 4), + PoolStart { + keep: true, + up: true + } + ); + for live in 1..4 { + assert_eq!( + pool_start(live, 4), + PoolStart { + keep: true, + up: false + }, + "{live} of 4 live" + ); + } + assert_eq!( + pool_start(0, 4), + PoolStart { + keep: false, + up: false + } + ); + let child = Child::EncryptWorkers; + assert_eq!(pool_start(4, 4).event(child), Event::SubstrateUp { child }); + assert_eq!( + pool_start(3, 4).event(child), + Event::SubstrateFailed { child } + ); + } + + #[test] + fn every_loss_is_reported_and_no_loss_is_not() { + assert!(lost_workers(4, 3)); + assert!(lost_workers(3, 0)); + assert!(!lost_workers(4, 4)); + assert!(!lost_workers(0, 0)); + } +} diff --git a/src/node/mod.rs b/src/node/mod.rs index f0dfc200..7eda46b7 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -33,6 +33,8 @@ pub(crate) mod stats_history; #[cfg(test)] mod tests; mod tree; +#[cfg(unix)] +pub(crate) mod worker_set; use self::peer_error_budget::PeerErrorBudget; use self::rate_limit::{HandshakeRateLimiter, LookupSignRateLimiter, SessionSetupRateLimiter}; diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index ba623b0e..7a233498 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -4816,3 +4816,244 @@ fn mesh_filter_resolves_the_live_tun_device_rather_than_the_configured_name() { node.tun_name = Some(loopback.to_string()); assert_eq!(node.mesh_ifindex(), Some(expected)); } + +/// The tick's crypto worker liveness sweep: a lost worker degrades the node and +/// is logged with the live count; a healthy pool changes nothing. +#[cfg(unix)] +mod worker_liveness { + use super::*; + use crate::node::decrypt_worker::DecryptWorkerPool; + use crate::node::encrypt_worker::EncryptWorkerPool; + use crate::node::lifecycle::supervisor::{Child, SupervisorFsm}; + use crate::node::worker_set::{TestWorker, WorkerLiveness, wait_for}; + use std::collections::HashSet; + use std::sync::mpsc; + + /// A node seeded straight into a full `Running` with both pools up. + fn running_node() -> Node { + let mut node = make_node(); + node.supervisor.fsm = SupervisorFsm::running_with([ + Child::Transport(TransportId::new(1)), + Child::EncryptWorkers, + Child::DecryptWorkers, + ]); + node.supervisor.state = NodeState::Running; + node + } + + /// Two workers that each exit when their sender is used. + fn two_workers() -> (Vec, [mpsc::Sender<()>; 2]) { + let (kill0, exit0) = mpsc::channel(); + let (kill1, exit1) = mpsc::channel(); + ( + vec![TestWorker::ExitOn(exit0), TestWorker::ExitOn(exit1)], + [kill0, kill1], + ) + } + + /// An encrypt pool of two workers that each exit when their sender is + /// used. + fn pool_of_two() -> (EncryptWorkerPool, [mpsc::Sender<()>; 2]) { + let (plan, kills) = two_workers(); + (EncryptWorkerPool::for_test(plan), kills) + } + + /// Signal one worker of `workers` to exit and wait until `live_after` + /// remain. + fn kill_worker(workers: &dyn WorkerLiveness, kill: &mpsc::Sender<()>, live_after: usize) { + kill.send(()).expect("worker gone before its signal"); + assert!( + wait_for(|| workers.live_workers() == live_after), + "worker never exited" + ); + } + + fn kill(node: &Node, kill: &mpsc::Sender<()>, live_after: usize) { + let pool = node.supervisor.encrypt_workers.as_ref().unwrap(); + kill_worker(pool.liveness(), kill, live_after); + } + + #[tokio::test] + async fn a_dead_worker_degrades_the_node_and_names_the_live_count() { + let mut node = running_node(); + let (pool, kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(pool); + kill(&node, &kills[1], 1); + + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert_eq!(node.state(), NodeState::Degraded); + assert!( + node.supervisor + .fsm + .degraded_children() + .contains(&Child::EncryptWorkers) + ); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"encrypt\""), "{warnings:?}"); + assert!(warnings[0].contains(" live=1"), "{warnings:?}"); + assert!(warnings[0].contains(" configured=2"), "{warnings:?}"); + + // Nothing new: no second report. + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert!(logs.warnings().is_empty(), "{:?}", logs.warnings()); + assert_eq!(node.state(), NodeState::Degraded); + + // The last worker goes: reported with the new count, and the node is + // still degraded, not failed. + kill(&node, &kills[0], 0); + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" live=0"), "{warnings:?}"); + assert_eq!(node.state(), NodeState::Degraded); + } + + #[tokio::test] + async fn a_dead_decrypt_worker_degrades_the_node_and_names_its_pool() { + let mut node = running_node(); + let (encrypt, _encrypt_kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(encrypt); + let (plan, kills) = two_workers(); + node.supervisor.decrypt_workers = Some(DecryptWorkerPool::for_test(plan)); + let pool = node.supervisor.decrypt_workers.as_ref().unwrap(); + kill_worker(pool.liveness(), &kills[0], 1); + + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert_eq!(node.state(), NodeState::Degraded); + assert_eq!( + node.supervisor.fsm.degraded_children(), + HashSet::from([Child::DecryptWorkers]) + ); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"decrypt\""), "{warnings:?}"); + assert!(warnings[0].contains(" live=1"), "{warnings:?}"); + assert!(warnings[0].contains(" configured=2"), "{warnings:?}"); + } + + #[tokio::test] + async fn healthy_pools_leave_the_node_running() { + let mut node = running_node(); + let (pool, _kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(pool); + let (plan, _decrypt_kills) = two_workers(); + node.supervisor.decrypt_workers = Some(DecryptWorkerPool::for_test(plan)); + + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert!(logs.warnings().is_empty(), "{:?}", logs.warnings()); + assert_eq!(node.state(), NodeState::Running); + } + + /// What start-up made of the pools a test staged for it. + struct StartOutcome { + state: NodeState, + degraded: HashSet, + encrypt_installed: bool, + decrypt_installed: bool, + } + + /// Start a node that installs `encrypt` and `decrypt` in place of the + /// pools it would spawn, then stop it. Fails if start-up did not take both + /// staged pools, since the outcome would then say nothing about them. + async fn start_with_staged( + encrypt: EncryptWorkerPool, + decrypt: DecryptWorkerPool, + ) -> StartOutcome { + let mut node = make_healthy_node(); + node.supervisor.staged_pools.encrypt = Some(encrypt); + node.supervisor.staged_pools.decrypt = Some(decrypt); + node.start().await.unwrap(); + let staged = &node.supervisor.staged_pools; + let taken = staged.encrypt.is_none() && staged.decrypt.is_none(); + let outcome = StartOutcome { + state: node.state(), + degraded: node.supervisor.fsm.degraded_children(), + encrypt_installed: node.supervisor.encrypt_workers.is_some(), + decrypt_installed: node.supervisor.decrypt_workers.is_some(), + }; + node.stop().await.unwrap(); + assert!(taken, "start-up did not install both staged pools"); + outcome + } + + #[tokio::test] + async fn a_worker_that_fails_to_start_leaves_the_node_degraded_with_its_pool_installed() { + let outcome = start_with_staged( + EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]), + DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + ) + .await; + assert_eq!(outcome.state, NodeState::Degraded); + assert_eq!(outcome.degraded, HashSet::from([Child::EncryptWorkers])); + assert!( + outcome.encrypt_installed, + "a pool with a live worker is kept" + ); + assert!(outcome.decrypt_installed); + } + + #[tokio::test] + async fn a_pool_with_no_worker_started_leaves_the_node_degraded_and_is_not_installed() { + let outcome = start_with_staged( + EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + DecryptWorkerPool::for_test(vec![TestWorker::FailSpawn, TestWorker::FailSpawn]), + ) + .await; + assert_eq!(outcome.state, NodeState::Degraded); + assert_eq!(outcome.degraded, HashSet::from([Child::DecryptWorkers])); + assert!(outcome.encrypt_installed); + assert!( + !outcome.decrypt_installed, + "a pool with no live worker is dropped for the main-loop path" + ); + } + + #[tokio::test] + async fn staged_pools_with_every_worker_started_leave_the_node_running() { + let outcome = start_with_staged( + EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + ) + .await; + assert_eq!(outcome.state, NodeState::Running); + assert!(outcome.degraded.is_empty(), "{:?}", outcome.degraded); + assert!(outcome.encrypt_installed && outcome.decrypt_installed); + } + + /// Start a node, swap in a pool of two, optionally kill one worker, and + /// drive the real rx loop past one tick. Returns the state it published. + async fn drive_with_pool(kill_one: bool) -> NodeState { + let mut node = make_healthy_node(); + node.start().await.unwrap(); + assert_eq!(node.state(), NodeState::Running); + let (pool, kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(pool); + if kill_one { + kill(&node, &kills[1], 1); + } + + let drive = tokio::time::timeout( + Duration::from_millis(1500), + node.run_rx_loop_with_shutdown(std::future::pending()), + ) + .await; + assert!( + drive.is_err(), + "the rx loop must still be running: {drive:?}" + ); + let state = node.state(); + node.stop().await.unwrap(); + state + } + + #[tokio::test] + async fn a_dead_worker_degrades_the_node_through_the_rx_loop() { + assert_eq!(drive_with_pool(true).await, NodeState::Degraded); + } + + #[tokio::test] + async fn a_healthy_pool_leaves_the_node_running_through_the_rx_loop() { + assert_eq!(drive_with_pool(false).await, NodeState::Running); + } +} diff --git a/src/node/worker_set.rs b/src/node/worker_set.rs new file mode 100644 index 00000000..3c7f17cb --- /dev/null +++ b/src/node/worker_set.rs @@ -0,0 +1,245 @@ +//! The worker threads behind one crypto worker pool, and what is known about +//! whether each is still running. +//! +//! Both pools (`encrypt_worker`, `decrypt_worker`) are a set of OS threads, +//! each reached through its own bounded channel. A worker that exits, a panic +//! unwinding included, drops its receiver, so its channel closes and every +//! later dispatch to it is refused. Nothing restarts it. This module keeps the +//! thread handles so that loss can be seen, and counts the dispatches it +//! refused. +//! +//! It reports facts only. Whether a loss degrades the node is decided by the +//! driver in `lifecycle`. + +use portable_atomic::AtomicU64; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::thread::JoinHandle; +use tracing::warn; + +/// One worker: the sending end of its channel and its thread, `None` when the +/// thread could not be started. +struct Worker { + sender: S, + thread: Option>, +} + +/// The workers of one pool. Held behind an `Arc` by the pool so that cloning +/// the pool per packet stays a reference-count bump. +pub(crate) struct WorkerSet { + workers: Box<[Worker]>, + /// Dispatches refused because the target worker had exited. + refused: AtomicU64, + /// The live-worker count the liveness sweep last reported. Starts at the + /// number of threads that started, so a worker that never started is + /// reported once, at spawn, and not again by the sweep. It lives with the + /// workers so that a pool replaced by another starts from its own count. + reported_live: AtomicUsize, +} + +impl WorkerSet { + /// Start `n` workers, at least one. + /// + /// `channel` builds each worker's sender and receiver. `spawn` starts the + /// thread that owns the receiver. A worker whose thread cannot be started + /// is logged and left dead rather than aborting the caller: the receiver + /// was moved into the failed spawn and is dropped with it, so that + /// worker's channel is closed and a dispatch to it is refused, not queued. + pub(crate) fn start( + pool: &'static str, + n: usize, + mut channel: impl FnMut() -> (S, R), + mut spawn: impl FnMut(usize, R) -> std::io::Result>, + ) -> Self { + let n = n.max(1); + let mut workers = Vec::with_capacity(n); + for idx in 0..n { + let (sender, receiver) = channel(); + let thread = match spawn(idx, receiver) { + Ok(handle) => Some(handle), + Err(error) => { + warn!(pool, worker = idx, %error, "Failed to start a crypto worker thread"); + None + } + }; + workers.push(Worker { sender, thread }); + } + let started = workers.iter().filter(|w| w.thread.is_some()).count(); + Self { + workers: workers.into(), + refused: AtomicU64::new(0), + reported_live: AtomicUsize::new(started), + } + } + + /// Number of workers, including dead ones. Never zero. + pub(crate) fn len(&self) -> usize { + self.workers.len() + } + + /// The sending end of worker `idx`'s channel. + pub(crate) fn sender(&self, idx: usize) -> &S { + &self.workers[idx].sender + } + + /// Count one dispatch refused by a dead worker, returning the count before + /// this one. + pub(crate) fn note_refused(&self) -> u64 { + self.refused.fetch_add(1, Ordering::Relaxed) + } +} + +/// What the liveness sweep and the tests read from a pool, without naming its +/// channel type. +pub(crate) trait WorkerLiveness { + /// Number of workers the pool was built with. + fn worker_count(&self) -> usize; + /// Workers whose thread started and has not finished. A worker thread + /// finishes only by unwinding or when every sender to it is dropped, and + /// the node holds the pool while it runs, so a finished thread is a dead + /// worker. + fn live_workers(&self) -> usize; + /// Indices of the workers that are not live. + fn dead_workers(&self) -> Vec; + /// Dispatches refused because their worker had exited. + #[cfg(test)] + fn refused_dispatches(&self) -> u64; + /// Record `live` as the count last reported, returning the previous one. + fn swap_reported_live(&self, live: usize) -> usize; +} + +impl WorkerLiveness for WorkerSet { + fn worker_count(&self) -> usize { + self.workers.len() + } + + fn live_workers(&self) -> usize { + self.workers.iter().filter(|w| is_live(w)).count() + } + + fn dead_workers(&self) -> Vec { + self.workers + .iter() + .enumerate() + .filter(|(_, w)| !is_live(w)) + .map(|(idx, _)| idx) + .collect() + } + + #[cfg(test)] + fn refused_dispatches(&self) -> u64 { + self.refused.load(Ordering::Relaxed) + } + + fn swap_reported_live(&self, live: usize) -> usize { + self.reported_live.swap(live, Ordering::Relaxed) + } +} + +fn is_live(worker: &Worker) -> bool { + worker.thread.as_ref().is_some_and(|t| !t.is_finished()) +} + +/// Whether the `n`th event (counting from zero) of a repeating condition is +/// logged: the first eight, then one in ten thousand. +pub(crate) fn worth_logging(n: u64) -> bool { + n < 8 || n.is_multiple_of(10_000) +} + +/// How a test wants one worker of a pool to behave. +#[cfg(test)] +pub(crate) enum TestWorker { + /// The production worker loop. + Run, + /// The thread fails to start. + FailSpawn, + /// The thread holds its receiver until signalled, then exits. + ExitOn(std::sync::mpsc::Receiver<()>), +} + +/// A spawner that starts each worker as `plan` says, using `run` for +/// [`TestWorker::Run`]. +#[cfg(test)] +pub(crate) fn test_spawner( + plan: Vec, + run: impl Fn(usize, R) -> std::io::Result>, +) -> impl FnMut(usize, R) -> std::io::Result> { + let mut plan: Vec> = plan.into_iter().map(Some).collect(); + move |idx, rx| match plan[idx].take().expect("each worker is started once") { + TestWorker::Run => run(idx, rx), + TestWorker::FailSpawn => { + drop(rx); + Err(std::io::Error::other("worker start refused by the test")) + } + TestWorker::ExitOn(signal) => std::thread::Builder::new().spawn(move || { + let _ = signal.recv(); + drop(rx); + }), + } +} + +/// Wait up to five seconds for `cond`, returning whether it came true. +#[cfg(test)] +pub(crate) fn wait_for(cond: impl Fn() -> bool) -> bool { + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + while !cond() { + if std::time::Instant::now() >= deadline { + return false; + } + std::thread::sleep(std::time::Duration::from_millis(5)); + } + true +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::mpsc; + + #[test] + fn a_worker_that_exits_is_no_longer_live_and_is_named_dead() { + let (stop_tx, stop_rx) = mpsc::channel::<()>(); + let mut stop_rx = Some(stop_rx); + let set = WorkerSet::start( + "test", + 2, + mpsc::channel::, + |idx, rx: mpsc::Receiver| { + let stop = if idx == 1 { stop_rx.take() } else { None }; + std::thread::Builder::new().spawn(move || match stop { + Some(stop) => { + let _ = stop.recv(); + drop(rx); + } + None => while rx.recv().is_ok() {}, + }) + }, + ); + assert_eq!(set.live_workers(), 2); + stop_tx.send(()).unwrap(); + assert!( + wait_for(|| set.live_workers() == 1), + "worker 1 never exited" + ); + assert_eq!(set.dead_workers(), vec![1]); + assert_eq!(set.worker_count(), 2); + } + + #[test] + fn the_reported_baseline_starts_at_the_workers_that_started() { + let set = WorkerSet::start("test", 3, mpsc::channel::, |idx, rx| { + if idx == 2 { + drop(rx); + return Err(std::io::Error::other("refused by the test")); + } + std::thread::Builder::new().spawn(move || while rx.recv().is_ok() {}) + }); + assert_eq!(set.swap_reported_live(2), 2); + assert_eq!(set.dead_workers(), vec![2]); + } + + #[test] + fn worth_logging_keeps_the_first_eight_then_one_in_ten_thousand() { + let logged: Vec = (0..20_001).filter(|n| worth_logging(*n)).collect(); + assert_eq!(logged, vec![0, 1, 2, 3, 4, 5, 6, 7, 10_000, 20_000]); + } +} From c38fb00f50380acc88f97cb461c27665dc61d616 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 2 Oct 2026 16:10:25 +0000 Subject: [PATCH 7/8] Encrypt packets on the main loop when their encrypt worker has exited A packet hashed to an encrypt worker that had exited was dropped after its counters were reserved and, on the session data path, after the link and session statistics had already counted it as sent. Every peer that hashed to the dead worker was cut off until the process restarted. Dispatch now hands a job its worker refused back to the caller instead of dropping it. Both send sites seal the job on the main loop with the counters it already reserved, so no counter is skipped and the receiver sees no gap, and send it through the transport the way the inline path does. On the link-message path the statistics count the bytes actually sent; on the session data path they were recorded before dispatch and now describe a packet that was sent. A job is handed back only when no worker holds it, so each reserved counter is still used at most once. In the macOS ordered sender the job's place in its flow is released before it is handed back, so later packets for that flow are not held behind it. The worker and the main loop share one seal function. It refuses a job whose offsets do not fit its buffer instead of indexing out of bounds, since a panic there would now end the node rather than one worker. --- src/node/encrypt_worker.rs | 494 +++++++++++++++++++++++++++------- src/node/handlers/session.rs | 37 ++- src/node/lifecycle/workers.rs | 3 +- src/node/mod.rs | 67 +++-- src/node/tests/handshake.rs | 206 ++++++++++++++ src/node/tests/mod.rs | 25 ++ src/node/tests/session.rs | 107 ++++++++ 7 files changed, 810 insertions(+), 129 deletions(-) diff --git a/src/node/encrypt_worker.rs b/src/node/encrypt_worker.rs index ba1881bc..c68019d3 100644 --- a/src/node/encrypt_worker.rs +++ b/src/node/encrypt_worker.rs @@ -57,7 +57,7 @@ use crate::proto::fmp::wire::ESTABLISHED_HEADER_SIZE; use crate::proto::fsp::wire::FSP_HEADER_SIZE; use crate::transport::udp::io::AsyncUdpSocket; #[cfg(not(target_os = "macos"))] -use crossbeam_channel::{Receiver, Sender, TrySendError, bounded}; +use crossbeam_channel::{Receiver, SendError, Sender, TrySendError, bounded}; use ring::aead::{Aad, LessSafeKey, Nonce}; #[cfg(any(target_os = "macos", test))] use std::collections::VecDeque; @@ -181,6 +181,24 @@ impl QueuedFmpSendJob { } } +/// A queued job its worker refused because the worker has exited. Boxed so +/// the dispatch path's `Result` stays small; the allocation happens only on a +/// refusal. +struct Refused(Box); + +impl Refused { + /// The job, with any macOS ordered-flow slot it held released as a skip. + fn into_job(self) -> Box { + #[cfg(target_os = "macos")] + let QueuedFmpSendJob { job, macos_ticket } = *self.0; + #[cfg(not(target_os = "macos"))] + let QueuedFmpSendJob { job } = *self.0; + #[cfg(target_os = "macos")] + drop(macos_ticket); + Box::new(job) + } +} + /// Handle to the encrypt worker pool. Dispatches jobs **hash-by- /// destination** across N worker tasks via per-worker bounded /// crossbeam channels. The bounded queue intentionally backpressures @@ -242,14 +260,17 @@ struct MacWorkerQueueState { closed: bool, } +/// Why `try_push` did not queue a job. Both variants hand the job back. #[cfg(any(target_os = "macos", test))] enum MacWorkerTryPushError { Full(Box), - Closed, + /// The receiver is gone: the worker has exited. + Closed(Box), } +/// `push_blocking` found the receiver gone; the job is handed back. #[cfg(any(target_os = "macos", test))] -struct MacWorkerPushError; +struct MacWorkerPushError(Box); #[cfg(any(target_os = "macos", test))] fn mac_worker_channel(cap: usize) -> (MacWorkerSender, MacWorkerReceiver) { @@ -280,10 +301,10 @@ impl MacWorkerSender { .lock() .expect("encrypt worker queue poisoned"); if state.closed { - // Outside the lock: dropping a sequenced job completes its slot. + // The caller drops or reuses the job, outside the lock: dropping a + // sequenced job completes its slot. drop(state); - drop(job); - return Err(MacWorkerTryPushError::Closed); + return Err(MacWorkerTryPushError::Closed(Box::new(job))); } if state.queue.len() >= self.inner.cap { return Err(MacWorkerTryPushError::Full(Box::new(job))); @@ -298,7 +319,7 @@ impl MacWorkerSender { Ok(()) } - fn push_blocking(&self, job: T) -> Result<(), MacWorkerPushError> { + fn push_blocking(&self, job: T) -> Result<(), MacWorkerPushError> { let mut state = self .inner .state @@ -307,8 +328,7 @@ impl MacWorkerSender { loop { if state.closed { drop(state); - drop(job); - return Err(MacWorkerPushError); + return Err(MacWorkerPushError(Box::new(job))); } if state.queue.len() < self.inner.cap { let was_empty = state.queue.is_empty(); @@ -503,8 +523,15 @@ impl EncryptWorkerPool { /// kernel `SocketAddr` lands on the same worker and stays in /// order — required for TCP's fast-retransmit logic above to /// behave on a single-flow run. The worker handles send errors - /// itself via stats counters. A job whose worker has exited is - /// dropped, counted, and logged at WARN. + /// itself via stats counters. + /// + /// A job whose worker has exited is never queued: it is counted, + /// logged at WARN, and handed back as `Err` so the caller can seal + /// and send it on its own path with the counters it already + /// reserved. A job is handed back only when no worker has it, so + /// each reserved counter is still used at most once. In the macOS + /// ordered mode the job's place in its flow is released as a skip + /// before it is handed back. /// /// Uses `try_send` for the common uncontended case, then blocks /// only when the bounded worker channel is full. These jobs carry @@ -512,19 +539,22 @@ impl EncryptWorkerPool { /// this internal queue makes TCP-over-TUN collapse with avoidable /// retransmits. Blocking here pushes back toward the TUN reader /// and lets the kernel/app TCP stack pace the flow instead. - pub fn dispatch(&self, job: FmpSendJob) { + #[must_use = "a job handed back was not sent; the caller must send it another way"] + pub fn dispatch(&self, job: FmpSendJob) -> Result<(), Box> { let (idx, job) = self.prepare_dispatch(job); - if !self.dispatch_to_worker(idx, job) { - let n = self.workers.note_refused(); - if worth_logging(n) { - warn!( - pool = "encrypt", - worker = idx, - refused = n + 1, - "Encrypt worker has exited; dropping packet" - ); - } + let Err(refused) = self.dispatch_to_worker(idx, job) else { + return Ok(()); + }; + let n = self.workers.note_refused(); + if worth_logging(n) { + warn!( + pool = "encrypt", + worker = idx, + refused = n + 1, + "Encrypt worker has exited; encrypting the packet on the main loop" + ); } + Err(refused.into_job()) } #[cfg(target_os = "macos")] @@ -573,13 +603,13 @@ impl EncryptWorkerPool { (h.finish() as usize) % self.workers.len() } - /// Queue `job` on worker `idx`, returning `false` when that worker has - /// exited and the job was not queued. + /// Queue `job` on worker `idx`, or hand it back when that worker has + /// exited. #[cfg(target_os = "macos")] - fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> bool { + fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> Result<(), Refused> { let sender = self.workers.sender(idx); match sender.try_push(job) { - Ok(()) => true, + Ok(()) => Ok(()), Err(MacWorkerTryPushError::Full(job)) => { static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed); @@ -590,19 +620,21 @@ impl EncryptWorkerPool { "EncryptWorker channel full; applying outbound backpressure" ); } - sender.push_blocking(*job).is_ok() + sender + .push_blocking(*job) + .map_err(|MacWorkerPushError(job)| Refused(job)) } - Err(MacWorkerTryPushError::Closed) => false, + Err(MacWorkerTryPushError::Closed(job)) => Err(Refused(job)), } } - /// Queue `job` on worker `idx`, returning `false` when that worker has - /// exited and the job was not queued. + /// Queue `job` on worker `idx`, or hand it back when that worker has + /// exited. #[cfg(not(target_os = "macos"))] - fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> bool { + fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> Result<(), Refused> { let sender = self.workers.sender(idx); match sender.try_send(job) { - Ok(()) => true, + Ok(()) => Ok(()), Err(TrySendError::Full(job)) => { static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed); @@ -613,9 +645,11 @@ impl EncryptWorkerPool { "EncryptWorker channel full; applying outbound backpressure" ); } - sender.send(job).is_ok() + sender + .send(job) + .map_err(|SendError(job)| Refused(Box::new(job))) } - Err(TrySendError::Disconnected(_)) => false, + Err(TrySendError::Disconnected(job)) => Err(Refused(Box::new(job))), } } } @@ -627,6 +661,21 @@ impl EncryptWorkerPool { pub(crate) fn for_test(plan: Vec) -> Self { Self::start_with(plan.len(), test_spawner(plan, spawn_worker)) } + + /// The worker a job to `dest` is dispatched to, where that depends on + /// `dest` alone. On macOS it also depends on the sending sockets, or on a + /// round-robin in the ordered mode, so there it is `None`. + pub(crate) fn worker_index_for_dest(&self, dest: SocketAddr) -> Option { + #[cfg(target_os = "macos")] + { + let _ = dest; + None + } + #[cfg(not(target_os = "macos"))] + { + Some(self.worker_index_for(dest)) + } + } } #[cfg(target_os = "macos")] @@ -1145,6 +1194,89 @@ fn run_worker_macos(idx: usize, rx: MacWorkerReceiver) { trace!(worker = idx, "FMP encrypt worker thread exiting"); } +/// Why a job could not be sealed. +#[derive(Debug, thiserror::Error)] +pub(crate) enum SealError { + /// The offsets the job carries do not fit its buffer. + #[error("job layout does not fit its buffer")] + Layout, + /// The AEAD refused to seal. + #[error("AEAD seal failed")] + Aead, +} + +impl FmpSendJob { + /// Seal this job on the calling thread, as its worker would have, and + /// return the wire packet. For a job its worker refused: the counters it + /// carries were reserved for it and are used here, once. + pub(crate) fn seal_inline(self) -> Result, SealError> { + let FmpSendJob { + cipher, + counter, + mut wire_buf, + fsp_seal, + .. + } = self; + seal_wire(&cipher, counter, &mut wire_buf, fsp_seal)?; + Ok(wire_buf) + } +} + +/// Seal one job's wire buffer in place: the inner FSP seal first when the job +/// carries one, then the outer FMP seal over `[16..]` with the header as AAD. +/// Each tag is appended into capacity the builder reserved, so the buffer +/// becomes the wire packet without reallocating. +/// +/// A layout that does not fit the buffer is refused rather than indexed out of +/// bounds: this runs on a worker thread and, for a job that worker refused, on +/// the rx loop, where a panic would end the node rather than one worker. +fn seal_wire( + cipher: &LessSafeKey, + counter: u64, + wire_buf: &mut Vec, + fsp_seal: Option, +) -> Result<(), SealError> { + if let Some(fsp) = fsp_seal { + let aad_end = fsp + .aad_offset + .checked_add(FSP_HEADER_SIZE) + .ok_or(SealError::Layout)?; + if aad_end > fsp.plaintext_offset || fsp.plaintext_offset > wire_buf.len() { + return Err(SealError::Layout); + } + + let mut nonce_bytes = [0u8; 12]; + nonce_bytes[4..12].copy_from_slice(&fsp.counter.to_le_bytes()); + let nonce = Nonce::assume_unique_for_key(nonce_bytes); + let (prefix, plaintext_slice) = wire_buf.split_at_mut(fsp.plaintext_offset); + let aad = &prefix[fsp.aad_offset..aad_end]; + let tag = fsp + .cipher + .seal_in_place_separate_tag(nonce, Aad::from(aad), plaintext_slice) + .map_err(|_| SealError::Aead)?; + wire_buf.extend_from_slice(tag.as_ref()); + } + + if wire_buf.len() < ESTABLISHED_HEADER_SIZE { + return Err(SealError::Layout); + } + let mut nonce_bytes = [0u8; 12]; + nonce_bytes[4..12].copy_from_slice(&counter.to_le_bytes()); + let nonce = Nonce::assume_unique_for_key(nonce_bytes); + // Split-borrow: AAD reads from header bytes [0..16], seal writes + // into the plaintext slice [16..]. ring::aead's `seal_in_place_ + // separate_tag` takes `&mut [u8]` so we can hand it the + // post-header slice while AAD references the header slice. + // `split_at_mut` is the standard way to do this safely. + let (header_slice, plaintext_slice) = wire_buf.split_at_mut(ESTABLISHED_HEADER_SIZE); + let tag = cipher + .seal_in_place_separate_tag(nonce, Aad::from(&*header_slice), plaintext_slice) + .map_err(|_| SealError::Aead)?; + // wire_buf already has `+16` capacity reserved → no realloc. + wire_buf.extend_from_slice(tag.as_ref()); + Ok(()) +} + /// Encrypt every job in `batch` in place, then issue one or more /// bulk-send syscalls grouped **by exact send target**. Clears /// `batch` on return. Sync version — operates directly on the raw @@ -1227,65 +1359,14 @@ fn flush_batch_sync( crate::perf_profile::Stage::FmpWorkerQueueWait, queued_at, ); - if let Some(fsp) = fsp_seal { - if fsp.aad_offset + FSP_HEADER_SIZE > fsp.plaintext_offset - || fsp.plaintext_offset > wire_buf.len() - { - #[cfg(target_os = "macos")] - if let Some(ticket) = macos_ticket { - push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip); - } - continue; + if seal_wire(&cipher, counter, &mut wire_buf, fsp_seal).is_err() { + #[cfg(target_os = "macos")] + if let Some(ticket) = macos_ticket { + push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip); } - - let mut nonce_bytes = [0u8; 12]; - nonce_bytes[4..12].copy_from_slice(&fsp.counter.to_le_bytes()); - let nonce = Nonce::assume_unique_for_key(nonce_bytes); - let (prefix, plaintext_slice) = wire_buf.split_at_mut(fsp.plaintext_offset); - let aad = &prefix[fsp.aad_offset..fsp.aad_offset + FSP_HEADER_SIZE]; - let tag = - match fsp - .cipher - .seal_in_place_separate_tag(nonce, Aad::from(aad), plaintext_slice) - { - Ok(tag) => tag, - Err(_) => { - #[cfg(target_os = "macos")] - if let Some(ticket) = macos_ticket { - push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip); - } - continue; - } - }; - wire_buf.extend_from_slice(tag.as_ref()); + continue; } - let mut nonce_bytes = [0u8; 12]; - nonce_bytes[4..12].copy_from_slice(&counter.to_le_bytes()); - let nonce = Nonce::assume_unique_for_key(nonce_bytes); - // Split-borrow: AAD reads from header bytes [0..16], seal writes - // into the plaintext slice [16..]. ring::aead's `seal_in_place_ - // separate_tag` takes `&mut [u8]` so we can hand it the - // post-header slice while AAD references the header slice. - // `split_at_mut` is the standard way to do this safely. - let (header_slice, plaintext_slice) = wire_buf.split_at_mut(ESTABLISHED_HEADER_SIZE); - let tag = match cipher.seal_in_place_separate_tag( - nonce, - Aad::from(&*header_slice), - plaintext_slice, - ) { - Ok(tag) => tag, - Err(_) => { - #[cfg(target_os = "macos")] - if let Some(ticket) = macos_ticket { - push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip); - } - continue; - } - }; - // wire_buf already has `+16` capacity reserved → no realloc. - wire_buf.extend_from_slice(tag.as_ref()); - #[cfg(target_os = "macos")] if let Some(ticket) = macos_ticket { push_mac_completion( @@ -2202,6 +2283,150 @@ mod unix_tests { assert_eq!(recovered_fsp_plaintext, fsp_plaintext); }); } + + /// The seal a refused job gets on the main loop produces a packet the + /// canonical receive-side decoders accept, inner FSP layer included. + #[test] + fn an_inline_seal_matches_the_worker_wire_layout() { + use crate::NodeAddr; + use crate::noise::TAG_SIZE; + use crate::proto::fmp::wire::{EncryptedHeader, build_established_header}; + use crate::proto::fsp::wire::build_fsp_header; + use crate::proto::link::{ + LinkMessageType, SESSION_DATAGRAM_HEADER_SIZE, SessionDatagramRef, + }; + use crate::utils::index::SessionIndex; + + let rt = tokio::runtime::Builder::new_current_thread() + .enable_io() + .build() + .expect("tokio rt"); + let _enter = rt.enter(); + let socket = UdpRawSocket::open("127.0.0.1:0".parse().unwrap(), 1 << 20, 1 << 20) + .expect("open send socket") + .into_async() + .expect("into_async"); + + let fmp_cipher = test_cipher(0x31); + let fsp_cipher = test_cipher(0x32); + let (fmp_counter, fsp_counter) = (900u64, 77u64); + let fsp_plaintext = b"sealed on the main loop".to_vec(); + let link_plaintext_len = + SESSION_DATAGRAM_HEADER_SIZE + FSP_HEADER_SIZE + fsp_plaintext.len(); + let fmp_inner_len = 4 + link_plaintext_len + TAG_SIZE; + let fsp_header = build_fsp_header(fsp_counter, 0, fsp_plaintext.len() as u16); + let fmp_header = + build_established_header(SessionIndex::new(5), fmp_counter, 0, fmp_inner_len as u16); + + let mut wire_buf = Vec::with_capacity(ESTABLISHED_HEADER_SIZE + fmp_inner_len + TAG_SIZE); + wire_buf.extend_from_slice(&fmp_header); + wire_buf.extend_from_slice(&7u32.to_le_bytes()); + wire_buf.push(LinkMessageType::SessionDatagram.to_byte()); + wire_buf.push(16); + wire_buf.extend_from_slice(&1280u16.to_le_bytes()); + wire_buf.extend_from_slice(NodeAddr::from_bytes([0xAA; 16]).as_bytes()); + wire_buf.extend_from_slice(NodeAddr::from_bytes([0xBB; 16]).as_bytes()); + let aad_offset = wire_buf.len(); + wire_buf.extend_from_slice(&fsp_header); + let plaintext_offset = wire_buf.len(); + wire_buf.extend_from_slice(&fsp_plaintext); + + let job = FmpSendJob { + cipher: fmp_cipher.clone(), + counter: fmp_counter, + wire_buf, + fsp_seal: Some(FspSealJob { + cipher: fsp_cipher.clone(), + counter: fsp_counter, + aad_offset, + plaintext_offset, + }), + socket: socket.clone(), + dest_addr: "127.0.0.1:9".parse().unwrap(), + #[cfg(any(target_os = "linux", target_os = "macos"))] + connected_socket: None, + drop_on_backpressure: true, + queued_at: None, + }; + let wire = job.seal_inline().expect("inline seal"); + + let parsed = EncryptedHeader::parse(&wire).expect("FMP header parses"); + assert_eq!(parsed.counter, fmp_counter); + let fmp_plaintext = crate::noise::open( + Some(&fmp_cipher), + fmp_counter, + &parsed.header_bytes, + &wire[ESTABLISHED_HEADER_SIZE..], + ) + .expect("FMP open"); + let datagram = SessionDatagramRef::decode(&fmp_plaintext[5..]).expect("datagram decodes"); + let inner = crate::noise::open( + Some(&fsp_cipher), + fsp_counter, + &datagram.payload[..FSP_HEADER_SIZE], + &datagram.payload[FSP_HEADER_SIZE..], + ) + .expect("FSP open"); + assert_eq!(inner, fsp_plaintext); + + // FMP-only twin: a link message carries no inner seal. + let link_plaintext = b"\x10link message".to_vec(); + let header = build_established_header( + SessionIndex::new(5), + fmp_counter + 1, + 0, + (4 + link_plaintext.len()) as u16, + ); + let mut wire_buf = Vec::with_capacity(ESTABLISHED_HEADER_SIZE + 4 + 64); + wire_buf.extend_from_slice(&header); + wire_buf.extend_from_slice(&9u32.to_le_bytes()); + wire_buf.extend_from_slice(&link_plaintext); + let job = FmpSendJob { + cipher: fmp_cipher.clone(), + counter: fmp_counter + 1, + wire_buf, + fsp_seal: None, + socket, + dest_addr: "127.0.0.1:9".parse().unwrap(), + #[cfg(any(target_os = "linux", target_os = "macos"))] + connected_socket: None, + drop_on_backpressure: false, + queued_at: None, + }; + let wire = job.seal_inline().expect("inline seal"); + let opened = crate::noise::open( + Some(&fmp_cipher), + fmp_counter + 1, + &wire[..ESTABLISHED_HEADER_SIZE], + &wire[ESTABLISHED_HEADER_SIZE..], + ) + .expect("FMP open"); + assert_eq!(&opened[4..], &link_plaintext[..]); + } + + /// A job whose offsets do not fit its buffer is refused, not indexed out + /// of bounds. On the main loop a panic here would end the node. + #[test] + fn a_job_whose_layout_does_not_fit_is_refused_not_a_panic() { + let cipher = test_cipher(9); + let mut short = vec![0u8; ESTABLISHED_HEADER_SIZE - 1]; + assert!(matches!( + seal_wire(&cipher, 1, &mut short, None), + Err(SealError::Layout) + )); + + let mut buf = vec![0u8; 64]; + let overflowing = FspSealJob { + cipher: test_cipher(8), + counter: 1, + aad_offset: usize::MAX - 2, + plaintext_offset: 40, + }; + assert!(matches!( + seal_wire(&cipher, 1, &mut buf, Some(overflowing)), + Err(SealError::Layout) + )); + } } /// Standalone tests for the GSO-eligibility predicate. The full @@ -2581,7 +2806,7 @@ mod mac_queue_tests { fn spawn_pusher( tx: MacWorkerSender, item: T, - ) -> mpsc::Receiver> { + ) -> mpsc::Receiver>> { let (done_tx, done_rx) = mpsc::channel(); thread::spawn(move || { let result = tx.push_blocking(item); @@ -2621,7 +2846,10 @@ mod mac_queue_tests { let result = done .recv_timeout(WAIT) .expect("push_blocking still blocked after the worker thread died"); - assert!(matches!(result, Err(MacWorkerPushError))); + match result { + Err(MacWorkerPushError(job)) => assert_eq!(*job, 3, "the refused job comes back"), + Ok(()) => panic!("push_blocking queued onto a dead worker"), + } assert!(worker.join().is_err(), "worker thread should have panicked"); } @@ -2629,12 +2857,18 @@ mod mac_queue_tests { fn try_push_returns_closed_after_receiver_dropped() { let (tx, rx) = mac_worker_channel::(2); drop(rx); - assert!(matches!(tx.try_push(1), Err(MacWorkerTryPushError::Closed))); + match tx.try_push(1) { + Err(MacWorkerTryPushError::Closed(job)) => assert_eq!(*job, 1), + _ => panic!("try_push on a closed queue should hand the job back"), + } let done = spawn_pusher(tx, 2); let result = done .recv_timeout(WAIT) .expect("push_blocking blocked on a queue whose receiver is gone"); - assert!(matches!(result, Err(MacWorkerPushError))); + match result { + Err(MacWorkerPushError(job)) => assert_eq!(*job, 2), + Ok(()) => panic!("push_blocking queued onto a closed queue"), + } } #[test] @@ -2824,9 +3058,11 @@ mod mac_ordered_tests { assert!(tx.try_push(rig.sequenced(1)).is_ok()); assert!(tx.try_push(rig.sequenced(2)).is_ok()); drop(rx); + // The refused job comes back and is dropped here, which releases its + // slot as the dispatcher's caller would. assert!(matches!( tx.try_push(rig.sequenced(3)), - Err(MacWorkerTryPushError::Closed) + Err(MacWorkerTryPushError::Closed(_)) )); let mut batch = vec![rig.sequenced(4)]; flush_batch_sync(&mut batch).expect("flush"); @@ -3049,7 +3285,10 @@ mod pool_tests { assert_eq!(pool.liveness().live_workers(), 1); let recv = receiver_on_worker(&pool, 0); - pool.dispatch(rig.job(recv.local_addr().unwrap(), 1)); + assert!( + pool.dispatch(rig.job(recv.local_addr().unwrap(), 1)) + .is_ok() + ); let mut buf = [0u8; 128]; recv.recv_from(&mut buf) .expect("the live worker did not send the job dispatched to it"); @@ -3062,9 +3301,9 @@ mod pool_tests { let rig = Rig::new(); let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); let recv = receiver_on_worker(&pool, 1); - let ((), logs) = crate::testutil::capture_logs(|| { - pool.dispatch(rig.job(recv.local_addr().unwrap(), 1)); - }); + let (dispatched, logs) = + crate::testutil::capture_logs(|| pool.dispatch(rig.job(recv.local_addr().unwrap(), 1))); + assert!(dispatched.is_err(), "a dead worker's job must come back"); let warnings = logs.warnings(); assert_eq!(warnings.len(), 1, "{warnings:?}"); assert!(warnings[0].contains(" worker=1"), "{warnings:?}"); @@ -3091,15 +3330,14 @@ mod pool_tests { let recv = UdpSocket::bind("127.0.0.1:0").expect("bind receiver"); let dest = recv.local_addr().unwrap(); for counter in 0..WORKER_CHANNEL_CAP as u64 { - pool.dispatch(rig.job(dest, counter)); + assert!(pool.dispatch(rig.job(dest, counter)).is_ok()); } - let (done_tx, done_rx) = mpsc::channel::<()>(); + let (done_tx, done_rx) = mpsc::channel::(); let blocked_pool = pool.clone(); let last = rig.job(dest, WORKER_CHANNEL_CAP as u64); std::thread::spawn(move || { - blocked_pool.dispatch(last); - let _ = done_tx.send(()); + let _ = done_tx.send(blocked_pool.dispatch(last).is_ok()); }); std::thread::sleep(Duration::from_millis(200)); assert!( @@ -3108,9 +3346,57 @@ mod pool_tests { ); release_tx.send(()).expect("worker gone before release"); - done_rx + let queued = done_rx .recv_timeout(Duration::from_secs(5)) .expect("dispatch still blocked after the worker drained"); + assert!(queued, "the blocked job was refused, not queued"); assert_eq!(pool.liveness().refused_dispatches(), 0); } + + /// A job refused by a dead worker comes back whole, with the counter and + /// buffer the caller reserved, whether the worker was already gone or + /// died while the dispatch waited on its full queue. + #[cfg(not(target_os = "macos"))] + #[test] + fn dispatch_to_an_exited_worker_hands_the_job_back() { + let rig = Rig::new(); + let dest: SocketAddr = "127.0.0.1:9".parse().unwrap(); + + let pool = EncryptWorkerPool::for_test(vec![TestWorker::FailSpawn]); + let job = rig.job(dest, 41); + let wire_buf = job.wire_buf.clone(); + let back = match pool.dispatch(job) { + Err(back) => back, + Ok(()) => panic!("a dead worker took the job"), + }; + assert_eq!(back.counter, 41); + assert_eq!(back.wire_buf, wire_buf); + + // Worker alive but not draining; it exits while a dispatch waits. + let (exit_tx, exit_rx) = mpsc::channel::<()>(); + let pool = EncryptWorkerPool::for_test(vec![TestWorker::ExitOn(exit_rx)]); + for counter in 0..WORKER_CHANNEL_CAP as u64 { + assert!(pool.dispatch(rig.job(dest, counter)).is_ok()); + } + let (done_tx, done_rx) = mpsc::channel::>(); + let blocked_pool = pool.clone(); + let last = rig.job(dest, 7_000); + std::thread::spawn(move || { + let _ = done_tx.send(blocked_pool.dispatch(last).err().map(|job| job.counter)); + }); + std::thread::sleep(Duration::from_millis(100)); + assert!( + matches!(done_rx.try_recv(), Err(mpsc::TryRecvError::Empty)), + "dispatch returned while the queue was full and the worker alive" + ); + exit_tx.send(()).expect("worker gone before its signal"); + let back = done_rx + .recv_timeout(Duration::from_secs(5)) + .expect("dispatch still blocked after the worker exited"); + assert_eq!( + back, + Some(7_000), + "the job blocked on a dying worker must come back" + ); + } } diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index 4d4c1f7b..932eea89 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -2875,7 +2875,7 @@ impl Node { entry.touch(send.now_ms); } - workers.dispatch(crate::node::encrypt_worker::FmpSendJob { + let dispatched = workers.dispatch(crate::node::encrypt_worker::FmpSendJob { cipher: fmp_cipher, counter: fmp_counter, wire_buf, @@ -2895,9 +2895,44 @@ impl Node { drop_on_backpressure: true, queued_at: None, }); + if let Err(job) = dispatched { + self.send_refused_job_inline(*job, transport_id, &remote_addr, next_hop_addr) + .await; + } Ok(true) } + /// Seal and send a job the encrypt worker for its next hop refused + /// because that worker has exited, using the FSP and FMP counters the + /// job already reserved so neither counter is skipped. Stats were + /// recorded before dispatch and now describe this packet. + /// + /// A failure is logged and swallowed, as the worker does with its own: + /// the caller sees the same result whichever of the two sent the packet. + #[cfg(unix)] + async fn send_refused_job_inline( + &self, + job: crate::node::encrypt_worker::FmpSendJob, + transport_id: crate::transport::TransportId, + remote_addr: &crate::transport::TransportAddr, + next_hop_addr: NodeAddr, + ) { + let wire = match job.seal_inline() { + Ok(wire) => wire, + Err(error) => { + debug!(next_hop = %next_hop_addr, %error, "Inline seal of session data failed"); + return; + } + }; + let Some(transport) = self.transports.get(&transport_id) else { + debug!(next_hop = %next_hop_addr, "Transport gone before inline send of session data"); + return; + }; + if let Err(error) = transport.send(remote_addr, &wire).await { + debug!(next_hop = %next_hop_addr, %error, "Inline send of session data failed"); + } + } + /// Send an IPv6 packet through the IPv6 shim (port 256) with header compression. /// /// Compresses the IPv6 header (format 0x00), then sends via `send_session_data` diff --git a/src/node/lifecycle/workers.rs b/src/node/lifecycle/workers.rs index 98f6fe8a..42a4558b 100644 --- a/src/node/lifecycle/workers.rs +++ b/src/node/lifecycle/workers.rs @@ -3,7 +3,8 @@ //! exited. //! //! The pools are a performance offload, and Windows never starts them at -//! all, so losing workers is `Degraded` at most and never fatal. Inbound +//! all, so losing workers is `Degraded` at most and never fatal. An outbound +//! packet for a missing encrypt worker is sealed on the main loop. Inbound //! packets for a session already held by a missing decrypt worker are //! dropped until the session rekeys or the link is re-established; a session //! that would be registered on the missing worker after the loss is decrypted diff --git a/src/node/mod.rs b/src/node/mod.rs index 7eda46b7..9827d083 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -3914,7 +3914,7 @@ impl Node { // Drop bulk endpoint data on UDP backpressure to // keep the queue moving; control frames retry. let drop_on_backpressure = plaintext.first().is_some_and(|t| *t == 0x00); - workers.dispatch(crate::node::encrypt_worker::FmpSendJob { + let dispatched = workers.dispatch(crate::node::encrypt_worker::FmpSendJob { cipher: fmp_cipher, counter, wire_buf, @@ -3926,12 +3926,27 @@ impl Node { drop_on_backpressure, queued_at: None, }); + let sent_bytes = match dispatched { + Ok(()) => predicted_bytes, + // The worker for this destination has exited. Seal + // here with the counter already reserved, so no + // counter is skipped, and send as the inline path does. + Err(job) => { + let wire = job.seal_inline().map_err(|e| NodeError::SendFailed { + node_addr: *node_addr, + reason: format!("encryption failed: {}", e), + })?; + transport + .send(&remote_addr, &wire) + .await + .map_err(|e| link_send_error(*node_addr, e))? + } + }; if let Some(peer) = self.peers.get_mut(node_addr) { - peer.link_stats_mut().record_sent(predicted_bytes); + peer.link_stats_mut().record_sent(sent_bytes); if let Some(mmp) = peer.mmp_mut() { - mmp.sender - .record_sent(counter, timestamp_ms, predicted_bytes); + mmp.sender.record_sent(counter, timestamp_ms, sent_bytes); } } return Ok(()); @@ -3989,25 +4004,7 @@ impl Node { let bytes_sent = transport .send(&remote_addr, &wire_packet) .await - .map_err(|e| match e { - TransportError::MtuExceeded { packet_size, mtu } => NodeError::MtuExceeded { - node_addr: *node_addr, - packet_size, - mtu, - }, - // Preserve the transport's own classification instead of - // flattening every non-MTU failure into one string. A caller - // that wants to keep its half-built state across an interface - // flap can only do that if the distinction survives to it. - other if other.is_transient() => NodeError::SendUnavailable { - node_addr: *node_addr, - reason: format!("transport send: {}", other), - }, - other => NodeError::SendFailed { - node_addr: *node_addr, - reason: format!("transport send: {}", other), - }, - })?; + .map_err(|e| link_send_error(*node_addr, e))?; // Update send statistics if let Some(peer) = self.peers.get_mut(node_addr) { @@ -4022,6 +4019,30 @@ impl Node { } } +/// Map a transport's refusal of an encrypted link frame to `node_addr` onto +/// the error the link-send path reports. +fn link_send_error(node_addr: NodeAddr, e: TransportError) -> NodeError { + match e { + TransportError::MtuExceeded { packet_size, mtu } => NodeError::MtuExceeded { + node_addr, + packet_size, + mtu, + }, + // Preserve the transport's own classification instead of + // flattening every non-MTU failure into one string. A caller + // that wants to keep its half-built state across an interface + // flap can only do that if the distinction survives to it. + other if other.is_transient() => NodeError::SendUnavailable { + node_addr, + reason: format!("transport send: {}", other), + }, + other => NodeError::SendFailed { + node_addr, + reason: format!("transport send: {}", other), + }, + } +} + /// Shell-side [`routing::RoutingView`] seam over live `Node` state — the sole /// routing read adapter the shell retains. It hands the sans-IO routing core /// borrowed peers plus raw `may_reach` / `link_cost` / `coords` diff --git a/src/node/tests/handshake.rs b/src/node/tests/handshake.rs index b735a9e6..5f17cb12 100644 --- a/src/node/tests/handshake.rs +++ b/src/node/tests/handshake.rs @@ -2377,3 +2377,209 @@ async fn a_transient_msg2_failure_on_the_restart_path_leaves_the_fresh_leg_pendi "a local interface flap must not be recorded as the peer's misbehaviour" ); } + +/// Link messages to a peer whose encrypt worker has exited are still sent, +/// sealed on the main loop with the counter the worker path reserved. +#[cfg(unix)] +mod dead_encrypt_worker { + use super::*; + use crate::config::UdpConfig; + use crate::node::encrypt_worker::EncryptWorkerPool; + use crate::node::tests::pool_with_dead_worker_for; + use crate::node::worker_set::TestWorker; + use crate::proto::fmp::wire::{EncryptedHeader, build_msg1}; + use crate::transport::ReceivedPacket; + use crate::transport::udp::UdpTransport; + use tokio::time::{Duration, timeout}; + + /// Two nodes on real UDP sockets with an established link from A to B. + struct Pair { + node_a: Node, + node_b: Node, + packet_rx_b: crate::transport::PacketRx, + peer_a: NodeAddr, + peer_b: NodeAddr, + addr_b: std::net::SocketAddr, + } + + async fn linked_pair() -> Pair { + let mut node_a = make_node(); + let mut node_b = make_node(); + let tid = TransportId::new(1); + let udp_config = UdpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + mtu: Some(1280), + ..Default::default() + }; + let (packet_tx_a, mut packet_rx_a) = packet_channel(64); + let (packet_tx_b, mut packet_rx_b) = packet_channel(64); + let mut transport_a = UdpTransport::new(tid, None, udp_config.clone(), packet_tx_a); + let mut transport_b = UdpTransport::new(tid, None, udp_config, packet_tx_b); + transport_a.start_async().await.unwrap(); + transport_b.start_async().await.unwrap(); + let addr_b = transport_b.local_addr().unwrap(); + let remote_addr_b = TransportAddr::from_string(&addr_b.to_string()); + node_a + .transports + .insert(tid, TransportHandle::Udp(transport_a)); + node_b + .transports + .insert(tid, TransportHandle::Udp(transport_b)); + + let peer_b_identity = PeerIdentity::from_pubkey_full(node_b.identity().pubkey_full()); + let peer_b = *peer_b_identity.node_addr(); + let peer_a = *PeerIdentity::from_pubkey_full(node_a.identity().pubkey_full()).node_addr(); + let link_id = node_a.allocate_link_id(); + let our_index = node_a.index_allocator.allocate().unwrap(); + node_a + .seed_handshake_machine( + HandshakeSeed::outbound(link_id, peer_b_identity, 1000) + .with_our_index(our_index) + .with_transport_id(tid) + .with_source_addr(remote_addr_b.clone()), + ) + .unwrap(); + let keypair = node_a.identity().keypair(); + let epoch = node_a.startup_epoch(); + let msg1 = node_a + .peer_machines + .get_mut(&link_id) + .unwrap() + .start_handshake(keypair, epoch, 1000) + .unwrap(); + node_a.links.insert( + link_id, + Link::connectionless( + link_id, + tid, + remote_addr_b.clone(), + LinkDirection::Outbound, + Duration::from_millis(100), + ), + ); + node_a + .pending_outbound + .insert((tid, our_index.as_u32()), link_id); + node_a + .transports + .get(&tid) + .unwrap() + .send(&remote_addr_b, &build_msg1(our_index, &msg1)) + .await + .expect("send msg1"); + + let msg1 = next_packet(&mut packet_rx_b).await; + node_b.handle_msg1(msg1).await; + let msg2 = next_packet(&mut packet_rx_a).await; + node_a.handle_msg2(msg2).await; + assert!(node_a.get_peer(&peer_b).is_some(), "A promoted B"); + assert!(node_b.get_peer(&peer_a).is_some(), "B promoted A"); + + Pair { + node_a, + node_b, + packet_rx_b, + peer_a, + peer_b, + addr_b, + } + } + + async fn next_packet(rx: &mut crate::transport::PacketRx) -> ReceivedPacket { + timeout(Duration::from_secs(5), rx.recv()) + .await + .expect("no packet within the bound") + .expect("packet channel closed") + } + + /// Send one link message from A to B through `pool` and have B process + /// it. Returns (the frame's FMP counter, A's send counter before the + /// send, A's packets-sent delta, B's packets-received delta, encrypt + /// WARN lines). + async fn send_through( + pair: &mut Pair, + pool: &EncryptWorkerPool, + ) -> (u64, u64, u64, u64, Vec) { + // Let B take whatever A sent on promotion before measuring. + while let Ok(Some(packet)) = + timeout(Duration::from_millis(100), pair.packet_rx_b.recv()).await + { + pair.node_b.handle_encrypted_frame(packet).await; + } + pair.node_a.supervisor.encrypt_workers = Some(pool.clone()); + let peer = pair.node_a.get_peer(&pair.peer_b).unwrap(); + let counter_before = peer.noise_session().unwrap().current_send_counter(); + let sent_before = peer.link_stats().packets_sent; + let recv_before = pair + .node_b + .get_peer(&pair.peer_a) + .unwrap() + .link_stats() + .packets_recv; + + let (logs, guard) = crate::testutil::capture_logs_scoped(); + pair.node_a + .send_encrypted_link_message(&pair.peer_b, b"\x10dead worker test") + .await + .expect("link message send"); + drop(guard); + + let packet = next_packet(&mut pair.packet_rx_b).await; + let frame_counter = EncryptedHeader::parse(&packet.data) + .expect("an established frame") + .counter; + pair.node_b.handle_encrypted_frame(packet).await; + + let sent = pair + .node_a + .get_peer(&pair.peer_b) + .unwrap() + .link_stats() + .packets_sent + - sent_before; + let recv = pair + .node_b + .get_peer(&pair.peer_a) + .unwrap() + .link_stats() + .packets_recv + - recv_before; + let warnings = logs + .warnings() + .into_iter() + .filter(|line| line.contains("pool=\"encrypt\"")) + .collect(); + (frame_counter, counter_before, sent, recv, warnings) + } + + #[tokio::test] + async fn a_link_message_for_a_dead_workers_peer_is_still_sent() { + let mut pair = linked_pair().await; + let pool = pool_with_dead_worker_for(pair.addr_b); + let (frame_counter, counter_before, sent, recv, warnings) = + send_through(&mut pair, &pool).await; + + assert_eq!(recv, 1, "B did not authenticate the frame"); + assert_eq!( + frame_counter, counter_before, + "the frame must carry the counter reserved for it, not a fresh one" + ); + assert_eq!(sent, 1, "A counted the packet other than once"); + assert_eq!(pool.liveness().refused_dispatches(), 1); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + } + + #[tokio::test] + async fn a_link_message_through_live_workers_is_sent_by_the_worker() { + let mut pair = linked_pair().await; + let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]); + let (frame_counter, counter_before, sent, recv, warnings) = + send_through(&mut pair, &pool).await; + + assert_eq!(recv, 1, "B did not authenticate the frame"); + assert_eq!(frame_counter, counter_before); + assert_eq!(sent, 1); + assert_eq!(pool.liveness().refused_dispatches(), 0); + assert!(warnings.is_empty(), "{warnings:?}"); + } +} diff --git a/src/node/tests/mod.rs b/src/node/tests/mod.rs index 22b838e7..e6d5a681 100644 --- a/src/node/tests/mod.rs +++ b/src/node/tests/mod.rs @@ -94,6 +94,31 @@ pub(super) fn install_connected_udp( .set_connected_udp(socket, drain); } +/// An encrypt pool of two whose worker for `dest` has exited. Where the worker +/// for a destination is not a function of the address alone (macOS), both +/// have. +#[cfg(unix)] +pub(super) fn pool_with_dead_worker_for( + dest: std::net::SocketAddr, +) -> crate::node::encrypt_worker::EncryptWorkerPool { + use crate::node::encrypt_worker::EncryptWorkerPool; + use crate::node::worker_set::TestWorker; + + let dead = EncryptWorkerPool::for_test(vec![TestWorker::FailSpawn, TestWorker::FailSpawn]) + .worker_index_for_dest(dest); + EncryptWorkerPool::for_test( + (0..2) + .map(|idx| { + if dead.is_none_or(|d| d == idx) { + TestWorker::FailSpawn + } else { + TestWorker::Run + } + }) + .collect(), + ) +} + /// Build a test node with an explicit `max_peers` limit (replaces the removed /// `set_max_peers` setter; resource limits are immutable post-construction). pub(super) fn make_node_with_max_peers(max_peers: usize) -> Node { diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index 8bb70718..3d22725f 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -9016,3 +9016,110 @@ async fn a_path_broken_flood_releases_the_stored_path_mtu_only_once_per_interval "a second release for the same destination inside the interval is refused" ); } + +/// Session data whose next hop's encrypt worker has exited is still +/// delivered, sealed on the main loop with the FSP and FMP counters the +/// worker path reserved. +#[cfg(unix)] +mod dead_encrypt_worker { + use super::*; + use crate::config::UdpConfig; + use crate::node::tests::pool_with_dead_worker_for; + use crate::transport::udp::UdpTransport; + use crate::transport::{TransportAddr, TransportHandle, TransportId, packet_channel}; + + /// A test node on a real UDP socket: the worker path needs one. + async fn make_test_node_udp() -> TestNode { + let mut node = make_node(); + let transport_id = TransportId::new(1); + let config = UdpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + mtu: Some(1280), + ..Default::default() + }; + let (packet_tx, packet_rx) = packet_channel(256); + let mut transport = UdpTransport::new(transport_id, None, config, packet_tx); + transport.start_async().await.unwrap(); + let addr = TransportAddr::from_string(&transport.local_addr().unwrap().to_string()); + node.transports + .insert(transport_id, TransportHandle::Udp(transport)); + TestNode { + node, + transport_id, + packet_rx: crate::node::tests::spanning_tree::bridge_to_unbounded(packet_rx), + addr, + } + } + + fn fsp_send_counter(node: &Node, dest: &NodeAddr) -> u64 { + match node.get_session(dest).expect("session").state() { + EndToEndState::Established(session) => session.current_send_counter(), + _ => panic!("session not established"), + } + } + + fn fmp_send_counter(node: &Node, peer: &NodeAddr) -> u64 { + node.get_peer(peer) + .expect("peer") + .noise_session() + .expect("link session") + .current_send_counter() + } + + #[tokio::test] + async fn session_data_for_a_dead_workers_next_hop_is_still_delivered() { + let mut nodes = vec![make_test_node_udp().await, make_test_node_udp().await]; + initiate_handshake(&mut nodes, 0, 1).await; + drain_all_packets(&mut nodes, false).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + establish_pair_session(&mut nodes).await; + drain_all_packets(&mut nodes, false).await; + + let node0 = *nodes[0].node.node_addr(); + let node1 = *nodes[1].node.node_addr(); + let dest: std::net::SocketAddr = nodes[1].addr.to_string().parse().unwrap(); + let pool = pool_with_dead_worker_for(dest); + nodes[0].node.supervisor.encrypt_workers = Some(pool.clone()); + + let fsp_before = fsp_send_counter(&nodes[0].node, &node1); + let fmp_before = fmp_send_counter(&nodes[0].node, &node1); + let recv_before = nodes[1] + .node + .get_session(&node0) + .unwrap() + .traffic_counters() + .1; + + nodes[0] + .node + .send_session_data(&node1, 0, 0, b"for a dead worker") + .await + .expect("send_session_data"); + // Read before anything else runs on A: one packet, one counter each. + assert_eq!(fsp_send_counter(&nodes[0].node, &node1), fsp_before + 1); + assert_eq!(fmp_send_counter(&nodes[0].node, &node1), fmp_before + 1); + assert_eq!(pool.liveness().refused_dispatches(), 1); + + let delivered = |nodes: &[TestNode]| { + nodes[1] + .node + .get_session(&node0) + .unwrap() + .traffic_counters() + .1 + }; + let deadline = tokio::time::Instant::now() + Duration::from_secs(5); + while delivered(&nodes) == recv_before && tokio::time::Instant::now() < deadline { + tokio::time::sleep(Duration::from_millis(10)).await; + process_available_packets(&mut nodes).await; + } + assert_eq!( + delivered(&nodes), + recv_before + 1, + "the payload never reached the destination" + ); + + cleanup_nodes(&mut nodes).await; + } +} From 3a6b7fcea1296d5c5c2e0d7e7b9807e9bae095c2 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 2 Oct 2026 16:17:51 +0000 Subject: [PATCH 8/8] Changelog: crypto worker deaths reported, and their packets sent from the main loop --- CHANGELOG.md | 21 +++++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index cf9e5865..67d41c34 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -475,6 +475,27 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 per-peer `connect()`-ed UDP socket stayed pinned to the old 5-tuple; the sibling path already cleared it. +- A crypto worker thread that exits now shows: the node reports `Degraded`, + and a warning names the pool and how many of its workers are still live, + logged again each time another one goes. A worker thread that cannot be + started is logged and leaves the node degraded instead of stopping start-up. + Losing workers never fails the node. Outbound packets for a lost encrypt + worker are sent from the main loop, as described below. Inbound packets for + a session held by a lost decrypt worker are still dropped, now with a + rate-limited warning, until the session rekeys or the link is re-established; + a session that would be placed on the lost worker after the loss is + decrypted on the main loop instead. + +- On every Unix platform, a packet whose encrypt worker thread has exited is + now encrypted and sent from the main loop. It was dropped, while the link + statistics counted it as sent, so every peer whose traffic went to that + worker was cut off until the daemon restarted. The packet keeps the counters + reserved for it, so the receiver sees no gap. Each such packet is counted and + logged at WARN, rate-limited, in place of a debug line. This includes macOS, + where the fix listed under macOS below dropped these packets; only packets + already queued to the exited worker are lost. With the macOS ordered sender, + a packet sent this way can arrive out of order with the rest of its flow. + #### Data plane and transports - A peer that stops reading can no longer stall the node. TCP, Tor, Nym and