diff --git a/CHANGELOG.md b/CHANGELOG.md index ece45994..7cd14768 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -632,6 +632,27 @@ with v0.5.x or earlier peers. per-peer `connect()`-ed UDP socket stayed pinned to the old 5-tuple; the sibling path already cleared it. +- A crypto worker thread that exits now shows: the node reports `Degraded`, + and a warning names the pool and how many of its workers are still live, + logged again each time another one goes. A worker thread that cannot be + started is logged and leaves the node degraded instead of stopping start-up. + Losing workers never fails the node. Outbound packets for a lost encrypt + worker are sent from the main loop, as described below. Inbound packets for + a session held by a lost decrypt worker are still dropped, now with a + rate-limited warning, until the session rekeys or the link is re-established; + a session that would be placed on the lost worker after the loss is + decrypted on the main loop instead. + +- On every Unix platform, a packet whose encrypt worker thread has exited is + now encrypted and sent from the main loop. It was dropped, while the link + statistics counted it as sent, so every peer whose traffic went to that + worker was cut off until the daemon restarted. The packet keeps the counters + reserved for it, so the receiver sees no gap. Each such packet is counted and + logged at WARN, rate-limited, in place of a debug line. This includes macOS, + where the fix listed under macOS below dropped these packets; only packets + already queued to the exited worker are lost. With the macOS ordered sender, + a packet sent this way can arrive out of order with the rest of its flow. + #### Data plane and transports - A peer that stops reading can no longer stall the node. TCP, Tor, Nym and diff --git a/src/instr/recorder.rs b/src/instr/recorder.rs index 6d136cf9..1dde89c0 100644 --- a/src/instr/recorder.rs +++ b/src/instr/recorder.rs @@ -41,7 +41,7 @@ impl Domain { /// /// `as usize` indexes the counter arrays, so the discriminants are dense and /// `WholeTick` is last (it defines `N_STEPS`). Variants are declared -/// unconditionally — see [`Step::emitted`] for how the two platform- and +/// unconditionally — see [`Step::emitted`] for how the three platform- and /// profile-conditional steps are kept out of the emitted table. #[derive(Copy, Clone, Debug, PartialEq, Eq)] #[repr(usize)] @@ -76,6 +76,7 @@ pub(crate) enum Step { PollTransportDiscovery, SampleTransportCongestion, ActivateConnectedUdpSessions, + PollWorkerLiveness, DebugAssertPeerMapsCoherent, /// The whole tick-arm body, from before `check_timeouts` to after the last /// step. Composes safely with the per-step spans because the macro @@ -116,6 +117,7 @@ pub(crate) const STEPS: [Step; N_STEPS] = [ Step::PollTransportDiscovery, Step::SampleTransportCongestion, Step::ActivateConnectedUdpSessions, + Step::PollWorkerLiveness, Step::DebugAssertPeerMapsCoherent, Step::WholeTick, ]; @@ -151,6 +153,7 @@ impl Step { Step::PollTransportDiscovery => "poll_transport_discovery", Step::SampleTransportCongestion => "sample_transport_congestion", Step::ActivateConnectedUdpSessions => "activate_connected_udp_sessions", + Step::PollWorkerLiveness => "poll_worker_liveness", Step::DebugAssertPeerMapsCoherent => "debug_assert_peer_maps_coherent", Step::WholeTick => "whole_tick", } @@ -158,7 +161,7 @@ impl Step { /// Whether this step gets a row in this build. /// - /// Two steps are conditionally compiled at their call sites. Emitting a row + /// Three steps are conditionally compiled at their call sites. Emitting a row /// for them in a build where the call site does not exist would publish a /// count that is structurally zero forever, which reads as "this step never /// runs" rather than "this step is not in this build". The predicates below @@ -169,6 +172,7 @@ impl Step { Step::ActivateConnectedUdpSessions => { cfg!(any(target_os = "linux", target_os = "macos")) } + Step::PollWorkerLiveness => cfg!(unix), Step::DebugAssertPeerMapsCoherent => cfg!(debug_assertions), _ => true, } @@ -388,7 +392,7 @@ mod tests { let emitted = STEPS.iter().filter(|s| s.emitted()).count(); // 27 unconditional subsystem steps on this line (26 shared with the // 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 // is what caught the extra step when the master-line instrumentation // 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")) { expected += 1; } + if cfg!(unix) { + expected += 1; + } if cfg!(debug_assertions) { expected += 1; } diff --git a/src/node/dataplane/forwarding.rs b/src/node/dataplane/forwarding.rs index 76164e66..2af800f7 100644 --- a/src/node/dataplane/forwarding.rs +++ b/src/node/dataplane/forwarding.rs @@ -55,13 +55,12 @@ impl Node { self.try_warm_coord_cache_ref(&datagram_ref, payload.len()); // Pre-resolve the next hop only for datagrams the core can actually - // forward: not locally destined, and carrying a TTL that survives the - // decrement (`ttl > 1` — the shell-side mirror of the core's - // would-leave-zero drop). This keeps `find_next_hop`'s coord-cache - // LRU-touch side effect scoped to genuine forwards, as it was when the - // TTL test ran inline ahead of it. Warming above has already run, so - // the resolution observes freshly cached coords. - let next_hop = if datagram_ref.dest_addr != my_addr && datagram_ref.ttl > 1 { + // forward: not locally destined, and passing `can_forward`, which is + // the core's own hop-limit rule. This keeps `find_next_hop`'s + // coord-cache LRU-touch side effect scoped to genuine forwards, as it + // was when the TTL test ran inline ahead of it. Warming above has + // already run, so the resolution observes freshly cached coords. + let next_hop = if datagram_ref.dest_addr != my_addr && datagram_ref.can_forward() { self.resolve_next_hop(&datagram_ref.dest_addr) } else { None diff --git a/src/node/dataplane/rx_loop.rs b/src/node/dataplane/rx_loop.rs index 0d0381ac..26948e4b 100644 --- a/src/node/dataplane/rx_loop.rs +++ b/src/node/dataplane/rx_loop.rs @@ -345,17 +345,7 @@ impl Node { // republishing health, so nothing outside the node can // observe an address the listener no longer answers on. self.retract_child_publications(child); - let actions = self - .supervisor - .fsm - .step(crate::node::lifecycle::supervisor::Event::ChildExited { child }); - for action in actions { - if let crate::node::lifecycle::supervisor::Action::PublishState(ns) = - action - { - self.supervisor.state = ns; - } - } + self.step_child_exited(child); // A transport child exiting leaves the bound set, so // it can be the one that was holding the node's egress // MTU down. `is_bound()` is `is_operational()` plus the @@ -585,6 +575,10 @@ impl Node { #[cfg(any(target_os = "linux", target_os = "macos"))] instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ActivateConnectedUdpSessions, self.activate_connected_udp_sessions().await); + // Crypto worker threads that exited since the last tick. + #[cfg(unix)] + instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PollWorkerLiveness, + self.poll_worker_liveness()); // Debug-build sweep of the peer-lifecycle map invariant // (leaked machines / machine-less legs); two map scans, // compiled out of release builds. diff --git a/src/node/decrypt_worker.rs b/src/node/decrypt_worker.rs index 842203c4..e489e992 100644 --- a/src/node/decrypt_worker.rs +++ b/src/node/decrypt_worker.rs @@ -36,6 +36,7 @@ #![cfg_attr(not(unix), allow(dead_code))] use crate::NodeAddr; +use crate::node::worker_set::{WorkerLiveness, WorkerSet, worth_logging}; use crate::transport::{TransportAddr, TransportId}; use crossbeam_channel::{Receiver, Sender, TrySendError, bounded}; use portable_atomic::{AtomicU64, Ordering}; @@ -220,26 +221,49 @@ pub(crate) enum WorkerMsg { /// shard. #[derive(Clone)] pub(crate) struct DecryptWorkerPool { - senders: Arc<[Sender]>, + workers: Arc>>, +} + +/// Start the production worker loop on `rx` in a named OS thread. +fn spawn_worker( + idx: usize, + rx: Receiver, +) -> std::io::Result> { + std::thread::Builder::new() + .name(format!("fips-decrypt-{idx}")) + .spawn(move || run_worker(idx, rx)) } impl DecryptWorkerPool { + /// Spawn `n` worker OS threads (at least one). A worker thread that + /// cannot be started is logged and left dead; the caller reads how + /// many started from [`Self::liveness`]. pub fn spawn(n: usize) -> Self { - let n = n.max(1); - let mut senders = Vec::with_capacity(n); - for i in 0..n { - let (tx, rx) = bounded::(WORKER_CHANNEL_CAP); - std::thread::Builder::new() - .name(format!("fips-decrypt-{i}")) - .spawn(move || run_worker(i, rx)) - .expect("failed to spawn fips-decrypt OS thread"); - senders.push(tx); - } + Self::start_with(n, spawn_worker) + } + + /// Build the pool, starting each worker with `spawn`. Production passes + /// [`spawn_worker`]; tests pass workers that fail to start or exit on cue. + fn start_with( + n: usize, + spawn: impl FnMut(usize, Receiver) -> std::io::Result>, + ) -> Self { Self { - senders: senders.into(), + workers: Arc::new(WorkerSet::start( + "decrypt", + n, + || bounded::(WORKER_CHANNEL_CAP), + spawn, + )), } } + /// Whether each worker is still running, and how many dispatches a dead + /// one refused. + pub(crate) fn liveness(&self) -> &dyn WorkerLiveness { + &*self.workers + } + /// Stable hash from session key → worker index. Same hash is used /// for session registration and per-packet dispatch so packets and /// registration arrive at the same shard. @@ -247,24 +271,37 @@ impl DecryptWorkerPool { use std::hash::{Hash, Hasher}; let mut h = std::collections::hash_map::DefaultHasher::new(); cache_key.hash(&mut h); - (h.finish() as usize) % self.senders.len() + (h.finish() as usize) % self.workers.len() + } + + /// Count a message refused by the exited worker `idx`, and log it at a + /// bounded rate. + fn note_refused(&self, idx: usize, what: &'static str) { + let n = self.workers.note_refused(); + if worth_logging(n) { + warn!( + pool = "decrypt", + worker = idx, + refused = n + 1, + what, + "Decrypt worker has exited; message refused" + ); + } } /// Dispatch a per-packet decrypt job. Drops if the per-worker /// channel is full (sustained rate overrun); the rx_loop's drain /// caps inbound at the same scale upstream so the cliff is - /// bounded. + /// bounded. A job for an exited worker is dropped, counted and + /// logged at WARN. pub fn dispatch_job(&self, job: DecryptJob) { - if self.senders.is_empty() { - return; - } let idx = self.worker_idx_for(job.cache_key); - match self.senders[idx].try_send(WorkerMsg::Job(job)) { + match self.workers.sender(idx).try_send(WorkerMsg::Job(job)) { Ok(()) => {} Err(TrySendError::Full(_)) => { static FULL_COUNT: AtomicU64 = AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, Ordering::Relaxed); - if n < 8 || n.is_multiple_of(10000) { + if worth_logging(n) { warn!( worker = idx, drops = n + 1, @@ -272,9 +309,7 @@ impl DecryptWorkerPool { ); } } - Err(TrySendError::Disconnected(_)) => { - debug!(worker = idx, "DecryptWorker thread gone; dropping job"); - } + Err(TrySendError::Disconnected(_)) => self.note_refused(idx, "inbound packet"), } } @@ -303,11 +338,12 @@ impl DecryptWorkerPool { cache_key: (TransportId, u32), state: OwnedSessionState, ) -> bool { - if self.senders.is_empty() { - return false; - } let idx = self.worker_idx_for(cache_key); - match self.senders[idx].try_send(WorkerMsg::RegisterSession { cache_key, state }) { + match self + .workers + .sender(idx) + .try_send(WorkerMsg::RegisterSession { cache_key, state }) + { Ok(()) => true, Err(TrySendError::Full(_)) => { warn!( @@ -317,10 +353,7 @@ impl DecryptWorkerPool { false } Err(TrySendError::Disconnected(_)) => { - debug!( - worker = idx, - "DecryptWorker thread gone; ignoring registration" - ); + self.note_refused(idx, "session registration"); false } } @@ -329,11 +362,11 @@ impl DecryptWorkerPool { /// Drop a session from its worker (rekey, peer removed). Fire and /// forget — if the worker is gone we don't care. pub fn unregister_session(&self, cache_key: (TransportId, u32)) { - if self.senders.is_empty() { - return; - } let idx = self.worker_idx_for(cache_key); - let _ = self.senders[idx].try_send(WorkerMsg::UnregisterSession { cache_key }); + let _ = self + .workers + .sender(idx) + .try_send(WorkerMsg::UnregisterSession { cache_key }); } } @@ -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) -> Self { + Self::start_with( + plan.len(), + crate::node::worker_set::test_spawner(plan, spawn_worker), + ) + } +} + +/// The pool's view of its workers: which are live and what a dead one does to +/// a dispatch or a registration. Every wait is bounded. +#[cfg(test)] +mod pool_tests { + use super::*; + use crate::node::worker_set::{TestWorker, wait_for}; + use ring::aead::UnboundKey; + use std::sync::mpsc; + + /// A session key the pool hashes to worker `idx`. + fn key_on_worker(pool: &DecryptWorkerPool, idx: usize) -> (TransportId, u32) { + (0u32..) + .map(|n| (TransportId::new(1), n)) + .find(|key| pool.worker_idx_for(*key) == idx) + .expect("some key hashes to every worker") + } + + fn session_state() -> OwnedSessionState { + let key = UnboundKey::new(&ring::aead::CHACHA20_POLY1305, &[0u8; 32]).unwrap(); + OwnedSessionState { + fmp_cipher: LessSafeKey::new(key), + fmp_replay: ReplayWindow::new(), + source_npub: None, + } + } + + fn job(cache_key: (TransportId, u32)) -> DecryptJob { + let (fallback_tx, _) = tokio::sync::mpsc::unbounded_channel::(); + DecryptJob { + packet_data: vec![0u8; 48], + cache_key, + _transport_id: cache_key.0, + _remote_addr: TransportAddr::from_string("127.0.0.1:1234"), + timestamp_ms: 1_000, + source_node_addr: NodeAddr::from_bytes([1u8; 16]), + fmp_counter: 1, + fmp_flags: 0, + fmp_header: [0u8; 16], + fmp_ciphertext_offset: 16, + fallback_tx, + } + } + + #[test] + fn a_decrypt_worker_that_panics_is_counted_dead() { + let (die_tx, die_rx) = mpsc::channel::<()>(); + let mut die_rx = Some(die_rx); + let pool = DecryptWorkerPool::start_with(2, |idx, rx| { + if idx != 1 { + return spawn_worker(idx, rx); + } + let die = die_rx.take().expect("worker 1 starts once"); + std::thread::Builder::new().spawn(move || { + let _rx = rx; + let _ = die.recv(); + panic!("simulated decrypt worker panic"); + }) + }); + assert_eq!(pool.liveness().live_workers(), 2); + die_tx.send(()).expect("worker 1 gone before its signal"); + assert!( + wait_for(|| pool.liveness().live_workers() == 1), + "a worker that panicked is still counted live" + ); + assert_eq!(pool.liveness().dead_workers(), vec![1]); + } + + #[test] + fn a_worker_that_fails_to_spawn_leaves_the_pool_serving_the_rest() { + let pool = DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + assert_eq!(pool.liveness().worker_count(), 2); + assert_eq!(pool.liveness().live_workers(), 1); + let key = key_on_worker(&pool, 0); + assert!( + pool.register_session(key, session_state()), + "the live worker refused a registration" + ); + pool.dispatch_job(job(key)); + assert_eq!(pool.liveness().refused_dispatches(), 0); + } + + #[test] + fn dispatch_to_an_exited_decrypt_worker_warns_and_counts() { + let pool = DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + let key = key_on_worker(&pool, 1); + let ((), logs) = crate::testutil::capture_logs(|| pool.dispatch_job(job(key))); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" worker=1"), "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"decrypt\""), "{warnings:?}"); + assert_eq!(pool.liveness().refused_dispatches(), 1); + } + + #[test] + fn registration_with_an_exited_decrypt_worker_warns_counts_and_is_refused() { + let pool = DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + let key = key_on_worker(&pool, 1); + let (registered, logs) = + crate::testutil::capture_logs(|| pool.register_session(key, session_state())); + assert!(!registered, "a dead worker cannot own a session"); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" worker=1"), "{warnings:?}"); + assert_eq!(pool.liveness().refused_dispatches(), 1); + } +} diff --git a/src/node/encrypt_worker.rs b/src/node/encrypt_worker.rs index d8cd4da7..c68019d3 100644 --- a/src/node/encrypt_worker.rs +++ b/src/node/encrypt_worker.rs @@ -50,6 +50,9 @@ // warnings rather than gate every function individually. #![cfg_attr(not(unix), allow(dead_code))] +#[cfg(test)] +use crate::node::worker_set::{TestWorker, test_spawner}; +use crate::node::worker_set::{WorkerLiveness, WorkerSet, worth_logging}; use crate::proto::fmp::wire::ESTABLISHED_HEADER_SIZE; use crate::proto::fsp::wire::FSP_HEADER_SIZE; use crate::transport::udp::io::AsyncUdpSocket; @@ -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); + +impl Refused { + /// The job, with any macOS ordered-flow slot it held released as a skip. + fn into_job(self) -> Box { + #[cfg(target_os = "macos")] + let QueuedFmpSendJob { job, macos_ticket } = *self.0; + #[cfg(not(target_os = "macos"))] + let QueuedFmpSendJob { job } = *self.0; + #[cfg(target_os = "macos")] + drop(macos_ticket); + Box::new(job) + } +} + /// Handle to the encrypt worker pool. Dispatches jobs **hash-by- /// destination** across N worker tasks via per-worker bounded /// crossbeam channels. The bounded queue intentionally backpressures @@ -239,14 +260,17 @@ struct MacWorkerQueueState { closed: bool, } +/// Why `try_push` did not queue a job. Both variants hand the job back. #[cfg(any(target_os = "macos", test))] enum MacWorkerTryPushError { Full(Box), - Closed, + /// The receiver is gone: the worker has exited. + Closed(Box), } +/// `push_blocking` found the receiver gone; the job is handed back. #[cfg(any(target_os = "macos", test))] -struct MacWorkerPushError; +struct MacWorkerPushError(Box); #[cfg(any(target_os = "macos", test))] fn mac_worker_channel(cap: usize) -> (MacWorkerSender, MacWorkerReceiver) { @@ -277,10 +301,10 @@ impl MacWorkerSender { .lock() .expect("encrypt worker queue poisoned"); if state.closed { - // Outside the lock: dropping a sequenced job completes its slot. + // The caller drops or reuses the job, outside the lock: dropping a + // sequenced job completes its slot. drop(state); - drop(job); - return Err(MacWorkerTryPushError::Closed); + return Err(MacWorkerTryPushError::Closed(Box::new(job))); } if state.queue.len() >= self.inner.cap { return Err(MacWorkerTryPushError::Full(Box::new(job))); @@ -295,7 +319,7 @@ impl MacWorkerSender { Ok(()) } - fn push_blocking(&self, job: T) -> Result<(), MacWorkerPushError> { + fn push_blocking(&self, job: T) -> Result<(), MacWorkerPushError> { let mut state = self .inner .state @@ -304,8 +328,7 @@ impl MacWorkerSender { loop { if state.closed { drop(state); - drop(job); - return Err(MacWorkerPushError); + return Err(MacWorkerPushError(Box::new(job))); } if state.queue.len() < self.inner.cap { let was_empty = state.queue.is_empty(); @@ -406,6 +429,36 @@ type WorkerSender = MacWorkerSender; #[cfg(not(target_os = "macos"))] type WorkerSender = Sender; +#[cfg(target_os = "macos")] +type WorkerReceiver = MacWorkerReceiver; + +#[cfg(not(target_os = "macos"))] +type WorkerReceiver = Receiver; + +fn worker_channel() -> (WorkerSender, WorkerReceiver) { + #[cfg(target_os = "macos")] + { + mac_worker_channel(WORKER_CHANNEL_CAP) + } + #[cfg(not(target_os = "macos"))] + { + bounded::(WORKER_CHANNEL_CAP) + } +} + +/// Start the production worker loop on `rx` in a named OS thread. +fn spawn_worker(idx: usize, rx: WorkerReceiver) -> std::io::Result> { + let builder = std::thread::Builder::new().name(format!("fips-encrypt-{idx}")); + #[cfg(target_os = "macos")] + { + builder.spawn(move || run_worker_macos(idx, rx)) + } + #[cfg(not(target_os = "macos"))] + { + builder.spawn(move || run_worker(idx, rx)) + } +} + /// Handle to the encrypt worker pool. /// /// Workers are **dedicated `std::thread`s** with **`crossbeam_channel`** @@ -425,7 +478,7 @@ type WorkerSender = Sender; /// destinations hash to different workers. #[derive(Clone)] pub(crate) struct EncryptWorkerPool { - senders: Arc<[WorkerSender]>, + workers: Arc>, #[cfg(target_os = "macos")] macos_senders: Arc, #[cfg(target_os = "macos")] @@ -437,31 +490,21 @@ impl EncryptWorkerPool { /// dispatches jobs hash-by-destination to them. The workers exit /// when all senders for their channel are dropped (i.e. when the /// returned `EncryptWorkerPool` and all clones go away). + /// + /// A worker thread that cannot be started is logged and left dead; + /// the caller reads how many started from [`Self::liveness`]. pub fn spawn(n: usize) -> Self { - let n = n.max(1); - let mut senders = Vec::with_capacity(n); - for i in 0..n { - #[cfg(target_os = "macos")] - { - let (tx, rx) = mac_worker_channel(WORKER_CHANNEL_CAP); - std::thread::Builder::new() - .name(format!("fips-encrypt-{i}")) - .spawn(move || run_worker_macos(i, rx)) - .expect("failed to spawn fips-encrypt OS thread"); - senders.push(tx); - } - #[cfg(not(target_os = "macos"))] - { - let (tx, rx) = bounded::(WORKER_CHANNEL_CAP); - std::thread::Builder::new() - .name(format!("fips-encrypt-{i}")) - .spawn(move || run_worker(i, rx)) - .expect("failed to spawn fips-encrypt OS thread"); - senders.push(tx); - } - } + Self::start_with(n, spawn_worker) + } + + /// Build the pool, starting each worker with `spawn`. Production passes + /// [`spawn_worker`]; tests pass workers that fail to start or exit on cue. + fn start_with( + n: usize, + spawn: impl FnMut(usize, WorkerReceiver) -> std::io::Result>, + ) -> Self { Self { - senders: senders.into(), + workers: Arc::new(WorkerSet::start("encrypt", n, worker_channel, spawn)), #[cfg(target_os = "macos")] macos_senders: Arc::new(MacSequencedSendFlows::default()), #[cfg(target_os = "macos")] @@ -469,12 +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. /// The hash is over `dest_addr` so every packet for one peer's /// kernel `SocketAddr` lands on the same worker and stays in /// order — required for TCP's fast-retransmit logic above to - /// behave on a single-flow run. Fire-and-forget — the worker - /// handles send errors itself via stats counters. + /// behave on a single-flow run. The worker handles send errors + /// itself via stats counters. + /// + /// A job whose worker has exited is never queued: it is counted, + /// logged at WARN, and handed back as `Err` so the caller can seal + /// and send it on its own path with the counters it already + /// reserved. A job is handed back only when no worker has it, so + /// each reserved counter is still used at most once. In the macOS + /// ordered mode the job's place in its flow is released as a skip + /// before it is handed back. /// /// Uses `try_send` for the common uncontended case, then blocks /// only when the bounded worker channel is full. These jobs carry @@ -482,13 +539,22 @@ impl EncryptWorkerPool { /// this internal queue makes TCP-over-TUN collapse with avoidable /// retransmits. Blocking here pushes back toward the TUN reader /// and lets the kernel/app TCP stack pace the flow instead. - pub fn dispatch(&self, job: FmpSendJob) { - if self.senders.is_empty() { - debug!("EncryptWorkerPool has no workers; dropping job"); - return; - } + #[must_use = "a job handed back was not sent; the caller must send it another way"] + pub fn dispatch(&self, job: FmpSendJob) -> Result<(), Box> { let (idx, job) = self.prepare_dispatch(job); - 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")] @@ -503,7 +569,7 @@ impl EncryptWorkerPool { }; let mut h = std::collections::hash_map::DefaultHasher::new(); key.hash(&mut h); - let idx = (h.finish() as usize) % self.senders.len(); + let idx = (h.finish() as usize) % self.workers.len(); return (idx, QueuedFmpSendJob::direct(job)); } @@ -518,64 +584,96 @@ impl EncryptWorkerPool { .next_worker .fetch_add(1, std::sync::atomic::Ordering::Relaxed) / macos_worker_stride(); - let idx = ticket % self.senders.len(); + let idx = ticket % self.workers.len(); (idx, QueuedFmpSendJob::macos_sequenced(job, flow)) } #[cfg(not(target_os = "macos"))] fn prepare_dispatch(&self, job: FmpSendJob) -> (usize, QueuedFmpSendJob) { - use std::hash::{Hash, Hasher}; - let mut h = std::collections::hash_map::DefaultHasher::new(); - job.dest_addr.hash(&mut h); - let idx = (h.finish() as usize) % self.senders.len(); + let idx = self.worker_index_for(job.dest_addr); (idx, QueuedFmpSendJob::direct(job)) } + /// The worker that owns `dest`'s flow. + #[cfg(not(target_os = "macos"))] + fn worker_index_for(&self, dest: SocketAddr) -> usize { + use std::hash::{Hash, Hasher}; + let mut h = std::collections::hash_map::DefaultHasher::new(); + dest.hash(&mut h); + (h.finish() as usize) % self.workers.len() + } + + /// Queue `job` on worker `idx`, or hand it back when that worker has + /// exited. #[cfg(target_os = "macos")] - fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) { - match self.senders[idx].try_push(job) { - Ok(()) => {} + fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> Result<(), Refused> { + let sender = self.workers.sender(idx); + match sender.try_push(job) { + Ok(()) => Ok(()), Err(MacWorkerTryPushError::Full(job)) => { static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed); - if n < 8 || n.is_multiple_of(10000) { + if worth_logging(n) { warn!( worker = idx, full_events = n + 1, "EncryptWorker channel full; applying outbound backpressure" ); } - if let Err(MacWorkerPushError) = self.senders[idx].push_blocking(*job) { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); - } - } - Err(MacWorkerTryPushError::Closed) => { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); + sender + .push_blocking(*job) + .map_err(|MacWorkerPushError(job)| Refused(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"))] - fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) { - match self.senders[idx].try_send(job) { - Ok(()) => {} + fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> Result<(), Refused> { + let sender = self.workers.sender(idx); + match sender.try_send(job) { + Ok(()) => Ok(()), Err(TrySendError::Full(job)) => { static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed); - if n < 8 || n.is_multiple_of(10000) { + if worth_logging(n) { warn!( worker = idx, full_events = n + 1, "EncryptWorker channel full; applying outbound backpressure" ); } - if let Err(SendError(_)) = self.senders[idx].send(job) { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); - } - } - Err(TrySendError::Disconnected(_)) => { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); + sender + .send(job) + .map_err(|SendError(job)| Refused(Box::new(job))) } + Err(TrySendError::Disconnected(job)) => Err(Refused(Box::new(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) -> Self { + Self::start_with(plan.len(), test_spawner(plan, spawn_worker)) + } + + /// The worker a job to `dest` is dispatched to, where that depends on + /// `dest` alone. On macOS it also depends on the sending sockets, or on a + /// round-robin in the ordered mode, so there it is `None`. + pub(crate) fn worker_index_for_dest(&self, dest: SocketAddr) -> Option { + #[cfg(target_os = "macos")] + { + let _ = dest; + None + } + #[cfg(not(target_os = "macos"))] + { + Some(self.worker_index_for(dest)) } } } @@ -1096,6 +1194,89 @@ fn run_worker_macos(idx: usize, rx: MacWorkerReceiver) { trace!(worker = idx, "FMP encrypt worker thread exiting"); } +/// Why a job could not be sealed. +#[derive(Debug, thiserror::Error)] +pub(crate) enum SealError { + /// The offsets the job carries do not fit its buffer. + #[error("job layout does not fit its buffer")] + Layout, + /// The AEAD refused to seal. + #[error("AEAD seal failed")] + Aead, +} + +impl FmpSendJob { + /// Seal this job on the calling thread, as its worker would have, and + /// return the wire packet. For a job its worker refused: the counters it + /// carries were reserved for it and are used here, once. + pub(crate) fn seal_inline(self) -> Result, SealError> { + let FmpSendJob { + cipher, + counter, + mut wire_buf, + fsp_seal, + .. + } = self; + seal_wire(&cipher, counter, &mut wire_buf, fsp_seal)?; + Ok(wire_buf) + } +} + +/// Seal one job's wire buffer in place: the inner FSP seal first when the job +/// carries one, then the outer FMP seal over `[16..]` with the header as AAD. +/// Each tag is appended into capacity the builder reserved, so the buffer +/// becomes the wire packet without reallocating. +/// +/// A layout that does not fit the buffer is refused rather than indexed out of +/// bounds: this runs on a worker thread and, for a job that worker refused, on +/// the rx loop, where a panic would end the node rather than one worker. +fn seal_wire( + cipher: &LessSafeKey, + counter: u64, + wire_buf: &mut Vec, + fsp_seal: Option, +) -> Result<(), SealError> { + if let Some(fsp) = fsp_seal { + let aad_end = fsp + .aad_offset + .checked_add(FSP_HEADER_SIZE) + .ok_or(SealError::Layout)?; + if aad_end > fsp.plaintext_offset || fsp.plaintext_offset > wire_buf.len() { + return Err(SealError::Layout); + } + + let mut nonce_bytes = [0u8; 12]; + nonce_bytes[4..12].copy_from_slice(&fsp.counter.to_le_bytes()); + let nonce = Nonce::assume_unique_for_key(nonce_bytes); + let (prefix, plaintext_slice) = wire_buf.split_at_mut(fsp.plaintext_offset); + let aad = &prefix[fsp.aad_offset..aad_end]; + let tag = fsp + .cipher + .seal_in_place_separate_tag(nonce, Aad::from(aad), plaintext_slice) + .map_err(|_| SealError::Aead)?; + wire_buf.extend_from_slice(tag.as_ref()); + } + + if wire_buf.len() < ESTABLISHED_HEADER_SIZE { + return Err(SealError::Layout); + } + let mut nonce_bytes = [0u8; 12]; + nonce_bytes[4..12].copy_from_slice(&counter.to_le_bytes()); + let nonce = Nonce::assume_unique_for_key(nonce_bytes); + // Split-borrow: AAD reads from header bytes [0..16], seal writes + // into the plaintext slice [16..]. ring::aead's `seal_in_place_ + // separate_tag` takes `&mut [u8]` so we can hand it the + // post-header slice while AAD references the header slice. + // `split_at_mut` is the standard way to do this safely. + let (header_slice, plaintext_slice) = wire_buf.split_at_mut(ESTABLISHED_HEADER_SIZE); + let tag = cipher + .seal_in_place_separate_tag(nonce, Aad::from(&*header_slice), plaintext_slice) + .map_err(|_| SealError::Aead)?; + // wire_buf already has `+16` capacity reserved → no realloc. + wire_buf.extend_from_slice(tag.as_ref()); + Ok(()) +} + /// Encrypt every job in `batch` in place, then issue one or more /// bulk-send syscalls grouped **by exact send target**. Clears /// `batch` on return. Sync version — operates directly on the raw @@ -1178,65 +1359,14 @@ fn flush_batch_sync( crate::perf_profile::Stage::FmpWorkerQueueWait, queued_at, ); - if let Some(fsp) = fsp_seal { - if fsp.aad_offset + FSP_HEADER_SIZE > fsp.plaintext_offset - || fsp.plaintext_offset > wire_buf.len() - { - #[cfg(target_os = "macos")] - if let Some(ticket) = macos_ticket { - push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip); - } - continue; + if seal_wire(&cipher, counter, &mut wire_buf, fsp_seal).is_err() { + #[cfg(target_os = "macos")] + if let Some(ticket) = macos_ticket { + push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip); } - - let mut nonce_bytes = [0u8; 12]; - nonce_bytes[4..12].copy_from_slice(&fsp.counter.to_le_bytes()); - let nonce = Nonce::assume_unique_for_key(nonce_bytes); - let (prefix, plaintext_slice) = wire_buf.split_at_mut(fsp.plaintext_offset); - let aad = &prefix[fsp.aad_offset..fsp.aad_offset + FSP_HEADER_SIZE]; - let tag = - match fsp - .cipher - .seal_in_place_separate_tag(nonce, Aad::from(aad), plaintext_slice) - { - Ok(tag) => tag, - Err(_) => { - #[cfg(target_os = "macos")] - if let Some(ticket) = macos_ticket { - push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip); - } - continue; - } - }; - wire_buf.extend_from_slice(tag.as_ref()); + continue; } - let mut nonce_bytes = [0u8; 12]; - nonce_bytes[4..12].copy_from_slice(&counter.to_le_bytes()); - let nonce = Nonce::assume_unique_for_key(nonce_bytes); - // Split-borrow: AAD reads from header bytes [0..16], seal writes - // into the plaintext slice [16..]. ring::aead's `seal_in_place_ - // separate_tag` takes `&mut [u8]` so we can hand it the - // post-header slice while AAD references the header slice. - // `split_at_mut` is the standard way to do this safely. - let (header_slice, plaintext_slice) = wire_buf.split_at_mut(ESTABLISHED_HEADER_SIZE); - let tag = match cipher.seal_in_place_separate_tag( - nonce, - Aad::from(&*header_slice), - plaintext_slice, - ) { - Ok(tag) => tag, - Err(_) => { - #[cfg(target_os = "macos")] - if let Some(ticket) = macos_ticket { - push_mac_completion(&mut macos_completions, ticket, MacSendItem::Skip); - } - continue; - } - }; - // wire_buf already has `+16` capacity reserved → no realloc. - wire_buf.extend_from_slice(tag.as_ref()); - #[cfg(target_os = "macos")] if let Some(ticket) = macos_ticket { push_mac_completion( @@ -2153,6 +2283,150 @@ mod unix_tests { assert_eq!(recovered_fsp_plaintext, fsp_plaintext); }); } + + /// The seal a refused job gets on the main loop produces a packet the + /// canonical receive-side decoders accept, inner FSP layer included. + #[test] + fn an_inline_seal_matches_the_worker_wire_layout() { + use crate::NodeAddr; + use crate::noise::TAG_SIZE; + use crate::proto::fmp::wire::{EncryptedHeader, build_established_header}; + use crate::proto::fsp::wire::build_fsp_header; + use crate::proto::link::{ + LinkMessageType, SESSION_DATAGRAM_HEADER_SIZE, SessionDatagramRef, + }; + use crate::utils::index::SessionIndex; + + let rt = tokio::runtime::Builder::new_current_thread() + .enable_io() + .build() + .expect("tokio rt"); + let _enter = rt.enter(); + let socket = UdpRawSocket::open("127.0.0.1:0".parse().unwrap(), 1 << 20, 1 << 20) + .expect("open send socket") + .into_async() + .expect("into_async"); + + let fmp_cipher = test_cipher(0x31); + let fsp_cipher = test_cipher(0x32); + let (fmp_counter, fsp_counter) = (900u64, 77u64); + let fsp_plaintext = b"sealed on the main loop".to_vec(); + let link_plaintext_len = + SESSION_DATAGRAM_HEADER_SIZE + FSP_HEADER_SIZE + fsp_plaintext.len(); + let fmp_inner_len = 4 + link_plaintext_len + TAG_SIZE; + let fsp_header = build_fsp_header(fsp_counter, 0, fsp_plaintext.len() as u16); + let fmp_header = + build_established_header(SessionIndex::new(5), fmp_counter, 0, fmp_inner_len as u16); + + let mut wire_buf = Vec::with_capacity(ESTABLISHED_HEADER_SIZE + fmp_inner_len + TAG_SIZE); + wire_buf.extend_from_slice(&fmp_header); + wire_buf.extend_from_slice(&7u32.to_le_bytes()); + wire_buf.push(LinkMessageType::SessionDatagram.to_byte()); + wire_buf.push(16); + wire_buf.extend_from_slice(&1280u16.to_le_bytes()); + wire_buf.extend_from_slice(NodeAddr::from_bytes([0xAA; 16]).as_bytes()); + wire_buf.extend_from_slice(NodeAddr::from_bytes([0xBB; 16]).as_bytes()); + let aad_offset = wire_buf.len(); + wire_buf.extend_from_slice(&fsp_header); + let plaintext_offset = wire_buf.len(); + wire_buf.extend_from_slice(&fsp_plaintext); + + let job = FmpSendJob { + cipher: fmp_cipher.clone(), + counter: fmp_counter, + wire_buf, + fsp_seal: Some(FspSealJob { + cipher: fsp_cipher.clone(), + counter: fsp_counter, + aad_offset, + plaintext_offset, + }), + socket: socket.clone(), + dest_addr: "127.0.0.1:9".parse().unwrap(), + #[cfg(any(target_os = "linux", target_os = "macos"))] + connected_socket: None, + drop_on_backpressure: true, + queued_at: None, + }; + let wire = job.seal_inline().expect("inline seal"); + + let parsed = EncryptedHeader::parse(&wire).expect("FMP header parses"); + assert_eq!(parsed.counter, fmp_counter); + let fmp_plaintext = crate::noise::open( + Some(&fmp_cipher), + fmp_counter, + &parsed.header_bytes, + &wire[ESTABLISHED_HEADER_SIZE..], + ) + .expect("FMP open"); + let datagram = SessionDatagramRef::decode(&fmp_plaintext[5..]).expect("datagram decodes"); + let inner = crate::noise::open( + Some(&fsp_cipher), + fsp_counter, + &datagram.payload[..FSP_HEADER_SIZE], + &datagram.payload[FSP_HEADER_SIZE..], + ) + .expect("FSP open"); + assert_eq!(inner, fsp_plaintext); + + // FMP-only twin: a link message carries no inner seal. + let link_plaintext = b"\x10link message".to_vec(); + let header = build_established_header( + SessionIndex::new(5), + fmp_counter + 1, + 0, + (4 + link_plaintext.len()) as u16, + ); + let mut wire_buf = Vec::with_capacity(ESTABLISHED_HEADER_SIZE + 4 + 64); + wire_buf.extend_from_slice(&header); + wire_buf.extend_from_slice(&9u32.to_le_bytes()); + wire_buf.extend_from_slice(&link_plaintext); + let job = FmpSendJob { + cipher: fmp_cipher.clone(), + counter: fmp_counter + 1, + wire_buf, + fsp_seal: None, + socket, + dest_addr: "127.0.0.1:9".parse().unwrap(), + #[cfg(any(target_os = "linux", target_os = "macos"))] + connected_socket: None, + drop_on_backpressure: false, + queued_at: None, + }; + let wire = job.seal_inline().expect("inline seal"); + let opened = crate::noise::open( + Some(&fmp_cipher), + fmp_counter + 1, + &wire[..ESTABLISHED_HEADER_SIZE], + &wire[ESTABLISHED_HEADER_SIZE..], + ) + .expect("FMP open"); + assert_eq!(&opened[4..], &link_plaintext[..]); + } + + /// A job whose offsets do not fit its buffer is refused, not indexed out + /// of bounds. On the main loop a panic here would end the node. + #[test] + fn a_job_whose_layout_does_not_fit_is_refused_not_a_panic() { + let cipher = test_cipher(9); + let mut short = vec![0u8; ESTABLISHED_HEADER_SIZE - 1]; + assert!(matches!( + seal_wire(&cipher, 1, &mut short, None), + Err(SealError::Layout) + )); + + let mut buf = vec![0u8; 64]; + let overflowing = FspSealJob { + cipher: test_cipher(8), + counter: 1, + aad_offset: usize::MAX - 2, + plaintext_offset: 40, + }; + assert!(matches!( + seal_wire(&cipher, 1, &mut buf, Some(overflowing)), + Err(SealError::Layout) + )); + } } /// Standalone tests for the GSO-eligibility predicate. The full @@ -2532,7 +2806,7 @@ mod mac_queue_tests { fn spawn_pusher( tx: MacWorkerSender, item: T, - ) -> mpsc::Receiver> { + ) -> mpsc::Receiver>> { let (done_tx, done_rx) = mpsc::channel(); thread::spawn(move || { let result = tx.push_blocking(item); @@ -2572,7 +2846,10 @@ mod mac_queue_tests { let result = done .recv_timeout(WAIT) .expect("push_blocking still blocked after the worker thread died"); - assert!(matches!(result, Err(MacWorkerPushError))); + match result { + Err(MacWorkerPushError(job)) => assert_eq!(*job, 3, "the refused job comes back"), + Ok(()) => panic!("push_blocking queued onto a dead worker"), + } assert!(worker.join().is_err(), "worker thread should have panicked"); } @@ -2580,12 +2857,18 @@ mod mac_queue_tests { fn try_push_returns_closed_after_receiver_dropped() { let (tx, rx) = mac_worker_channel::(2); drop(rx); - assert!(matches!(tx.try_push(1), Err(MacWorkerTryPushError::Closed))); + match tx.try_push(1) { + Err(MacWorkerTryPushError::Closed(job)) => assert_eq!(*job, 1), + _ => panic!("try_push on a closed queue should hand the job back"), + } let done = spawn_pusher(tx, 2); let result = done .recv_timeout(WAIT) .expect("push_blocking blocked on a queue whose receiver is gone"); - assert!(matches!(result, Err(MacWorkerPushError))); + match result { + Err(MacWorkerPushError(job)) => assert_eq!(*job, 2), + Ok(()) => panic!("push_blocking queued onto a closed queue"), + } } #[test] @@ -2775,9 +3058,11 @@ mod mac_ordered_tests { assert!(tx.try_push(rig.sequenced(1)).is_ok()); assert!(tx.try_push(rig.sequenced(2)).is_ok()); drop(rx); + // The refused job comes back and is dropped here, which releases its + // slot as the dispatcher's caller would. assert!(matches!( tx.try_push(rig.sequenced(3)), - Err(MacWorkerTryPushError::Closed) + Err(MacWorkerTryPushError::Closed(_)) )); let mut batch = vec![rig.sequenced(4)]; flush_batch_sync(&mut batch).expect("flush"); @@ -2892,3 +3177,226 @@ mod mac_ordered_tests { drop(gap); } } + +/// The pool's view of its workers: which are live, what a dead one does to a +/// dispatch, and that a live but full queue still blocks. Every wait is +/// bounded so a regression fails instead of hanging. +#[cfg(test)] +mod pool_tests { + use super::*; + use crate::node::worker_set::wait_for; + #[cfg(not(target_os = "macos"))] + use crate::transport::udp::io::UdpRawSocket; + #[cfg(not(target_os = "macos"))] + use ring::aead::UnboundKey; + #[cfg(not(target_os = "macos"))] + use std::net::UdpSocket; + use std::sync::mpsc; + #[cfg(not(target_os = "macos"))] + use std::time::Duration; + + #[cfg(not(target_os = "macos"))] + struct Rig { + _rt: tokio::runtime::Runtime, + socket: AsyncUdpSocket, + } + + #[cfg(not(target_os = "macos"))] + impl Rig { + fn new() -> Self { + let rt = tokio::runtime::Builder::new_current_thread() + .enable_io() + .build() + .expect("tokio rt"); + let enter = rt.enter(); + let socket = UdpRawSocket::open("127.0.0.1:0".parse().unwrap(), 1 << 20, 1 << 20) + .expect("open send socket") + .into_async() + .expect("into_async"); + drop(enter); + Self { _rt: rt, socket } + } + + fn job(&self, dest: SocketAddr, counter: u64) -> FmpSendJob { + let key = UnboundKey::new(&ring::aead::CHACHA20_POLY1305, &[3u8; 32]).expect("key"); + let mut wire_buf = + Vec::with_capacity(ESTABLISHED_HEADER_SIZE + 8 + crate::noise::TAG_SIZE); + wire_buf.extend_from_slice(&[0x5A; ESTABLISHED_HEADER_SIZE]); + wire_buf.extend_from_slice(&counter.to_le_bytes()); + FmpSendJob { + cipher: LessSafeKey::new(key), + counter, + wire_buf, + fsp_seal: None, + socket: self.socket.clone(), + dest_addr: dest, + #[cfg(any(target_os = "linux", target_os = "macos"))] + connected_socket: None, + drop_on_backpressure: false, + queued_at: None, + } + } + } + + /// A bound receiver whose address the pool dispatches to worker `idx`. + #[cfg(not(target_os = "macos"))] + fn receiver_on_worker(pool: &EncryptWorkerPool, idx: usize) -> UdpSocket { + loop { + let sock = UdpSocket::bind("127.0.0.1:0").expect("bind receiver"); + if pool.worker_index_for(sock.local_addr().unwrap()) == idx { + sock.set_read_timeout(Some(Duration::from_secs(5))) + .expect("read timeout"); + return sock; + } + } + } + + #[test] + fn an_encrypt_worker_that_panics_is_counted_dead() { + let (die_tx, die_rx) = mpsc::channel::<()>(); + let mut die_rx = Some(die_rx); + let pool = EncryptWorkerPool::start_with(2, |idx, rx| { + if idx != 1 { + return spawn_worker(idx, rx); + } + let die = die_rx.take().expect("worker 1 starts once"); + std::thread::Builder::new().spawn(move || { + let _rx = rx; + let _ = die.recv(); + panic!("simulated encrypt worker panic"); + }) + }); + assert_eq!(pool.liveness().live_workers(), 2); + die_tx.send(()).expect("worker 1 gone before its signal"); + assert!( + wait_for(|| pool.liveness().live_workers() == 1), + "a worker that panicked is still counted live" + ); + assert_eq!(pool.liveness().dead_workers(), vec![1]); + assert_eq!(pool.liveness().worker_count(), 2); + } + + #[cfg(not(target_os = "macos"))] + #[test] + fn a_worker_that_fails_to_spawn_leaves_the_pool_serving_the_rest() { + let rig = Rig::new(); + let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + assert_eq!(pool.liveness().worker_count(), 2); + assert_eq!(pool.liveness().live_workers(), 1); + + let recv = receiver_on_worker(&pool, 0); + 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::(); + 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::>(); + let blocked_pool = pool.clone(); + let last = rig.job(dest, 7_000); + std::thread::spawn(move || { + let _ = done_tx.send(blocked_pool.dispatch(last).err().map(|job| job.counter)); + }); + std::thread::sleep(Duration::from_millis(100)); + assert!( + matches!(done_rx.try_recv(), Err(mpsc::TryRecvError::Empty)), + "dispatch returned while the queue was full and the worker alive" + ); + exit_tx.send(()).expect("worker gone before its signal"); + let back = done_rx + .recv_timeout(Duration::from_secs(5)) + .expect("dispatch still blocked after the worker exited"); + assert_eq!( + back, + Some(7_000), + "the job blocked on a dying worker must come back" + ); + } +} diff --git a/src/node/handlers/lookup.rs b/src/node/handlers/lookup.rs index aa6cf6dd..37beeaae 100644 --- a/src/node/handlers/lookup.rs +++ b/src/node/handlers/lookup.rs @@ -7,6 +7,7 @@ use crate::node::Node; use crate::node::reject::DiscoveryReject; +use crate::proto::fsp::should_apply_path_mtu; use crate::proto::lookup::{ LookupAction, LookupRequest, LookupResponse, MAX_RECENT_LOOKUP_REQUESTS, }; @@ -431,7 +432,9 @@ impl Node { let fips_addr = crate::FipsAddress::from_node_addr(&target); match self.path_mtu_lookup.write() { Ok(mut map) => match map.get(&fips_addr).copied() { - Some(existing) if existing.mtu <= path_mtu => { + Some(existing) + if !should_apply_path_mtu(Some(existing.mtu), path_mtu) => + { // Keep the tighter learned value; never loosen // the clamp. A reactive MtuExceeded or // PathMtuNotification tighten takes precedence @@ -439,12 +442,15 @@ impl Node { // (cross-carrier keep-tighter). // // This arm deliberately leaves `learned_ms` - // alone. That is what bounds a replayed - // response: the replay of a value already - // stored takes this arm, so the entry still - // expires at first-write plus the TTL rather - // than being pushed out again on every - // injection. Refreshing the stamp here would + // alone. A later answered lookup that reports + // the value already stored takes this arm, so + // the entry still expires at first-write plus + // the TTL rather than being pushed out again + // by every answer of the same value. (A + // replayed response never gets here: the + // pending lookup is gone once the first answer + // is accepted, so the copy is dropped as + // unsolicited.) Refreshing the stamp here would // read as a tidy-up and would silently restore // indefinite pinning. debug!( diff --git a/src/node/handlers/mmp.rs b/src/node/handlers/mmp.rs index db9bd17b..15051c03 100644 --- a/src/node/handlers/mmp.rs +++ b/src/node/handlers/mmp.rs @@ -15,7 +15,7 @@ use crate::proto::mmp::{ LinkReportKind, LinkReportSnapshot, MmpAction, PeerLivenessSnapshot, ReceiverReport, RrLog, SenderReport, }; -use crate::proto::stp::ParentEval; +use crate::proto::stp::{Stp, TreeDecision}; use crate::transport::{TransportAddr, TransportId}; use std::time::{Duration, Instant}; use tracing::{debug, info, trace, warn}; @@ -258,75 +258,88 @@ impl Node { // Compute the flap-dampening / hold-down veto at the edge; a mandatory // switch bypasses it, a discretionary one is taken only if not suppressed. let switch_suppressed = self.tree_state.is_switch_suppressed(mono_now_ms); - let new_parent = match self.tree_state.evaluate_parent(&peer_costs, &skip) { - ParentEval::Mandatory(p) => Some(p), - ParentEval::Discretionary(p) if !switch_suppressed => Some(p), - ParentEval::Discretionary(_) | ParentEval::None => None, - }; - if let Some(new_parent) = new_parent { - let new_seq = self.tree_state.my_declaration().sequence() + 1; - let flap_dampened = - self.tree_state - .set_parent(new_parent, new_seq, now_secs, mono_now_ms); - self.tree_state.recompute_coords(); - // Clone identity once: sign_declaration borrows &mut tree_state while - // the identity() accessor borrows all of &self, so an owned copy avoids - // the split-borrow conflict on this infrequent parent-switch path. - let our_identity = self.identity().clone(); - if let Err(e) = - sign_declaration(self.tree_state.my_declaration_mut(), &our_identity) - { - warn!(error = %e, "Failed to sign declaration after first-RTT parent eval"); - self.metrics() - .tree - .record_reject(TreeReject::OutboundSignFailed); - return; + match Stp::classify_periodic(&self.tree_state, &peer_costs, &skip, switch_suppressed) { + TreeDecision::Switch { + new_parent, + new_seq, + } => { + let flap_dampened = + self.tree_state + .set_parent(new_parent, new_seq, now_secs, mono_now_ms); + self.tree_state.recompute_coords(); + // Clone identity once: sign_declaration borrows &mut tree_state while + // the identity() accessor borrows all of &self, so an owned copy avoids + // the split-borrow conflict on this infrequent parent-switch path. + let our_identity = self.identity().clone(); + if let Err(e) = + sign_declaration(self.tree_state.my_declaration_mut(), &our_identity) + { + warn!(error = %e, "Failed to sign declaration after first-RTT parent eval"); + self.metrics() + .tree + .record_reject(TreeReject::OutboundSignFailed); + return; + } + // Surgical invalidation — see CoordCache::invalidate_via_node doc. + self.coord_cache + .invalidate_via_node(our_identity.node_addr()); + self.reset_lookup_backoff(); + self.metrics().tree.parent_switches.inc(); + info!( + new_parent = %self.peer_display_name(&new_parent), + new_seq = new_seq, + new_root = %self.tree_state.root(), + depth = self.tree_state.my_coords().depth(), + trigger = "first-rtt", + "Parent switched after first RTT measurement" + ); + if flap_dampened { + self.note_flap("first-rtt"); + } + self.send_tree_announce_to_all().await; + let all_peers: Vec = self.peers.keys().copied().collect(); + self.bloom_state.mark_all_updates_needed(all_peers); } - // Surgical invalidation — see CoordCache::invalidate_via_node doc. - self.coord_cache - .invalidate_via_node(our_identity.node_addr()); - self.reset_lookup_backoff(); - self.metrics().tree.parent_switches.inc(); - info!( - new_parent = %self.peer_display_name(&new_parent), - new_seq = new_seq, - new_root = %self.tree_state.root(), - depth = self.tree_state.my_coords().depth(), - trigger = "first-rtt", - "Parent switched after first RTT measurement" - ); - if flap_dampened { - self.note_flap("first-rtt"); + TreeDecision::SelfRoot => { + self.tree_state.become_root(now_secs); + // Clone identity once (see the parent-switch branch above for why). + let our_identity = self.identity().clone(); + if let Err(e) = + sign_declaration(self.tree_state.my_declaration_mut(), &our_identity) + { + warn!(error = %e, "Failed to sign self-root declaration after first-RTT"); + self.metrics() + .tree + .record_reject(TreeReject::OutboundSignFailed); + return; + } + // Surgical invalidation — see CoordCache::invalidate_other_roots doc. + self.coord_cache + .invalidate_other_roots(our_identity.node_addr()); + self.reset_lookup_backoff(); + self.metrics().tree.parent_switches.inc(); + info!( + new_root = %self.tree_state.root(), + trigger = "first-rtt", + "Self-promoted to root after first RTT: smallest visible NodeAddr" + ); + self.send_tree_announce_to_all().await; + let all_peers: Vec = self.peers.keys().copied().collect(); + self.bloom_state.mark_all_updates_needed(all_peers); } - self.send_tree_announce_to_all().await; - let all_peers: Vec = self.peers.keys().copied().collect(); - self.bloom_state.mark_all_updates_needed(all_peers); - } else if !self.tree_state.is_root() && self.tree_state.should_be_root() { - self.tree_state.become_root(now_secs); - // Clone identity once (see the parent-switch branch above for why). - let our_identity = self.identity().clone(); - if let Err(e) = - sign_declaration(self.tree_state.my_declaration_mut(), &our_identity) - { - warn!(error = %e, "Failed to sign self-root declaration after first-RTT"); - self.metrics() - .tree - .record_reject(TreeReject::OutboundSignFailed); - return; + // Nothing changed. The periodic tick rebroadcasts; this path does not. + TreeDecision::PeriodicRebroadcast => {} + // classify_periodic never yields these: there is no announcing + // peer, so the loop-drop / ancestry-update arms cannot arise, and + // ParentLost is the removal drive's outcome. + TreeDecision::LoopDrop + | TreeDecision::AncestryUpdate { .. } + | TreeDecision::ParentLost + | TreeDecision::NoChange => { + unreachable!( + "classify_periodic yields only Switch / SelfRoot / PeriodicRebroadcast" + ) } - // Surgical invalidation — see CoordCache::invalidate_other_roots doc. - self.coord_cache - .invalidate_other_roots(our_identity.node_addr()); - self.reset_lookup_backoff(); - self.metrics().tree.parent_switches.inc(); - info!( - new_root = %self.tree_state.root(), - trigger = "first-rtt", - "Self-promoted to root after first RTT: smallest visible NodeAddr" - ); - self.send_tree_announce_to_all().await; - let all_peers: Vec = self.peers.keys().copied().collect(); - self.bloom_state.mark_all_updates_needed(all_peers); } } } diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index 25e13560..f20372a0 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -3085,7 +3085,7 @@ impl Node { entry.touch(send.now_ms); } - workers.dispatch(crate::node::encrypt_worker::FmpSendJob { + let dispatched = workers.dispatch(crate::node::encrypt_worker::FmpSendJob { cipher: fmp_cipher, counter: fmp_counter, wire_buf, @@ -3105,9 +3105,44 @@ impl Node { drop_on_backpressure: true, queued_at: None, }); + if let Err(job) = dispatched { + self.send_refused_job_inline(*job, transport_id, &remote_addr, next_hop_addr) + .await; + } Ok(true) } + /// Seal and send a job the encrypt worker for its next hop refused + /// because that worker has exited, using the FSP and FMP counters the + /// job already reserved so neither counter is skipped. Stats were + /// recorded before dispatch and now describe this packet. + /// + /// A failure is logged and swallowed, as the worker does with its own: + /// the caller sees the same result whichever of the two sent the packet. + #[cfg(unix)] + async fn send_refused_job_inline( + &self, + job: crate::node::encrypt_worker::FmpSendJob, + transport_id: crate::transport::TransportId, + remote_addr: &crate::transport::TransportAddr, + next_hop_addr: NodeAddr, + ) { + let wire = match job.seal_inline() { + Ok(wire) => wire, + Err(error) => { + debug!(next_hop = %next_hop_addr, %error, "Inline seal of session data failed"); + return; + } + }; + let Some(transport) = self.transports.get(&transport_id) else { + debug!(next_hop = %next_hop_addr, "Transport gone before inline send of session data"); + return; + }; + if let Err(error) = transport.send(remote_addr, &wire).await { + debug!(next_hop = %next_hop_addr, %error, "Inline send of session data failed"); + } + } + /// Send an IPv6 packet through the IPv6 shim (port 256) with header compression. /// /// Compresses the IPv6 header (format 0x00), then sends via `send_session_data` diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index 9fa489df..f3fb2af3 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -1,6 +1,7 @@ //! Node lifecycle management: start, stop, and peer connection initiation. pub(crate) mod supervisor; +mod workers; use super::{Node, NodeError, NodeState}; use supervisor::{Action, Child, Event, PeeringDesired, SupervisorFsm}; @@ -1845,44 +1846,33 @@ impl Node { } } Child::EncryptWorkers => { - // Hash-by-destination pins a TCP flow to one worker - // (preserves wire ordering); additional workers light up - // under multi-flow load. Infallible → always up. + // A worker that cannot be started degrades the node; it + // never stops start-up. #[cfg(unix)] - { - self.supervisor.encrypt_workers = Some( - super::encrypt_worker::EncryptWorkerPool::spawn(encrypt_worker_count), - ); - info!( - workers = encrypt_worker_count, - "Spawned FMP-encrypt worker pool" - ); + let start = self.start_encrypt_workers(encrypt_worker_count); + #[cfg(not(unix))] + let start = workers::PoolStart::ALL_LIVE; - // `FIPS_DECRYPT_WORKERS=0` disables the pool entirely - // and forces the in-line rx_loop decrypt path. When 0 - // no DecryptWorkers child is emitted, so this info! - // sits here — exactly where the decrypt spawn would be - // in today's sequence (after the encrypt spawn+info, - // before nostr). - if decrypt_worker_count == 0 { - info!("FIPS_DECRYPT_WORKERS=0 → in-line decrypt in rx_loop"); - } + // `FIPS_DECRYPT_WORKERS=0` disables the pool entirely + // and forces the in-line rx_loop decrypt path. When 0 + // no DecryptWorkers child is emitted, so this info! + // sits here — exactly where the decrypt spawn would be + // in today's sequence (after the encrypt spawn+info, + // before nostr). + #[cfg(unix)] + if decrypt_worker_count == 0 { + info!("FIPS_DECRYPT_WORKERS=0 → in-line decrypt in rx_loop"); } - Event::SubstrateUp { child } + start.event(child) } Child::DecryptWorkers => { - // Shard-owned decrypt pool. Infallible → always up. + // Shard-owned decrypt pool. A worker that cannot be + // started degrades the node; it never stops start-up. #[cfg(unix)] - { - self.supervisor.decrypt_workers = Some( - super::decrypt_worker::DecryptWorkerPool::spawn(decrypt_worker_count), - ); - info!( - workers = decrypt_worker_count, - "Spawned FMP-decrypt worker pool" - ); - } - Event::SubstrateUp { child } + let start = self.start_decrypt_workers(decrypt_worker_count); + #[cfg(not(unix))] + let start = workers::PoolStart::ALL_LIVE; + start.event(child) } Child::Nostr => { match NostrRendezvous::start( @@ -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, /// returning the last [`NodeState`] it asked to publish (if any). /// diff --git a/src/node/lifecycle/supervisor.rs b/src/node/lifecycle/supervisor.rs index 301ad22c..edb4e658 100644 --- a/src/node/lifecycle/supervisor.rs +++ b/src/node/lifecycle/supervisor.rs @@ -63,10 +63,11 @@ //! - the degenerate no-children path now resolves to `Failed` (zero transports), //! **not** the old immediate-`Running`. //! -//! Runtime child-liveness monitoring (a `ChildExited` event re-routing health -//! when a task/thread dies at runtime) is **deferred**: start-completion health -//! resolution is start-framed, and liveness monitoring is a substantial unbuilt -//! mechanism. This commit is start-time health only. +//! Runtime child-liveness (a `ChildExited` event re-routing health when a +//! task/thread dies at runtime) came later. Its producers are the children that +//! report their own exit (TUN threads, the DNS task, the mDNS/Nostr monitor) +//! and, for the crypto worker pools, the rx loop tick's liveness sweep, which +//! reports a pool's child as exited when it loses a worker. //! //! ## Scope: interface presence, and `Degraded` as a level (this commit) //! @@ -298,7 +299,8 @@ pub(crate) enum Health { Full, /// ≥1 transport is up, but one or more configured optional children failed /// to start (a transport beyond the first, Nostr, mDNS, TUN, DNS, or a - /// worker-pool spawn). The node is operational (serving) but degraded. + /// worker pool that did not start every worker). The node is operational + /// (serving) but degraded. Degraded { /// The configured children that failed to start. reasons: HashSet, @@ -860,8 +862,8 @@ pub(crate) struct Supervisor { /// Off-task FMP-encrypt + UDP-send worker pool. Unix-only — /// the worker issues direct sendmmsg(2) / sendmsg+UDP_GSO calls - /// on raw fds via `AsRawFd`. None on Windows or when the worker - /// pool failed to spawn. + /// on raw fds via `AsRawFd`. None on Windows or when no worker + /// thread could be started. #[cfg(unix)] pub(crate) encrypt_workers: Option, @@ -869,9 +871,16 @@ pub(crate) struct Supervisor { /// `encrypt_workers`. Workers are shards: each owns its session /// state directly in a thread-local `HashMap` (no `RwLock`, /// no `Mutex` per packet). Hash-by-cache-key dispatch. + /// None on Windows, with `FIPS_DECRYPT_WORKERS=0`, or when no worker + /// thread could be started. #[cfg(unix)] pub(crate) decrypt_workers: Option, + /// Pools a test has built for start-up to install in place of spawning + /// its own. + #[cfg(all(test, unix))] + pub(in crate::node) staged_pools: StagedPools, + /// Transport-medium change detection: the receiver the rx loop drains and /// the detector task behind it. /// @@ -889,6 +898,15 @@ pub(crate) struct Supervisor { pub(in crate::node) fsm: SupervisorFsm, } +/// Crypto worker pools a test hands to start-up, each taken by the first +/// start of its pool's child. +#[cfg(all(test, unix))] +#[derive(Default)] +pub(in crate::node) struct StagedPools { + pub encrypt: Option, + pub decrypt: Option, +} + impl Supervisor { /// A fresh supervisor with all handles empty and the FSM in `Created`, /// matching the field initializers `Node::new` previously used. @@ -914,6 +932,8 @@ impl Supervisor { encrypt_workers: None, #[cfg(unix)] decrypt_workers: None, + #[cfg(all(test, unix))] + staged_pools: StagedPools::default(), netmon_rx: None, netmon_task: None, fsm: SupervisorFsm::new(), @@ -1796,6 +1816,39 @@ mod tests { assert!(s.absent().is_empty()); } + /// The crypto worker pools are a performance offload with a main-loop path + /// behind them, so losing one degrades the node and can never be what + /// fails it. + #[test] + fn a_worker_pool_exit_is_degraded_and_never_failed() { + let mut s = SupervisorFsm::running_with([ + Child::Transport(tid(1)), + Child::EncryptWorkers, + Child::DecryptWorkers, + ]); + assert_eq!( + s.step(Event::ChildExited { + child: Child::EncryptWorkers + }), + vec![Action::PublishState(NodeState::Degraded)] + ); + assert_eq!( + s.step(Event::ChildExited { + child: Child::DecryptWorkers + }), + vec![Action::PublishState(NodeState::Degraded)] + ); + assert!(matches!( + s.state(), + SupState::Running { + health: Health::Degraded { .. } + } + )); + let degraded = s.degraded_children(); + assert!(degraded.contains(&Child::EncryptWorkers)); + assert!(degraded.contains(&Child::DecryptWorkers)); + } + #[test] fn degraded_children_is_the_union_of_both_reason_sets() { let mut s = SupervisorFsm::new(); diff --git a/src/node/lifecycle/workers.rs b/src/node/lifecycle/workers.rs new file mode 100644 index 00000000..42a4558b --- /dev/null +++ b/src/node/lifecycle/workers.rs @@ -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)); + } +} diff --git a/src/node/mod.rs b/src/node/mod.rs index 20ac6574..4a29c997 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -33,6 +33,8 @@ pub(crate) mod stats_history; #[cfg(test)] mod tests; mod tree; +#[cfg(unix)] +pub(crate) mod worker_set; use self::peer_error_budget::PeerErrorBudget; use self::rate_limit::{HandshakeRateLimiter, LookupSignRateLimiter, SessionSetupRateLimiter}; @@ -4045,7 +4047,7 @@ impl Node { // Drop bulk endpoint data on UDP backpressure to // keep the queue moving; control frames retry. let drop_on_backpressure = plaintext.first().is_some_and(|t| *t == 0x00); - workers.dispatch(crate::node::encrypt_worker::FmpSendJob { + let dispatched = workers.dispatch(crate::node::encrypt_worker::FmpSendJob { cipher: fmp_cipher, counter, wire_buf, @@ -4057,12 +4059,27 @@ impl Node { drop_on_backpressure, queued_at: None, }); + let sent_bytes = match dispatched { + Ok(()) => predicted_bytes, + // The worker for this destination has exited. Seal + // here with the counter already reserved, so no + // counter is skipped, and send as the inline path does. + Err(job) => { + let wire = job.seal_inline().map_err(|e| NodeError::SendFailed { + node_addr: *node_addr, + reason: format!("encryption failed: {}", e), + })?; + transport + .send(&remote_addr, &wire) + .await + .map_err(|e| link_send_error(*node_addr, e))? + } + }; if let Some(peer) = self.peers.get_mut(node_addr) { - peer.link_stats_mut().record_sent(predicted_bytes); + peer.link_stats_mut().record_sent(sent_bytes); if let Some(mmp) = peer.mmp_mut() { - mmp.sender - .record_sent(counter, timestamp_ms, predicted_bytes); + mmp.sender.record_sent(counter, timestamp_ms, sent_bytes); } } return Ok(()); @@ -4120,25 +4137,7 @@ impl Node { let bytes_sent = transport .send(&remote_addr, &wire_packet) .await - .map_err(|e| match e { - TransportError::MtuExceeded { packet_size, mtu } => NodeError::MtuExceeded { - node_addr: *node_addr, - packet_size, - mtu, - }, - // Preserve the transport's own classification instead of - // flattening every non-MTU failure into one string. A caller - // that wants to keep its half-built state across an interface - // flap can only do that if the distinction survives to it. - other if other.is_transient() => NodeError::SendUnavailable { - node_addr: *node_addr, - reason: format!("transport send: {}", other), - }, - other => NodeError::SendFailed { - node_addr: *node_addr, - reason: format!("transport send: {}", other), - }, - })?; + .map_err(|e| link_send_error(*node_addr, e))?; // Update send statistics if let Some(peer) = self.peers.get_mut(node_addr) { @@ -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 /// routing read adapter the shell retains. It hands the sans-IO routing core /// borrowed peers plus raw `may_reach` / `link_cost` / `coords` diff --git a/src/node/netmon/mod.rs b/src/node/netmon/mod.rs index b67d8f16..74cec120 100644 --- a/src/node/netmon/mod.rs +++ b/src/node/netmon/mod.rs @@ -1,11 +1,10 @@ //! Transport-medium change detection. //! -//! A node that moves between media (WLAN → LAN, WLAN → 5G, a BLE adapter -//! coming or going) would otherwise learn about it only as *silence*: the peer -//! sits in the table until `node.link_dead_timeout_secs` reaps it, and the -//! reconnect then waits out whatever backoff the old medium had already -//! accumulated. The host kernel knew within milliseconds; the node would find -//! out half a minute later. +//! A node that moves between IP media (WLAN → LAN, WLAN → 5G) would otherwise +//! learn about it only as *silence*: the peer sits in the table until +//! `node.link_dead_timeout_secs` reaps it, and the reconnect then waits out +//! whatever backoff the old medium had already accumulated. The host kernel +//! knew within milliseconds; the node would find out half a minute later. //! //! This module closes that gap. It samples a coarse [`NetFingerprint`] of the //! host's network attachment and publishes a [`NetChange`] on the channel the @@ -110,9 +109,13 @@ //! is correct — there is nothing bound to the old path to repair. //! //! **A BLE adapter's state** is invisible here, as it was before: it is not an -//! IP attachment at all. That signal comes from the radio (BlueZ properties, -//! the Android callback) and belongs on this same channel, pushed by the BLE -//! transport rather than sampled here. +//! IP attachment at all, and nothing routes it onto this channel. The detector +//! below is the only source of a [`NetChange`]; the BLE transport publishes +//! none. On Android the embedder-facing half does exist: the app installs and +//! clears its radio through the `BleRadioSlot` returned by +//! `Node::enable_app_owned_ble_radio`, and the BLE transport re-resolves the +//! slot when it changes. That reaches only the BLE transport. The rest of the +//! node sees a radio going away as the loss of the links it carried. //! //! # Where the peer list comes from //! diff --git a/src/node/tests/discovery.rs b/src/node/tests/discovery.rs index 57ab19ee..f5a39ae1 100644 --- a/src/node/tests/discovery.rs +++ b/src/node/tests/discovery.rs @@ -1644,10 +1644,11 @@ async fn test_lookup_response_path_mtu_expires_without_a_session() { #[tokio::test] async fn test_replayed_lookup_response_does_not_extend_the_path_mtu_deadline() { - // The response carries no replay dedupe, so a captured one can be - // re-injected indefinitely. What bounds the damage is that a replay of a - // value already stored takes the keep-tighter arm, which does not touch - // the learn time: each injection buys one TTL, not one per packet. + // The response carries no replay dedupe of its own, so a captured one can + // be re-injected indefinitely. Accepting the first response clears the + // pending lookup, so each replay is dropped as unsolicited before it + // reaches the path-MTU write: each injection buys one TTL, not one per + // packet. The equal-value arm of that write is pinned by the next test. let mut node = make_node(); let from = make_node_addr(0xAA); @@ -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 // ============================================================================ diff --git a/src/node/tests/forwarding.rs b/src/node/tests/forwarding.rs index c092e2d8..80023fc9 100644 --- a/src/node/tests/forwarding.rs +++ b/src/node/tests/forwarding.rs @@ -155,6 +155,53 @@ async fn test_forwarding_ttl_two_transit_clears_the_gate() { ); } +/// The next hop is resolved only for a datagram the core can forward, and that +/// resolution refreshes the destination's cached coordinates. A last-hop +/// transit datagram (ttl=1) is dropped by the core, so it must not refresh +/// them; a ttl=2 datagram reaches the resolution and does. +#[tokio::test] +async fn test_forwarding_last_hop_transit_does_not_refresh_destination_coords() { + let mut node = make_node(); + let from = make_node_addr(0xAA); + let src = make_node_addr(0x01); + let dest = make_node_addr(0x02); + let root = make_node_addr(0xF0); + let coords = TreeCoordinate::from_addrs(vec![dest, root]).unwrap(); + + let now_ms = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis() as u64; + let stamped_ms = now_ms - node.coord_cache().default_ttl_ms() / 2; + node.coord_cache_mut() + .insert_verified(dest, coords, stamped_ms); + let last_used = |node: &Node| node.coord_cache().get_entry(&dest).unwrap().last_used(); + assert_eq!( + last_used(&node), + stamped_ms, + "precondition: the entry carries the past stamp" + ); + + for ttl in [1u8, 2] { + let dg = SessionDatagram::new(src, dest, vec![0x10, 0x00, 0x00, 0x00]).with_ttl(ttl); + let encoded = dg.encode(); + node.handle_session_datagram(&from, &encoded[1..], false) + .await; + if ttl == 1 { + assert_eq!( + last_used(&node), + stamped_ms, + "a transit ttl=1 datagram is dropped, so it must not refresh the destination's coords" + ); + } else { + assert!( + last_used(&node) > stamped_ms, + "a transit ttl=2 datagram reaches next-hop resolution, which refreshes the coords" + ); + } + } +} + // --- Local delivery --- #[tokio::test] diff --git a/src/node/tests/handshake.rs b/src/node/tests/handshake.rs index eec1600c..00bfd564 100644 --- a/src/node/tests/handshake.rs +++ b/src/node/tests/handshake.rs @@ -5680,3 +5680,212 @@ async fn a_terminal_rekey_msg1_send_failure_is_recorded_as_a_reject() { "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) { + // Let B take whatever A sent on promotion before measuring. + while let Ok(Some(packet)) = + timeout(Duration::from_millis(100), pair.packet_rx_b.recv()).await + { + pair.node_b.handle_encrypted_frame(packet).await; + } + pair.node_a.supervisor.encrypt_workers = Some(pool.clone()); + let peer = pair.node_a.get_peer(&pair.peer_b).unwrap(); + let counter_before = peer.noise_session().unwrap().current_send_counter(); + let sent_before = peer.link_stats().packets_sent; + let recv_before = pair + .node_b + .get_peer(&pair.peer_a) + .unwrap() + .link_stats() + .packets_recv; + + let (logs, guard) = crate::testutil::capture_logs_scoped(); + pair.node_a + .send_encrypted_link_message(&pair.peer_b, b"\x10dead worker test") + .await + .expect("link message send"); + drop(guard); + + let packet = next_packet(&mut pair.packet_rx_b).await; + let frame_counter = EncryptedHeader::parse(&packet.data) + .expect("an established frame") + .counter; + pair.node_b.handle_encrypted_frame(packet).await; + + let sent = pair + .node_a + .get_peer(&pair.peer_b) + .unwrap() + .link_stats() + .packets_sent + - sent_before; + let recv = pair + .node_b + .get_peer(&pair.peer_a) + .unwrap() + .link_stats() + .packets_recv + - recv_before; + let warnings = logs + .warnings() + .into_iter() + .filter(|line| line.contains("pool=\"encrypt\"")) + .collect(); + (frame_counter, counter_before, sent, recv, warnings) + } + + #[tokio::test] + async fn a_link_message_for_a_dead_workers_peer_is_still_sent() { + let mut pair = linked_pair().await; + let pool = pool_with_dead_worker_for(pair.addr_b); + let (frame_counter, counter_before, sent, recv, warnings) = + send_through(&mut pair, &pool).await; + + assert_eq!(recv, 1, "B did not authenticate the frame"); + assert_eq!( + frame_counter, counter_before, + "the frame must carry the counter reserved for it, not a fresh one" + ); + assert_eq!(sent, 1, "A counted the packet other than once"); + assert_eq!(pool.liveness().refused_dispatches(), 1); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + } + + #[tokio::test] + async fn a_link_message_through_live_workers_is_sent_by_the_worker() { + let mut pair = linked_pair().await; + let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]); + let (frame_counter, counter_before, sent, recv, warnings) = + send_through(&mut pair, &pool).await; + + assert_eq!(recv, 1, "B did not authenticate the frame"); + assert_eq!(frame_counter, counter_before); + assert_eq!(sent, 1); + assert_eq!(pool.liveness().refused_dispatches(), 0); + assert!(warnings.is_empty(), "{warnings:?}"); + } +} diff --git a/src/node/tests/mmp_chartests.rs b/src/node/tests/mmp_chartests.rs index 3a028e00..da772916 100644 --- a/src/node/tests/mmp_chartests.rs +++ b/src/node/tests/mmp_chartests.rs @@ -33,7 +33,7 @@ use crate::node::session::{EndToEndState, SessionEntry}; use crate::noise::HandshakeState; use crate::peer::ActivePeer; use crate::proto::mmp::{MmpMode, ReceiverReport}; -use crate::proto::stp::{ParentDeclaration, TreeCoordinate}; +use crate::proto::stp::{CoordEntry, ParentDeclaration, TreeCoordinate}; // =========================================================================== // Helpers @@ -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" ); } + +/// Insert a peer whose NodeAddr is strictly larger than the node's own, with +/// link MMP but no RTT yet, aged so a crafted ReceiverReport yields a first +/// RTT sample. The caller registers the peer's tree position. +fn insert_larger_unmeasured_peer(node: &mut Node) -> NodeAddr { + let my_addr = *node.node_addr(); + let (identity, addr) = loop { + let id = make_peer_identity(); + let a = *id.node_addr(); + if a > my_addr { + break (id, a); + } + }; + let mut peer = ActivePeer::new(identity, LinkId::new(1), 0); + peer.test_init_mmp(MmpMode::Full); + peer.test_backdate_session_start(std::time::Duration::from_secs(10)); + node.peers.insert(addr, peer); + addr +} + +/// Sum of the tree-announce fan-out counters. Every attempt to send an +/// announce to a peer moves exactly one of them. +fn tree_announce_attempts(node: &Node) -> u64 { + let tree = &node.metrics().tree; + tree.sent.get() + tree.send_failed.get() + tree.rate_limited.get() +} + +/// A root node whose only peer has a larger address has nothing to change on +/// the first RTT sample: it stays root, records no switch, and sends no +/// TreeAnnounce. The periodic tick rebroadcasts; the first-RTT path does not. +#[tokio::test] +async fn first_rtt_on_a_root_with_only_larger_peers_sends_no_tree_announce() { + let mut node = make_node(); + let addr = insert_larger_unmeasured_peer(&mut node); + node.tree_state_mut().update_peer( + ParentDeclaration::self_root(addr, 1, 0), + TreeCoordinate::root(addr), + ); + + assert!( + node.tree_state().is_root(), + "precondition: node starts as its own root" + ); + let switches_before = node.metrics().tree.parent_switches.get(); + let attempts_before = tree_announce_attempts(&node); + + node.handle_receiver_report(&addr, &craft_rr_payload(10, 5, 500)) + .await; + + assert!( + node.get_peer(&addr).unwrap().has_srtt(), + "precondition: this report was the peer's first RTT sample" + ); + assert!(node.tree_state().is_root(), "node remains its own root"); + assert_eq!( + node.metrics().tree.parent_switches.get(), + switches_before, + "no parent switch is recorded" + ); + assert_eq!( + tree_announce_attempts(&node), + attempts_before, + "the first-RTT path does not rebroadcast an unchanged declaration" + ); +} + +/// A node holding a parent whose tree has since re-rooted at a larger address +/// than its own promotes itself to root on the first RTT sample. +#[tokio::test] +async fn first_rtt_self_promotes_when_no_visible_root_is_smaller() { + let mut node = make_node(); + let addr = insert_larger_unmeasured_peer(&mut node); + + // The peer first sits under a root smaller than us, and we take it as + // parent, so our root is that smaller node. + let far_root = NodeAddr::from_bytes([0u8; 16]); + assert!( + far_root < *node.node_addr(), + "precondition: the far root is smaller than the node" + ); + node.tree_state_mut().update_peer( + ParentDeclaration::new(addr, far_root, 1, 0), + TreeCoordinate::new(vec![ + CoordEntry::new(addr, 1, 0), + CoordEntry::new(far_root, 1, 0), + ]) + .unwrap(), + ); + let seq = node.tree_state().my_declaration().sequence() + 1; + node.tree_state_mut() + .set_parent(addr, seq, 0, crate::time::mono_ms()); + node.tree_state_mut().recompute_coords(); + assert_eq!( + node.tree_state().root(), + &far_root, + "precondition: the node sits under the far root" + ); + + // The peer then re-roots at itself, larger than us: no visible root is + // smaller than the node any more. + node.tree_state_mut().update_peer( + ParentDeclaration::self_root(addr, 2, 0), + TreeCoordinate::root(addr), + ); + assert!( + !node.tree_state().is_root() && node.tree_state().should_be_root(), + "precondition: not root, but should be" + ); + let switches_before = node.metrics().tree.parent_switches.get(); + + node.handle_receiver_report(&addr, &craft_rr_payload(10, 5, 500)) + .await; + + assert!( + node.get_peer(&addr).unwrap().has_srtt(), + "precondition: this report was the peer's first RTT sample" + ); + assert!( + node.tree_state().is_root(), + "the first-RTT path promoted the node to root" + ); + assert_eq!( + node.metrics().tree.parent_switches.get(), + switches_before + 1, + "the self-promotion is recorded as one parent switch" + ); +} diff --git a/src/node/tests/mod.rs b/src/node/tests/mod.rs index f5efde00..06e18165 100644 --- a/src/node/tests/mod.rs +++ b/src/node/tests/mod.rs @@ -93,6 +93,31 @@ pub(super) fn install_connected_udp( .set_connected_udp(socket, drain); } +/// An encrypt pool of two whose worker for `dest` has exited. Where the worker +/// for a destination is not a function of the address alone (macOS), both +/// have. +#[cfg(unix)] +pub(super) fn pool_with_dead_worker_for( + dest: std::net::SocketAddr, +) -> crate::node::encrypt_worker::EncryptWorkerPool { + use crate::node::encrypt_worker::EncryptWorkerPool; + use crate::node::worker_set::TestWorker; + + let dead = EncryptWorkerPool::for_test(vec![TestWorker::FailSpawn, TestWorker::FailSpawn]) + .worker_index_for_dest(dest); + EncryptWorkerPool::for_test( + (0..2) + .map(|idx| { + if dead.is_none_or(|d| d == idx) { + TestWorker::FailSpawn + } else { + TestWorker::Run + } + }) + .collect(), + ) +} + /// Build a test node with an explicit `max_peers` limit (replaces the removed /// `set_max_peers` setter; resource limits are immutable post-construction). pub(super) fn make_node_with_max_peers(max_peers: usize) -> Node { diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index 66f27379..301df635 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -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" ); } + +/// 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; + } +} diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index a1e2a951..f07d8efb 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -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" ); } + +/// The tick's crypto worker liveness sweep: a lost worker degrades the node and +/// is logged with the live count; a healthy pool changes nothing. +#[cfg(unix)] +mod worker_liveness { + use super::*; + use crate::node::decrypt_worker::DecryptWorkerPool; + use crate::node::encrypt_worker::EncryptWorkerPool; + use crate::node::lifecycle::supervisor::{Child, SupervisorFsm}; + use crate::node::worker_set::{TestWorker, WorkerLiveness, wait_for}; + use std::collections::HashSet; + use std::sync::mpsc; + + /// A node seeded straight into a full `Running` with both pools up. + fn running_node() -> Node { + let mut node = make_node(); + node.supervisor.fsm = SupervisorFsm::running_with([ + Child::Transport(TransportId::new(1)), + Child::EncryptWorkers, + Child::DecryptWorkers, + ]); + node.supervisor.state = NodeState::Running; + node + } + + /// Two workers that each exit when their sender is used. + fn two_workers() -> (Vec, [mpsc::Sender<()>; 2]) { + let (kill0, exit0) = mpsc::channel(); + let (kill1, exit1) = mpsc::channel(); + ( + vec![TestWorker::ExitOn(exit0), TestWorker::ExitOn(exit1)], + [kill0, kill1], + ) + } + + /// An encrypt pool of two workers that each exit when their sender is + /// used. + fn pool_of_two() -> (EncryptWorkerPool, [mpsc::Sender<()>; 2]) { + let (plan, kills) = two_workers(); + (EncryptWorkerPool::for_test(plan), kills) + } + + /// Signal one worker of `workers` to exit and wait until `live_after` + /// remain. + fn kill_worker(workers: &dyn WorkerLiveness, kill: &mpsc::Sender<()>, live_after: usize) { + kill.send(()).expect("worker gone before its signal"); + assert!( + wait_for(|| workers.live_workers() == live_after), + "worker never exited" + ); + } + + fn kill(node: &Node, kill: &mpsc::Sender<()>, live_after: usize) { + let pool = node.supervisor.encrypt_workers.as_ref().unwrap(); + kill_worker(pool.liveness(), kill, live_after); + } + + #[tokio::test] + async fn a_dead_worker_degrades_the_node_and_names_the_live_count() { + let mut node = running_node(); + let (pool, kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(pool); + kill(&node, &kills[1], 1); + + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert_eq!(node.state(), NodeState::Degraded); + assert!( + node.supervisor + .fsm + .degraded_children() + .contains(&Child::EncryptWorkers) + ); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"encrypt\""), "{warnings:?}"); + assert!(warnings[0].contains(" live=1"), "{warnings:?}"); + assert!(warnings[0].contains(" configured=2"), "{warnings:?}"); + + // Nothing new: no second report. + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert!(logs.warnings().is_empty(), "{:?}", logs.warnings()); + assert_eq!(node.state(), NodeState::Degraded); + + // The last worker goes: reported with the new count, and the node is + // still degraded, not failed. + kill(&node, &kills[0], 0); + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" live=0"), "{warnings:?}"); + assert_eq!(node.state(), NodeState::Degraded); + } + + #[tokio::test] + async fn a_dead_decrypt_worker_degrades_the_node_and_names_its_pool() { + let mut node = running_node(); + let (encrypt, _encrypt_kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(encrypt); + let (plan, kills) = two_workers(); + node.supervisor.decrypt_workers = Some(DecryptWorkerPool::for_test(plan)); + let pool = node.supervisor.decrypt_workers.as_ref().unwrap(); + kill_worker(pool.liveness(), &kills[0], 1); + + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert_eq!(node.state(), NodeState::Degraded); + assert_eq!( + node.supervisor.fsm.degraded_children(), + HashSet::from([Child::DecryptWorkers]) + ); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"decrypt\""), "{warnings:?}"); + assert!(warnings[0].contains(" live=1"), "{warnings:?}"); + assert!(warnings[0].contains(" configured=2"), "{warnings:?}"); + } + + #[tokio::test] + async fn healthy_pools_leave_the_node_running() { + let mut node = running_node(); + let (pool, _kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(pool); + let (plan, _decrypt_kills) = two_workers(); + node.supervisor.decrypt_workers = Some(DecryptWorkerPool::for_test(plan)); + + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert!(logs.warnings().is_empty(), "{:?}", logs.warnings()); + assert_eq!(node.state(), NodeState::Running); + } + + /// What start-up made of the pools a test staged for it. + struct StartOutcome { + state: NodeState, + degraded: HashSet, + encrypt_installed: bool, + decrypt_installed: bool, + } + + /// Start a node that installs `encrypt` and `decrypt` in place of the + /// pools it would spawn, then stop it. Fails if start-up did not take both + /// staged pools, since the outcome would then say nothing about them. + async fn start_with_staged( + encrypt: EncryptWorkerPool, + decrypt: DecryptWorkerPool, + ) -> StartOutcome { + let mut node = make_healthy_node(); + node.supervisor.staged_pools.encrypt = Some(encrypt); + node.supervisor.staged_pools.decrypt = Some(decrypt); + node.start().await.unwrap(); + let staged = &node.supervisor.staged_pools; + let taken = staged.encrypt.is_none() && staged.decrypt.is_none(); + let outcome = StartOutcome { + state: node.state(), + degraded: node.supervisor.fsm.degraded_children(), + encrypt_installed: node.supervisor.encrypt_workers.is_some(), + decrypt_installed: node.supervisor.decrypt_workers.is_some(), + }; + node.stop().await.unwrap(); + assert!(taken, "start-up did not install both staged pools"); + outcome + } + + #[tokio::test] + async fn a_worker_that_fails_to_start_leaves_the_node_degraded_with_its_pool_installed() { + let outcome = start_with_staged( + EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]), + DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + ) + .await; + assert_eq!(outcome.state, NodeState::Degraded); + assert_eq!(outcome.degraded, HashSet::from([Child::EncryptWorkers])); + assert!( + outcome.encrypt_installed, + "a pool with a live worker is kept" + ); + assert!(outcome.decrypt_installed); + } + + #[tokio::test] + async fn a_pool_with_no_worker_started_leaves_the_node_degraded_and_is_not_installed() { + let outcome = start_with_staged( + EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + DecryptWorkerPool::for_test(vec![TestWorker::FailSpawn, TestWorker::FailSpawn]), + ) + .await; + assert_eq!(outcome.state, NodeState::Degraded); + assert_eq!(outcome.degraded, HashSet::from([Child::DecryptWorkers])); + assert!(outcome.encrypt_installed); + assert!( + !outcome.decrypt_installed, + "a pool with no live worker is dropped for the main-loop path" + ); + } + + #[tokio::test] + async fn staged_pools_with_every_worker_started_leave_the_node_running() { + let outcome = start_with_staged( + EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + ) + .await; + assert_eq!(outcome.state, NodeState::Running); + assert!(outcome.degraded.is_empty(), "{:?}", outcome.degraded); + assert!(outcome.encrypt_installed && outcome.decrypt_installed); + } + + /// Start a node, swap in a pool of two, optionally kill one worker, and + /// drive the real rx loop past one tick. Returns the state it published. + async fn drive_with_pool(kill_one: bool) -> NodeState { + let mut node = make_healthy_node(); + node.start().await.unwrap(); + assert_eq!(node.state(), NodeState::Running); + let (pool, kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(pool); + if kill_one { + kill(&node, &kills[1], 1); + } + + let drive = tokio::time::timeout( + Duration::from_millis(1500), + node.run_rx_loop_with_shutdown(std::future::pending()), + ) + .await; + assert!( + drive.is_err(), + "the rx loop must still be running: {drive:?}" + ); + let state = node.state(); + node.stop().await.unwrap(); + state + } + + #[tokio::test] + async fn a_dead_worker_degrades_the_node_through_the_rx_loop() { + assert_eq!(drive_with_pool(true).await, NodeState::Degraded); + } + + #[tokio::test] + async fn a_healthy_pool_leaves_the_node_running_through_the_rx_loop() { + assert_eq!(drive_with_pool(false).await, NodeState::Running); + } +} diff --git a/src/node/worker_set.rs b/src/node/worker_set.rs new file mode 100644 index 00000000..3c7f17cb --- /dev/null +++ b/src/node/worker_set.rs @@ -0,0 +1,245 @@ +//! The worker threads behind one crypto worker pool, and what is known about +//! whether each is still running. +//! +//! Both pools (`encrypt_worker`, `decrypt_worker`) are a set of OS threads, +//! each reached through its own bounded channel. A worker that exits, a panic +//! unwinding included, drops its receiver, so its channel closes and every +//! later dispatch to it is refused. Nothing restarts it. This module keeps the +//! thread handles so that loss can be seen, and counts the dispatches it +//! refused. +//! +//! It reports facts only. Whether a loss degrades the node is decided by the +//! driver in `lifecycle`. + +use portable_atomic::AtomicU64; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::thread::JoinHandle; +use tracing::warn; + +/// One worker: the sending end of its channel and its thread, `None` when the +/// thread could not be started. +struct Worker { + sender: S, + thread: Option>, +} + +/// The workers of one pool. Held behind an `Arc` by the pool so that cloning +/// the pool per packet stays a reference-count bump. +pub(crate) struct WorkerSet { + workers: Box<[Worker]>, + /// Dispatches refused because the target worker had exited. + refused: AtomicU64, + /// The live-worker count the liveness sweep last reported. Starts at the + /// number of threads that started, so a worker that never started is + /// reported once, at spawn, and not again by the sweep. It lives with the + /// workers so that a pool replaced by another starts from its own count. + reported_live: AtomicUsize, +} + +impl WorkerSet { + /// Start `n` workers, at least one. + /// + /// `channel` builds each worker's sender and receiver. `spawn` starts the + /// thread that owns the receiver. A worker whose thread cannot be started + /// is logged and left dead rather than aborting the caller: the receiver + /// was moved into the failed spawn and is dropped with it, so that + /// worker's channel is closed and a dispatch to it is refused, not queued. + pub(crate) fn start( + pool: &'static str, + n: usize, + mut channel: impl FnMut() -> (S, R), + mut spawn: impl FnMut(usize, R) -> std::io::Result>, + ) -> Self { + let n = n.max(1); + let mut workers = Vec::with_capacity(n); + for idx in 0..n { + let (sender, receiver) = channel(); + let thread = match spawn(idx, receiver) { + Ok(handle) => Some(handle), + Err(error) => { + warn!(pool, worker = idx, %error, "Failed to start a crypto worker thread"); + None + } + }; + workers.push(Worker { sender, thread }); + } + let started = workers.iter().filter(|w| w.thread.is_some()).count(); + Self { + workers: workers.into(), + refused: AtomicU64::new(0), + reported_live: AtomicUsize::new(started), + } + } + + /// Number of workers, including dead ones. Never zero. + pub(crate) fn len(&self) -> usize { + self.workers.len() + } + + /// The sending end of worker `idx`'s channel. + pub(crate) fn sender(&self, idx: usize) -> &S { + &self.workers[idx].sender + } + + /// Count one dispatch refused by a dead worker, returning the count before + /// this one. + pub(crate) fn note_refused(&self) -> u64 { + self.refused.fetch_add(1, Ordering::Relaxed) + } +} + +/// What the liveness sweep and the tests read from a pool, without naming its +/// channel type. +pub(crate) trait WorkerLiveness { + /// Number of workers the pool was built with. + fn worker_count(&self) -> usize; + /// Workers whose thread started and has not finished. A worker thread + /// finishes only by unwinding or when every sender to it is dropped, and + /// the node holds the pool while it runs, so a finished thread is a dead + /// worker. + fn live_workers(&self) -> usize; + /// Indices of the workers that are not live. + fn dead_workers(&self) -> Vec; + /// Dispatches refused because their worker had exited. + #[cfg(test)] + fn refused_dispatches(&self) -> u64; + /// Record `live` as the count last reported, returning the previous one. + fn swap_reported_live(&self, live: usize) -> usize; +} + +impl WorkerLiveness for WorkerSet { + fn worker_count(&self) -> usize { + self.workers.len() + } + + fn live_workers(&self) -> usize { + self.workers.iter().filter(|w| is_live(w)).count() + } + + fn dead_workers(&self) -> Vec { + self.workers + .iter() + .enumerate() + .filter(|(_, w)| !is_live(w)) + .map(|(idx, _)| idx) + .collect() + } + + #[cfg(test)] + fn refused_dispatches(&self) -> u64 { + self.refused.load(Ordering::Relaxed) + } + + fn swap_reported_live(&self, live: usize) -> usize { + self.reported_live.swap(live, Ordering::Relaxed) + } +} + +fn is_live(worker: &Worker) -> bool { + worker.thread.as_ref().is_some_and(|t| !t.is_finished()) +} + +/// Whether the `n`th event (counting from zero) of a repeating condition is +/// logged: the first eight, then one in ten thousand. +pub(crate) fn worth_logging(n: u64) -> bool { + n < 8 || n.is_multiple_of(10_000) +} + +/// How a test wants one worker of a pool to behave. +#[cfg(test)] +pub(crate) enum TestWorker { + /// The production worker loop. + Run, + /// The thread fails to start. + FailSpawn, + /// The thread holds its receiver until signalled, then exits. + ExitOn(std::sync::mpsc::Receiver<()>), +} + +/// A spawner that starts each worker as `plan` says, using `run` for +/// [`TestWorker::Run`]. +#[cfg(test)] +pub(crate) fn test_spawner( + plan: Vec, + run: impl Fn(usize, R) -> std::io::Result>, +) -> impl FnMut(usize, R) -> std::io::Result> { + let mut plan: Vec> = plan.into_iter().map(Some).collect(); + move |idx, rx| match plan[idx].take().expect("each worker is started once") { + TestWorker::Run => run(idx, rx), + TestWorker::FailSpawn => { + drop(rx); + Err(std::io::Error::other("worker start refused by the test")) + } + TestWorker::ExitOn(signal) => std::thread::Builder::new().spawn(move || { + let _ = signal.recv(); + drop(rx); + }), + } +} + +/// Wait up to five seconds for `cond`, returning whether it came true. +#[cfg(test)] +pub(crate) fn wait_for(cond: impl Fn() -> bool) -> bool { + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + while !cond() { + if std::time::Instant::now() >= deadline { + return false; + } + std::thread::sleep(std::time::Duration::from_millis(5)); + } + true +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::mpsc; + + #[test] + fn a_worker_that_exits_is_no_longer_live_and_is_named_dead() { + let (stop_tx, stop_rx) = mpsc::channel::<()>(); + let mut stop_rx = Some(stop_rx); + let set = WorkerSet::start( + "test", + 2, + mpsc::channel::, + |idx, rx: mpsc::Receiver| { + let stop = if idx == 1 { stop_rx.take() } else { None }; + std::thread::Builder::new().spawn(move || match stop { + Some(stop) => { + let _ = stop.recv(); + drop(rx); + } + None => while rx.recv().is_ok() {}, + }) + }, + ); + assert_eq!(set.live_workers(), 2); + stop_tx.send(()).unwrap(); + assert!( + wait_for(|| set.live_workers() == 1), + "worker 1 never exited" + ); + assert_eq!(set.dead_workers(), vec![1]); + assert_eq!(set.worker_count(), 2); + } + + #[test] + fn the_reported_baseline_starts_at_the_workers_that_started() { + let set = WorkerSet::start("test", 3, mpsc::channel::, |idx, rx| { + if idx == 2 { + drop(rx); + return Err(std::io::Error::other("refused by the test")); + } + std::thread::Builder::new().spawn(move || while rx.recv().is_ok() {}) + }); + assert_eq!(set.swap_reported_live(2), 2); + assert_eq!(set.dead_workers(), vec![2]); + } + + #[test] + fn worth_logging_keeps_the_first_eight_then_one_in_ten_thousand() { + let logged: Vec = (0..20_001).filter(|n| worth_logging(*n)).collect(); + assert_eq!(logged, vec![0, 1, 2, 3, 4, 5, 6, 7, 10_000, 20_000]); + } +} diff --git a/src/proto/fsp/core.rs b/src/proto/fsp/core.rs index f45526a0..8a3b25c4 100644 --- a/src/proto/fsp/core.rs +++ b/src/proto/fsp/core.rs @@ -483,10 +483,10 @@ impl Fsp { } /// Decide whether a path-MTU update should tighten the shared lookup: emit - /// `TightenPathMtuLookup` only when `candidate` is at least as tight as the - /// `existing` value (keep-tighter, never loosen). The `existing` read and - /// the applied write are performed shell-side under one `path_mtu_lookup` - /// write guard, so the decision stays atomic. + /// `TightenPathMtuLookup` only when there is no `existing` value or + /// `candidate` is strictly tighter than it (keep-tighter, never loosen). + /// The `existing` read and the applied write are performed shell-side + /// under one `path_mtu_lookup` write guard, so the decision stays atomic. pub(crate) fn plan_path_mtu_tighten( &self, fips_addr: FipsAddress, @@ -524,7 +524,8 @@ pub(crate) fn initiation_winner(our_node_addr: &NodeAddr, their_node_addr: &Node /// Decide whether a path-MTU update should be applied to the shared /// `FipsAddress`-keyed lookup: keep the tighter of existing-or-candidate, never /// loosen. Returns `true` when `candidate` should be written (there is no -/// existing value, or the candidate is at least as tight). +/// existing value, or the candidate is strictly tighter). An equal candidate is +/// not written, so the stored entry, and any learn time it carries, is kept. pub(crate) fn should_apply_path_mtu(existing: Option, candidate: u16) -> bool { !matches!(existing, Some(existing) if existing <= candidate) } diff --git a/src/proto/fsp/mod.rs b/src/proto/fsp/mod.rs index 7c9b37f5..15df2ff5 100644 --- a/src/proto/fsp/mod.rs +++ b/src/proto/fsp/mod.rs @@ -33,7 +33,7 @@ mod tests; pub(crate) use core::{ DecryptSlot, EpochReaction, Fsp, FspAction, InitialMsg3ResendSnapshot, RekeyCfg, RekeyMsg3ResendSnapshot, SessionSnapshot, cutover_timer_elapsed, initiation_winner, - mark_ipv6_ecn_ce, push_bounded_pending, + mark_ipv6_ecn_ce, push_bounded_pending, should_apply_path_mtu, }; pub use wire::{ FspInnerFlags, SessionAck, SessionFlags, SessionMessageType, SessionMsg3, SessionSetup, diff --git a/src/proto/link.rs b/src/proto/link.rs index 51b1c777..95f7c226 100644 --- a/src/proto/link.rs +++ b/src/proto/link.rs @@ -144,6 +144,20 @@ pub struct SessionDatagramRef<'a> { pub payload: &'a [u8], } +/// The TTL a transit datagram leaves this node with, or `None` when it may not +/// be transmitted because it would leave with zero. +/// +/// Follows IP semantics: the decrement comes first, and `saturating_sub` folds +/// an already-exhausted arrival (TTL 0) into the same outcome as a last-hop +/// arrival (TTL 1). This is the one rule behind the routing core's hop-limit +/// drop and both `can_forward`s. +pub(crate) fn ttl_after_hop(ttl: u8) -> Option { + match ttl.saturating_sub(1) { + 0 => None, + left => Some(left), + } +} + /// SessionDatagram fixed header size: msg_type(1) + ttl(1) + path_mtu(2) + src_addr(16) + dest_addr(16). pub const SESSION_DATAGRAM_HEADER_SIZE: usize = 36; @@ -191,7 +205,7 @@ impl SessionDatagram { /// True only at TTL 2 or more: at TTL 1 the decrement leaves zero, so the /// datagram is dropped rather than forwarded. pub fn can_forward(&self) -> bool { - self.ttl > 1 + ttl_after_hop(self.ttl).is_some() } /// Encode as link-layer message (msg_type + ttl + path_mtu + src_addr + dest_addr + payload). @@ -239,6 +253,12 @@ impl<'a> SessionDatagramRef<'a> { }) } + /// Check whether this datagram would survive a transit hop, by the same + /// rule the routing core drops on (`ttl_after_hop`). + pub fn can_forward(&self) -> bool { + ttl_after_hop(self.ttl).is_some() + } + /// Materialize an owned datagram for forwarding/re-encoding paths. pub fn into_owned(self) -> SessionDatagram { SessionDatagram { @@ -424,6 +444,36 @@ mod tests { assert!(dg.with_ttl(255).can_forward()); } + #[test] + fn ttl_after_hop_drops_only_what_would_leave_at_zero() { + assert_eq!( + ttl_after_hop(0), + None, + "an exhausted arrival leaves at zero" + ); + assert_eq!(ttl_after_hop(1), None, "a last-hop arrival leaves at zero"); + assert_eq!(ttl_after_hop(2), Some(1)); + assert_eq!(ttl_after_hop(255), Some(254)); + + let dg = SessionDatagram::new(make_node_addr(1), make_node_addr(2), vec![0x42]); + for ttl in 0..=u8::MAX { + assert_eq!( + ttl_after_hop(ttl).is_some(), + ttl > 1, + "ttl={ttl}: only a TTL of 2 or more survives the hop" + ); + let owned = dg.clone().with_ttl(ttl); + let encoded = owned.encode(); + let view = SessionDatagramRef::decode(&encoded[1..]).unwrap(); + assert_eq!( + view.can_forward(), + owned.can_forward(), + "ttl={ttl}: the borrowed and owned views must agree" + ); + assert_eq!(view.can_forward(), ttl_after_hop(ttl).is_some()); + } + } + #[test] fn test_session_datagram_decrement_ttl() { let base = SessionDatagram::new(make_node_addr(1), make_node_addr(2), vec![0x42]); diff --git a/src/proto/routing/core.rs b/src/proto/routing/core.rs index 8d30480d..d7e4108a 100644 --- a/src/proto/routing/core.rs +++ b/src/proto/routing/core.rs @@ -17,7 +17,7 @@ use super::limits::LimitVerdict; use super::state::Router; use super::wire::{CoordsRequired, MtuExceeded, PathBroken}; -use crate::proto::link::{SessionDatagram, SessionDatagramRef}; +use crate::proto::link::{SessionDatagram, SessionDatagramRef, ttl_after_hop}; use crate::{NodeAddr, TreeCoordinate}; /// Read-only view of routing state the routing core needs. @@ -116,7 +116,8 @@ impl Router { /// datagram that would leave with a TTL of zero is not transmitted. /// /// The shell pre-resolves `next_hop` only for datagrams this can actually - /// forward (dest not local and TTL surviving the decrement), so + /// forward (dest not local, and `SessionDatagramRef::can_forward`, which + /// applies the same [`ttl_after_hop`] rule this drops on), so /// `find_next_hop`'s LRU-touch side effect stays scoped to genuine /// forwards. `route` still re-checks local delivery and the TTL /// authoritatively. @@ -136,15 +137,14 @@ impl Router { } // TTL enforcement on the transit path: decrement first, then drop if - // the datagram would leave with a TTL of zero. `saturating_sub` folds - // the already-exhausted arrival (ttl=0) into the same test as the - // last-hop arrival (ttl=1); neither is transmitted. - let forwarded_ttl = dg.ttl.saturating_sub(1); - if forwarded_ttl == 0 { + // the datagram would leave with a TTL of zero. The already-exhausted + // arrival (ttl=0) and the last-hop arrival (ttl=1) are both dropped; + // neither is transmitted. + let Some(forwarded_ttl) = ttl_after_hop(dg.ttl) else { return RouteOutcome::Drop { reason: DropReason::TtlExhausted, }; - } + }; let nh = match next_hop { Some(nh) => nh, diff --git a/src/proto/routing/tests/core.rs b/src/proto/routing/tests/core.rs index 6eeb945f..10fd97f4 100644 --- a/src/proto/routing/tests/core.rs +++ b/src/proto/routing/tests/core.rs @@ -1,7 +1,7 @@ //! Tests for the sans-IO routing decision core. use super::util::{MockPeer, MockRoutingView, make_coords, make_datagram_ref, make_next_hop}; -use crate::proto::link::SessionDatagramRef; +use crate::proto::link::{SessionDatagramRef, ttl_after_hop}; use crate::proto::routing::RoutingSignalType; use crate::proto::routing::{ DropReason, LimitVerdict, RouteAction, RouteOutcome, Router, RoutingView, select_best_candidate, @@ -560,3 +560,45 @@ fn candidate_selection_excludes_non_full_peers() { "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"), + } + } +} diff --git a/src/proto/stp/core.rs b/src/proto/stp/core.rs index 8181e94f..04e2fffe 100644 --- a/src/proto/stp/core.rs +++ b/src/proto/stp/core.rs @@ -165,6 +165,8 @@ impl Stp { /// `classify_announce`, the periodic path has no same-parent loop-drop / /// ancestry-update arms — a periodic tick has no announcing peer, so those cases /// never arise; the no-change tail is a re-broadcast rather than a true no-op. + /// The first-RTT re-evaluation in `node::handlers::mmp` is driven by it too, + /// and ignores `PeriodicRebroadcast`. pub(crate) fn classify_periodic( tree: &TreeState, peer_costs: &BTreeMap, diff --git a/src/proto/stp/mod.rs b/src/proto/stp/mod.rs index 2c679e24..677bfd96 100644 --- a/src/proto/stp/mod.rs +++ b/src/proto/stp/mod.rs @@ -34,7 +34,11 @@ pub use crate::proto::coord::{CoordEntry, CoordError, TreeCoordinate}; pub(crate) use crate::proto::coord::{ coords_wire_size, decode_coords, decode_optional_coords, encode_coords, encode_empty_coords, }; -pub(crate) use core::{ParentEval, Stp, TreeDecision}; +// Callers outside this module take the decision from `Stp`; only the tests +// name the parent evaluation itself. +#[cfg(test)] +pub(crate) use core::ParentEval; +pub(crate) use core::{Stp, TreeDecision}; pub use declaration::ParentDeclaration; pub use state::TreeState; pub use wire::TreeAnnounce; diff --git a/testing/chaos/sim/runner.py b/testing/chaos/sim/runner.py index 7af8660a..957bd828 100644 --- a/testing/chaos/sim/runner.py +++ b/testing/chaos/sim/runner.py @@ -48,8 +48,9 @@ from .veth import VethManager log = logging.getLogger(__name__) -# The final snapshot waits for this many identical consecutive tree reads, -# taken this far apart, so the tree must hold still for two intervals. +# The final snapshot waits until every node answers and this many identical +# consecutive tree reads, taken this far apart, agree, so the tree must hold +# still, with every node in it, for two intervals. SETTLE_READS = 3 SETTLE_INTERVAL_SECS = 5 SETTLE_TIMEOUT_SECS = 90 @@ -744,7 +745,8 @@ class SimRunner: # Take final tree snapshot while nodes are still running, once the # tree has stopped moving. A node restored a moment ago is its own # root until it re-parents, so a snapshot taken straight after the - # restore reads a mesh still converging. + # restore reads a mesh still converging. It may not answer at all + # yet either, and the settle waits for it. self._settle_tree() self._take_snapshot("final") @@ -893,29 +895,39 @@ class SimRunner: return result def _settle_tree(self): - """Wait until consecutive tree reads agree, or the settle time runs out. + """Wait until every node answers and consecutive tree reads agree. - Compares each answering node's root and parent. A fixed delay would - either waste time on a mesh that settled at once or cut off one that - had not. Running out is logged and is not a failure in itself: the - final snapshot is taken anyway, and the assertions judge what it - shows. + Returns once SETTLE_READS reads in a row, SETTLE_INTERVAL_SECS apart, + have each been answered by every node in the topology and show the + same root and parent for each, or once SETTLE_TIMEOUT_SECS has passed. + The bound is checked after each read, so the settle can return up to + one interval and one read past it. A fixed delay would either waste + time on a mesh that settled at once or cut off one that had not. - A read that no node answered never counts toward agreement. It does - not catch a node that stays its own root for longer than the reads - span, which is a tree that is stable and wrong, and is left to the - assertions. + A read that any node did not answer never counts toward agreement. A + node restored at teardown may not have opened its control socket yet, + and a tree that holds still without it is not the tree the final + snapshot is meant to record. + + Running out is logged, naming any node that still does not answer, + and is not a failure in itself: the final snapshot is taken anyway + and the assertions judge what it shows. A node that never comes back + is therefore reported as absent by the assertions that count nodes. + This does not catch a node that answers but stays its own root for + longer than the reads span, which is a tree that is stable and + wrong, and is left to the assertions. """ started = time.time() previous = None agreeing = 0 while not self._interrupted: trees = snapshot_all_trees(self.topology) + missing = sorted(set(self.topology.nodes) - trees.keys()) shape = { nid: (data.get("root"), data.get("parent")) for nid, data in trees.items() } - if not shape: + if missing: agreeing = 0 else: agreeing = agreeing + 1 if shape == previous else 1 @@ -925,10 +937,19 @@ class SimRunner: log.info("Tree settled after %.0fs", waited) return if waited >= SETTLE_TIMEOUT_SECS: - log.warning( - "Tree still changing after %.0fs; taking the final snapshot anyway", - waited, - ) + if missing: + log.warning( + "%s still not answering after %.0fs; " + "taking the final snapshot anyway", + ", ".join(missing), + waited, + ) + else: + log.warning( + "Tree still changing after %.0fs; " + "taking the final snapshot anyway", + waited, + ) return self._sleep(SETTLE_INTERVAL_SECS)