diff --git a/CHANGELOG.md b/CHANGELOG.md index 9ac36e73..3edbabf9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -445,6 +445,21 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 rarely dialed name in a very large peer list may re-resolve more often; the cost of a wrong eviction is one DNS lookup, not a failed dial. +- macOS: stopping an Ethernet transport under load no longer hangs the + process. The BPF reader thread handed each frame to the async consumer with + `blocking_send`, which parks with no way to be woken. Stopping the transport + aborts the consumer first, so nothing drains the 1024-frame channel, and the + socket's `Drop` then joined a thread that could never return: on a busy + interface the daemon had to be killed. The socket now drops the receiver + before joining, which releases a parked send at once, and the reader thread + sends through a helper that watches the same shutdown pipe its `select()` + already honours, so a send waiting for room cannot outlive a shutdown + request. The helper yields before it sleeps, so the saturated-path handoff + rate is unchanged. **Not covered by CI**: the reader thread is macOS-only + and Linux CI compiles none of it. What the tests prove is that the helper + the thread now waits in is cancellable; that a real BPF thread exits under + load still needs a manual check on a Mac. + - A failed private-key write no longer leaves a node silently running an ephemeral identity. Six write results in the identity path were discarded, and the sharpest was in `persistent` mode: a failed write to `fips.key` fell diff --git a/src/transport/ethernet/socket.rs b/src/transport/ethernet/socket.rs index 622d7f6c..a58e4f85 100644 --- a/src/transport/ethernet/socket.rs +++ b/src/transport/ethernet/socket.rs @@ -21,6 +21,90 @@ mod platform; #[cfg(unix)] pub use platform::PacketSocket; +/// Outcome of `send_frame`. +#[cfg(unix)] +#[cfg_attr(not(target_os = "macos"), allow(dead_code))] +pub(crate) enum SendOutcome { + Sent, + Stop, +} + +/// Retry iterations spent yielding before the send loop starts sleeping. +/// +/// A transiently full channel drains in microseconds, so yielding keeps the +/// saturated-path handoff rate uncapped, which is the whole reason this +/// module has a dedicated reader thread. Raising it burns more CPU against a +/// genuinely stuck consumer; lowering it puts a sleep in the common case. +#[cfg(unix)] +#[cfg_attr(not(target_os = "macos"), allow(dead_code))] +const SEND_YIELD_SPINS: u32 = 64; + +/// Longest the send loop sleeps between attempts on a full channel. +/// +/// This bounds only how quickly a parked send notices a shutdown request that +/// closing the receiver has not already covered. Raising it delays that +/// notice; lowering it costs more wakeups under sustained backpressure. +#[cfg(unix)] +#[cfg_attr(not(target_os = "macos"), allow(dead_code))] +const SEND_RETRY_MAX: std::time::Duration = std::time::Duration::from_millis(1); + +/// Send one item, waiting out a full channel but waking on `shutdown_fd`. +/// +/// Returns `Stop` when the receiver is gone or shutdown has been requested, +/// which is the reader thread's cue to exit. Unlike `blocking_send` this +/// cannot park past a shutdown request, so the `join()` in `Drop` always +/// returns. The caller must keep the socket owning `shutdown_fd` alive across +/// the call; `poll` on a closed fd reports `POLLNVAL` rather than `POLLIN`, so +/// even a lifetime mistake degrades to waiting rather than to a false stop. +/// +/// Compiled on every unix so Linux CI exercises the tests below; only the +/// macOS reader thread calls it. +#[cfg(unix)] +#[cfg_attr(not(target_os = "macos"), allow(dead_code))] +pub(crate) fn send_frame( + tx: &tokio::sync::mpsc::Sender, + item: T, + shutdown_fd: std::os::unix::io::RawFd, +) -> SendOutcome { + use tokio::sync::mpsc::error::TrySendError; + + let mut item = item; + let mut spins = 0u32; + let mut backoff = std::time::Duration::from_micros(50); + loop { + match tx.try_send(item) { + Ok(()) => return SendOutcome::Sent, + Err(TrySendError::Closed(_)) => return SendOutcome::Stop, + Err(TrySendError::Full(returned)) => { + if fd_is_readable(shutdown_fd) { + return SendOutcome::Stop; + } + item = returned; + if spins < SEND_YIELD_SPINS { + spins += 1; + std::thread::yield_now(); + } else { + std::thread::sleep(backoff); + backoff = (backoff * 2).min(SEND_RETRY_MAX); + } + } + } + } +} + +/// True if `fd` has data ready, tested without blocking. +#[cfg(unix)] +#[cfg_attr(not(target_os = "macos"), allow(dead_code))] +pub(crate) fn fd_is_readable(fd: std::os::unix::io::RawFd) -> bool { + let mut pfd = libc::pollfd { + fd, + events: libc::POLLIN, + revents: 0, + }; + let ret = unsafe { libc::poll(&mut pfd, 1, 0) }; + ret > 0 && (pfd.revents & libc::POLLIN) != 0 +} + // ============================================================================= // Linux: AsyncFd-based async wrapper // ============================================================================= @@ -111,7 +195,9 @@ mod async_impl { pub struct AsyncPacketSocket { inner: Arc, - rx: tokio::sync::Mutex>, + /// `None` once shutdown has taken the receiver, which is what makes + /// a reader thread parked on a full channel return at once. + rx: tokio::sync::Mutex>>, reader_thread: Option>, } @@ -146,7 +232,13 @@ mod async_impl { match result { Ok((n, mac)) => { let data = read_buf[..n].to_vec(); - if tx.blocking_send((data, mac)).is_err() { + // Not blocking_send: a send parked on a + // full channel must still notice shutdown, + // or Drop's join() never returns. + if matches!( + super::send_frame(&tx, (data, mac), shutdown_fd), + super::SendOutcome::Stop + ) { return; } } @@ -207,7 +299,7 @@ mod async_impl { Ok(Self { inner, - rx: tokio::sync::Mutex::new(rx), + rx: tokio::sync::Mutex::new(Some(rx)), reader_thread: Some(reader_thread), }) } @@ -230,7 +322,10 @@ mod async_impl { } pub async fn recv_from(&self, buf: &mut [u8]) -> Result<(usize, [u8; 6]), TransportError> { - let mut rx = self.rx.lock().await; + let mut guard = self.rx.lock().await; + let Some(rx) = guard.as_mut() else { + return Err(TransportError::RecvFailed("reader thread stopped".into())); + }; match rx.recv().await { Some((data, mac)) => { let n = data.len().min(buf.len()); @@ -247,15 +342,25 @@ mod async_impl { /// Signal the reader thread to stop. /// - /// Sets the shutdown flag; the reader thread checks it after - /// each BPF read timeout (~250ms) and exits. + /// Drops the receiver where it can, which makes a send parked on a + /// full channel fail immediately, then writes the shutdown pipe that + /// the thread's `select()` and `send_frame` both watch. The receiver + /// is unavailable while a `recv_from` holds the lock; `Drop` takes it + /// unconditionally, so the pipe is what covers that window. pub fn shutdown(&self) { + if let Ok(mut guard) = self.rx.try_lock() { + guard.take(); + } self.inner.request_shutdown(); } } impl Drop for AsyncPacketSocket { fn drop(&mut self) { + // Drop the receiver before joining: a send parked on a full + // channel then returns at once, with no polling and no latency + // added to the steady-state path. + self.rx.get_mut().take(); self.inner.request_shutdown(); if let Some(handle) = self.reader_thread.take() { let _ = handle.join(); @@ -284,3 +389,133 @@ pub struct PacketSocket; #[cfg(windows)] pub struct AsyncPacketSocket; + +// ============================================================================= +// Tests +// ============================================================================= + +#[cfg(all(test, unix))] +mod tests { + use super::{SendOutcome, fd_is_readable, send_frame}; + use std::sync::mpsc; + use std::time::Duration; + + /// A pipe, as the shutdown signal, returned as (read fd, write fd). + /// + /// Leaked deliberately: these live for the length of one test and closing + /// them mid-poll is exactly the confusion the test is meant to avoid. + fn shutdown_pipe() -> (std::os::unix::io::RawFd, std::os::unix::io::RawFd) { + let mut fds = [0i32; 2]; + let ret = unsafe { libc::pipe(fds.as_mut_ptr()) }; + assert_eq!(ret, 0, "pipe() failed"); + (fds[0], fds[1]) + } + + fn signal(write_fd: std::os::unix::io::RawFd) { + let byte = [1u8]; + let ret = unsafe { libc::write(write_fd, byte.as_ptr() as *const libc::c_void, 1) }; + assert_eq!(ret, 1, "write() to shutdown pipe failed"); + } + + #[test] + fn fd_is_readable_is_false_for_an_unwritten_pipe_and_true_after_a_write() { + let (read_fd, write_fd) = shutdown_pipe(); + assert!(!fd_is_readable(read_fd)); + signal(write_fd); + assert!(fd_is_readable(read_fd)); + } + + #[test] + fn send_frame_delivers_when_the_channel_has_room() { + let (tx, mut rx) = tokio::sync::mpsc::channel::>(1); + let (read_fd, _write_fd) = shutdown_pipe(); + + assert!(matches!( + send_frame(&tx, vec![1u8, 2, 3], read_fd), + SendOutcome::Sent + )); + assert_eq!(rx.try_recv().unwrap(), vec![1u8, 2, 3]); + } + + #[test] + fn send_frame_returns_stop_when_the_receiver_is_gone() { + let (tx, rx) = tokio::sync::mpsc::channel::>(1); + let (read_fd, _write_fd) = shutdown_pipe(); + drop(rx); + + assert!(matches!( + send_frame(&tx, vec![0u8], read_fd), + SendOutcome::Stop + )); + } + + #[test] + fn send_frame_returns_stop_when_the_receiver_is_dropped_while_the_channel_is_full() { + // The mechanism `Drop` relies on: closing the channel releases a + // sender that is waiting for room. + let (tx, rx) = tokio::sync::mpsc::channel::>(1); + let (read_fd, _write_fd) = shutdown_pipe(); + tx.try_send(vec![0u8]).unwrap(); + + let (done_tx, done_rx) = mpsc::channel(); + let sender = std::thread::spawn(move || { + let outcome = send_frame(&tx, vec![1u8], read_fd); + done_tx.send(matches!(outcome, SendOutcome::Stop)).unwrap(); + }); + // The send is parked on a full channel; only the drop frees it. + assert!(done_rx.recv_timeout(Duration::from_millis(50)).is_err()); + drop(rx); + + let stopped = done_rx + .recv_timeout(Duration::from_secs(5)) + .expect("send_frame did not return after the receiver was dropped"); + sender.join().unwrap(); + assert!(stopped); + } + + #[test] + fn send_frame_returns_stop_when_shutdown_is_requested_and_the_channel_is_full() { + // The defect: `blocking_send` on a full channel nobody is draining + // parks forever, so the reader thread never sees shutdown and the + // `join()` in `Drop` never returns. See the ignored test below for + // the same fixture against `blocking_send`. + let (tx, _rx) = tokio::sync::mpsc::channel::>(1); + let (read_fd, write_fd) = shutdown_pipe(); + tx.try_send(vec![0u8]).unwrap(); + signal(write_fd); + + let (done_tx, done_rx) = mpsc::channel(); + let sender = std::thread::spawn(move || { + let outcome = send_frame(&tx, vec![1u8], read_fd); + done_tx.send(matches!(outcome, SendOutcome::Stop)).unwrap(); + }); + + let stopped = done_rx + .recv_timeout(Duration::from_secs(5)) + .expect("send_frame parked past a shutdown request"); + sender.join().unwrap(); + assert!(stopped); + } + + #[test] + #[ignore = "demonstrates the defect: blocking_send never returns, so this hangs"] + fn blocking_send_parks_past_a_shutdown_request_when_the_channel_is_full() { + // Run with `--ignored` to watch the old send site hang. Kept as the + // observed red-before for the test above, which cannot itself fail + // against the old code because `send_frame` did not exist then. + let (tx, _rx) = tokio::sync::mpsc::channel::>(1); + let (_read_fd, write_fd) = shutdown_pipe(); + tx.try_send(vec![0u8]).unwrap(); + signal(write_fd); + + let (done_tx, done_rx) = mpsc::channel(); + std::thread::spawn(move || { + let _ = tx.blocking_send(vec![1u8]); + done_tx.send(()).unwrap(); + }); + + done_rx + .recv_timeout(Duration::from_secs(5)) + .expect("blocking_send returned, so the send site was already cancellable"); + } +}