Merge branch 'master' into next

This commit is contained in:
Johnathan Corgan
2026-10-02 18:45:24 +00:00
30 changed files with 2598 additions and 368 deletions
+21
View File
@@ -632,6 +632,27 @@ with v0.5.x or earlier peers.
per-peer `connect()`-ed UDP socket stayed pinned to the old 5-tuple; per-peer `connect()`-ed UDP socket stayed pinned to the old 5-tuple;
the sibling path already cleared it. 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 #### Data plane and transports
- A peer that stops reading can no longer stall the node. TCP, Tor, Nym and - A peer that stops reading can no longer stall the node. TCP, Tor, Nym and
+10 -3
View File
@@ -41,7 +41,7 @@ impl Domain {
/// ///
/// `as usize` indexes the counter arrays, so the discriminants are dense and /// `as usize` indexes the counter arrays, so the discriminants are dense and
/// `WholeTick` is last (it defines `N_STEPS`). Variants are declared /// `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. /// profile-conditional steps are kept out of the emitted table.
#[derive(Copy, Clone, Debug, PartialEq, Eq)] #[derive(Copy, Clone, Debug, PartialEq, Eq)]
#[repr(usize)] #[repr(usize)]
@@ -76,6 +76,7 @@ pub(crate) enum Step {
PollTransportDiscovery, PollTransportDiscovery,
SampleTransportCongestion, SampleTransportCongestion,
ActivateConnectedUdpSessions, ActivateConnectedUdpSessions,
PollWorkerLiveness,
DebugAssertPeerMapsCoherent, DebugAssertPeerMapsCoherent,
/// The whole tick-arm body, from before `check_timeouts` to after the last /// The whole tick-arm body, from before `check_timeouts` to after the last
/// step. Composes safely with the per-step spans because the macro /// step. Composes safely with the per-step spans because the macro
@@ -116,6 +117,7 @@ pub(crate) const STEPS: [Step; N_STEPS] = [
Step::PollTransportDiscovery, Step::PollTransportDiscovery,
Step::SampleTransportCongestion, Step::SampleTransportCongestion,
Step::ActivateConnectedUdpSessions, Step::ActivateConnectedUdpSessions,
Step::PollWorkerLiveness,
Step::DebugAssertPeerMapsCoherent, Step::DebugAssertPeerMapsCoherent,
Step::WholeTick, Step::WholeTick,
]; ];
@@ -151,6 +153,7 @@ impl Step {
Step::PollTransportDiscovery => "poll_transport_discovery", Step::PollTransportDiscovery => "poll_transport_discovery",
Step::SampleTransportCongestion => "sample_transport_congestion", Step::SampleTransportCongestion => "sample_transport_congestion",
Step::ActivateConnectedUdpSessions => "activate_connected_udp_sessions", Step::ActivateConnectedUdpSessions => "activate_connected_udp_sessions",
Step::PollWorkerLiveness => "poll_worker_liveness",
Step::DebugAssertPeerMapsCoherent => "debug_assert_peer_maps_coherent", Step::DebugAssertPeerMapsCoherent => "debug_assert_peer_maps_coherent",
Step::WholeTick => "whole_tick", Step::WholeTick => "whole_tick",
} }
@@ -158,7 +161,7 @@ impl Step {
/// Whether this step gets a row in this build. /// 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 /// 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 /// 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 /// runs" rather than "this step is not in this build". The predicates below
@@ -169,6 +172,7 @@ impl Step {
Step::ActivateConnectedUdpSessions => { Step::ActivateConnectedUdpSessions => {
cfg!(any(target_os = "linux", target_os = "macos")) cfg!(any(target_os = "linux", target_os = "macos"))
} }
Step::PollWorkerLiveness => cfg!(unix),
Step::DebugAssertPeerMapsCoherent => cfg!(debug_assertions), Step::DebugAssertPeerMapsCoherent => cfg!(debug_assertions),
_ => true, _ => true,
} }
@@ -388,7 +392,7 @@ mod tests {
let emitted = STEPS.iter().filter(|s| s.emitted()).count(); let emitted = STEPS.iter().filter(|s| s.emitted()).count();
// 27 unconditional subsystem steps on this line (26 shared with the // 27 unconditional subsystem steps on this line (26 shared with the
// master line, plus `resend_pending_fmp_rekey_msg3`, which exists only // master line, plus `resend_pending_fmp_rekey_msg3`, which exists only
// here) + the whole-tick span, plus the two conditionally-compiled // here) + the whole-tick span, plus the three conditionally-compiled
// steps where this build has them. The count is pinned deliberately: it // steps where this build has them. The count is pinned deliberately: it
// is what caught the extra step when the master-line instrumentation // is what caught the extra step when the master-line instrumentation
// was merged up, rather than letting the tables silently disagree — and // was merged up, rather than letting the tables silently disagree — and
@@ -402,6 +406,9 @@ mod tests {
if cfg!(any(target_os = "linux", target_os = "macos")) { if cfg!(any(target_os = "linux", target_os = "macos")) {
expected += 1; expected += 1;
} }
if cfg!(unix) {
expected += 1;
}
if cfg!(debug_assertions) { if cfg!(debug_assertions) {
expected += 1; expected += 1;
} }
+6 -7
View File
@@ -55,13 +55,12 @@ impl Node {
self.try_warm_coord_cache_ref(&datagram_ref, payload.len()); self.try_warm_coord_cache_ref(&datagram_ref, payload.len());
// Pre-resolve the next hop only for datagrams the core can actually // Pre-resolve the next hop only for datagrams the core can actually
// forward: not locally destined, and carrying a TTL that survives the // forward: not locally destined, and passing `can_forward`, which is
// decrement (`ttl > 1` — the shell-side mirror of the core's // the core's own hop-limit rule. This keeps `find_next_hop`'s
// would-leave-zero drop). This keeps `find_next_hop`'s coord-cache // coord-cache LRU-touch side effect scoped to genuine forwards, as it
// LRU-touch side effect scoped to genuine forwards, as it was when the // was when the TTL test ran inline ahead of it. Warming above has
// TTL test ran inline ahead of it. Warming above has already run, so // already run, so the resolution observes freshly cached coords.
// the resolution observes freshly cached coords. let next_hop = if datagram_ref.dest_addr != my_addr && datagram_ref.can_forward() {
let next_hop = if datagram_ref.dest_addr != my_addr && datagram_ref.ttl > 1 {
self.resolve_next_hop(&datagram_ref.dest_addr) self.resolve_next_hop(&datagram_ref.dest_addr)
} else { } else {
None None
+5 -11
View File
@@ -345,17 +345,7 @@ impl Node {
// republishing health, so nothing outside the node can // republishing health, so nothing outside the node can
// observe an address the listener no longer answers on. // observe an address the listener no longer answers on.
self.retract_child_publications(child); self.retract_child_publications(child);
let actions = self self.step_child_exited(child);
.supervisor
.fsm
.step(crate::node::lifecycle::supervisor::Event::ChildExited { child });
for action in actions {
if let crate::node::lifecycle::supervisor::Action::PublishState(ns) =
action
{
self.supervisor.state = ns;
}
}
// A transport child exiting leaves the bound set, so // A transport child exiting leaves the bound set, so
// it can be the one that was holding the node's egress // it can be the one that was holding the node's egress
// MTU down. `is_bound()` is `is_operational()` plus the // MTU down. `is_bound()` is `is_operational()` plus the
@@ -585,6 +575,10 @@ impl Node {
#[cfg(any(target_os = "linux", target_os = "macos"))] #[cfg(any(target_os = "linux", target_os = "macos"))]
instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ActivateConnectedUdpSessions, instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ActivateConnectedUdpSessions,
self.activate_connected_udp_sessions().await); 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 // Debug-build sweep of the peer-lifecycle map invariant
// (leaked machines / machine-less legs); two map scans, // (leaked machines / machine-less legs); two map scans,
// compiled out of release builds. // compiled out of release builds.
+185 -33
View File
@@ -36,6 +36,7 @@
#![cfg_attr(not(unix), allow(dead_code))] #![cfg_attr(not(unix), allow(dead_code))]
use crate::NodeAddr; use crate::NodeAddr;
use crate::node::worker_set::{WorkerLiveness, WorkerSet, worth_logging};
use crate::transport::{TransportAddr, TransportId}; use crate::transport::{TransportAddr, TransportId};
use crossbeam_channel::{Receiver, Sender, TrySendError, bounded}; use crossbeam_channel::{Receiver, Sender, TrySendError, bounded};
use portable_atomic::{AtomicU64, Ordering}; use portable_atomic::{AtomicU64, Ordering};
@@ -220,26 +221,49 @@ pub(crate) enum WorkerMsg {
/// shard. /// shard.
#[derive(Clone)] #[derive(Clone)]
pub(crate) struct DecryptWorkerPool { pub(crate) struct DecryptWorkerPool {
senders: Arc<[Sender<WorkerMsg>]>, workers: Arc<WorkerSet<Sender<WorkerMsg>>>,
}
/// Start the production worker loop on `rx` in a named OS thread.
fn spawn_worker(
idx: usize,
rx: Receiver<WorkerMsg>,
) -> std::io::Result<std::thread::JoinHandle<()>> {
std::thread::Builder::new()
.name(format!("fips-decrypt-{idx}"))
.spawn(move || run_worker(idx, rx))
} }
impl DecryptWorkerPool { 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 { pub fn spawn(n: usize) -> Self {
let n = n.max(1); Self::start_with(n, spawn_worker)
let mut senders = Vec::with_capacity(n);
for i in 0..n {
let (tx, rx) = bounded::<WorkerMsg>(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);
} }
/// 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<WorkerMsg>) -> std::io::Result<std::thread::JoinHandle<()>>,
) -> Self {
Self { Self {
senders: senders.into(), workers: Arc::new(WorkerSet::start(
"decrypt",
n,
|| bounded::<WorkerMsg>(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 /// Stable hash from session key → worker index. Same hash is used
/// for session registration and per-packet dispatch so packets and /// for session registration and per-packet dispatch so packets and
/// registration arrive at the same shard. /// registration arrive at the same shard.
@@ -247,24 +271,37 @@ impl DecryptWorkerPool {
use std::hash::{Hash, Hasher}; use std::hash::{Hash, Hasher};
let mut h = std::collections::hash_map::DefaultHasher::new(); let mut h = std::collections::hash_map::DefaultHasher::new();
cache_key.hash(&mut h); 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 /// Dispatch a per-packet decrypt job. Drops if the per-worker
/// channel is full (sustained rate overrun); the rx_loop's drain /// channel is full (sustained rate overrun); the rx_loop's drain
/// caps inbound at the same scale upstream so the cliff is /// 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) { pub fn dispatch_job(&self, job: DecryptJob) {
if self.senders.is_empty() {
return;
}
let idx = self.worker_idx_for(job.cache_key); 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(()) => {} Ok(()) => {}
Err(TrySendError::Full(_)) => { Err(TrySendError::Full(_)) => {
static FULL_COUNT: AtomicU64 = AtomicU64::new(0); static FULL_COUNT: AtomicU64 = AtomicU64::new(0);
let n = FULL_COUNT.fetch_add(1, Ordering::Relaxed); let n = FULL_COUNT.fetch_add(1, Ordering::Relaxed);
if n < 8 || n.is_multiple_of(10000) { if worth_logging(n) {
warn!( warn!(
worker = idx, worker = idx,
drops = n + 1, drops = n + 1,
@@ -272,9 +309,7 @@ impl DecryptWorkerPool {
); );
} }
} }
Err(TrySendError::Disconnected(_)) => { Err(TrySendError::Disconnected(_)) => self.note_refused(idx, "inbound packet"),
debug!(worker = idx, "DecryptWorker thread gone; dropping job");
}
} }
} }
@@ -303,11 +338,12 @@ impl DecryptWorkerPool {
cache_key: (TransportId, u32), cache_key: (TransportId, u32),
state: OwnedSessionState, state: OwnedSessionState,
) -> bool { ) -> bool {
if self.senders.is_empty() {
return false;
}
let idx = self.worker_idx_for(cache_key); 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, Ok(()) => true,
Err(TrySendError::Full(_)) => { Err(TrySendError::Full(_)) => {
warn!( warn!(
@@ -317,10 +353,7 @@ impl DecryptWorkerPool {
false false
} }
Err(TrySendError::Disconnected(_)) => { Err(TrySendError::Disconnected(_)) => {
debug!( self.note_refused(idx, "session registration");
worker = idx,
"DecryptWorker thread gone; ignoring registration"
);
false false
} }
} }
@@ -329,11 +362,11 @@ impl DecryptWorkerPool {
/// Drop a session from its worker (rekey, peer removed). Fire and /// Drop a session from its worker (rekey, peer removed). Fire and
/// forget — if the worker is gone we don't care. /// forget — if the worker is gone we don't care.
pub fn unregister_session(&self, cache_key: (TransportId, u32)) { pub fn unregister_session(&self, cache_key: (TransportId, u32)) {
if self.senders.is_empty() {
return;
}
let idx = self.worker_idx_for(cache_key); 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 });
} }
} }
@@ -773,3 +806,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<crate::node::worker_set::TestWorker>) -> 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::<DecryptWorkerEvent>();
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);
}
}
+628 -120
View File
@@ -50,6 +50,9 @@
// warnings rather than gate every function individually. // warnings rather than gate every function individually.
#![cfg_attr(not(unix), allow(dead_code))] #![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::fmp::wire::ESTABLISHED_HEADER_SIZE;
use crate::proto::fsp::wire::FSP_HEADER_SIZE; use crate::proto::fsp::wire::FSP_HEADER_SIZE;
use crate::transport::udp::io::AsyncUdpSocket; use crate::transport::udp::io::AsyncUdpSocket;
@@ -178,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<QueuedFmpSendJob>);
impl Refused {
/// The job, with any macOS ordered-flow slot it held released as a skip.
fn into_job(self) -> Box<FmpSendJob> {
#[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- /// Handle to the encrypt worker pool. Dispatches jobs **hash-by-
/// destination** across N worker tasks via per-worker bounded /// destination** across N worker tasks via per-worker bounded
/// crossbeam channels. The bounded queue intentionally backpressures /// crossbeam channels. The bounded queue intentionally backpressures
@@ -239,14 +260,17 @@ struct MacWorkerQueueState<T> {
closed: bool, closed: bool,
} }
/// Why `try_push` did not queue a job. Both variants hand the job back.
#[cfg(any(target_os = "macos", test))] #[cfg(any(target_os = "macos", test))]
enum MacWorkerTryPushError<T> { enum MacWorkerTryPushError<T> {
Full(Box<T>), Full(Box<T>),
Closed, /// The receiver is gone: the worker has exited.
Closed(Box<T>),
} }
/// `push_blocking` found the receiver gone; the job is handed back.
#[cfg(any(target_os = "macos", test))] #[cfg(any(target_os = "macos", test))]
struct MacWorkerPushError; struct MacWorkerPushError<T>(Box<T>);
#[cfg(any(target_os = "macos", test))] #[cfg(any(target_os = "macos", test))]
fn mac_worker_channel<T>(cap: usize) -> (MacWorkerSender<T>, MacWorkerReceiver<T>) { fn mac_worker_channel<T>(cap: usize) -> (MacWorkerSender<T>, MacWorkerReceiver<T>) {
@@ -277,10 +301,10 @@ impl<T> MacWorkerSender<T> {
.lock() .lock()
.expect("encrypt worker queue poisoned"); .expect("encrypt worker queue poisoned");
if state.closed { 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(state);
drop(job); return Err(MacWorkerTryPushError::Closed(Box::new(job)));
return Err(MacWorkerTryPushError::Closed);
} }
if state.queue.len() >= self.inner.cap { if state.queue.len() >= self.inner.cap {
return Err(MacWorkerTryPushError::Full(Box::new(job))); return Err(MacWorkerTryPushError::Full(Box::new(job)));
@@ -295,7 +319,7 @@ impl<T> MacWorkerSender<T> {
Ok(()) Ok(())
} }
fn push_blocking(&self, job: T) -> Result<(), MacWorkerPushError> { fn push_blocking(&self, job: T) -> Result<(), MacWorkerPushError<T>> {
let mut state = self let mut state = self
.inner .inner
.state .state
@@ -304,8 +328,7 @@ impl<T> MacWorkerSender<T> {
loop { loop {
if state.closed { if state.closed {
drop(state); drop(state);
drop(job); return Err(MacWorkerPushError(Box::new(job)));
return Err(MacWorkerPushError);
} }
if state.queue.len() < self.inner.cap { if state.queue.len() < self.inner.cap {
let was_empty = state.queue.is_empty(); let was_empty = state.queue.is_empty();
@@ -406,6 +429,36 @@ type WorkerSender = MacWorkerSender<QueuedFmpSendJob>;
#[cfg(not(target_os = "macos"))] #[cfg(not(target_os = "macos"))]
type WorkerSender = Sender<QueuedFmpSendJob>; type WorkerSender = Sender<QueuedFmpSendJob>;
#[cfg(target_os = "macos")]
type WorkerReceiver = MacWorkerReceiver<QueuedFmpSendJob>;
#[cfg(not(target_os = "macos"))]
type WorkerReceiver = Receiver<QueuedFmpSendJob>;
fn worker_channel() -> (WorkerSender, WorkerReceiver) {
#[cfg(target_os = "macos")]
{
mac_worker_channel(WORKER_CHANNEL_CAP)
}
#[cfg(not(target_os = "macos"))]
{
bounded::<QueuedFmpSendJob>(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<std::thread::JoinHandle<()>> {
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. /// Handle to the encrypt worker pool.
/// ///
/// Workers are **dedicated `std::thread`s** with **`crossbeam_channel`** /// Workers are **dedicated `std::thread`s** with **`crossbeam_channel`**
@@ -425,7 +478,7 @@ type WorkerSender = Sender<QueuedFmpSendJob>;
/// destinations hash to different workers. /// destinations hash to different workers.
#[derive(Clone)] #[derive(Clone)]
pub(crate) struct EncryptWorkerPool { pub(crate) struct EncryptWorkerPool {
senders: Arc<[WorkerSender]>, workers: Arc<WorkerSet<WorkerSender>>,
#[cfg(target_os = "macos")] #[cfg(target_os = "macos")]
macos_senders: Arc<MacSequencedSendFlows>, macos_senders: Arc<MacSequencedSendFlows>,
#[cfg(target_os = "macos")] #[cfg(target_os = "macos")]
@@ -437,31 +490,21 @@ impl EncryptWorkerPool {
/// dispatches jobs hash-by-destination to them. The workers exit /// dispatches jobs hash-by-destination to them. The workers exit
/// when all senders for their channel are dropped (i.e. when the /// when all senders for their channel are dropped (i.e. when the
/// returned `EncryptWorkerPool` and all clones go away). /// 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 { pub fn spawn(n: usize) -> Self {
let n = n.max(1); Self::start_with(n, spawn_worker)
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::<QueuedFmpSendJob>(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);
}
} }
/// 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<std::thread::JoinHandle<()>>,
) -> Self {
Self { Self {
senders: senders.into(), workers: Arc::new(WorkerSet::start("encrypt", n, worker_channel, spawn)),
#[cfg(target_os = "macos")] #[cfg(target_os = "macos")]
macos_senders: Arc::new(MacSequencedSendFlows::default()), macos_senders: Arc::new(MacSequencedSendFlows::default()),
#[cfg(target_os = "macos")] #[cfg(target_os = "macos")]
@@ -469,12 +512,26 @@ 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. /// Dispatch a job to the worker that owns its destination flow.
/// The hash is over `dest_addr` so every packet for one peer's /// The hash is over `dest_addr` so every packet for one peer's
/// kernel `SocketAddr` lands on the same worker and stays in /// kernel `SocketAddr` lands on the same worker and stays in
/// order — required for TCP's fast-retransmit logic above to /// order — required for TCP's fast-retransmit logic above to
/// behave on a single-flow run. Fire-and-forget — the worker /// behave on a single-flow run. The worker handles send errors
/// handles send errors itself via stats counters. /// 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 /// Uses `try_send` for the common uncontended case, then blocks
/// only when the bounded worker channel is full. These jobs carry /// only when the bounded worker channel is full. These jobs carry
@@ -482,13 +539,22 @@ impl EncryptWorkerPool {
/// this internal queue makes TCP-over-TUN collapse with avoidable /// this internal queue makes TCP-over-TUN collapse with avoidable
/// retransmits. Blocking here pushes back toward the TUN reader /// retransmits. Blocking here pushes back toward the TUN reader
/// and lets the kernel/app TCP stack pace the flow instead. /// 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"]
if self.senders.is_empty() { pub fn dispatch(&self, job: FmpSendJob) -> Result<(), Box<FmpSendJob>> {
debug!("EncryptWorkerPool has no workers; dropping job");
return;
}
let (idx, job) = self.prepare_dispatch(job); let (idx, job) = self.prepare_dispatch(job);
self.dispatch_to_worker(idx, job); 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")] #[cfg(target_os = "macos")]
@@ -503,7 +569,7 @@ impl EncryptWorkerPool {
}; };
let mut h = std::collections::hash_map::DefaultHasher::new(); let mut h = std::collections::hash_map::DefaultHasher::new();
key.hash(&mut h); 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)); return (idx, QueuedFmpSendJob::direct(job));
} }
@@ -518,64 +584,96 @@ impl EncryptWorkerPool {
.next_worker .next_worker
.fetch_add(1, std::sync::atomic::Ordering::Relaxed) .fetch_add(1, std::sync::atomic::Ordering::Relaxed)
/ macos_worker_stride(); / macos_worker_stride();
let idx = ticket % self.senders.len(); let idx = ticket % self.workers.len();
(idx, QueuedFmpSendJob::macos_sequenced(job, flow)) (idx, QueuedFmpSendJob::macos_sequenced(job, flow))
} }
#[cfg(not(target_os = "macos"))] #[cfg(not(target_os = "macos"))]
fn prepare_dispatch(&self, job: FmpSendJob) -> (usize, QueuedFmpSendJob) { fn prepare_dispatch(&self, job: FmpSendJob) -> (usize, QueuedFmpSendJob) {
use std::hash::{Hash, Hasher}; let idx = self.worker_index_for(job.dest_addr);
let mut h = std::collections::hash_map::DefaultHasher::new();
job.dest_addr.hash(&mut h);
let idx = (h.finish() as usize) % self.senders.len();
(idx, QueuedFmpSendJob::direct(job)) (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`, or hand it back when that worker has
/// exited.
#[cfg(target_os = "macos")] #[cfg(target_os = "macos")]
fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) { fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> Result<(), Refused> {
match self.senders[idx].try_push(job) { let sender = self.workers.sender(idx);
Ok(()) => {} match sender.try_push(job) {
Ok(()) => Ok(()),
Err(MacWorkerTryPushError::Full(job)) => { Err(MacWorkerTryPushError::Full(job)) => {
static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0); static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0);
let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed); 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!( warn!(
worker = idx, worker = idx,
full_events = n + 1, full_events = n + 1,
"EncryptWorker channel full; applying outbound backpressure" "EncryptWorker channel full; applying outbound backpressure"
); );
} }
if let Err(MacWorkerPushError) = self.senders[idx].push_blocking(*job) { sender
debug!(worker = idx, "EncryptWorker thread gone; dropping job"); .push_blocking(*job)
} .map_err(|MacWorkerPushError(job)| Refused(job))
}
Err(MacWorkerTryPushError::Closed) => {
debug!(worker = idx, "EncryptWorker thread gone; dropping job");
} }
Err(MacWorkerTryPushError::Closed(job)) => Err(Refused(job)),
} }
} }
/// Queue `job` on worker `idx`, or hand it back when that worker has
/// exited.
#[cfg(not(target_os = "macos"))] #[cfg(not(target_os = "macos"))]
fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) { fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> Result<(), Refused> {
match self.senders[idx].try_send(job) { let sender = self.workers.sender(idx);
Ok(()) => {} match sender.try_send(job) {
Ok(()) => Ok(()),
Err(TrySendError::Full(job)) => { Err(TrySendError::Full(job)) => {
static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0); static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0);
let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed); 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!( warn!(
worker = idx, worker = idx,
full_events = n + 1, full_events = n + 1,
"EncryptWorker channel full; applying outbound backpressure" "EncryptWorker channel full; applying outbound backpressure"
); );
} }
if let Err(SendError(_)) = self.senders[idx].send(job) { sender
debug!(worker = idx, "EncryptWorker thread gone; dropping job"); .send(job)
.map_err(|SendError(job)| Refused(Box::new(job)))
}
Err(TrySendError::Disconnected(job)) => Err(Refused(Box::new(job))),
} }
} }
Err(TrySendError::Disconnected(_)) => { }
debug!(worker = idx, "EncryptWorker thread gone; dropping job");
#[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<TestWorker>) -> 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<usize> {
#[cfg(target_os = "macos")]
{
let _ = dest;
None
}
#[cfg(not(target_os = "macos"))]
{
Some(self.worker_index_for(dest))
} }
} }
} }
@@ -1096,6 +1194,89 @@ fn run_worker_macos(idx: usize, rx: MacWorkerReceiver<QueuedFmpSendJob>) {
trace!(worker = idx, "FMP encrypt worker thread exiting"); 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<Vec<u8>, 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<u8>,
fsp_seal: Option<FspSealJob>,
) -> 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 /// Encrypt every job in `batch` in place, then issue one or more
/// bulk-send syscalls grouped **by exact send target**. Clears /// bulk-send syscalls grouped **by exact send target**. Clears
/// `batch` on return. Sync version — operates directly on the raw /// `batch` on return. Sync version — operates directly on the raw
@@ -1178,10 +1359,7 @@ fn flush_batch_sync(
crate::perf_profile::Stage::FmpWorkerQueueWait, crate::perf_profile::Stage::FmpWorkerQueueWait,
queued_at, queued_at,
); );
if let Some(fsp) = fsp_seal { if seal_wire(&cipher, counter, &mut wire_buf, fsp_seal).is_err() {
if fsp.aad_offset + FSP_HEADER_SIZE > fsp.plaintext_offset
|| fsp.plaintext_offset > wire_buf.len()
{
#[cfg(target_os = "macos")] #[cfg(target_os = "macos")]
if let Some(ticket) = macos_ticket { if let Some(ticket) = macos_ticket {
push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip); push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip);
@@ -1189,54 +1367,6 @@ fn flush_batch_sync(
continue; continue;
} }
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());
}
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")] #[cfg(target_os = "macos")]
if let Some(ticket) = macos_ticket { if let Some(ticket) = macos_ticket {
push_mac_completion( push_mac_completion(
@@ -2153,6 +2283,150 @@ mod unix_tests {
assert_eq!(recovered_fsp_plaintext, fsp_plaintext); 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 /// Standalone tests for the GSO-eligibility predicate. The full
@@ -2532,7 +2806,7 @@ mod mac_queue_tests {
fn spawn_pusher<T: Send + 'static>( fn spawn_pusher<T: Send + 'static>(
tx: MacWorkerSender<T>, tx: MacWorkerSender<T>,
item: T, item: T,
) -> mpsc::Receiver<Result<(), MacWorkerPushError>> { ) -> mpsc::Receiver<Result<(), MacWorkerPushError<T>>> {
let (done_tx, done_rx) = mpsc::channel(); let (done_tx, done_rx) = mpsc::channel();
thread::spawn(move || { thread::spawn(move || {
let result = tx.push_blocking(item); let result = tx.push_blocking(item);
@@ -2572,7 +2846,10 @@ mod mac_queue_tests {
let result = done let result = done
.recv_timeout(WAIT) .recv_timeout(WAIT)
.expect("push_blocking still blocked after the worker thread died"); .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"); assert!(worker.join().is_err(), "worker thread should have panicked");
} }
@@ -2580,12 +2857,18 @@ mod mac_queue_tests {
fn try_push_returns_closed_after_receiver_dropped() { fn try_push_returns_closed_after_receiver_dropped() {
let (tx, rx) = mac_worker_channel::<u32>(2); let (tx, rx) = mac_worker_channel::<u32>(2);
drop(rx); 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 done = spawn_pusher(tx, 2);
let result = done let result = done
.recv_timeout(WAIT) .recv_timeout(WAIT)
.expect("push_blocking blocked on a queue whose receiver is gone"); .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] #[test]
@@ -2775,9 +3058,11 @@ mod mac_ordered_tests {
assert!(tx.try_push(rig.sequenced(1)).is_ok()); assert!(tx.try_push(rig.sequenced(1)).is_ok());
assert!(tx.try_push(rig.sequenced(2)).is_ok()); assert!(tx.try_push(rig.sequenced(2)).is_ok());
drop(rx); drop(rx);
// The refused job comes back and is dropped here, which releases its
// slot as the dispatcher's caller would.
assert!(matches!( assert!(matches!(
tx.try_push(rig.sequenced(3)), tx.try_push(rig.sequenced(3)),
Err(MacWorkerTryPushError::Closed) Err(MacWorkerTryPushError::Closed(_))
)); ));
let mut batch = vec![rig.sequenced(4)]; let mut batch = vec![rig.sequenced(4)];
flush_batch_sync(&mut batch).expect("flush"); flush_batch_sync(&mut batch).expect("flush");
@@ -2892,3 +3177,226 @@ mod mac_ordered_tests {
drop(gap); 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);
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");
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 (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:?}");
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 {
assert!(pool.dispatch(rig.job(dest, counter)).is_ok());
}
let (done_tx, done_rx) = mpsc::channel::<bool>();
let blocked_pool = pool.clone();
let last = rig.job(dest, WORKER_CHANNEL_CAP as u64);
std::thread::spawn(move || {
let _ = done_tx.send(blocked_pool.dispatch(last).is_ok());
});
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");
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::<Option<u64>>();
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"
);
}
}
+13 -7
View File
@@ -7,6 +7,7 @@
use crate::node::Node; use crate::node::Node;
use crate::node::reject::DiscoveryReject; use crate::node::reject::DiscoveryReject;
use crate::proto::fsp::should_apply_path_mtu;
use crate::proto::lookup::{ use crate::proto::lookup::{
LookupAction, LookupRequest, LookupResponse, MAX_RECENT_LOOKUP_REQUESTS, LookupAction, LookupRequest, LookupResponse, MAX_RECENT_LOOKUP_REQUESTS,
}; };
@@ -431,7 +432,9 @@ impl Node {
let fips_addr = crate::FipsAddress::from_node_addr(&target); let fips_addr = crate::FipsAddress::from_node_addr(&target);
match self.path_mtu_lookup.write() { match self.path_mtu_lookup.write() {
Ok(mut map) => match map.get(&fips_addr).copied() { 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 // Keep the tighter learned value; never loosen
// the clamp. A reactive MtuExceeded or // the clamp. A reactive MtuExceeded or
// PathMtuNotification tighten takes precedence // PathMtuNotification tighten takes precedence
@@ -439,12 +442,15 @@ impl Node {
// (cross-carrier keep-tighter). // (cross-carrier keep-tighter).
// //
// This arm deliberately leaves `learned_ms` // This arm deliberately leaves `learned_ms`
// alone. That is what bounds a replayed // alone. A later answered lookup that reports
// response: the replay of a value already // the value already stored takes this arm, so
// stored takes this arm, so the entry still // the entry still expires at first-write plus
// expires at first-write plus the TTL rather // the TTL rather than being pushed out again
// than being pushed out again on every // by every answer of the same value. (A
// injection. Refreshing the stamp here would // 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 // read as a tidy-up and would silently restore
// indefinite pinning. // indefinite pinning.
debug!( debug!(
+22 -9
View File
@@ -15,7 +15,7 @@ use crate::proto::mmp::{
LinkReportKind, LinkReportSnapshot, MmpAction, PeerLivenessSnapshot, ReceiverReport, RrLog, LinkReportKind, LinkReportSnapshot, MmpAction, PeerLivenessSnapshot, ReceiverReport, RrLog,
SenderReport, SenderReport,
}; };
use crate::proto::stp::ParentEval; use crate::proto::stp::{Stp, TreeDecision};
use crate::transport::{TransportAddr, TransportId}; use crate::transport::{TransportAddr, TransportId};
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
use tracing::{debug, info, trace, warn}; use tracing::{debug, info, trace, warn};
@@ -258,13 +258,11 @@ impl Node {
// Compute the flap-dampening / hold-down veto at the edge; a mandatory // Compute the flap-dampening / hold-down veto at the edge; a mandatory
// switch bypasses it, a discretionary one is taken only if not suppressed. // 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 switch_suppressed = self.tree_state.is_switch_suppressed(mono_now_ms);
let new_parent = match self.tree_state.evaluate_parent(&peer_costs, &skip) { match Stp::classify_periodic(&self.tree_state, &peer_costs, &skip, switch_suppressed) {
ParentEval::Mandatory(p) => Some(p), TreeDecision::Switch {
ParentEval::Discretionary(p) if !switch_suppressed => Some(p), new_parent,
ParentEval::Discretionary(_) | ParentEval::None => None, new_seq,
}; } => {
if let Some(new_parent) = new_parent {
let new_seq = self.tree_state.my_declaration().sequence() + 1;
let flap_dampened = let flap_dampened =
self.tree_state self.tree_state
.set_parent(new_parent, new_seq, now_secs, mono_now_ms); .set_parent(new_parent, new_seq, now_secs, mono_now_ms);
@@ -301,7 +299,8 @@ impl Node {
self.send_tree_announce_to_all().await; self.send_tree_announce_to_all().await;
let all_peers: Vec<crate::NodeAddr> = self.peers.keys().copied().collect(); let all_peers: Vec<crate::NodeAddr> = self.peers.keys().copied().collect();
self.bloom_state.mark_all_updates_needed(all_peers); self.bloom_state.mark_all_updates_needed(all_peers);
} else if !self.tree_state.is_root() && self.tree_state.should_be_root() { }
TreeDecision::SelfRoot => {
self.tree_state.become_root(now_secs); self.tree_state.become_root(now_secs);
// Clone identity once (see the parent-switch branch above for why). // Clone identity once (see the parent-switch branch above for why).
let our_identity = self.identity().clone(); let our_identity = self.identity().clone();
@@ -328,6 +327,20 @@ impl Node {
let all_peers: Vec<crate::NodeAddr> = self.peers.keys().copied().collect(); let all_peers: Vec<crate::NodeAddr> = self.peers.keys().copied().collect();
self.bloom_state.mark_all_updates_needed(all_peers); self.bloom_state.mark_all_updates_needed(all_peers);
} }
// 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"
)
}
}
} }
} }
+36 -1
View File
@@ -3085,7 +3085,7 @@ impl Node {
entry.touch(send.now_ms); entry.touch(send.now_ms);
} }
workers.dispatch(crate::node::encrypt_worker::FmpSendJob { let dispatched = workers.dispatch(crate::node::encrypt_worker::FmpSendJob {
cipher: fmp_cipher, cipher: fmp_cipher,
counter: fmp_counter, counter: fmp_counter,
wire_buf, wire_buf,
@@ -3105,9 +3105,44 @@ impl Node {
drop_on_backpressure: true, drop_on_backpressure: true,
queued_at: None, queued_at: None,
}); });
if let Err(job) = dispatched {
self.send_refused_job_inline(*job, transport_id, &remote_addr, next_hop_addr)
.await;
}
Ok(true) 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. /// 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` /// Compresses the IPv6 header (format 0x00), then sends via `send_session_data`
+26 -24
View File
@@ -1,6 +1,7 @@
//! Node lifecycle management: start, stop, and peer connection initiation. //! Node lifecycle management: start, stop, and peer connection initiation.
pub(crate) mod supervisor; pub(crate) mod supervisor;
mod workers;
use super::{Node, NodeError, NodeState}; use super::{Node, NodeError, NodeState};
use supervisor::{Action, Child, Event, PeeringDesired, SupervisorFsm}; use supervisor::{Action, Child, Event, PeeringDesired, SupervisorFsm};
@@ -1845,18 +1846,12 @@ impl Node {
} }
} }
Child::EncryptWorkers => { Child::EncryptWorkers => {
// Hash-by-destination pins a TCP flow to one worker // A worker that cannot be started degrades the node; it
// (preserves wire ordering); additional workers light up // never stops start-up.
// under multi-flow load. Infallible → always up.
#[cfg(unix)] #[cfg(unix)]
{ let start = self.start_encrypt_workers(encrypt_worker_count);
self.supervisor.encrypt_workers = Some( #[cfg(not(unix))]
super::encrypt_worker::EncryptWorkerPool::spawn(encrypt_worker_count), let start = workers::PoolStart::ALL_LIVE;
);
info!(
workers = encrypt_worker_count,
"Spawned FMP-encrypt worker pool"
);
// `FIPS_DECRYPT_WORKERS=0` disables the pool entirely // `FIPS_DECRYPT_WORKERS=0` disables the pool entirely
// and forces the in-line rx_loop decrypt path. When 0 // and forces the in-line rx_loop decrypt path. When 0
@@ -1864,25 +1859,20 @@ impl Node {
// sits here — exactly where the decrypt spawn would be // sits here — exactly where the decrypt spawn would be
// in today's sequence (after the encrypt spawn+info, // in today's sequence (after the encrypt spawn+info,
// before nostr). // before nostr).
#[cfg(unix)]
if decrypt_worker_count == 0 { if decrypt_worker_count == 0 {
info!("FIPS_DECRYPT_WORKERS=0 → in-line decrypt in rx_loop"); info!("FIPS_DECRYPT_WORKERS=0 → in-line decrypt in rx_loop");
} }
} start.event(child)
Event::SubstrateUp { child }
} }
Child::DecryptWorkers => { 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)] #[cfg(unix)]
{ let start = self.start_decrypt_workers(decrypt_worker_count);
self.supervisor.decrypt_workers = Some( #[cfg(not(unix))]
super::decrypt_worker::DecryptWorkerPool::spawn(decrypt_worker_count), let start = workers::PoolStart::ALL_LIVE;
); start.event(child)
info!(
workers = decrypt_worker_count,
"Spawned FMP-decrypt worker pool"
);
}
Event::SubstrateUp { child }
} }
Child::Nostr => { Child::Nostr => {
match NostrRendezvous::start( match NostrRendezvous::start(
@@ -2624,6 +2614,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, /// Feed every queued interface-presence edge to the supervisor FSM,
/// returning the last [`NodeState`] it asked to publish (if any). /// returning the last [`NodeState`] it asked to publish (if any).
/// ///
+60 -7
View File
@@ -63,10 +63,11 @@
//! - the degenerate no-children path now resolves to `Failed` (zero transports), //! - the degenerate no-children path now resolves to `Failed` (zero transports),
//! **not** the old immediate-`Running`. //! **not** the old immediate-`Running`.
//! //!
//! Runtime child-liveness monitoring (a `ChildExited` event re-routing health //! Runtime child-liveness (a `ChildExited` event re-routing health when a
//! when a task/thread dies at runtime) is **deferred**: start-completion health //! task/thread dies at runtime) came later. Its producers are the children that
//! resolution is start-framed, and liveness monitoring is a substantial unbuilt //! report their own exit (TUN threads, the DNS task, the mDNS/Nostr monitor)
//! mechanism. This commit is start-time health only. //! 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) //! ## Scope: interface presence, and `Degraded` as a level (this commit)
//! //!
@@ -298,7 +299,8 @@ pub(crate) enum Health {
Full, Full,
/// ≥1 transport is up, but one or more configured optional children failed /// ≥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 /// 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 { Degraded {
/// The configured children that failed to start. /// The configured children that failed to start.
reasons: HashSet<Child>, reasons: HashSet<Child>,
@@ -860,8 +862,8 @@ pub(crate) struct Supervisor {
/// Off-task FMP-encrypt + UDP-send worker pool. Unix-only — /// Off-task FMP-encrypt + UDP-send worker pool. Unix-only —
/// the worker issues direct sendmmsg(2) / sendmsg+UDP_GSO calls /// the worker issues direct sendmmsg(2) / sendmsg+UDP_GSO calls
/// on raw fds via `AsRawFd`. None on Windows or when the worker /// on raw fds via `AsRawFd`. None on Windows or when no worker
/// pool failed to spawn. /// thread could be started.
#[cfg(unix)] #[cfg(unix)]
pub(crate) encrypt_workers: Option<crate::node::encrypt_worker::EncryptWorkerPool>, pub(crate) encrypt_workers: Option<crate::node::encrypt_worker::EncryptWorkerPool>,
@@ -869,9 +871,16 @@ pub(crate) struct Supervisor {
/// `encrypt_workers`. Workers are shards: each owns its session /// `encrypt_workers`. Workers are shards: each owns its session
/// state directly in a thread-local `HashMap` (no `RwLock`, /// state directly in a thread-local `HashMap` (no `RwLock`,
/// no `Mutex` per packet). Hash-by-cache-key dispatch. /// 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)] #[cfg(unix)]
pub(crate) decrypt_workers: Option<crate::node::decrypt_worker::DecryptWorkerPool>, pub(crate) decrypt_workers: Option<crate::node::decrypt_worker::DecryptWorkerPool>,
/// 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 /// Transport-medium change detection: the receiver the rx loop drains and
/// the detector task behind it. /// the detector task behind it.
/// ///
@@ -889,6 +898,15 @@ pub(crate) struct Supervisor {
pub(in crate::node) fsm: SupervisorFsm, 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<crate::node::encrypt_worker::EncryptWorkerPool>,
pub decrypt: Option<crate::node::decrypt_worker::DecryptWorkerPool>,
}
impl Supervisor { impl Supervisor {
/// A fresh supervisor with all handles empty and the FSM in `Created`, /// A fresh supervisor with all handles empty and the FSM in `Created`,
/// matching the field initializers `Node::new` previously used. /// matching the field initializers `Node::new` previously used.
@@ -914,6 +932,8 @@ impl Supervisor {
encrypt_workers: None, encrypt_workers: None,
#[cfg(unix)] #[cfg(unix)]
decrypt_workers: None, decrypt_workers: None,
#[cfg(all(test, unix))]
staged_pools: StagedPools::default(),
netmon_rx: None, netmon_rx: None,
netmon_task: None, netmon_task: None,
fsm: SupervisorFsm::new(), fsm: SupervisorFsm::new(),
@@ -1796,6 +1816,39 @@ mod tests {
assert!(s.absent().is_empty()); 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] #[test]
fn degraded_children_is_the_union_of_both_reason_sets() { fn degraded_children_is_the_union_of_both_reason_sets() {
let mut s = SupervisorFsm::new(); let mut s = SupervisorFsm::new();
+241
View File
@@ -0,0 +1,241 @@
//! 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. 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
//! 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));
}
}
+46 -23
View File
@@ -33,6 +33,8 @@ pub(crate) mod stats_history;
#[cfg(test)] #[cfg(test)]
mod tests; mod tests;
mod tree; mod tree;
#[cfg(unix)]
pub(crate) mod worker_set;
use self::peer_error_budget::PeerErrorBudget; use self::peer_error_budget::PeerErrorBudget;
use self::rate_limit::{HandshakeRateLimiter, LookupSignRateLimiter, SessionSetupRateLimiter}; use self::rate_limit::{HandshakeRateLimiter, LookupSignRateLimiter, SessionSetupRateLimiter};
@@ -4045,7 +4047,7 @@ impl Node {
// Drop bulk endpoint data on UDP backpressure to // Drop bulk endpoint data on UDP backpressure to
// keep the queue moving; control frames retry. // keep the queue moving; control frames retry.
let drop_on_backpressure = plaintext.first().is_some_and(|t| *t == 0x00); 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, cipher: fmp_cipher,
counter, counter,
wire_buf, wire_buf,
@@ -4057,12 +4059,27 @@ impl Node {
drop_on_backpressure, drop_on_backpressure,
queued_at: None, 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) { 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() { if let Some(mmp) = peer.mmp_mut() {
mmp.sender mmp.sender.record_sent(counter, timestamp_ms, sent_bytes);
.record_sent(counter, timestamp_ms, predicted_bytes);
} }
} }
return Ok(()); return Ok(());
@@ -4120,25 +4137,7 @@ impl Node {
let bytes_sent = transport let bytes_sent = transport
.send(&remote_addr, &wire_packet) .send(&remote_addr, &wire_packet)
.await .await
.map_err(|e| match e { .map_err(|e| link_send_error(*node_addr, 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),
},
})?;
// Update send statistics // Update send statistics
if let Some(peer) = self.peers.get_mut(node_addr) { if let Some(peer) = self.peers.get_mut(node_addr) {
@@ -4211,6 +4210,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 /// Shell-side [`routing::RoutingView`] seam over live `Node` state — the sole
/// routing read adapter the shell retains. It hands the sans-IO routing core /// routing read adapter the shell retains. It hands the sans-IO routing core
/// borrowed peers plus raw `may_reach` / `link_cost` / `coords` /// borrowed peers plus raw `may_reach` / `link_cost` / `coords`
+12 -9
View File
@@ -1,11 +1,10 @@
//! Transport-medium change detection. //! Transport-medium change detection.
//! //!
//! A node that moves between media (WLAN → LAN, WLAN → 5G, a BLE adapter //! A node that moves between IP media (WLAN → LAN, WLAN → 5G) would otherwise
//! coming or going) would otherwise learn about it only as *silence*: the peer //! learn about it only as *silence*: the peer sits in the table until
//! sits in the table until `node.link_dead_timeout_secs` reaps it, and the //! `node.link_dead_timeout_secs` reaps it, and the reconnect then waits out
//! reconnect then waits out whatever backoff the old medium had already //! whatever backoff the old medium had already accumulated. The host kernel
//! accumulated. The host kernel knew within milliseconds; the node would find //! knew within milliseconds; the node would find out half a minute later.
//! out half a minute later.
//! //!
//! This module closes that gap. It samples a coarse [`NetFingerprint`] of the //! This module closes that gap. It samples a coarse [`NetFingerprint`] of the
//! host's network attachment and publishes a [`NetChange`] on the channel 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. //! 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 //! **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, //! IP attachment at all, and nothing routes it onto this channel. The detector
//! the Android callback) and belongs on this same channel, pushed by the BLE //! below is the only source of a [`NetChange`]; the BLE transport publishes
//! transport rather than sampled here. //! 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 //! # Where the peer list comes from
//! //!
+56 -4
View File
@@ -1644,10 +1644,11 @@ async fn test_lookup_response_path_mtu_expires_without_a_session() {
#[tokio::test] #[tokio::test]
async fn test_replayed_lookup_response_does_not_extend_the_path_mtu_deadline() { 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 // The response carries no replay dedupe of its own, so a captured one can
// re-injected indefinitely. What bounds the damage is that a replay of a // be re-injected indefinitely. Accepting the first response clears the
// value already stored takes the keep-tighter arm, which does not touch // pending lookup, so each replay is dropped as unsolicited before it
// the learn time: each injection buys one TTL, not one per packet. // 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 mut node = make_node();
let from = make_node_addr(0xAA); let from = make_node_addr(0xAA);
@@ -1685,6 +1686,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"
);
}
// ============================================================================ // ============================================================================
// Integration Tests — min_mtu transit pruning // Integration Tests — min_mtu transit pruning
// ============================================================================ // ============================================================================
+47
View File
@@ -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 --- // --- Local delivery ---
#[tokio::test] #[tokio::test]
+209
View File
@@ -5680,3 +5680,212 @@ async fn a_terminal_rekey_msg1_send_failure_is_recorded_as_a_reject() {
"a terminal send failure is still counted" "a terminal send failure is still counted"
); );
} }
/// 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");
// XX: B promotes only on msg3.
let msg3 = next_packet(&mut packet_rx_b).await;
node_b.handle_msg3(msg3).await;
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<String>) {
// 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:?}");
}
}
+128 -1
View File
@@ -33,7 +33,7 @@ use crate::node::session::{EndToEndState, SessionEntry};
use crate::noise::HandshakeState; use crate::noise::HandshakeState;
use crate::peer::ActivePeer; use crate::peer::ActivePeer;
use crate::proto::mmp::{MmpMode, ReceiverReport}; use crate::proto::mmp::{MmpMode, ReceiverReport};
use crate::proto::stp::{ParentDeclaration, TreeCoordinate}; use crate::proto::stp::{CoordEntry, ParentDeclaration, TreeCoordinate};
// =========================================================================== // ===========================================================================
// Helpers // Helpers
@@ -434,3 +434,130 @@ async fn non_first_receiver_report_does_not_retrigger_tree() {
"a non-first ReceiverReport does not re-enter the first-RTT tree branch" "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"
);
}
+25
View File
@@ -93,6 +93,31 @@ pub(super) fn install_connected_udp(
.set_connected_udp(socket, drain); .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 /// Build a test node with an explicit `max_peers` limit (replaces the removed
/// `set_max_peers` setter; resource limits are immutable post-construction). /// `set_max_peers` setter; resource limits are immutable post-construction).
pub(super) fn make_node_with_max_peers(max_peers: usize) -> Node { pub(super) fn make_node_with_max_peers(max_peers: usize) -> Node {
+107
View File
@@ -9068,3 +9068,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" "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;
}
}
+241
View File
@@ -5010,3 +5010,244 @@ async fn msg1_handler_holds_its_pending_slot_while_the_handler_runs() {
"each handler released its own slot exactly once on the way out" "each handler released its own slot exactly once on the way out"
); );
} }
/// 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<TestWorker>, [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<Child>,
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);
}
}
+245
View File
@@ -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<S> {
sender: S,
thread: Option<JoinHandle<()>>,
}
/// 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<S> {
workers: Box<[Worker<S>]>,
/// 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<S> WorkerSet<S> {
/// 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<R>(
pool: &'static str,
n: usize,
mut channel: impl FnMut() -> (S, R),
mut spawn: impl FnMut(usize, R) -> std::io::Result<JoinHandle<()>>,
) -> 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<usize>;
/// 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<S> WorkerLiveness for WorkerSet<S> {
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<usize> {
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<S>(worker: &Worker<S>) -> 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<R: Send + 'static>(
plan: Vec<TestWorker>,
run: impl Fn(usize, R) -> std::io::Result<JoinHandle<()>>,
) -> impl FnMut(usize, R) -> std::io::Result<JoinHandle<()>> {
let mut plan: Vec<Option<TestWorker>> = 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::<u32>,
|idx, rx: mpsc::Receiver<u32>| {
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::<u32>, |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<u64> = (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]);
}
}
+6 -5
View File
@@ -483,10 +483,10 @@ impl Fsp {
} }
/// Decide whether a path-MTU update should tighten the shared lookup: emit /// Decide whether a path-MTU update should tighten the shared lookup: emit
/// `TightenPathMtuLookup` only when `candidate` is at least as tight as the /// `TightenPathMtuLookup` only when there is no `existing` value or
/// `existing` value (keep-tighter, never loosen). The `existing` read and /// `candidate` is strictly tighter than it (keep-tighter, never loosen).
/// the applied write are performed shell-side under one `path_mtu_lookup` /// The `existing` read and the applied write are performed shell-side
/// write guard, so the decision stays atomic. /// under one `path_mtu_lookup` write guard, so the decision stays atomic.
pub(crate) fn plan_path_mtu_tighten( pub(crate) fn plan_path_mtu_tighten(
&self, &self,
fips_addr: FipsAddress, 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 /// Decide whether a path-MTU update should be applied to the shared
/// `FipsAddress`-keyed lookup: keep the tighter of existing-or-candidate, never /// `FipsAddress`-keyed lookup: keep the tighter of existing-or-candidate, never
/// loosen. Returns `true` when `candidate` should be written (there is no /// 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<u16>, candidate: u16) -> bool { pub(crate) fn should_apply_path_mtu(existing: Option<u16>, candidate: u16) -> bool {
!matches!(existing, Some(existing) if existing <= candidate) !matches!(existing, Some(existing) if existing <= candidate)
} }
+1 -1
View File
@@ -33,7 +33,7 @@ mod tests;
pub(crate) use core::{ pub(crate) use core::{
DecryptSlot, EpochReaction, Fsp, FspAction, InitialMsg3ResendSnapshot, RekeyCfg, DecryptSlot, EpochReaction, Fsp, FspAction, InitialMsg3ResendSnapshot, RekeyCfg,
RekeyMsg3ResendSnapshot, SessionSnapshot, cutover_timer_elapsed, initiation_winner, 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::{ pub use wire::{
FspInnerFlags, SessionAck, SessionFlags, SessionMessageType, SessionMsg3, SessionSetup, FspInnerFlags, SessionAck, SessionFlags, SessionMessageType, SessionMsg3, SessionSetup,
+51 -1
View File
@@ -144,6 +144,20 @@ pub struct SessionDatagramRef<'a> {
pub payload: &'a [u8], 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<u8> {
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). /// 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; 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 /// True only at TTL 2 or more: at TTL 1 the decrement leaves zero, so the
/// datagram is dropped rather than forwarded. /// datagram is dropped rather than forwarded.
pub fn can_forward(&self) -> bool { 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). /// 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. /// Materialize an owned datagram for forwarding/re-encoding paths.
pub fn into_owned(self) -> SessionDatagram { pub fn into_owned(self) -> SessionDatagram {
SessionDatagram { SessionDatagram {
@@ -424,6 +444,36 @@ mod tests {
assert!(dg.with_ttl(255).can_forward()); 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] #[test]
fn test_session_datagram_decrement_ttl() { fn test_session_datagram_decrement_ttl() {
let base = SessionDatagram::new(make_node_addr(1), make_node_addr(2), vec![0x42]); let base = SessionDatagram::new(make_node_addr(1), make_node_addr(2), vec![0x42]);
+8 -8
View File
@@ -17,7 +17,7 @@
use super::limits::LimitVerdict; use super::limits::LimitVerdict;
use super::state::Router; use super::state::Router;
use super::wire::{CoordsRequired, MtuExceeded, PathBroken}; 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}; use crate::{NodeAddr, TreeCoordinate};
/// Read-only view of routing state the routing core needs. /// Read-only view of routing state the routing core needs.
@@ -116,7 +116,8 @@ impl Router {
/// datagram that would leave with a TTL of zero is not transmitted. /// datagram that would leave with a TTL of zero is not transmitted.
/// ///
/// The shell pre-resolves `next_hop` only for datagrams this can actually /// 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 /// `find_next_hop`'s LRU-touch side effect stays scoped to genuine
/// forwards. `route` still re-checks local delivery and the TTL /// forwards. `route` still re-checks local delivery and the TTL
/// authoritatively. /// authoritatively.
@@ -136,15 +137,14 @@ impl Router {
} }
// TTL enforcement on the transit path: decrement first, then drop if // TTL enforcement on the transit path: decrement first, then drop if
// the datagram would leave with a TTL of zero. `saturating_sub` folds // the datagram would leave with a TTL of zero. The already-exhausted
// the already-exhausted arrival (ttl=0) into the same test as the // arrival (ttl=0) and the last-hop arrival (ttl=1) are both dropped;
// last-hop arrival (ttl=1); neither is transmitted. // neither is transmitted.
let forwarded_ttl = dg.ttl.saturating_sub(1); let Some(forwarded_ttl) = ttl_after_hop(dg.ttl) else {
if forwarded_ttl == 0 {
return RouteOutcome::Drop { return RouteOutcome::Drop {
reason: DropReason::TtlExhausted, reason: DropReason::TtlExhausted,
}; };
} };
let nh = match next_hop { let nh = match next_hop {
Some(nh) => nh, Some(nh) => nh,
+43 -1
View File
@@ -1,7 +1,7 @@
//! Tests for the sans-IO routing decision core. //! Tests for the sans-IO routing decision core.
use super::util::{MockPeer, MockRoutingView, make_coords, make_datagram_ref, make_next_hop}; 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::RoutingSignalType;
use crate::proto::routing::{ use crate::proto::routing::{
DropReason, LimitVerdict, RouteAction, RouteOutcome, Router, RoutingView, select_best_candidate, DropReason, LimitVerdict, RouteAction, RouteOutcome, Router, RoutingView, select_best_candidate,
@@ -560,3 +560,45 @@ fn candidate_selection_excludes_non_full_peers() {
"with no Full candidate there is no bloom next hop at all" "with no Full candidate there is no bloom next hop at all"
); );
} }
/// 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"),
}
}
}
+2
View File
@@ -165,6 +165,8 @@ impl Stp {
/// `classify_announce`, the periodic path has no same-parent loop-drop / /// `classify_announce`, the periodic path has no same-parent loop-drop /
/// ancestry-update arms — a periodic tick has no announcing peer, so those cases /// 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. /// 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( pub(crate) fn classify_periodic(
tree: &TreeState, tree: &TreeState,
peer_costs: &BTreeMap<NodeAddr, f64>, peer_costs: &BTreeMap<NodeAddr, f64>,
+5 -1
View File
@@ -34,7 +34,11 @@ pub use crate::proto::coord::{CoordEntry, CoordError, TreeCoordinate};
pub(crate) use crate::proto::coord::{ pub(crate) use crate::proto::coord::{
coords_wire_size, decode_coords, decode_optional_coords, encode_coords, encode_empty_coords, 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 declaration::ParentDeclaration;
pub use state::TreeState; pub use state::TreeState;
pub use wire::TreeAnnounce; pub use wire::TreeAnnounce;
+36 -15
View File
@@ -48,8 +48,9 @@ from .veth import VethManager
log = logging.getLogger(__name__) log = logging.getLogger(__name__)
# The final snapshot waits for this many identical consecutive tree reads, # The final snapshot waits until every node answers and this many identical
# taken this far apart, so the tree must hold still for two intervals. # 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_READS = 3
SETTLE_INTERVAL_SECS = 5 SETTLE_INTERVAL_SECS = 5
SETTLE_TIMEOUT_SECS = 90 SETTLE_TIMEOUT_SECS = 90
@@ -744,7 +745,8 @@ class SimRunner:
# Take final tree snapshot while nodes are still running, once the # Take final tree snapshot while nodes are still running, once the
# tree has stopped moving. A node restored a moment ago is its own # 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 # 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._settle_tree()
self._take_snapshot("final") self._take_snapshot("final")
@@ -893,29 +895,39 @@ class SimRunner:
return result return result
def _settle_tree(self): 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 Returns once SETTLE_READS reads in a row, SETTLE_INTERVAL_SECS apart,
either waste time on a mesh that settled at once or cut off one that have each been answered by every node in the topology and show the
had not. Running out is logged and is not a failure in itself: the same root and parent for each, or once SETTLE_TIMEOUT_SECS has passed.
final snapshot is taken anyway, and the assertions judge what it The bound is checked after each read, so the settle can return up to
shows. 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 A read that any node did not answer never counts toward agreement. A
not catch a node that stays its own root for longer than the reads node restored at teardown may not have opened its control socket yet,
span, which is a tree that is stable and wrong, and is left to the and a tree that holds still without it is not the tree the final
assertions. 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() started = time.time()
previous = None previous = None
agreeing = 0 agreeing = 0
while not self._interrupted: while not self._interrupted:
trees = snapshot_all_trees(self.topology) trees = snapshot_all_trees(self.topology)
missing = sorted(set(self.topology.nodes) - trees.keys())
shape = { shape = {
nid: (data.get("root"), data.get("parent")) nid: (data.get("root"), data.get("parent"))
for nid, data in trees.items() for nid, data in trees.items()
} }
if not shape: if missing:
agreeing = 0 agreeing = 0
else: else:
agreeing = agreeing + 1 if shape == previous else 1 agreeing = agreeing + 1 if shape == previous else 1
@@ -925,8 +937,17 @@ class SimRunner:
log.info("Tree settled after %.0fs", waited) log.info("Tree settled after %.0fs", waited)
return return
if waited >= SETTLE_TIMEOUT_SECS: if waited >= SETTLE_TIMEOUT_SECS:
if missing:
log.warning( log.warning(
"Tree still changing after %.0fs; taking the final snapshot anyway", "%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, waited,
) )
return return