Files
smarm/src/io.rs
T

462 lines
19 KiB
Rust

//! Off-scheduler IO: blocking-work offload and epoll-based fd readiness.
//!
//! `block_on_io(closure)` runs `closure` on a dedicated worker OS thread,
//! parks the calling actor in the meantime, and returns the closure's
//! value when it completes. Lets actors call into blocking C libraries,
//! synchronous file IO, or anything else that doesn't fit the readiness
//! model.
//!
//! `wait_readable(fd)` / `wait_writable(fd)` register interest in an fd
//! with epoll and park the calling actor. When the fd becomes ready, the
//! epoll thread unparks the actor. The actual `read(2)`/`write(2)` syscall
//! runs back on the scheduler thread, *inside* the actor — buffer never
//! leaves the actor, no copying through an intermediary thread. Built on
//! these are the conveniences `read(fd, &mut buf)` and `write(fd, &buf)`.
//!
//! Architecture (RFC 018: driver-enqueues)
//! =======================================
//! Per `run()`, two OS threads, each a *producer* behind the runtime's
//! two-call contract — make the actor runnable (`unpark_at`, whose enqueue
//! tail wakes a parked scheduler), nothing else:
//!
//! - **epoll thread**: owns `epoll_wait` on the epollfd. On a ready fd it
//! removes the parked waiter from the shared `waiters` map and DELs the
//! fd (both under the waiters lock — see below), then unparks the
//! actor directly. On the shutdown pipe (also registered in the
//! epollfd), exits.
//! - **pool thread**: blocks on the request mpsc. Runs the closure inside
//! `catch_unwind`, stashes the result in the actor's slot
//! (`pending_io_result`, under the cold lock, generation-checked),
//! decrements the runtime's `io_outstanding`, and unparks the actor.
//!
//! There is no shared completion queue and no wake pipe: each producer
//! routes its own completion, so the whole byte-vs-completion visibility
//! discipline of the drain era — and the stranded-completion hazards it
//! defended against — is unrepresentable. Producers reach the runtime
//! through a `Weak<RuntimeInner>`: upgraded per completion (the path is
//! syscall-bound; the refcount op is noise) and avoiding an Arc cycle
//! through `RuntimeInner::io`.
//!
//! `epoll_ctl` (register fd interest) is called by the scheduler thread
//! directly on the epollfd. That's well-defined per `epoll_ctl(2)`: a
//! thread may be calling `epoll_wait` on the epollfd while another thread
//! calls `epoll_ctl`.
//!
//! Epoll mode
//! ==========
//! Level-triggered with EPOLLONESHOT. After a wakeup the kernel
//! auto-disarms the fd, so we never get two wakeups for one
//! `wait_readable` call. The epoll thread explicitly `EPOLL_CTL_DEL`s the
//! fd on readiness to free the slot for re-registration. Net effect: each
//! `wait_readable(fd)` is one ADD, one wakeup, one DEL — symmetric and
//! stateless between calls.
//!
//! ## The waiters lock is the ADD/DEL serialization
//!
//! Registration (scheduler thread: check-vacant, defensive DEL, ADD,
//! insert) and readiness consumption (epoll thread: remove, DEL) each run
//! entirely under the `waiters` mutex. This is what makes the
//! oneshot-rearm race unrepresentable: a woken actor re-registering the
//! same fd cannot interleave with the epoll thread's DEL for the *previous*
//! registration — whichever takes the lock second sees a consistent
//! kernel-side state. Lock order: `io` (the runtime's outer mutex, held by
//! scheduler-side callers) → `waiters` → slot/queue leaves via `unpark_at`.
//! The epoll thread takes `waiters` without `io` — it must never take
//! `io`, both for lock-order hygiene and because teardown holds `io` while
//! joining it.
//!
//! Fd hygiene
//! ==========
//! An actor stopped while waiting on an fd unwinds out of `wait_fd`'s park;
//! a drop guard there (armed after a successful register, forgotten on a
//! normal wake) calls [`IoThread::cancel_waiter`], which removes the
//! `waiters` entry iff it is still that wait's `(pid, epoch)` and only then
//! `EPOLL_CTL_DEL`s the fd — an entry already consumed by the epoll thread
//! means the fd may carry someone else's fresh registration, which must be
//! left alone. `epoll_register` keeps a defensive bare DEL before ADD as
//! belt-and-braces.
//!
//! Buffers used with `read`/`write` should be on fds opened with
//! `O_NONBLOCK`. If they aren't, the syscall may block the scheduler
//! thread despite the readiness notification (the fd reporting readable
//! doesn't guarantee the syscall completes without blocking — e.g. a
//! signal could be delivered). Documented; not enforced.
//!
//! Panic handling
//! ==============
//! The pool worker runs the closure inside `catch_unwind` and ships either
//! the return value or the panic payload back to the scheduler.
//! `block_on_io` resumes the panic on the calling actor's stack, so the
//! actor's supervisor sees a real `Signal::Panic` as if the work had run
//! inline. Fd-wait primitives don't run user code on the IO thread, so
//! they have no equivalent panic-propagation path.
use crate::pid::Pid;
use crate::runtime::RuntimeInner;
use std::any::Any;
use std::collections::HashMap;
use std::io;
use std::os::fd::RawFd;
use std::panic;
use std::sync::atomic::Ordering;
use std::sync::{mpsc, Arc, Mutex, Weak};
use std::thread::JoinHandle as OsJoinHandle;
// ---------------------------------------------------------------------------
// Wire types
// ---------------------------------------------------------------------------
/// What the pool stores while computing a result. `Ok` is the closure's
/// return value (boxed as `Any`); `Err` is the panic payload.
pub type IoResult = Result<Box<dyn Any + Send>, Box<dyn Any + Send>>;
struct Request {
/// The submitter's park-epoch — the eventual wake is epoch-matched.
epoch: u32,
pid: Pid,
/// The work to perform. Returns the wire-form result directly.
work: Box<dyn FnOnce() -> IoResult + Send>,
}
/// The parked-waiter map, shared between scheduler-side registration and
/// the epoll thread's readiness consumption. See the module docs on why
/// this single lock is the ADD/DEL serialization.
type Waiters = Arc<Mutex<HashMap<RawFd, (Pid, u32)>>>;
// ---------------------------------------------------------------------------
// IoThread — created per `run()`, owned by `RuntimeInner::io`.
// ---------------------------------------------------------------------------
pub struct IoThread {
/// Submission queue into the blocking-work pool.
tx: mpsc::Sender<Request>,
/// One parked actor per registered fd. Populated by `epoll_register`,
/// consumed by the epoll thread on readiness or `cancel_waiter` on an
/// unwound wait.
waiters: Waiters,
// ----- Epoll machinery -----
/// The epollfd, owned by `IoThread`. Callable cross-thread via
/// `epoll_ctl` per the man page.
epollfd: RawFd,
/// Pipe used to signal the epoll thread to exit. Registered inside the
/// epollfd so a single `epoll_wait` covers both fd readiness and
/// shutdown.
shutdown_read: RawFd,
shutdown_write: RawFd,
// ----- Threads -----
pool_thread: Option<OsJoinHandle<()>>,
epoll_thread: Option<OsJoinHandle<()>>,
}
impl IoThread {
/// Start the pool and epoll threads. `rt` is the producers' route back
/// into the runtime (slot table + unpark protocol); a `Weak` so the
/// `RuntimeInner → IoThread → RuntimeInner` cycle never forms.
pub(crate) fn start(rt: Weak<RuntimeInner>) -> io::Result<Self> {
// Pool submission channel.
let (tx, rx) = mpsc::channel::<Request>();
let waiters: Waiters = Arc::new(Mutex::new(HashMap::new()));
// Epoll machinery.
let epollfd = unsafe { libc::epoll_create1(libc::EPOLL_CLOEXEC) };
if epollfd < 0 {
return Err(io::Error::last_os_error());
}
let (shutdown_read, shutdown_write) = match make_pipe() {
Ok(p) => p,
Err(e) => {
unsafe {
libc::close(epollfd);
}
return Err(e);
}
};
// Register the shutdown pipe in epollfd. We use a sentinel `data`
// value to recognise shutdown events. RawFd values are non-negative,
// so u64::MAX is unambiguously not a real fd-data encoding.
let mut shutdown_ev = libc::epoll_event {
events: libc::EPOLLIN as u32,
u64: SHUTDOWN_EPOLL_TOKEN,
};
if unsafe {
libc::epoll_ctl(
epollfd,
libc::EPOLL_CTL_ADD,
shutdown_read,
&mut shutdown_ev as *mut _,
)
} < 0
{
let e = io::Error::last_os_error();
unsafe {
libc::close(epollfd);
libc::close(shutdown_read);
libc::close(shutdown_write);
}
return Err(e);
}
// Spawn pool thread.
let pool_rt = rt.clone();
let pool_thread = std::thread::Builder::new()
.name("smarm-io-pool".into())
.spawn(move || pool_loop(rx, pool_rt))?;
// Spawn epoll thread.
let epoll_waiters = waiters.clone();
let epoll_thread = std::thread::Builder::new()
.name("smarm-io-epoll".into())
.spawn(move || epoll_loop(epollfd, epoll_waiters, rt))?;
Ok(Self {
tx,
waiters,
epollfd,
shutdown_read,
shutdown_write,
pool_thread: Some(pool_thread),
epoll_thread: Some(epoll_thread),
})
}
/// Hand a request to the pool. The caller (scheduler.rs) increments
/// `io_outstanding` BEFORE calling — the pool decrements on completion,
/// and an increment that trailed the completion would underflow.
pub fn submit(&mut self, pid: Pid, epoch: u32, work: Box<dyn FnOnce() -> IoResult + Send>) {
// Send can only fail if the pool has hung up, which only happens
// on shutdown. submit during shutdown is a bug.
if self.tx.send(Request { pid, epoch, work }).is_err() {
panic!("smarm: io pool hung up unexpectedly (submit during shutdown)");
}
}
/// Register interest in `fd` becoming readable/writable; record `pid`
/// as the parked waiter. The epoll thread unparks it on readiness.
/// The caller increments `io_fd_waiters` BEFORE calling (mirror of
/// `submit`'s contract) and decrements it again if this errors.
///
/// EPOLLONESHOT: one wakeup per registration; the epoll thread DELs on
/// readiness, `cancel_waiter` DELs on an unwound wait.
pub fn epoll_register(
&mut self,
fd: RawFd,
pid: Pid,
epoch: u32,
readable: bool,
writable: bool,
) -> io::Result<()> {
let mut waiters = match self.waiters.lock() {
Ok(g) => g,
Err(e) => panic!("smarm: io waiters lock poisoned (core corrupt): {e}"),
};
// Two actors waiting on the same fd would be a misuse: the kernel
// delivers exactly one EPOLLONESHOT wakeup, so the second waiter
// would hang. Reject up front.
if waiters.contains_key(&fd) {
return Err(io::Error::new(
io::ErrorKind::AlreadyExists,
"fd already has a parked waiter",
));
}
// Belt-and-braces: `cancel_waiter` is responsible for cleaning up a
// stopped waiter's registration, but a bare DEL is harmless if the
// fd isn't registered (ENOENT) and removes any leak a path we
// haven't thought of might leave behind.
unsafe {
libc::epoll_ctl(self.epollfd, libc::EPOLL_CTL_DEL, fd, std::ptr::null_mut());
}
let mut events: u32 = libc::EPOLLONESHOT as u32;
if readable {
events |= libc::EPOLLIN as u32;
}
if writable {
events |= libc::EPOLLOUT as u32;
}
let mut ev = libc::epoll_event {
events,
u64: fd as u64,
};
let r =
unsafe { libc::epoll_ctl(self.epollfd, libc::EPOLL_CTL_ADD, fd, &mut ev as *mut _) };
if r < 0 {
return Err(io::Error::last_os_error());
}
waiters.insert(fd, (pid, epoch));
Ok(())
}
/// Remove `fd`'s waiter iff it is still `(pid, epoch)`, DELing the fd
/// from the epollfd in the same critical section. Returns whether the
/// entry was removed (the caller then decrements `io_fd_waiters`).
/// `false` means the epoll thread consumed the registration first —
/// the fd may already carry someone else's fresh ADD; hands off.
pub fn cancel_waiter(&mut self, fd: RawFd, pid: Pid, epoch: u32) -> bool {
let mut waiters = match self.waiters.lock() {
Ok(g) => g,
Err(e) => panic!("smarm: io waiters lock poisoned (core corrupt): {e}"),
};
if waiters.get(&fd) == Some(&(pid, epoch)) {
waiters.remove(&fd);
// EPOLL_CTL_DEL of an already-removed fd returns ENOENT; ignore.
unsafe {
libc::epoll_ctl(self.epollfd, libc::EPOLL_CTL_DEL, fd, std::ptr::null_mut());
}
true
} else {
false
}
}
}
impl Drop for IoThread {
fn drop(&mut self) {
// 1. Signal the epoll thread to exit by writing the shutdown pipe.
unsafe {
let buf: [u8; 1] = [0];
// Single byte; we don't care about EINTR retry here — worst
// case the epoll thread blocks until process exit, which is
// fine because we then close fds out from under it.
libc::write(self.shutdown_write, buf.as_ptr() as *const _, 1);
}
// 2. Hang up the pool's request channel so the pool thread exits.
let (dead_tx, _) = mpsc::channel::<Request>();
let real_tx = std::mem::replace(&mut self.tx, dead_tx);
drop(real_tx);
// 3. Join both threads. Safe even while the caller holds the
// runtime's `io` mutex: neither thread ever takes it (they reach
// the runtime through a Weak they upgrade per completion, and
// the epoll thread's only lock is `waiters`).
if let Some(h) = self.epoll_thread.take() {
let _ = h.join();
}
if let Some(h) = self.pool_thread.take() {
let _ = h.join();
}
// 4. Close fds.
unsafe {
libc::close(self.epollfd);
libc::close(self.shutdown_read);
libc::close(self.shutdown_write);
}
}
}
/// Sentinel `epoll_event.u64` distinguishing the shutdown pipe from
/// registered actor fds. RawFd values fit in i32, so the high bits are
/// available for a marker; we use u64::MAX which can't be a valid fd.
const SHUTDOWN_EPOLL_TOKEN: u64 = u64::MAX;
// ---------------------------------------------------------------------------
// Pool loop (producer: Blocking completions)
// ---------------------------------------------------------------------------
fn pool_loop(rx: mpsc::Receiver<Request>, rt: Weak<RuntimeInner>) {
while let Ok(Request { pid, epoch, work }) = rx.recv() {
let result: IoResult = match panic::catch_unwind(panic::AssertUnwindSafe(work)) {
Ok(r) => r,
Err(payload) => Err(payload),
};
let Some(inner) = rt.upgrade() else { return };
// Stash the result under the cold lock (generation-checked: an
// actor stopped with the op in flight discards it), decrement the
// in-flight count, then wake through the epoch-matched unpark. The
// unpark's enqueue tail wakes a parked scheduler; the actor stays
// `live` until it resumes and finalizes, so the decrement's
// ordering against the termination verdict is not load-bearing.
if let Some(slot) = inner.slot_at(pid) {
let mut cold = slot.cold.lock();
if slot.generation() == pid.generation() {
cold.pending_io_result = Some(result);
}
}
inner.io_outstanding.fetch_sub(1, Ordering::AcqRel);
inner.unpark_at(pid, epoch);
}
}
// ---------------------------------------------------------------------------
// Epoll loop (producer: FdReady completions)
// ---------------------------------------------------------------------------
fn epoll_loop(epollfd: RawFd, waiters: Waiters, rt: Weak<RuntimeInner>) {
// Buffer for epoll_wait. 64 is plenty for our scale; if a real load
// appears that needs more, this is a one-line change.
const MAX_EVENTS: usize = 64;
let mut events: [libc::epoll_event; MAX_EVENTS] = unsafe { std::mem::zeroed() };
loop {
let n = unsafe {
libc::epoll_wait(epollfd, events.as_mut_ptr(), MAX_EVENTS as libc::c_int, -1)
};
if n < 0 {
let e = unsafe { *libc::__errno_location() };
if e == libc::EINTR {
continue;
}
// Anything else here is a programming error (EBADF on epollfd
// after we've closed it from Drop — the close races with us).
// Treat as shutdown.
return;
}
let mut shutdown_requested = false;
for ev in events.iter().take(n as usize) {
if ev.u64 == SHUTDOWN_EPOLL_TOKEN {
shutdown_requested = true;
continue;
}
let fd = ev.u64 as RawFd;
// Consume the registration: remove + DEL under the waiters
// lock (the ADD/DEL serialization — see module docs). A
// vanished entry means `cancel_waiter` beat us: the wake is
// already moot.
let entry = {
let mut w = match waiters.lock() {
Ok(g) => g,
Err(e) => {
panic!("smarm: io waiters lock poisoned (core corrupt): {e}")
}
};
let entry = w.remove(&fd);
if entry.is_some() {
unsafe {
libc::epoll_ctl(epollfd, libc::EPOLL_CTL_DEL, fd, std::ptr::null_mut());
}
}
entry
};
if let Some((pid, epoch)) = entry {
let Some(inner) = rt.upgrade() else { return };
inner.io_fd_waiters.fetch_sub(1, Ordering::AcqRel);
inner.unpark_at(pid, epoch);
}
}
if shutdown_requested {
return;
}
}
}
// ---------------------------------------------------------------------------
// Pipe helper
// ---------------------------------------------------------------------------
fn make_pipe() -> io::Result<(RawFd, RawFd)> {
let mut fds: [libc::c_int; 2] = [0; 2];
let r = unsafe { libc::pipe2(fds.as_mut_ptr(), libc::O_CLOEXEC | libc::O_NONBLOCK) };
if r != 0 {
return Err(io::Error::last_os_error());
}
Ok((fds[0], fds[1]))
}