Stop the macOS BPF reader from parking past a shutdown request

The 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.

Drop the receiver before joining. A send parked on a full channel then fails
at once, which costs nothing in steady state. Do the same take in shutdown()
where the lock is free, since stop_async calls it before the abort.

Send through a helper that watches the same shutdown pipe the thread's
select() already honours, so a send waiting for room cannot outlive a
shutdown request even if the receiver is still held. The helper yields before
it sleeps, so the saturated-path handoff rate is unchanged; that matters here
because this thread exists to lift a throughput ceiling. It lives at module
scope rather than inside the macOS-gated module so Linux CI compiles and
tests it.
This commit is contained in:
Johnathan Corgan
2026-08-23 11:45:07 +01:00
parent 759fee50f5
commit 7f5813f305
2 changed files with 256 additions and 6 deletions
+15
View File
@@ -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
+241 -6
View File
@@ -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<T>(
tx: &tokio::sync::mpsc::Sender<T>,
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<PacketSocket>,
rx: tokio::sync::Mutex<tokio::sync::mpsc::Receiver<Frame>>,
/// `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<Option<tokio::sync::mpsc::Receiver<Frame>>>,
reader_thread: Option<std::thread::JoinHandle<()>>,
}
@@ -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::<Vec<u8>>(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::<Vec<u8>>(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::<Vec<u8>>(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::<Vec<u8>>(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::<Vec<u8>>(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");
}
}