From a9423f801ac6faf8d3ac835a9710a94ece650535 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Fri, 2 Oct 2026 15:48:20 +0000 Subject: [PATCH] Report crypto worker deaths as degraded health instead of hiding them The encrypt and decrypt worker pools discarded their thread handles, so a worker that panicked left the node reporting full health while every packet hashed to that worker was dropped behind a DEBUG line. A worker thread that could not be started panicked start-up instead. The pools now keep each worker's thread handle and report how many are live. A worker that cannot be started is logged and left dead, and the node starts degraded instead of aborting; a pool with no live worker is not installed, so its traffic takes the main-loop path. A sweep on the rx loop tick notices a worker that exits at runtime, logs a warning naming the pool and the live and configured counts, and reports the pool to the supervisor, which publishes Degraded. Losing workers never fails the node: the pools are an offload, and Windows runs without them. A dispatch refused by an exited worker is now counted and logged at WARN, rate-limited, in place of the DEBUG line. --- src/instr/recorder.rs | 15 +- src/node/dataplane/rx_loop.rs | 16 +- src/node/decrypt_worker.rs | 220 ++++++++++++++++---- src/node/encrypt_worker.rs | 338 +++++++++++++++++++++++++------ src/node/lifecycle/mod.rs | 66 +++--- src/node/lifecycle/supervisor.rs | 67 +++++- src/node/lifecycle/workers.rs | 240 ++++++++++++++++++++++ src/node/mod.rs | 2 + src/node/tests/unit.rs | 241 ++++++++++++++++++++++ src/node/worker_set.rs | 245 ++++++++++++++++++++++ 10 files changed, 1304 insertions(+), 146 deletions(-) create mode 100644 src/node/lifecycle/workers.rs create mode 100644 src/node/worker_set.rs diff --git a/src/instr/recorder.rs b/src/instr/recorder.rs index 93bd51b6..98b430d4 100644 --- a/src/instr/recorder.rs +++ b/src/instr/recorder.rs @@ -41,7 +41,7 @@ impl Domain { /// /// `as usize` indexes the counter arrays, so the discriminants are dense and /// `WholeTick` is last (it defines `N_STEPS`). Variants are declared -/// unconditionally — see [`Step::emitted`] for how the two platform- and +/// unconditionally — see [`Step::emitted`] for how the three platform- and /// profile-conditional steps are kept out of the emitted table. #[derive(Copy, Clone, Debug, PartialEq, Eq)] #[repr(usize)] @@ -73,6 +73,7 @@ pub(crate) enum Step { PollTransportDiscovery, SampleTransportCongestion, ActivateConnectedUdpSessions, + PollWorkerLiveness, DebugAssertPeerMapsCoherent, /// The whole tick-arm body, from before `check_timeouts` to after the last /// step. Composes safely with the per-step spans because the macro @@ -112,6 +113,7 @@ pub(crate) const STEPS: [Step; N_STEPS] = [ Step::PollTransportDiscovery, Step::SampleTransportCongestion, Step::ActivateConnectedUdpSessions, + Step::PollWorkerLiveness, Step::DebugAssertPeerMapsCoherent, Step::WholeTick, ]; @@ -146,6 +148,7 @@ impl Step { Step::PollTransportDiscovery => "poll_transport_discovery", Step::SampleTransportCongestion => "sample_transport_congestion", Step::ActivateConnectedUdpSessions => "activate_connected_udp_sessions", + Step::PollWorkerLiveness => "poll_worker_liveness", Step::DebugAssertPeerMapsCoherent => "debug_assert_peer_maps_coherent", Step::WholeTick => "whole_tick", } @@ -153,7 +156,7 @@ impl Step { /// Whether this step gets a row in this build. /// - /// Two steps are conditionally compiled at their call sites. Emitting a row + /// Three steps are conditionally compiled at their call sites. Emitting a row /// for them in a build where the call site does not exist would publish a /// count that is structurally zero forever, which reads as "this step never /// runs" rather than "this step is not in this build". The predicates below @@ -164,6 +167,7 @@ impl Step { Step::ActivateConnectedUdpSessions => { cfg!(any(target_os = "linux", target_os = "macos")) } + Step::PollWorkerLiveness => cfg!(unix), Step::DebugAssertPeerMapsCoherent => cfg!(debug_assertions), _ => true, } @@ -381,12 +385,15 @@ mod tests { #[test] fn emitted_row_count_matches_build() { let emitted = STEPS.iter().filter(|s| s.emitted()).count(); - // 26 unconditional subsystem steps + the whole-tick span, plus the two - // conditionally-compiled steps where this build has them. + // 26 unconditional subsystem steps + the whole-tick span, plus the + // three conditionally-compiled steps where this build has them. let mut expected = 27; if cfg!(any(target_os = "linux", target_os = "macos")) { expected += 1; } + if cfg!(unix) { + expected += 1; + } if cfg!(debug_assertions) { expected += 1; } diff --git a/src/node/dataplane/rx_loop.rs b/src/node/dataplane/rx_loop.rs index bb26d78a..425e7c8d 100644 --- a/src/node/dataplane/rx_loop.rs +++ b/src/node/dataplane/rx_loop.rs @@ -345,17 +345,7 @@ impl Node { // republishing health, so nothing outside the node can // observe an address the listener no longer answers on. self.retract_child_publications(child); - let actions = self - .supervisor - .fsm - .step(crate::node::lifecycle::supervisor::Event::ChildExited { child }); - for action in actions { - if let crate::node::lifecycle::supervisor::Action::PublishState(ns) = - action - { - self.supervisor.state = ns; - } - } + self.step_child_exited(child); // A transport child exiting leaves the bound set, so // it can be the one that was holding the node's egress // MTU down. `is_bound()` is `is_operational()` plus the @@ -583,6 +573,10 @@ impl Node { #[cfg(any(target_os = "linux", target_os = "macos"))] instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ActivateConnectedUdpSessions, self.activate_connected_udp_sessions().await); + // Crypto worker threads that exited since the last tick. + #[cfg(unix)] + instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PollWorkerLiveness, + self.poll_worker_liveness()); // Debug-build sweep of the peer-lifecycle map invariant // (leaked machines / machine-less legs); two map scans, // compiled out of release builds. diff --git a/src/node/decrypt_worker.rs b/src/node/decrypt_worker.rs index d4d5173c..63449205 100644 --- a/src/node/decrypt_worker.rs +++ b/src/node/decrypt_worker.rs @@ -36,6 +36,7 @@ #![cfg_attr(not(unix), allow(dead_code))] use crate::NodeAddr; +use crate::node::worker_set::{WorkerLiveness, WorkerSet, worth_logging}; use crate::transport::{TransportAddr, TransportId}; use crossbeam_channel::{Receiver, Sender, TrySendError, bounded}; use portable_atomic::{AtomicU64, Ordering}; @@ -220,26 +221,49 @@ pub(crate) enum WorkerMsg { /// shard. #[derive(Clone)] pub(crate) struct DecryptWorkerPool { - senders: Arc<[Sender]>, + workers: Arc>>, +} + +/// Start the production worker loop on `rx` in a named OS thread. +fn spawn_worker( + idx: usize, + rx: Receiver, +) -> std::io::Result> { + std::thread::Builder::new() + .name(format!("fips-decrypt-{idx}")) + .spawn(move || run_worker(idx, rx)) } impl DecryptWorkerPool { + /// Spawn `n` worker OS threads (at least one). A worker thread that + /// cannot be started is logged and left dead; the caller reads how + /// many started from [`Self::liveness`]. pub fn spawn(n: usize) -> Self { - let n = n.max(1); - let mut senders = Vec::with_capacity(n); - for i in 0..n { - let (tx, rx) = bounded::(WORKER_CHANNEL_CAP); - std::thread::Builder::new() - .name(format!("fips-decrypt-{i}")) - .spawn(move || run_worker(i, rx)) - .expect("failed to spawn fips-decrypt OS thread"); - senders.push(tx); - } + Self::start_with(n, spawn_worker) + } + + /// Build the pool, starting each worker with `spawn`. Production passes + /// [`spawn_worker`]; tests pass workers that fail to start or exit on cue. + fn start_with( + n: usize, + spawn: impl FnMut(usize, Receiver) -> std::io::Result>, + ) -> Self { Self { - senders: senders.into(), + workers: Arc::new(WorkerSet::start( + "decrypt", + n, + || bounded::(WORKER_CHANNEL_CAP), + spawn, + )), } } + /// Whether each worker is still running, and how many dispatches a dead + /// one refused. + pub(crate) fn liveness(&self) -> &dyn WorkerLiveness { + &*self.workers + } + /// Stable hash from session key → worker index. Same hash is used /// for session registration and per-packet dispatch so packets and /// registration arrive at the same shard. @@ -247,24 +271,37 @@ impl DecryptWorkerPool { use std::hash::{Hash, Hasher}; let mut h = std::collections::hash_map::DefaultHasher::new(); cache_key.hash(&mut h); - (h.finish() as usize) % self.senders.len() + (h.finish() as usize) % self.workers.len() + } + + /// Count a message refused by the exited worker `idx`, and log it at a + /// bounded rate. + fn note_refused(&self, idx: usize, what: &'static str) { + let n = self.workers.note_refused(); + if worth_logging(n) { + warn!( + pool = "decrypt", + worker = idx, + refused = n + 1, + what, + "Decrypt worker has exited; message refused" + ); + } } /// Dispatch a per-packet decrypt job. Drops if the per-worker /// channel is full (sustained rate overrun); the rx_loop's drain /// caps inbound at the same scale upstream so the cliff is - /// bounded. + /// bounded. A job for an exited worker is dropped, counted and + /// logged at WARN. pub fn dispatch_job(&self, job: DecryptJob) { - if self.senders.is_empty() { - return; - } let idx = self.worker_idx_for(job.cache_key); - match self.senders[idx].try_send(WorkerMsg::Job(job)) { + match self.workers.sender(idx).try_send(WorkerMsg::Job(job)) { Ok(()) => {} Err(TrySendError::Full(_)) => { static FULL_COUNT: AtomicU64 = AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, Ordering::Relaxed); - if n < 8 || n.is_multiple_of(10000) { + if worth_logging(n) { warn!( worker = idx, drops = n + 1, @@ -272,9 +309,7 @@ impl DecryptWorkerPool { ); } } - Err(TrySendError::Disconnected(_)) => { - debug!(worker = idx, "DecryptWorker thread gone; dropping job"); - } + Err(TrySendError::Disconnected(_)) => self.note_refused(idx, "inbound packet"), } } @@ -303,11 +338,12 @@ impl DecryptWorkerPool { cache_key: (TransportId, u32), state: OwnedSessionState, ) -> bool { - if self.senders.is_empty() { - return false; - } let idx = self.worker_idx_for(cache_key); - match self.senders[idx].try_send(WorkerMsg::RegisterSession { cache_key, state }) { + match self + .workers + .sender(idx) + .try_send(WorkerMsg::RegisterSession { cache_key, state }) + { Ok(()) => true, Err(TrySendError::Full(_)) => { warn!( @@ -317,10 +353,7 @@ impl DecryptWorkerPool { false } Err(TrySendError::Disconnected(_)) => { - debug!( - worker = idx, - "DecryptWorker thread gone; ignoring registration" - ); + self.note_refused(idx, "session registration"); false } } @@ -329,11 +362,11 @@ impl DecryptWorkerPool { /// Drop a session from its worker (rekey, peer removed). Fire and /// forget — if the worker is gone we don't care. pub fn unregister_session(&self, cache_key: (TransportId, u32)) { - if self.senders.is_empty() { - return; - } let idx = self.worker_idx_for(cache_key); - let _ = self.senders[idx].try_send(WorkerMsg::UnregisterSession { cache_key }); + let _ = self + .workers + .sender(idx) + .try_send(WorkerMsg::UnregisterSession { cache_key }); } } @@ -774,3 +807,122 @@ mod tests { } } } + +#[cfg(test)] +impl DecryptWorkerPool { + /// A pool whose workers behave as `plan` says; `Run` workers are the + /// production loop. + pub(crate) fn for_test(plan: Vec) -> Self { + Self::start_with( + plan.len(), + crate::node::worker_set::test_spawner(plan, spawn_worker), + ) + } +} + +/// The pool's view of its workers: which are live and what a dead one does to +/// a dispatch or a registration. Every wait is bounded. +#[cfg(test)] +mod pool_tests { + use super::*; + use crate::node::worker_set::{TestWorker, wait_for}; + use ring::aead::UnboundKey; + use std::sync::mpsc; + + /// A session key the pool hashes to worker `idx`. + fn key_on_worker(pool: &DecryptWorkerPool, idx: usize) -> (TransportId, u32) { + (0u32..) + .map(|n| (TransportId::new(1), n)) + .find(|key| pool.worker_idx_for(*key) == idx) + .expect("some key hashes to every worker") + } + + fn session_state() -> OwnedSessionState { + let key = UnboundKey::new(&ring::aead::CHACHA20_POLY1305, &[0u8; 32]).unwrap(); + OwnedSessionState { + fmp_cipher: LessSafeKey::new(key), + fmp_replay: ReplayWindow::new(), + source_npub: None, + } + } + + fn job(cache_key: (TransportId, u32)) -> DecryptJob { + let (fallback_tx, _) = tokio::sync::mpsc::unbounded_channel::(); + DecryptJob { + packet_data: vec![0u8; 48], + cache_key, + _transport_id: cache_key.0, + _remote_addr: TransportAddr::from_string("127.0.0.1:1234"), + timestamp_ms: 1_000, + source_node_addr: NodeAddr::from_bytes([1u8; 16]), + fmp_counter: 1, + fmp_flags: 0, + fmp_header: [0u8; 16], + fmp_ciphertext_offset: 16, + fallback_tx, + } + } + + #[test] + fn a_decrypt_worker_that_panics_is_counted_dead() { + let (die_tx, die_rx) = mpsc::channel::<()>(); + let mut die_rx = Some(die_rx); + let pool = DecryptWorkerPool::start_with(2, |idx, rx| { + if idx != 1 { + return spawn_worker(idx, rx); + } + let die = die_rx.take().expect("worker 1 starts once"); + std::thread::Builder::new().spawn(move || { + let _rx = rx; + let _ = die.recv(); + panic!("simulated decrypt worker panic"); + }) + }); + assert_eq!(pool.liveness().live_workers(), 2); + die_tx.send(()).expect("worker 1 gone before its signal"); + assert!( + wait_for(|| pool.liveness().live_workers() == 1), + "a worker that panicked is still counted live" + ); + assert_eq!(pool.liveness().dead_workers(), vec![1]); + } + + #[test] + fn a_worker_that_fails_to_spawn_leaves_the_pool_serving_the_rest() { + let pool = DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + assert_eq!(pool.liveness().worker_count(), 2); + assert_eq!(pool.liveness().live_workers(), 1); + let key = key_on_worker(&pool, 0); + assert!( + pool.register_session(key, session_state()), + "the live worker refused a registration" + ); + pool.dispatch_job(job(key)); + assert_eq!(pool.liveness().refused_dispatches(), 0); + } + + #[test] + fn dispatch_to_an_exited_decrypt_worker_warns_and_counts() { + let pool = DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + let key = key_on_worker(&pool, 1); + let ((), logs) = crate::testutil::capture_logs(|| pool.dispatch_job(job(key))); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" worker=1"), "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"decrypt\""), "{warnings:?}"); + assert_eq!(pool.liveness().refused_dispatches(), 1); + } + + #[test] + fn registration_with_an_exited_decrypt_worker_warns_counts_and_is_refused() { + let pool = DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + let key = key_on_worker(&pool, 1); + let (registered, logs) = + crate::testutil::capture_logs(|| pool.register_session(key, session_state())); + assert!(!registered, "a dead worker cannot own a session"); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" worker=1"), "{warnings:?}"); + assert_eq!(pool.liveness().refused_dispatches(), 1); + } +} diff --git a/src/node/encrypt_worker.rs b/src/node/encrypt_worker.rs index d8cd4da7..ba1881bc 100644 --- a/src/node/encrypt_worker.rs +++ b/src/node/encrypt_worker.rs @@ -50,11 +50,14 @@ // warnings rather than gate every function individually. #![cfg_attr(not(unix), allow(dead_code))] +#[cfg(test)] +use crate::node::worker_set::{TestWorker, test_spawner}; +use crate::node::worker_set::{WorkerLiveness, WorkerSet, worth_logging}; use crate::proto::fmp::wire::ESTABLISHED_HEADER_SIZE; use crate::proto::fsp::wire::FSP_HEADER_SIZE; use crate::transport::udp::io::AsyncUdpSocket; #[cfg(not(target_os = "macos"))] -use crossbeam_channel::{Receiver, SendError, Sender, TrySendError, bounded}; +use crossbeam_channel::{Receiver, Sender, TrySendError, bounded}; use ring::aead::{Aad, LessSafeKey, Nonce}; #[cfg(any(target_os = "macos", test))] use std::collections::VecDeque; @@ -406,6 +409,36 @@ type WorkerSender = MacWorkerSender; #[cfg(not(target_os = "macos"))] type WorkerSender = Sender; +#[cfg(target_os = "macos")] +type WorkerReceiver = MacWorkerReceiver; + +#[cfg(not(target_os = "macos"))] +type WorkerReceiver = Receiver; + +fn worker_channel() -> (WorkerSender, WorkerReceiver) { + #[cfg(target_os = "macos")] + { + mac_worker_channel(WORKER_CHANNEL_CAP) + } + #[cfg(not(target_os = "macos"))] + { + bounded::(WORKER_CHANNEL_CAP) + } +} + +/// Start the production worker loop on `rx` in a named OS thread. +fn spawn_worker(idx: usize, rx: WorkerReceiver) -> std::io::Result> { + let builder = std::thread::Builder::new().name(format!("fips-encrypt-{idx}")); + #[cfg(target_os = "macos")] + { + builder.spawn(move || run_worker_macos(idx, rx)) + } + #[cfg(not(target_os = "macos"))] + { + builder.spawn(move || run_worker(idx, rx)) + } +} + /// Handle to the encrypt worker pool. /// /// Workers are **dedicated `std::thread`s** with **`crossbeam_channel`** @@ -425,7 +458,7 @@ type WorkerSender = Sender; /// destinations hash to different workers. #[derive(Clone)] pub(crate) struct EncryptWorkerPool { - senders: Arc<[WorkerSender]>, + workers: Arc>, #[cfg(target_os = "macos")] macos_senders: Arc, #[cfg(target_os = "macos")] @@ -437,31 +470,21 @@ impl EncryptWorkerPool { /// dispatches jobs hash-by-destination to them. The workers exit /// when all senders for their channel are dropped (i.e. when the /// returned `EncryptWorkerPool` and all clones go away). + /// + /// A worker thread that cannot be started is logged and left dead; + /// the caller reads how many started from [`Self::liveness`]. pub fn spawn(n: usize) -> Self { - let n = n.max(1); - let mut senders = Vec::with_capacity(n); - for i in 0..n { - #[cfg(target_os = "macos")] - { - let (tx, rx) = mac_worker_channel(WORKER_CHANNEL_CAP); - std::thread::Builder::new() - .name(format!("fips-encrypt-{i}")) - .spawn(move || run_worker_macos(i, rx)) - .expect("failed to spawn fips-encrypt OS thread"); - senders.push(tx); - } - #[cfg(not(target_os = "macos"))] - { - let (tx, rx) = bounded::(WORKER_CHANNEL_CAP); - std::thread::Builder::new() - .name(format!("fips-encrypt-{i}")) - .spawn(move || run_worker(i, rx)) - .expect("failed to spawn fips-encrypt OS thread"); - senders.push(tx); - } - } + Self::start_with(n, spawn_worker) + } + + /// Build the pool, starting each worker with `spawn`. Production passes + /// [`spawn_worker`]; tests pass workers that fail to start or exit on cue. + fn start_with( + n: usize, + spawn: impl FnMut(usize, WorkerReceiver) -> std::io::Result>, + ) -> Self { Self { - senders: senders.into(), + workers: Arc::new(WorkerSet::start("encrypt", n, worker_channel, spawn)), #[cfg(target_os = "macos")] macos_senders: Arc::new(MacSequencedSendFlows::default()), #[cfg(target_os = "macos")] @@ -469,12 +492,19 @@ impl EncryptWorkerPool { } } + /// Whether each worker is still running, and how many dispatches a dead + /// one refused. + pub(crate) fn liveness(&self) -> &dyn WorkerLiveness { + &*self.workers + } + /// Dispatch a job to the worker that owns its destination flow. /// The hash is over `dest_addr` so every packet for one peer's /// kernel `SocketAddr` lands on the same worker and stays in /// order — required for TCP's fast-retransmit logic above to - /// behave on a single-flow run. Fire-and-forget — the worker - /// handles send errors itself via stats counters. + /// behave on a single-flow run. The worker handles send errors + /// itself via stats counters. A job whose worker has exited is + /// dropped, counted, and logged at WARN. /// /// Uses `try_send` for the common uncontended case, then blocks /// only when the bounded worker channel is full. These jobs carry @@ -483,12 +513,18 @@ impl EncryptWorkerPool { /// retransmits. Blocking here pushes back toward the TUN reader /// and lets the kernel/app TCP stack pace the flow instead. pub fn dispatch(&self, job: FmpSendJob) { - if self.senders.is_empty() { - debug!("EncryptWorkerPool has no workers; dropping job"); - return; - } let (idx, job) = self.prepare_dispatch(job); - self.dispatch_to_worker(idx, job); + if !self.dispatch_to_worker(idx, job) { + let n = self.workers.note_refused(); + if worth_logging(n) { + warn!( + pool = "encrypt", + worker = idx, + refused = n + 1, + "Encrypt worker has exited; dropping packet" + ); + } + } } #[cfg(target_os = "macos")] @@ -503,7 +539,7 @@ impl EncryptWorkerPool { }; let mut h = std::collections::hash_map::DefaultHasher::new(); key.hash(&mut h); - let idx = (h.finish() as usize) % self.senders.len(); + let idx = (h.finish() as usize) % self.workers.len(); return (idx, QueuedFmpSendJob::direct(job)); } @@ -518,68 +554,81 @@ impl EncryptWorkerPool { .next_worker .fetch_add(1, std::sync::atomic::Ordering::Relaxed) / macos_worker_stride(); - let idx = ticket % self.senders.len(); + let idx = ticket % self.workers.len(); (idx, QueuedFmpSendJob::macos_sequenced(job, flow)) } #[cfg(not(target_os = "macos"))] fn prepare_dispatch(&self, job: FmpSendJob) -> (usize, QueuedFmpSendJob) { - use std::hash::{Hash, Hasher}; - let mut h = std::collections::hash_map::DefaultHasher::new(); - job.dest_addr.hash(&mut h); - let idx = (h.finish() as usize) % self.senders.len(); + let idx = self.worker_index_for(job.dest_addr); (idx, QueuedFmpSendJob::direct(job)) } + /// The worker that owns `dest`'s flow. + #[cfg(not(target_os = "macos"))] + fn worker_index_for(&self, dest: SocketAddr) -> usize { + use std::hash::{Hash, Hasher}; + let mut h = std::collections::hash_map::DefaultHasher::new(); + dest.hash(&mut h); + (h.finish() as usize) % self.workers.len() + } + + /// Queue `job` on worker `idx`, returning `false` when that worker has + /// exited and the job was not queued. #[cfg(target_os = "macos")] - fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) { - match self.senders[idx].try_push(job) { - Ok(()) => {} + fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> bool { + let sender = self.workers.sender(idx); + match sender.try_push(job) { + Ok(()) => true, Err(MacWorkerTryPushError::Full(job)) => { static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed); - if n < 8 || n.is_multiple_of(10000) { + if worth_logging(n) { warn!( worker = idx, full_events = n + 1, "EncryptWorker channel full; applying outbound backpressure" ); } - if let Err(MacWorkerPushError) = self.senders[idx].push_blocking(*job) { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); - } - } - Err(MacWorkerTryPushError::Closed) => { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); + sender.push_blocking(*job).is_ok() } + Err(MacWorkerTryPushError::Closed) => false, } } + /// Queue `job` on worker `idx`, returning `false` when that worker has + /// exited and the job was not queued. #[cfg(not(target_os = "macos"))] - fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) { - match self.senders[idx].try_send(job) { - Ok(()) => {} + fn dispatch_to_worker(&self, idx: usize, job: QueuedFmpSendJob) -> bool { + let sender = self.workers.sender(idx); + match sender.try_send(job) { + Ok(()) => true, Err(TrySendError::Full(job)) => { static FULL_COUNT: portable_atomic::AtomicU64 = portable_atomic::AtomicU64::new(0); let n = FULL_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed); - if n < 8 || n.is_multiple_of(10000) { + if worth_logging(n) { warn!( worker = idx, full_events = n + 1, "EncryptWorker channel full; applying outbound backpressure" ); } - if let Err(SendError(_)) = self.senders[idx].send(job) { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); - } - } - Err(TrySendError::Disconnected(_)) => { - debug!(worker = idx, "EncryptWorker thread gone; dropping job"); + sender.send(job).is_ok() } + Err(TrySendError::Disconnected(_)) => false, } } } +#[cfg(test)] +impl EncryptWorkerPool { + /// A pool whose workers behave as `plan` says; `Run` workers are the + /// production loop. + pub(crate) fn for_test(plan: Vec) -> Self { + Self::start_with(plan.len(), test_spawner(plan, spawn_worker)) + } +} + #[cfg(target_os = "macos")] #[derive(Clone, Copy, Debug, Hash, PartialEq, Eq)] struct MacSendFlowKey { @@ -2892,3 +2941,176 @@ mod mac_ordered_tests { drop(gap); } } + +/// The pool's view of its workers: which are live, what a dead one does to a +/// dispatch, and that a live but full queue still blocks. Every wait is +/// bounded so a regression fails instead of hanging. +#[cfg(test)] +mod pool_tests { + use super::*; + use crate::node::worker_set::wait_for; + #[cfg(not(target_os = "macos"))] + use crate::transport::udp::io::UdpRawSocket; + #[cfg(not(target_os = "macos"))] + use ring::aead::UnboundKey; + #[cfg(not(target_os = "macos"))] + use std::net::UdpSocket; + use std::sync::mpsc; + #[cfg(not(target_os = "macos"))] + use std::time::Duration; + + #[cfg(not(target_os = "macos"))] + struct Rig { + _rt: tokio::runtime::Runtime, + socket: AsyncUdpSocket, + } + + #[cfg(not(target_os = "macos"))] + impl Rig { + fn new() -> Self { + let rt = tokio::runtime::Builder::new_current_thread() + .enable_io() + .build() + .expect("tokio rt"); + let enter = rt.enter(); + let socket = UdpRawSocket::open("127.0.0.1:0".parse().unwrap(), 1 << 20, 1 << 20) + .expect("open send socket") + .into_async() + .expect("into_async"); + drop(enter); + Self { _rt: rt, socket } + } + + fn job(&self, dest: SocketAddr, counter: u64) -> FmpSendJob { + let key = UnboundKey::new(&ring::aead::CHACHA20_POLY1305, &[3u8; 32]).expect("key"); + let mut wire_buf = + Vec::with_capacity(ESTABLISHED_HEADER_SIZE + 8 + crate::noise::TAG_SIZE); + wire_buf.extend_from_slice(&[0x5A; ESTABLISHED_HEADER_SIZE]); + wire_buf.extend_from_slice(&counter.to_le_bytes()); + FmpSendJob { + cipher: LessSafeKey::new(key), + counter, + wire_buf, + fsp_seal: None, + socket: self.socket.clone(), + dest_addr: dest, + #[cfg(any(target_os = "linux", target_os = "macos"))] + connected_socket: None, + drop_on_backpressure: false, + queued_at: None, + } + } + } + + /// A bound receiver whose address the pool dispatches to worker `idx`. + #[cfg(not(target_os = "macos"))] + fn receiver_on_worker(pool: &EncryptWorkerPool, idx: usize) -> UdpSocket { + loop { + let sock = UdpSocket::bind("127.0.0.1:0").expect("bind receiver"); + if pool.worker_index_for(sock.local_addr().unwrap()) == idx { + sock.set_read_timeout(Some(Duration::from_secs(5))) + .expect("read timeout"); + return sock; + } + } + } + + #[test] + fn an_encrypt_worker_that_panics_is_counted_dead() { + let (die_tx, die_rx) = mpsc::channel::<()>(); + let mut die_rx = Some(die_rx); + let pool = EncryptWorkerPool::start_with(2, |idx, rx| { + if idx != 1 { + return spawn_worker(idx, rx); + } + let die = die_rx.take().expect("worker 1 starts once"); + std::thread::Builder::new().spawn(move || { + let _rx = rx; + let _ = die.recv(); + panic!("simulated encrypt worker panic"); + }) + }); + assert_eq!(pool.liveness().live_workers(), 2); + die_tx.send(()).expect("worker 1 gone before its signal"); + assert!( + wait_for(|| pool.liveness().live_workers() == 1), + "a worker that panicked is still counted live" + ); + assert_eq!(pool.liveness().dead_workers(), vec![1]); + assert_eq!(pool.liveness().worker_count(), 2); + } + + #[cfg(not(target_os = "macos"))] + #[test] + fn a_worker_that_fails_to_spawn_leaves_the_pool_serving_the_rest() { + let rig = Rig::new(); + let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + assert_eq!(pool.liveness().worker_count(), 2); + assert_eq!(pool.liveness().live_workers(), 1); + + let recv = receiver_on_worker(&pool, 0); + pool.dispatch(rig.job(recv.local_addr().unwrap(), 1)); + let mut buf = [0u8; 128]; + recv.recv_from(&mut buf) + .expect("the live worker did not send the job dispatched to it"); + assert_eq!(pool.liveness().refused_dispatches(), 0); + } + + #[cfg(not(target_os = "macos"))] + #[test] + fn dispatch_to_an_exited_encrypt_worker_warns_and_counts() { + let rig = Rig::new(); + let pool = EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]); + let recv = receiver_on_worker(&pool, 1); + let ((), logs) = crate::testutil::capture_logs(|| { + pool.dispatch(rig.job(recv.local_addr().unwrap(), 1)); + }); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" worker=1"), "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"encrypt\""), "{warnings:?}"); + assert_eq!(pool.liveness().refused_dispatches(), 1); + } + + /// A live worker that has fallen behind must hold the rx loop back, not + /// have its packets dropped: these are tunnelled packets, and a drop here + /// reads as loss to TCP inside the tunnel. + #[cfg(not(target_os = "macos"))] + #[test] + fn a_full_live_encrypt_queue_still_blocks_dispatch() { + let rig = Rig::new(); + let (release_tx, release_rx) = mpsc::channel::<()>(); + let mut release_rx = Some(release_rx); + let pool = EncryptWorkerPool::start_with(1, |idx, rx| { + let release = release_rx.take().expect("one worker"); + std::thread::Builder::new().spawn(move || { + let _ = release.recv(); + run_worker(idx, rx); + }) + }); + let recv = UdpSocket::bind("127.0.0.1:0").expect("bind receiver"); + let dest = recv.local_addr().unwrap(); + for counter in 0..WORKER_CHANNEL_CAP as u64 { + pool.dispatch(rig.job(dest, counter)); + } + + let (done_tx, done_rx) = mpsc::channel::<()>(); + let blocked_pool = pool.clone(); + let last = rig.job(dest, WORKER_CHANNEL_CAP as u64); + std::thread::spawn(move || { + blocked_pool.dispatch(last); + let _ = done_tx.send(()); + }); + std::thread::sleep(Duration::from_millis(200)); + assert!( + matches!(done_rx.try_recv(), Err(mpsc::TryRecvError::Empty)), + "dispatch returned while the live worker's queue was full" + ); + + release_tx.send(()).expect("worker gone before release"); + done_rx + .recv_timeout(Duration::from_secs(5)) + .expect("dispatch still blocked after the worker drained"); + assert_eq!(pool.liveness().refused_dispatches(), 0); + } +} diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index 10b9b930..38e469e8 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -1,6 +1,7 @@ //! Node lifecycle management: start, stop, and peer connection initiation. pub(crate) mod supervisor; +mod workers; use super::{Node, NodeError, NodeState}; use supervisor::{Action, Child, Event, PeeringDesired, SupervisorFsm}; @@ -1609,44 +1610,33 @@ impl Node { } } Child::EncryptWorkers => { - // Hash-by-destination pins a TCP flow to one worker - // (preserves wire ordering); additional workers light up - // under multi-flow load. Infallible → always up. + // A worker that cannot be started degrades the node; it + // never stops start-up. #[cfg(unix)] - { - self.supervisor.encrypt_workers = Some( - super::encrypt_worker::EncryptWorkerPool::spawn(encrypt_worker_count), - ); - info!( - workers = encrypt_worker_count, - "Spawned FMP-encrypt worker pool" - ); + let start = self.start_encrypt_workers(encrypt_worker_count); + #[cfg(not(unix))] + let start = workers::PoolStart::ALL_LIVE; - // `FIPS_DECRYPT_WORKERS=0` disables the pool entirely - // and forces the in-line rx_loop decrypt path. When 0 - // no DecryptWorkers child is emitted, so this info! - // sits here — exactly where the decrypt spawn would be - // in today's sequence (after the encrypt spawn+info, - // before nostr). - if decrypt_worker_count == 0 { - info!("FIPS_DECRYPT_WORKERS=0 → in-line decrypt in rx_loop"); - } + // `FIPS_DECRYPT_WORKERS=0` disables the pool entirely + // and forces the in-line rx_loop decrypt path. When 0 + // no DecryptWorkers child is emitted, so this info! + // sits here — exactly where the decrypt spawn would be + // in today's sequence (after the encrypt spawn+info, + // before nostr). + #[cfg(unix)] + if decrypt_worker_count == 0 { + info!("FIPS_DECRYPT_WORKERS=0 → in-line decrypt in rx_loop"); } - Event::SubstrateUp { child } + start.event(child) } Child::DecryptWorkers => { - // Shard-owned decrypt pool. Infallible → always up. + // Shard-owned decrypt pool. A worker that cannot be + // started degrades the node; it never stops start-up. #[cfg(unix)] - { - self.supervisor.decrypt_workers = Some( - super::decrypt_worker::DecryptWorkerPool::spawn(decrypt_worker_count), - ); - info!( - workers = decrypt_worker_count, - "Spawned FMP-decrypt worker pool" - ); - } - Event::SubstrateUp { child } + let start = self.start_decrypt_workers(decrypt_worker_count); + #[cfg(not(unix))] + let start = workers::PoolStart::ALL_LIVE; + start.event(child) } Child::Nostr => { match NostrRendezvous::start( @@ -2388,6 +2378,18 @@ impl Node { } } + /// Tell the supervisor FSM that `child` exited on its own at runtime, and + /// publish the health it resolves. Shared by the child-exit channel's arm + /// in the rx loop and the tick's worker-liveness sweep, which steps the FSM + /// directly because it runs on the loop that drains that channel. + pub(in crate::node) fn step_child_exited(&mut self, child: Child) { + for action in self.supervisor.fsm.step(Event::ChildExited { child }) { + if let Action::PublishState(ns) = action { + self.supervisor.state = ns; + } + } + } + /// Feed every queued interface-presence edge to the supervisor FSM, /// returning the last [`NodeState`] it asked to publish (if any). /// diff --git a/src/node/lifecycle/supervisor.rs b/src/node/lifecycle/supervisor.rs index 301ad22c..edb4e658 100644 --- a/src/node/lifecycle/supervisor.rs +++ b/src/node/lifecycle/supervisor.rs @@ -63,10 +63,11 @@ //! - the degenerate no-children path now resolves to `Failed` (zero transports), //! **not** the old immediate-`Running`. //! -//! Runtime child-liveness monitoring (a `ChildExited` event re-routing health -//! when a task/thread dies at runtime) is **deferred**: start-completion health -//! resolution is start-framed, and liveness monitoring is a substantial unbuilt -//! mechanism. This commit is start-time health only. +//! Runtime child-liveness (a `ChildExited` event re-routing health when a +//! task/thread dies at runtime) came later. Its producers are the children that +//! report their own exit (TUN threads, the DNS task, the mDNS/Nostr monitor) +//! and, for the crypto worker pools, the rx loop tick's liveness sweep, which +//! reports a pool's child as exited when it loses a worker. //! //! ## Scope: interface presence, and `Degraded` as a level (this commit) //! @@ -298,7 +299,8 @@ pub(crate) enum Health { Full, /// ≥1 transport is up, but one or more configured optional children failed /// to start (a transport beyond the first, Nostr, mDNS, TUN, DNS, or a - /// worker-pool spawn). The node is operational (serving) but degraded. + /// worker pool that did not start every worker). The node is operational + /// (serving) but degraded. Degraded { /// The configured children that failed to start. reasons: HashSet, @@ -860,8 +862,8 @@ pub(crate) struct Supervisor { /// Off-task FMP-encrypt + UDP-send worker pool. Unix-only — /// the worker issues direct sendmmsg(2) / sendmsg+UDP_GSO calls - /// on raw fds via `AsRawFd`. None on Windows or when the worker - /// pool failed to spawn. + /// on raw fds via `AsRawFd`. None on Windows or when no worker + /// thread could be started. #[cfg(unix)] pub(crate) encrypt_workers: Option, @@ -869,9 +871,16 @@ pub(crate) struct Supervisor { /// `encrypt_workers`. Workers are shards: each owns its session /// state directly in a thread-local `HashMap` (no `RwLock`, /// no `Mutex` per packet). Hash-by-cache-key dispatch. + /// None on Windows, with `FIPS_DECRYPT_WORKERS=0`, or when no worker + /// thread could be started. #[cfg(unix)] pub(crate) decrypt_workers: Option, + /// Pools a test has built for start-up to install in place of spawning + /// its own. + #[cfg(all(test, unix))] + pub(in crate::node) staged_pools: StagedPools, + /// Transport-medium change detection: the receiver the rx loop drains and /// the detector task behind it. /// @@ -889,6 +898,15 @@ pub(crate) struct Supervisor { pub(in crate::node) fsm: SupervisorFsm, } +/// Crypto worker pools a test hands to start-up, each taken by the first +/// start of its pool's child. +#[cfg(all(test, unix))] +#[derive(Default)] +pub(in crate::node) struct StagedPools { + pub encrypt: Option, + pub decrypt: Option, +} + impl Supervisor { /// A fresh supervisor with all handles empty and the FSM in `Created`, /// matching the field initializers `Node::new` previously used. @@ -914,6 +932,8 @@ impl Supervisor { encrypt_workers: None, #[cfg(unix)] decrypt_workers: None, + #[cfg(all(test, unix))] + staged_pools: StagedPools::default(), netmon_rx: None, netmon_task: None, fsm: SupervisorFsm::new(), @@ -1796,6 +1816,39 @@ mod tests { assert!(s.absent().is_empty()); } + /// The crypto worker pools are a performance offload with a main-loop path + /// behind them, so losing one degrades the node and can never be what + /// fails it. + #[test] + fn a_worker_pool_exit_is_degraded_and_never_failed() { + let mut s = SupervisorFsm::running_with([ + Child::Transport(tid(1)), + Child::EncryptWorkers, + Child::DecryptWorkers, + ]); + assert_eq!( + s.step(Event::ChildExited { + child: Child::EncryptWorkers + }), + vec![Action::PublishState(NodeState::Degraded)] + ); + assert_eq!( + s.step(Event::ChildExited { + child: Child::DecryptWorkers + }), + vec![Action::PublishState(NodeState::Degraded)] + ); + assert!(matches!( + s.state(), + SupState::Running { + health: Health::Degraded { .. } + } + )); + let degraded = s.degraded_children(); + assert!(degraded.contains(&Child::EncryptWorkers)); + assert!(degraded.contains(&Child::DecryptWorkers)); + } + #[test] fn degraded_children_is_the_union_of_both_reason_sets() { let mut s = SupervisorFsm::new(); diff --git a/src/node/lifecycle/workers.rs b/src/node/lifecycle/workers.rs new file mode 100644 index 00000000..98f6fe8a --- /dev/null +++ b/src/node/lifecycle/workers.rs @@ -0,0 +1,240 @@ +//! Supervision of the crypto worker pools: what a pool that started short +//! means for the node, and the tick's sweep for workers that have since +//! exited. +//! +//! The pools are a performance offload, and Windows never starts them at +//! all, so losing workers is `Degraded` at most and never fatal. Inbound +//! packets for a session already held by a missing decrypt worker are +//! dropped until the session rekeys or the link is re-established; a session +//! that would be registered on the missing worker after the loss is decrypted +//! on the main loop instead. +//! The pools report facts (how many workers, how many live); the decisions +//! below turn them into supervisor events. + +use super::supervisor::{Child, Event}; +#[cfg(unix)] +use crate::node::Node; +#[cfg(unix)] +use crate::node::worker_set::WorkerLiveness; +#[cfg(unix)] +use tracing::{info, warn}; + +/// What a freshly started pool means for the node. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(in crate::node) struct PoolStart { + /// Keep the pool. A pool with no live worker is dropped, so every packet + /// takes the main-loop path directly instead of being refused first. + pub keep: bool, + /// Report the pool's child as up. Anything short of every worker live is + /// reported as failed to start, so start-completion health resolves to + /// `Degraded` on the first publish rather than `Full` corrected a tick + /// later. + pub up: bool, +} + +impl PoolStart { + /// Every worker started. Off Unix the pools are never started, and their + /// children report up as before. + #[cfg(not(unix))] + pub(in crate::node) const ALL_LIVE: Self = Self { + keep: true, + up: true, + }; + + /// The supervisor event that reports this outcome for `child`. + pub(in crate::node) fn event(self, child: Child) -> Event { + if self.up { + Event::SubstrateUp { child } + } else { + Event::SubstrateFailed { child } + } + } +} + +/// Classify a pool that started `live` of `configured` workers. +#[cfg(unix)] +pub(in crate::node) fn pool_start(live: usize, configured: usize) -> PoolStart { + PoolStart { + keep: live > 0, + up: live >= configured, + } +} + +/// Whether a pool now at `live` workers has lost any since it last reported +/// `reported`. Every loss is reported, the first one and each after it, so the +/// count of live workers stays visible as it falls. Nothing restarts a worker, +/// so the count never rises. +#[cfg(unix)] +pub(in crate::node) fn lost_workers(reported: usize, live: usize) -> bool { + live < reported +} + +#[cfg(unix)] +impl Node { + /// Start the encrypt pool and install it, or the main-loop path when no + /// worker started. Under test, a pool staged in + /// [`StagedPools`](super::supervisor::StagedPools) is installed instead of + /// spawning one, so a pool that starts short can be driven through + /// start-up. + pub(in crate::node) fn start_encrypt_workers(&mut self, n: usize) -> PoolStart { + #[cfg(test)] + if let Some(pool) = self.supervisor.staged_pools.encrypt.take() { + return self.install_encrypt_workers(pool); + } + self.install_encrypt_workers(crate::node::encrypt_worker::EncryptWorkerPool::spawn(n)) + } + + /// Install a started encrypt pool, or the main-loop path when none of + /// its workers started. + fn install_encrypt_workers( + &mut self, + pool: crate::node::encrypt_worker::EncryptWorkerPool, + ) -> PoolStart { + let start = report_pool_start("encrypt", pool.liveness()); + if start.up { + info!( + workers = pool.liveness().worker_count(), + "Spawned FMP-encrypt worker pool" + ); + } + self.supervisor.encrypt_workers = start.keep.then_some(pool); + start + } + + /// Start the decrypt pool and install it, or the main-loop path when no + /// worker started. Under test, a pool staged in + /// [`StagedPools`](super::supervisor::StagedPools) is installed instead of + /// spawning one, so a pool that starts short can be driven through + /// start-up. + pub(in crate::node) fn start_decrypt_workers(&mut self, n: usize) -> PoolStart { + #[cfg(test)] + if let Some(pool) = self.supervisor.staged_pools.decrypt.take() { + return self.install_decrypt_workers(pool); + } + self.install_decrypt_workers(crate::node::decrypt_worker::DecryptWorkerPool::spawn(n)) + } + + /// Install a started decrypt pool, or the main-loop path when none of + /// its workers started. + fn install_decrypt_workers( + &mut self, + pool: crate::node::decrypt_worker::DecryptWorkerPool, + ) -> PoolStart { + let start = report_pool_start("decrypt", pool.liveness()); + if start.up { + info!( + workers = pool.liveness().worker_count(), + "Spawned FMP-decrypt worker pool" + ); + } + self.supervisor.decrypt_workers = start.keep.then_some(pool); + start + } + + /// The tick's worker-liveness sweep. For each pool that has lost a worker + /// since the last sweep, log the loss with the live count and report the + /// pool's child as exited. The FSM republishes `Degraded` the first time; + /// a pool already out of its up-set republishes nothing, but each further + /// loss is still logged with the new count. + /// + /// Steps the FSM directly rather than through the child-exit channel: that + /// channel is drained by the same loop that runs this sweep, so a send into + /// a full channel from here would never complete. + pub(in crate::node) fn poll_worker_liveness(&mut self) { + let mut exited = Vec::with_capacity(2); + if let Some(pool) = &self.supervisor.encrypt_workers + && report_worker_loss("encrypt", pool.liveness()) + { + exited.push(Child::EncryptWorkers); + } + if let Some(pool) = &self.supervisor.decrypt_workers + && report_worker_loss("decrypt", pool.liveness()) + { + exited.push(Child::DecryptWorkers); + } + for child in exited { + self.step_child_exited(child); + } + } +} + +/// Classify how a pool started, logging a pool that started short. +#[cfg(unix)] +fn report_pool_start(pool: &'static str, workers: &dyn WorkerLiveness) -> PoolStart { + let live = workers.live_workers(); + let configured = workers.worker_count(); + let start = pool_start(live, configured); + if !start.up { + warn!( + pool, + live, configured, "Crypto worker pool started with fewer workers than configured" + ); + } + start +} + +/// Record a pool's live count, logging and returning `true` when it has lost a +/// worker since the last call. +#[cfg(unix)] +fn report_worker_loss(pool: &'static str, workers: &dyn WorkerLiveness) -> bool { + let live = workers.live_workers(); + let reported = workers.swap_reported_live(live); + if !lost_workers(reported, live) { + return false; + } + warn!( + pool, + live, + configured = workers.worker_count(), + dead = ?workers.dead_workers(), + "Crypto worker thread exited" + ); + true +} + +#[cfg(all(test, unix))] +mod tests { + use super::*; + + #[test] + fn worker_spawn_outcome_maps_k_of_n() { + assert_eq!( + pool_start(4, 4), + PoolStart { + keep: true, + up: true + } + ); + for live in 1..4 { + assert_eq!( + pool_start(live, 4), + PoolStart { + keep: true, + up: false + }, + "{live} of 4 live" + ); + } + assert_eq!( + pool_start(0, 4), + PoolStart { + keep: false, + up: false + } + ); + let child = Child::EncryptWorkers; + assert_eq!(pool_start(4, 4).event(child), Event::SubstrateUp { child }); + assert_eq!( + pool_start(3, 4).event(child), + Event::SubstrateFailed { child } + ); + } + + #[test] + fn every_loss_is_reported_and_no_loss_is_not() { + assert!(lost_workers(4, 3)); + assert!(lost_workers(3, 0)); + assert!(!lost_workers(4, 4)); + assert!(!lost_workers(0, 0)); + } +} diff --git a/src/node/mod.rs b/src/node/mod.rs index f0dfc200..7eda46b7 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -33,6 +33,8 @@ pub(crate) mod stats_history; #[cfg(test)] mod tests; mod tree; +#[cfg(unix)] +pub(crate) mod worker_set; use self::peer_error_budget::PeerErrorBudget; use self::rate_limit::{HandshakeRateLimiter, LookupSignRateLimiter, SessionSetupRateLimiter}; diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index ba623b0e..7a233498 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -4816,3 +4816,244 @@ fn mesh_filter_resolves_the_live_tun_device_rather_than_the_configured_name() { node.tun_name = Some(loopback.to_string()); assert_eq!(node.mesh_ifindex(), Some(expected)); } + +/// The tick's crypto worker liveness sweep: a lost worker degrades the node and +/// is logged with the live count; a healthy pool changes nothing. +#[cfg(unix)] +mod worker_liveness { + use super::*; + use crate::node::decrypt_worker::DecryptWorkerPool; + use crate::node::encrypt_worker::EncryptWorkerPool; + use crate::node::lifecycle::supervisor::{Child, SupervisorFsm}; + use crate::node::worker_set::{TestWorker, WorkerLiveness, wait_for}; + use std::collections::HashSet; + use std::sync::mpsc; + + /// A node seeded straight into a full `Running` with both pools up. + fn running_node() -> Node { + let mut node = make_node(); + node.supervisor.fsm = SupervisorFsm::running_with([ + Child::Transport(TransportId::new(1)), + Child::EncryptWorkers, + Child::DecryptWorkers, + ]); + node.supervisor.state = NodeState::Running; + node + } + + /// Two workers that each exit when their sender is used. + fn two_workers() -> (Vec, [mpsc::Sender<()>; 2]) { + let (kill0, exit0) = mpsc::channel(); + let (kill1, exit1) = mpsc::channel(); + ( + vec![TestWorker::ExitOn(exit0), TestWorker::ExitOn(exit1)], + [kill0, kill1], + ) + } + + /// An encrypt pool of two workers that each exit when their sender is + /// used. + fn pool_of_two() -> (EncryptWorkerPool, [mpsc::Sender<()>; 2]) { + let (plan, kills) = two_workers(); + (EncryptWorkerPool::for_test(plan), kills) + } + + /// Signal one worker of `workers` to exit and wait until `live_after` + /// remain. + fn kill_worker(workers: &dyn WorkerLiveness, kill: &mpsc::Sender<()>, live_after: usize) { + kill.send(()).expect("worker gone before its signal"); + assert!( + wait_for(|| workers.live_workers() == live_after), + "worker never exited" + ); + } + + fn kill(node: &Node, kill: &mpsc::Sender<()>, live_after: usize) { + let pool = node.supervisor.encrypt_workers.as_ref().unwrap(); + kill_worker(pool.liveness(), kill, live_after); + } + + #[tokio::test] + async fn a_dead_worker_degrades_the_node_and_names_the_live_count() { + let mut node = running_node(); + let (pool, kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(pool); + kill(&node, &kills[1], 1); + + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert_eq!(node.state(), NodeState::Degraded); + assert!( + node.supervisor + .fsm + .degraded_children() + .contains(&Child::EncryptWorkers) + ); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"encrypt\""), "{warnings:?}"); + assert!(warnings[0].contains(" live=1"), "{warnings:?}"); + assert!(warnings[0].contains(" configured=2"), "{warnings:?}"); + + // Nothing new: no second report. + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert!(logs.warnings().is_empty(), "{:?}", logs.warnings()); + assert_eq!(node.state(), NodeState::Degraded); + + // The last worker goes: reported with the new count, and the node is + // still degraded, not failed. + kill(&node, &kills[0], 0); + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" live=0"), "{warnings:?}"); + assert_eq!(node.state(), NodeState::Degraded); + } + + #[tokio::test] + async fn a_dead_decrypt_worker_degrades_the_node_and_names_its_pool() { + let mut node = running_node(); + let (encrypt, _encrypt_kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(encrypt); + let (plan, kills) = two_workers(); + node.supervisor.decrypt_workers = Some(DecryptWorkerPool::for_test(plan)); + let pool = node.supervisor.decrypt_workers.as_ref().unwrap(); + kill_worker(pool.liveness(), &kills[0], 1); + + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert_eq!(node.state(), NodeState::Degraded); + assert_eq!( + node.supervisor.fsm.degraded_children(), + HashSet::from([Child::DecryptWorkers]) + ); + let warnings = logs.warnings(); + assert_eq!(warnings.len(), 1, "{warnings:?}"); + assert!(warnings[0].contains(" pool=\"decrypt\""), "{warnings:?}"); + assert!(warnings[0].contains(" live=1"), "{warnings:?}"); + assert!(warnings[0].contains(" configured=2"), "{warnings:?}"); + } + + #[tokio::test] + async fn healthy_pools_leave_the_node_running() { + let mut node = running_node(); + let (pool, _kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(pool); + let (plan, _decrypt_kills) = two_workers(); + node.supervisor.decrypt_workers = Some(DecryptWorkerPool::for_test(plan)); + + let ((), logs) = crate::testutil::capture_logs(|| node.poll_worker_liveness()); + assert!(logs.warnings().is_empty(), "{:?}", logs.warnings()); + assert_eq!(node.state(), NodeState::Running); + } + + /// What start-up made of the pools a test staged for it. + struct StartOutcome { + state: NodeState, + degraded: HashSet, + encrypt_installed: bool, + decrypt_installed: bool, + } + + /// Start a node that installs `encrypt` and `decrypt` in place of the + /// pools it would spawn, then stop it. Fails if start-up did not take both + /// staged pools, since the outcome would then say nothing about them. + async fn start_with_staged( + encrypt: EncryptWorkerPool, + decrypt: DecryptWorkerPool, + ) -> StartOutcome { + let mut node = make_healthy_node(); + node.supervisor.staged_pools.encrypt = Some(encrypt); + node.supervisor.staged_pools.decrypt = Some(decrypt); + node.start().await.unwrap(); + let staged = &node.supervisor.staged_pools; + let taken = staged.encrypt.is_none() && staged.decrypt.is_none(); + let outcome = StartOutcome { + state: node.state(), + degraded: node.supervisor.fsm.degraded_children(), + encrypt_installed: node.supervisor.encrypt_workers.is_some(), + decrypt_installed: node.supervisor.decrypt_workers.is_some(), + }; + node.stop().await.unwrap(); + assert!(taken, "start-up did not install both staged pools"); + outcome + } + + #[tokio::test] + async fn a_worker_that_fails_to_start_leaves_the_node_degraded_with_its_pool_installed() { + let outcome = start_with_staged( + EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::FailSpawn]), + DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + ) + .await; + assert_eq!(outcome.state, NodeState::Degraded); + assert_eq!(outcome.degraded, HashSet::from([Child::EncryptWorkers])); + assert!( + outcome.encrypt_installed, + "a pool with a live worker is kept" + ); + assert!(outcome.decrypt_installed); + } + + #[tokio::test] + async fn a_pool_with_no_worker_started_leaves_the_node_degraded_and_is_not_installed() { + let outcome = start_with_staged( + EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + DecryptWorkerPool::for_test(vec![TestWorker::FailSpawn, TestWorker::FailSpawn]), + ) + .await; + assert_eq!(outcome.state, NodeState::Degraded); + assert_eq!(outcome.degraded, HashSet::from([Child::DecryptWorkers])); + assert!(outcome.encrypt_installed); + assert!( + !outcome.decrypt_installed, + "a pool with no live worker is dropped for the main-loop path" + ); + } + + #[tokio::test] + async fn staged_pools_with_every_worker_started_leave_the_node_running() { + let outcome = start_with_staged( + EncryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + DecryptWorkerPool::for_test(vec![TestWorker::Run, TestWorker::Run]), + ) + .await; + assert_eq!(outcome.state, NodeState::Running); + assert!(outcome.degraded.is_empty(), "{:?}", outcome.degraded); + assert!(outcome.encrypt_installed && outcome.decrypt_installed); + } + + /// Start a node, swap in a pool of two, optionally kill one worker, and + /// drive the real rx loop past one tick. Returns the state it published. + async fn drive_with_pool(kill_one: bool) -> NodeState { + let mut node = make_healthy_node(); + node.start().await.unwrap(); + assert_eq!(node.state(), NodeState::Running); + let (pool, kills) = pool_of_two(); + node.supervisor.encrypt_workers = Some(pool); + if kill_one { + kill(&node, &kills[1], 1); + } + + let drive = tokio::time::timeout( + Duration::from_millis(1500), + node.run_rx_loop_with_shutdown(std::future::pending()), + ) + .await; + assert!( + drive.is_err(), + "the rx loop must still be running: {drive:?}" + ); + let state = node.state(); + node.stop().await.unwrap(); + state + } + + #[tokio::test] + async fn a_dead_worker_degrades_the_node_through_the_rx_loop() { + assert_eq!(drive_with_pool(true).await, NodeState::Degraded); + } + + #[tokio::test] + async fn a_healthy_pool_leaves_the_node_running_through_the_rx_loop() { + assert_eq!(drive_with_pool(false).await, NodeState::Running); + } +} diff --git a/src/node/worker_set.rs b/src/node/worker_set.rs new file mode 100644 index 00000000..3c7f17cb --- /dev/null +++ b/src/node/worker_set.rs @@ -0,0 +1,245 @@ +//! The worker threads behind one crypto worker pool, and what is known about +//! whether each is still running. +//! +//! Both pools (`encrypt_worker`, `decrypt_worker`) are a set of OS threads, +//! each reached through its own bounded channel. A worker that exits, a panic +//! unwinding included, drops its receiver, so its channel closes and every +//! later dispatch to it is refused. Nothing restarts it. This module keeps the +//! thread handles so that loss can be seen, and counts the dispatches it +//! refused. +//! +//! It reports facts only. Whether a loss degrades the node is decided by the +//! driver in `lifecycle`. + +use portable_atomic::AtomicU64; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::thread::JoinHandle; +use tracing::warn; + +/// One worker: the sending end of its channel and its thread, `None` when the +/// thread could not be started. +struct Worker { + sender: S, + thread: Option>, +} + +/// The workers of one pool. Held behind an `Arc` by the pool so that cloning +/// the pool per packet stays a reference-count bump. +pub(crate) struct WorkerSet { + workers: Box<[Worker]>, + /// Dispatches refused because the target worker had exited. + refused: AtomicU64, + /// The live-worker count the liveness sweep last reported. Starts at the + /// number of threads that started, so a worker that never started is + /// reported once, at spawn, and not again by the sweep. It lives with the + /// workers so that a pool replaced by another starts from its own count. + reported_live: AtomicUsize, +} + +impl WorkerSet { + /// Start `n` workers, at least one. + /// + /// `channel` builds each worker's sender and receiver. `spawn` starts the + /// thread that owns the receiver. A worker whose thread cannot be started + /// is logged and left dead rather than aborting the caller: the receiver + /// was moved into the failed spawn and is dropped with it, so that + /// worker's channel is closed and a dispatch to it is refused, not queued. + pub(crate) fn start( + pool: &'static str, + n: usize, + mut channel: impl FnMut() -> (S, R), + mut spawn: impl FnMut(usize, R) -> std::io::Result>, + ) -> Self { + let n = n.max(1); + let mut workers = Vec::with_capacity(n); + for idx in 0..n { + let (sender, receiver) = channel(); + let thread = match spawn(idx, receiver) { + Ok(handle) => Some(handle), + Err(error) => { + warn!(pool, worker = idx, %error, "Failed to start a crypto worker thread"); + None + } + }; + workers.push(Worker { sender, thread }); + } + let started = workers.iter().filter(|w| w.thread.is_some()).count(); + Self { + workers: workers.into(), + refused: AtomicU64::new(0), + reported_live: AtomicUsize::new(started), + } + } + + /// Number of workers, including dead ones. Never zero. + pub(crate) fn len(&self) -> usize { + self.workers.len() + } + + /// The sending end of worker `idx`'s channel. + pub(crate) fn sender(&self, idx: usize) -> &S { + &self.workers[idx].sender + } + + /// Count one dispatch refused by a dead worker, returning the count before + /// this one. + pub(crate) fn note_refused(&self) -> u64 { + self.refused.fetch_add(1, Ordering::Relaxed) + } +} + +/// What the liveness sweep and the tests read from a pool, without naming its +/// channel type. +pub(crate) trait WorkerLiveness { + /// Number of workers the pool was built with. + fn worker_count(&self) -> usize; + /// Workers whose thread started and has not finished. A worker thread + /// finishes only by unwinding or when every sender to it is dropped, and + /// the node holds the pool while it runs, so a finished thread is a dead + /// worker. + fn live_workers(&self) -> usize; + /// Indices of the workers that are not live. + fn dead_workers(&self) -> Vec; + /// Dispatches refused because their worker had exited. + #[cfg(test)] + fn refused_dispatches(&self) -> u64; + /// Record `live` as the count last reported, returning the previous one. + fn swap_reported_live(&self, live: usize) -> usize; +} + +impl WorkerLiveness for WorkerSet { + fn worker_count(&self) -> usize { + self.workers.len() + } + + fn live_workers(&self) -> usize { + self.workers.iter().filter(|w| is_live(w)).count() + } + + fn dead_workers(&self) -> Vec { + self.workers + .iter() + .enumerate() + .filter(|(_, w)| !is_live(w)) + .map(|(idx, _)| idx) + .collect() + } + + #[cfg(test)] + fn refused_dispatches(&self) -> u64 { + self.refused.load(Ordering::Relaxed) + } + + fn swap_reported_live(&self, live: usize) -> usize { + self.reported_live.swap(live, Ordering::Relaxed) + } +} + +fn is_live(worker: &Worker) -> bool { + worker.thread.as_ref().is_some_and(|t| !t.is_finished()) +} + +/// Whether the `n`th event (counting from zero) of a repeating condition is +/// logged: the first eight, then one in ten thousand. +pub(crate) fn worth_logging(n: u64) -> bool { + n < 8 || n.is_multiple_of(10_000) +} + +/// How a test wants one worker of a pool to behave. +#[cfg(test)] +pub(crate) enum TestWorker { + /// The production worker loop. + Run, + /// The thread fails to start. + FailSpawn, + /// The thread holds its receiver until signalled, then exits. + ExitOn(std::sync::mpsc::Receiver<()>), +} + +/// A spawner that starts each worker as `plan` says, using `run` for +/// [`TestWorker::Run`]. +#[cfg(test)] +pub(crate) fn test_spawner( + plan: Vec, + run: impl Fn(usize, R) -> std::io::Result>, +) -> impl FnMut(usize, R) -> std::io::Result> { + let mut plan: Vec> = plan.into_iter().map(Some).collect(); + move |idx, rx| match plan[idx].take().expect("each worker is started once") { + TestWorker::Run => run(idx, rx), + TestWorker::FailSpawn => { + drop(rx); + Err(std::io::Error::other("worker start refused by the test")) + } + TestWorker::ExitOn(signal) => std::thread::Builder::new().spawn(move || { + let _ = signal.recv(); + drop(rx); + }), + } +} + +/// Wait up to five seconds for `cond`, returning whether it came true. +#[cfg(test)] +pub(crate) fn wait_for(cond: impl Fn() -> bool) -> bool { + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + while !cond() { + if std::time::Instant::now() >= deadline { + return false; + } + std::thread::sleep(std::time::Duration::from_millis(5)); + } + true +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::mpsc; + + #[test] + fn a_worker_that_exits_is_no_longer_live_and_is_named_dead() { + let (stop_tx, stop_rx) = mpsc::channel::<()>(); + let mut stop_rx = Some(stop_rx); + let set = WorkerSet::start( + "test", + 2, + mpsc::channel::, + |idx, rx: mpsc::Receiver| { + let stop = if idx == 1 { stop_rx.take() } else { None }; + std::thread::Builder::new().spawn(move || match stop { + Some(stop) => { + let _ = stop.recv(); + drop(rx); + } + None => while rx.recv().is_ok() {}, + }) + }, + ); + assert_eq!(set.live_workers(), 2); + stop_tx.send(()).unwrap(); + assert!( + wait_for(|| set.live_workers() == 1), + "worker 1 never exited" + ); + assert_eq!(set.dead_workers(), vec![1]); + assert_eq!(set.worker_count(), 2); + } + + #[test] + fn the_reported_baseline_starts_at_the_workers_that_started() { + let set = WorkerSet::start("test", 3, mpsc::channel::, |idx, rx| { + if idx == 2 { + drop(rx); + return Err(std::io::Error::other("refused by the test")); + } + std::thread::Builder::new().spawn(move || while rx.recv().is_ok() {}) + }); + assert_eq!(set.swap_reported_live(2), 2); + assert_eq!(set.dead_workers(), vec![2]); + } + + #[test] + fn worth_logging_keeps_the_first_eight_then_one_in_ten_thousand() { + let logged: Vec = (0..20_001).filter(|n| worth_logging(*n)).collect(); + assert_eq!(logged, vec![0, 1, 2, 3, 4, 5, 6, 7, 10_000, 20_000]); + } +}