Compare commits
4
Commits
8c764e9169
...
master
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d4839f1d81 | ||
|
|
2854b560d6 | ||
|
|
7b026cfe56 | ||
|
|
006a3283e7 |
@@ -2,7 +2,23 @@
|
||||
# smarm pre-commit gate: clippy the library (src/) with warnings as errors.
|
||||
# unwrap_used / expect_used are denied (Cargo.toml [lints.clippy]): library
|
||||
# code must not hide a panic behind unwrap/expect. Tests/examples are not gated.
|
||||
#
|
||||
# Toolchain resolution: prefer an installed cargo-clippy; on machines whose
|
||||
# rust comes without the clippy component (e.g. NixOS home-manager), fall
|
||||
# back to an ephemeral nix-shell toolchain. The fallback uses its own target
|
||||
# dir (target/clippy) because the shell's rustc version may differ from the
|
||||
# default toolchain's — mixed-compiler artifacts in one target dir are an
|
||||
# E0514 hard error. MSRV (Cargo.toml rust-version) keeps the older shell
|
||||
# toolchain a legitimate gate.
|
||||
set -eu
|
||||
[ -f "$HOME/.cargo/env" ] && . "$HOME/.cargo/env"
|
||||
cd "$(git rev-parse --show-toplevel)"
|
||||
if cargo clippy --version >/dev/null 2>&1; then
|
||||
cargo clippy --lib -- -D warnings
|
||||
elif command -v nix-shell >/dev/null 2>&1; then
|
||||
nix-shell -p clippy -p cargo -p rustc \
|
||||
--run 'CARGO_TARGET_DIR=target/clippy cargo clippy --lib -- -D warnings'
|
||||
else
|
||||
echo "pre-commit: cargo clippy unavailable and no nix-shell fallback" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
@@ -13,44 +13,68 @@
|
||||
//! leaves the actor, no copying through an intermediary thread. Built on
|
||||
//! these are the conveniences `read(fd, &mut buf)` and `write(fd, &buf)`.
|
||||
//!
|
||||
//! Architecture
|
||||
//! ============
|
||||
//! Per `run()`, two OS threads:
|
||||
//! - **epoll thread**: owns the epollfd. Loops in `epoll_wait`. On a
|
||||
//! ready fd, pushes `Completion::FdReady { pid, fd, events }` to the
|
||||
//! shared completion queue and writes the scheduler-wake pipe. On the
|
||||
//! shutdown pipe (also registered in epollfd), exits.
|
||||
//! - **pool thread**: blocks on the request mpsc. Runs the closure
|
||||
//! inside `catch_unwind`, pushes `Completion::Blocking { pid, result }`,
|
||||
//! writes the scheduler-wake pipe.
|
||||
//! 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:
|
||||
//!
|
||||
//! Both threads share a single `completions: Arc<Mutex<VecDeque<Completion>>>`
|
||||
//! and the same scheduler-wake pipe.
|
||||
//! - **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.
|
||||
//!
|
||||
//! `epoll_ctl` (register/unregister 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`. Avoids needing a second mpsc
|
||||
//! and a second wake mechanism.
|
||||
//! 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 scheduler explicitly `EPOLL_CTL_DEL`s the fd
|
||||
//! on completion to free the slot for re-registration. Net effect: each
|
||||
//! `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) 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 a racing `FdReady` 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.
|
||||
//! 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
|
||||
@@ -68,13 +92,14 @@
|
||||
//! they have no equivalent panic-propagation path.
|
||||
|
||||
use crate::pid::Pid;
|
||||
use crate::runtime::RuntimeInner;
|
||||
use std::any::Any;
|
||||
use std::collections::{HashMap, VecDeque};
|
||||
use std::collections::HashMap;
|
||||
use std::io;
|
||||
use std::os::fd::RawFd;
|
||||
use std::panic;
|
||||
use std::sync::mpsc;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::sync::atomic::Ordering;
|
||||
use std::sync::{mpsc, Arc, Mutex, Weak};
|
||||
use std::thread::JoinHandle as OsJoinHandle;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -86,42 +111,29 @@ use std::thread::JoinHandle as OsJoinHandle;
|
||||
pub type IoResult = Result<Box<dyn Any + Send>, Box<dyn Any + Send>>;
|
||||
|
||||
struct Request {
|
||||
/// The submitter's park-epoch — carried through to the `Blocking`
|
||||
/// completion so the wake is epoch-matched.
|
||||
/// 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>,
|
||||
}
|
||||
|
||||
/// Completion message from either IO thread back to the scheduler.
|
||||
pub enum Completion {
|
||||
/// A `block_on_io` closure has finished (Ok = return value, Err = panic
|
||||
/// payload).
|
||||
Blocking { pid: Pid, epoch: u32, result: IoResult },
|
||||
/// An fd registered via `wait_readable`/`wait_writable` is ready. The
|
||||
/// scheduler looks up the parked pid in `waiters`, unparks it, and
|
||||
/// removes the entry. `pid` isn't in this variant because the epoll
|
||||
/// thread doesn't have access to the `waiters` map; the scheduler
|
||||
/// thread owns that.
|
||||
FdReady { fd: RawFd, events: u32 },
|
||||
}
|
||||
/// 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 `SchedulerState`.
|
||||
// IoThread — created per `run()`, owned by `RuntimeInner::io`.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
pub struct IoThread {
|
||||
// ----- Channels & queues -----
|
||||
|
||||
/// Submission queue into the blocking-work pool.
|
||||
tx: mpsc::Sender<Request>,
|
||||
/// Shared completion queue, fed by both the pool and the epoll thread.
|
||||
completions: Arc<Mutex<VecDeque<Completion>>>,
|
||||
/// Pipe the scheduler polls in its idle path. Both IO threads write to
|
||||
/// `wake_write` after pushing a completion.
|
||||
wake_read: RawFd,
|
||||
wake_write: RawFd,
|
||||
/// 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 -----
|
||||
|
||||
@@ -133,39 +145,25 @@ pub struct IoThread {
|
||||
/// shutdown.
|
||||
shutdown_read: RawFd,
|
||||
shutdown_write: RawFd,
|
||||
/// One parked actor per registered fd. Populated by `wait_readable` /
|
||||
/// `wait_writable` and drained by the scheduler when a `FdReady`
|
||||
/// completion is processed.
|
||||
pub waiters: HashMap<RawFd, (Pid, u32)>,
|
||||
|
||||
// ----- Threads -----
|
||||
|
||||
pool_thread: Option<OsJoinHandle<()>>,
|
||||
epoll_thread: Option<OsJoinHandle<()>>,
|
||||
|
||||
/// Number of `block_on_io` requests in-flight. Used by the scheduler's
|
||||
/// idle path to decide whether to wait on the pipe or exit. Fd waits
|
||||
/// are not counted here; they're counted by `waiters.len()`.
|
||||
pub outstanding: u32,
|
||||
}
|
||||
|
||||
impl IoThread {
|
||||
pub fn start() -> io::Result<Self> {
|
||||
// Scheduler-facing wake pipe.
|
||||
let (wake_read, wake_write) = make_pipe()?;
|
||||
// Pool submission channel + shared completion queue.
|
||||
/// 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 completions: Arc<Mutex<VecDeque<Completion>>> =
|
||||
Arc::new(Mutex::new(VecDeque::new()));
|
||||
let waiters: Waiters = Arc::new(Mutex::new(HashMap::new()));
|
||||
|
||||
// Epoll machinery.
|
||||
let epollfd = unsafe { libc::epoll_create1(libc::EPOLL_CLOEXEC) };
|
||||
if epollfd < 0 {
|
||||
// Best-effort fd cleanup before bailing.
|
||||
unsafe {
|
||||
libc::close(wake_read);
|
||||
libc::close(wake_write);
|
||||
}
|
||||
return Err(io::Error::last_os_error());
|
||||
}
|
||||
|
||||
@@ -174,8 +172,6 @@ impl IoThread {
|
||||
Err(e) => {
|
||||
unsafe {
|
||||
libc::close(epollfd);
|
||||
libc::close(wake_read);
|
||||
libc::close(wake_write);
|
||||
}
|
||||
return Err(e);
|
||||
}
|
||||
@@ -202,42 +198,37 @@ impl IoThread {
|
||||
libc::close(epollfd);
|
||||
libc::close(shutdown_read);
|
||||
libc::close(shutdown_write);
|
||||
libc::close(wake_read);
|
||||
libc::close(wake_write);
|
||||
}
|
||||
return Err(e);
|
||||
}
|
||||
|
||||
// Spawn pool thread.
|
||||
let pool_comps = completions.clone();
|
||||
let pool_rt = rt.clone();
|
||||
let pool_thread = std::thread::Builder::new()
|
||||
.name("smarm-io-pool".into())
|
||||
.spawn(move || pool_loop(rx, pool_comps, wake_write))?;
|
||||
.spawn(move || pool_loop(rx, pool_rt))?;
|
||||
|
||||
// Spawn epoll thread.
|
||||
let epoll_comps = completions.clone();
|
||||
let epoll_waiters = waiters.clone();
|
||||
let epoll_thread = std::thread::Builder::new()
|
||||
.name("smarm-io-epoll".into())
|
||||
.spawn(move || epoll_loop(epollfd, epoll_comps, wake_write))?;
|
||||
.spawn(move || epoll_loop(epollfd, epoll_waiters, rt))?;
|
||||
|
||||
Ok(Self {
|
||||
tx,
|
||||
completions,
|
||||
wake_read,
|
||||
wake_write,
|
||||
waiters,
|
||||
epollfd,
|
||||
shutdown_read,
|
||||
shutdown_write,
|
||||
waiters: HashMap::new(),
|
||||
pool_thread: Some(pool_thread),
|
||||
epoll_thread: Some(epoll_thread),
|
||||
outstanding: 0,
|
||||
})
|
||||
}
|
||||
|
||||
/// Hand a request to the pool. Increments `outstanding`.
|
||||
/// 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>) {
|
||||
self.outstanding += 1;
|
||||
// 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() {
|
||||
@@ -245,39 +236,13 @@ impl IoThread {
|
||||
}
|
||||
}
|
||||
|
||||
/// Drain every available completion. Caller (the scheduler) routes the
|
||||
/// results and updates `outstanding` / `waiters` accordingly.
|
||||
pub fn drain_completions(&mut self) -> Vec<Completion> {
|
||||
let mut q = match self.completions.lock() {
|
||||
Ok(g) => g,
|
||||
Err(e) => panic!("smarm: io completions lock poisoned (core corrupt): {e}"),
|
||||
};
|
||||
let mut out = Vec::with_capacity(q.len());
|
||||
while let Some(c) = q.pop_front() {
|
||||
out.push(c);
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
pub fn wake_fd(&self) -> RawFd {
|
||||
self.wake_read
|
||||
}
|
||||
|
||||
/// Write the wake pipe directly: rouse every scheduler thread blocked in
|
||||
/// its idle `poll_wake`. Used by the terminal (AllDone) path — an idle
|
||||
/// sibling may be blocked on a snapshot that nothing will ever refresh
|
||||
/// (an orphaned timer deadline, or `io_outstanding` from a waiter that
|
||||
/// was stop-cancelled and so never produces a completion).
|
||||
pub fn wake(&self) {
|
||||
wake_scheduler(self.wake_write);
|
||||
}
|
||||
|
||||
/// Register interest in `fd` becoming readable/writable; record `pid`
|
||||
/// as the parked waiter. The epoll thread will push a `FdReady`
|
||||
/// completion when the kernel signals.
|
||||
/// 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 scheduler must
|
||||
/// `epoll_del` on completion to free the slot for re-registration.
|
||||
/// 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,
|
||||
@@ -286,20 +251,24 @@ impl IoThread {
|
||||
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 self.waiters.contains_key(&fd) {
|
||||
if waiters.contains_key(&fd) {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::AlreadyExists,
|
||||
"fd already has a parked waiter",
|
||||
));
|
||||
}
|
||||
|
||||
// Belt-and-braces: the unwind guard in `wait_fd` 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.
|
||||
// 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());
|
||||
}
|
||||
@@ -321,20 +290,30 @@ impl IoThread {
|
||||
if r < 0 {
|
||||
return Err(io::Error::last_os_error());
|
||||
}
|
||||
self.waiters.insert(fd, (pid, epoch));
|
||||
waiters.insert(fd, (pid, epoch));
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Remove `fd` from the epollfd. Called by the scheduler after a
|
||||
/// `FdReady` completion, so the next `wait_readable(fd)` can ADD again.
|
||||
///
|
||||
/// Does NOT touch `waiters` — that's the scheduler's bookkeeping; this
|
||||
/// is purely the kernel-side cleanup.
|
||||
pub fn epoll_deregister(&mut self, fd: RawFd) {
|
||||
/// 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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -354,7 +333,10 @@ impl Drop for IoThread {
|
||||
let real_tx = std::mem::replace(&mut self.tx, dead_tx);
|
||||
drop(real_tx);
|
||||
|
||||
// 3. Join both threads.
|
||||
// 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();
|
||||
}
|
||||
@@ -367,8 +349,6 @@ impl Drop for IoThread {
|
||||
libc::close(self.epollfd);
|
||||
libc::close(self.shutdown_read);
|
||||
libc::close(self.shutdown_write);
|
||||
libc::close(self.wake_read);
|
||||
libc::close(self.wake_write);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -379,36 +359,38 @@ impl Drop for IoThread {
|
||||
const SHUTDOWN_EPOLL_TOKEN: u64 = u64::MAX;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Pool loop
|
||||
// Pool loop (producer: Blocking completions)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn pool_loop(
|
||||
rx: mpsc::Receiver<Request>,
|
||||
completions: Arc<Mutex<VecDeque<Completion>>>,
|
||||
wake_write: RawFd,
|
||||
) {
|
||||
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),
|
||||
};
|
||||
match completions.lock() {
|
||||
Ok(mut g) => g.push_back(Completion::Blocking { pid, epoch, result }),
|
||||
Err(e) => panic!("smarm: io completions lock poisoned (core corrupt): {e}"),
|
||||
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);
|
||||
}
|
||||
wake_scheduler(wake_write);
|
||||
}
|
||||
inner.io_outstanding.fetch_sub(1, Ordering::AcqRel);
|
||||
inner.unpark_at(pid, epoch);
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Epoll loop
|
||||
// Epoll loop (producer: FdReady completions)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn epoll_loop(
|
||||
epollfd: RawFd,
|
||||
completions: Arc<Mutex<VecDeque<Completion>>>,
|
||||
wake_write: RawFd,
|
||||
) {
|
||||
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;
|
||||
@@ -436,29 +418,41 @@ fn epoll_loop(
|
||||
}
|
||||
|
||||
let mut shutdown_requested = false;
|
||||
let mut pushed_any = false;
|
||||
{
|
||||
let mut q = match completions.lock() {
|
||||
Ok(g) => g,
|
||||
Err(e) => panic!("smarm: io completions lock poisoned (core corrupt): {e}"),
|
||||
};
|
||||
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;
|
||||
let evs = ev.events;
|
||||
q.push_back(Completion::FdReady {
|
||||
// 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,
|
||||
events: evs,
|
||||
});
|
||||
pushed_any = true;
|
||||
std::ptr::null_mut(),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
if pushed_any {
|
||||
wake_scheduler(wake_write);
|
||||
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;
|
||||
@@ -466,27 +460,8 @@ fn epoll_loop(
|
||||
}
|
||||
}
|
||||
|
||||
/// Write one byte to the scheduler's wake pipe. Retries on EINTR; ignores
|
||||
/// EAGAIN (pipe full means there's already an outstanding wake we haven't
|
||||
/// consumed yet, which is sufficient).
|
||||
fn wake_scheduler(wake_write: RawFd) {
|
||||
let buf: [u8; 1] = [0];
|
||||
unsafe {
|
||||
loop {
|
||||
let n = libc::write(wake_write, buf.as_ptr() as *const _, 1);
|
||||
if n < 0 {
|
||||
let e = *libc::__errno_location();
|
||||
if e == libc::EINTR {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Pipe helpers (unchanged from v0.2)
|
||||
// Pipe helper
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn make_pipe() -> io::Result<(RawFd, RawFd)> {
|
||||
@@ -497,50 +472,3 @@ fn make_pipe() -> io::Result<(RawFd, RawFd)> {
|
||||
}
|
||||
Ok((fds[0], fds[1]))
|
||||
}
|
||||
|
||||
/// Drain pending bytes from the wake pipe. Nonblocking (pipe is O_NONBLOCK).
|
||||
///
|
||||
/// DISCIPLINE: called only by the phase-1 drain-lock winner, immediately
|
||||
/// before `drain_completions`. Bytes are the notification channel for
|
||||
/// completions; consuming one anywhere else can strand the completion it
|
||||
/// announces (see the lost-wakeup note at the call site in `schedule_loop`).
|
||||
pub fn drain_wake_pipe(fd: RawFd) {
|
||||
let mut buf = [0u8; 64];
|
||||
loop {
|
||||
let n = unsafe { libc::read(fd, buf.as_mut_ptr() as *mut _, buf.len()) };
|
||||
if n <= 0 {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Block on `fd` for up to `timeout`, returning when either there's data
|
||||
/// to read or the timeout elapses. `None` for `timeout` means wait forever.
|
||||
pub fn poll_wake(fd: RawFd, timeout: Option<std::time::Duration>) {
|
||||
let timeout_ms: libc::c_int = match timeout {
|
||||
None => -1,
|
||||
Some(d) => {
|
||||
let ms = d.as_millis();
|
||||
if ms > i32::MAX as u128 {
|
||||
i32::MAX
|
||||
} else {
|
||||
ms as i32
|
||||
}
|
||||
}
|
||||
};
|
||||
let mut pfd = libc::pollfd {
|
||||
fd,
|
||||
events: libc::POLLIN,
|
||||
revents: 0,
|
||||
};
|
||||
loop {
|
||||
let r = unsafe { libc::poll(&mut pfd as *mut _, 1, timeout_ms) };
|
||||
if r < 0 {
|
||||
let e = unsafe { *libc::__errno_location() };
|
||||
if e == libc::EINTR {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,6 +32,7 @@ pub mod introspect;
|
||||
#[cfg(feature = "observer")]
|
||||
pub mod observer;
|
||||
pub mod runtime;
|
||||
pub(crate) mod park;
|
||||
pub(crate) mod raw_mutex;
|
||||
pub(crate) mod slot_state;
|
||||
pub(crate) mod sync_shim;
|
||||
|
||||
+994
@@ -0,0 +1,994 @@
|
||||
//! Scheduler park/wake coordination layer (RFC 018).
|
||||
//!
|
||||
//! Schedulers never touch an fd to sleep: they park on a per-thread
|
||||
//! [`Parker`] and are woken through an idle-mask protocol the runtime owns
|
||||
//! outright. IO backends (epoll today, io_uring later) are *producers*
|
||||
//! behind a two-call contract — make actors runnable, then wake — which is
|
||||
//! what makes backend selection tractable (RFC 018 §step-back).
|
||||
//!
|
||||
//! Three pieces, all in [`Coordinator`]:
|
||||
//!
|
||||
//! - **Parker** (one per scheduler): permit semantics, `std::thread::park`
|
||||
//! shaped — an unpark delivered before the park sets a permit; the next
|
||||
//! park consumes it and returns immediately. This single property closes
|
||||
//! the check-then-park race. Linux: `futex(2)` `FUTEX_WAIT`/`FUTEX_WAKE`
|
||||
//! (private) with a nanosecond-precision relative `timespec` — the
|
||||
//! `as_millis` truncation defect of the retired wake pipe is
|
||||
//! unrepresentable here. Loom / non-Linux: `Mutex<bool>` + `Condvar`
|
||||
//! (the loom models run against this build).
|
||||
//! - **Idle mask**: an `AtomicU64` bitmask of parked scheduler ids
|
||||
//! (construction asserts ≤ 64 schedulers). Park protocol: set own bit,
|
||||
//! run the caller's mandatory post-publish re-check, then wait. A
|
||||
//! producer that published work before observing our bit has left us
|
||||
//! work the re-check finds; one that observes the bit wakes us.
|
||||
//!
|
||||
//! The publish/re-check pair is a store-buffer (Dekker) shape. Two sound
|
||||
//! resolutions coexist here, chosen per call-site cost profile: the
|
||||
//! **fence handshake** on the hot producer path
|
||||
//! ([`Coordinator::wake_one_if_idle`]: publish work; `fence(SeqCst)`;
|
||||
//! one *Relaxed* mask load — pairing with the consumer's `fetch_or(bit)`;
|
||||
//! `fence(SeqCst)`; re-check inside [`Coordinator::park`]), so the
|
||||
//! pure-compute hot path (everyone busy, mask 0) never takes the shared
|
||||
//! mask line exclusive — the RFC's "one relaxed load" fast path, made
|
||||
//! sound; and the **same-location-RMW read** (`fetch_or(0)`, in
|
||||
//! [`Coordinator::wake_one`] / [`Coordinator::idle_mask`]) for the rare
|
||||
//! paths (chain rule — at most one per wake) where reading the latest
|
||||
//! mask by modification-order coherence is worth an RMW. The loom models
|
||||
//! drive the fence pattern end to end (a fence-less plain-load draft
|
||||
//! would — and did — fail model 1/2 with a lost wake, as it must).
|
||||
//! - **Timekeeper**: at most one parked scheduler holds the timer
|
||||
//! deadline (RFC 018 §timers) so a timer expiry wakes one scheduler,
|
||||
//! not a herd. The role is an atomic `(holder id, armed deadline)`
|
||||
//! pair; a timer insertion with an earlier deadline wakes the holder to
|
||||
//! re-peek. Arm / insert-check MUST be serialized by the caller (the
|
||||
//! timers mutex in the runtime) — the atomics exist so the *busy-path
|
||||
//! due-check* (one Relaxed load, [`Coordinator::armed_deadline_nanos`])
|
||||
//! and the wake stay lock-free. All races are biased over-wake: a
|
||||
//! spurious permit costs one failed pop; a missed wake would cost a
|
||||
//! stranded actor, and is unrepresentable under the serialization rule.
|
||||
//!
|
||||
//! Every wake here is *at most one* futex round-trip and wakes *exactly
|
||||
//! one* scheduler by construction (`wake_one` CASes a bit clear before
|
||||
//! unparking its owner) — there is no shared level-triggered anything
|
||||
//! left to herd on.
|
||||
//!
|
||||
//! Standalone until the runtime swap (RFC 018 commit 2): nothing outside
|
||||
//! tests constructs a [`Coordinator`] yet.
|
||||
|
||||
use crate::sync_shim::{fence, AtomicU64, Ordering};
|
||||
use std::time::Instant;
|
||||
|
||||
/// Sentinel for "no timekeeper" in the holder word.
|
||||
const NO_TIMEKEEPER: u64 = u64::MAX;
|
||||
/// Sentinel for "no armed deadline" in the deadline word. Also what the
|
||||
/// busy-path due-check compares against: `now_nanos < armed` is one branch.
|
||||
pub(crate) const NO_DEADLINE: u64 = u64::MAX;
|
||||
|
||||
/// Outcome of a [`Coordinator::park`] call.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub(crate) enum ParkResult {
|
||||
/// A permit was consumed (wake delivered before or during the park).
|
||||
Woken,
|
||||
/// The deadline passed with no wake. Only possible when a deadline was
|
||||
/// supplied (and never under loom, which has no time — see `park`).
|
||||
TimedOut,
|
||||
/// The post-publish re-check found work; the thread never blocked.
|
||||
/// A permit may still be pending (a racing `wake_one` picked us after
|
||||
/// the bit was set) — it will surface as one spurious `Woken` on a
|
||||
/// later park. Benign: over-wake by design.
|
||||
WorkFound,
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Parker — permit-semantics thread parking
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Linux, non-loom: futex on a state word.
|
||||
///
|
||||
/// States: EMPTY (no permit, nobody waiting), PARKED (a thread is, or is
|
||||
/// about to be, in `futex_wait`), NOTIFIED (permit pending). The classic
|
||||
/// std-parker protocol: `unpark` swaps to NOTIFIED and futex-wakes iff it
|
||||
/// displaced PARKED; `park` CASes EMPTY→PARKED, waits, and consumes
|
||||
/// NOTIFIED on every exit path.
|
||||
#[cfg(all(target_os = "linux", not(loom)))]
|
||||
mod parker {
|
||||
use std::sync::atomic::{AtomicU32, Ordering};
|
||||
use std::time::Instant;
|
||||
|
||||
const EMPTY: u32 = 0;
|
||||
const PARKED: u32 = 1;
|
||||
const NOTIFIED: u32 = 2;
|
||||
|
||||
pub(super) struct Parker {
|
||||
state: AtomicU32,
|
||||
}
|
||||
|
||||
impl Parker {
|
||||
pub(super) fn new() -> Self {
|
||||
Self { state: AtomicU32::new(EMPTY) }
|
||||
}
|
||||
|
||||
/// Returns `true` = woken (permit consumed), `false` = timed out.
|
||||
pub(super) fn park(&self, deadline: Option<Instant>) -> bool {
|
||||
// Fast path: consume a pending permit without blocking.
|
||||
if self
|
||||
.state
|
||||
.compare_exchange(EMPTY, PARKED, Ordering::AcqRel, Ordering::Acquire)
|
||||
.is_err()
|
||||
{
|
||||
// Only NOTIFIED can be here (one thread parks at a time).
|
||||
self.state.store(EMPTY, Ordering::Release);
|
||||
return true;
|
||||
}
|
||||
loop {
|
||||
let timeout = match deadline {
|
||||
None => None,
|
||||
Some(d) => {
|
||||
let now = Instant::now();
|
||||
if d <= now {
|
||||
// Deadline passed: cancel the park. The swap
|
||||
// races a concurrent unpark — if it delivered
|
||||
// NOTIFIED first, report Woken (never lose a
|
||||
// permit).
|
||||
return self.state.swap(EMPTY, Ordering::AcqRel) == NOTIFIED;
|
||||
}
|
||||
Some(d - now)
|
||||
}
|
||||
};
|
||||
futex_wait(&self.state, PARKED, timeout);
|
||||
if self
|
||||
.state
|
||||
.compare_exchange(NOTIFIED, EMPTY, Ordering::AcqRel, Ordering::Acquire)
|
||||
.is_ok()
|
||||
{
|
||||
return true;
|
||||
}
|
||||
// Spurious wake or timeout with state still PARKED: loop —
|
||||
// the deadline check at the top decides.
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn unpark(&self) {
|
||||
if self.state.swap(NOTIFIED, Ordering::AcqRel) == PARKED {
|
||||
futex_wake(&self.state, 1);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// `FUTEX_WAIT` with a *relative* nanosecond timeout (`CLOCK_MONOTONIC`
|
||||
/// per futex(2) for relative waits). No millisecond conversion anywhere:
|
||||
/// the timespec carries the full sub-ms remainder (RFC 018 kills the
|
||||
/// `as_millis` truncation structurally).
|
||||
fn futex_wait(word: &AtomicU32, expected: u32, timeout: Option<std::time::Duration>) {
|
||||
let ts;
|
||||
let ts_ptr: *const libc::timespec = match timeout {
|
||||
Some(d) => {
|
||||
ts = libc::timespec {
|
||||
tv_sec: d.as_secs() as libc::time_t,
|
||||
tv_nsec: d.subsec_nanos() as libc::c_long,
|
||||
};
|
||||
&ts
|
||||
}
|
||||
None => std::ptr::null(),
|
||||
};
|
||||
// Errors (EAGAIN: word changed; ETIMEDOUT; EINTR) all mean "return
|
||||
// and let the caller's state machine decide" — deliberately ignored.
|
||||
unsafe {
|
||||
libc::syscall(
|
||||
libc::SYS_futex,
|
||||
word.as_ptr(),
|
||||
libc::FUTEX_WAIT | libc::FUTEX_PRIVATE_FLAG,
|
||||
expected,
|
||||
ts_ptr,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
fn futex_wake(word: &AtomicU32, n: u32) {
|
||||
unsafe {
|
||||
libc::syscall(
|
||||
libc::SYS_futex,
|
||||
word.as_ptr(),
|
||||
libc::FUTEX_WAKE | libc::FUTEX_PRIVATE_FLAG,
|
||||
n,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Loom / non-Linux: `Mutex<bool>` permit + `Condvar` — loom's own model
|
||||
/// of a parker, and the portable fallback. Under loom the deadline is
|
||||
/// ignored (loom has no time); the models exercise wake paths only.
|
||||
#[cfg(any(loom, not(target_os = "linux")))]
|
||||
mod parker {
|
||||
use crate::sync_shim::{Condvar, Mutex};
|
||||
use std::time::Instant;
|
||||
|
||||
pub(super) struct Parker {
|
||||
permit: Mutex<bool>,
|
||||
cv: Condvar,
|
||||
}
|
||||
|
||||
impl Parker {
|
||||
pub(super) fn new() -> Self {
|
||||
Self { permit: Mutex::new(false), cv: Condvar::new() }
|
||||
}
|
||||
|
||||
/// Returns `true` = woken (permit consumed), `false` = timed out.
|
||||
pub(super) fn park(&self, deadline: Option<Instant>) -> bool {
|
||||
let mut permit = match self.permit.lock() {
|
||||
Ok(g) => g,
|
||||
Err(_) => panic!("smarm: parker permit lock poisoned (core corrupt)"),
|
||||
};
|
||||
loop {
|
||||
if *permit {
|
||||
*permit = false;
|
||||
return true;
|
||||
}
|
||||
#[cfg(loom)]
|
||||
{
|
||||
// Loom has no clock: block until a wake. Models must
|
||||
// deliver one (a park nobody wakes is a real deadlock
|
||||
// and loom reports it as such).
|
||||
let _ = deadline;
|
||||
permit = match self.cv.wait(permit) {
|
||||
Ok(g) => g,
|
||||
Err(_) => panic!("smarm: parker cv poisoned (core corrupt)"),
|
||||
};
|
||||
}
|
||||
#[cfg(not(loom))]
|
||||
{
|
||||
match deadline {
|
||||
None => {
|
||||
permit = match self.cv.wait(permit) {
|
||||
Ok(g) => g,
|
||||
Err(_) => panic!("smarm: parker cv poisoned (core corrupt)"),
|
||||
};
|
||||
}
|
||||
Some(d) => {
|
||||
let now = Instant::now();
|
||||
if d <= now {
|
||||
return false;
|
||||
}
|
||||
permit = match self.cv.wait_timeout(permit, d - now) {
|
||||
Ok((g, _)) => g,
|
||||
Err(_) => panic!("smarm: parker cv poisoned (core corrupt)"),
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn unpark(&self) {
|
||||
let mut permit = match self.permit.lock() {
|
||||
Ok(g) => g,
|
||||
Err(_) => panic!("smarm: parker permit lock poisoned (core corrupt)"),
|
||||
};
|
||||
*permit = true;
|
||||
self.cv.notify_one();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
use parker::Parker;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Coordinator — idle mask + wake protocol + timekeeper
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
pub(crate) struct Coordinator {
|
||||
parkers: Box<[Parker]>,
|
||||
/// Bit `i` set = scheduler `i` is parked or committed to parking (set
|
||||
/// before the re-check; cleared by `wake_one`'s CAS or by the parker
|
||||
/// itself on return). AcqRel same-location-RMW handshake — see module
|
||||
/// docs (no SeqCst needed: every producer-side read is an RMW).
|
||||
idle: AtomicU64,
|
||||
/// Timekeeper holder id, or `NO_TIMEKEEPER`. Written under the
|
||||
/// caller's timer serialization (arm/disarm/insert-check); read
|
||||
/// lock-free by the insert wake path.
|
||||
tk_holder: AtomicU64,
|
||||
/// Armed deadline as nanos since `origin`, or `NO_DEADLINE`. Written
|
||||
/// only by the timekeeper arm/disarm protocol.
|
||||
tk_armed: AtomicU64,
|
||||
/// Earliest KNOWN timer deadline (nanos since `origin`), or
|
||||
/// `NO_DEADLINE` — independent of whether any scheduler is parked,
|
||||
/// which is what the timekeeper's `tk_armed` cannot give: under
|
||||
/// saturation nobody parks and nobody arms, yet due timers must still
|
||||
/// fire (ratified design point (a)). Maintained under the caller's
|
||||
/// timers mutex (`note_deadline` on insert, `refresh_deadline` after a
|
||||
/// pop/peek); read lock-free by the busy-path due-check.
|
||||
next_deadline: AtomicU64,
|
||||
/// Time origin for the nanos encoding.
|
||||
origin: Instant,
|
||||
}
|
||||
|
||||
impl Coordinator {
|
||||
pub(crate) fn new(schedulers: usize) -> Self {
|
||||
assert!(
|
||||
(1..=64).contains(&schedulers),
|
||||
"smarm: scheduler count must be 1..=64 (idle mask is one u64); got {schedulers}"
|
||||
);
|
||||
Self {
|
||||
parkers: (0..schedulers).map(|_| Parker::new()).collect(),
|
||||
idle: AtomicU64::new(0),
|
||||
tk_holder: AtomicU64::new(NO_TIMEKEEPER),
|
||||
tk_armed: AtomicU64::new(NO_DEADLINE),
|
||||
next_deadline: AtomicU64::new(NO_DEADLINE),
|
||||
origin: Instant::now(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Encode a deadline for the armed snapshot / busy-path compare.
|
||||
/// Saturating: a deadline at-or-before `origin` encodes as 0 (always
|
||||
/// due), one beyond ~584 years as `NO_DEADLINE - 1`.
|
||||
pub(crate) fn deadline_nanos(&self, deadline: Instant) -> u64 {
|
||||
let nanos = deadline
|
||||
.checked_duration_since(self.origin)
|
||||
.map(|d| d.as_nanos())
|
||||
.unwrap_or(0);
|
||||
if nanos >= NO_DEADLINE as u128 {
|
||||
NO_DEADLINE - 1
|
||||
} else {
|
||||
nanos as u64
|
||||
}
|
||||
}
|
||||
|
||||
/// Park scheduler `id` until a wake, the deadline, or a positive
|
||||
/// re-check. Protocol: (1) publish own idle bit (SeqCst), (2) run
|
||||
/// `recheck` — it MUST re-read the work source with an ordering that
|
||||
/// pairs with the producer's publish (SeqCst load, or take the mutex
|
||||
/// the producer publishes under); a `true` aborts the park, (3) block.
|
||||
pub(crate) fn park(
|
||||
&self,
|
||||
id: usize,
|
||||
deadline: Option<Instant>,
|
||||
recheck: impl FnOnce() -> bool,
|
||||
) -> ParkResult {
|
||||
debug_assert!(id < self.parkers.len(), "park: scheduler id out of range");
|
||||
let bit = 1u64 << id;
|
||||
// (1) publish. AcqRel RMW: the acquire half is the handshake — if
|
||||
// this lands after a producer's mask RMW in modification order, we
|
||||
// read-from it and the producer's earlier work publication is
|
||||
// visible to the re-check below (see module docs).
|
||||
let prev = self.idle.fetch_or(bit, Ordering::AcqRel);
|
||||
debug_assert_eq!(prev & bit, 0, "park: idle bit already set for this id");
|
||||
// Fence half of the producer handshake (see `wake_one_if_idle`):
|
||||
// orders our bit-publish before the re-check's loads, so it pairs
|
||||
// with the producer's publish→fence→mask-load — at least one side
|
||||
// must see the other's store, whichever queue backend is in play.
|
||||
fence(Ordering::SeqCst);
|
||||
// (2) the mandatory post-publish re-check.
|
||||
if recheck() {
|
||||
self.idle.fetch_and(!bit, Ordering::AcqRel);
|
||||
return ParkResult::WorkFound;
|
||||
}
|
||||
// (3) block. The permit closes the window between the re-check and
|
||||
// the futex wait: a wake_one that picked us in that window has
|
||||
// already CASed our bit clear and set the permit.
|
||||
let woken = self.parkers[id].park(deadline);
|
||||
// Clear own bit — a no-op when a waker already CASed it clear.
|
||||
self.idle.fetch_and(!bit, Ordering::AcqRel);
|
||||
if woken {
|
||||
ParkResult::Woken
|
||||
} else {
|
||||
ParkResult::TimedOut
|
||||
}
|
||||
}
|
||||
|
||||
/// Wake exactly one parked scheduler, if any: pick the highest set idle
|
||||
/// bit (LIFO-ish — warmest cache), CAS it clear, deliver a permit.
|
||||
/// Empty mask = no-op (everyone is awake and will find work by
|
||||
/// popping). Returns whether a scheduler was woken.
|
||||
pub(crate) fn wake_one(&self) -> bool {
|
||||
// RMW read, not a load: reads the latest mask by modification-order
|
||||
// coherence, closing the Dekker race with a parking consumer (see
|
||||
// module docs). The release side of the RMW is what a later-parking
|
||||
// consumer's fetch_or acquires to make its re-check sound.
|
||||
let mut mask = self.idle.fetch_or(0, Ordering::AcqRel);
|
||||
loop {
|
||||
if mask == 0 {
|
||||
return false;
|
||||
}
|
||||
let id = 63 - mask.leading_zeros() as usize; // highest set bit
|
||||
let bit = 1u64 << id;
|
||||
// The CAS is the exactly-one guarantee: whoever clears the bit
|
||||
// owns the wake; a racing wake_one retries on the observed value
|
||||
// (coherence: a failed CAS can never read older than `mask`).
|
||||
match self.idle.compare_exchange(
|
||||
mask,
|
||||
mask & !bit,
|
||||
Ordering::AcqRel,
|
||||
Ordering::Acquire,
|
||||
) {
|
||||
Ok(_) => {
|
||||
self.parkers[id].unpark();
|
||||
return true;
|
||||
}
|
||||
Err(m) => mask = m,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The producer-side wake tail (`enqueue`'s fast path, RFC 018 "enqueue
|
||||
/// wakes"). The caller has just published work (queue push); we fence,
|
||||
/// then read the mask with ONE Relaxed load — 0 means every scheduler
|
||||
/// is awake and the pure-compute hot path pays no RMW on the shared
|
||||
/// mask line. Soundness is the fence handshake (module docs): our
|
||||
/// fence orders the caller's push before the mask load; the consumer's
|
||||
/// fence (in `park`) orders its bit-publish before its re-check — at
|
||||
/// least one side must observe the other's store, so a consumer we
|
||||
/// miss here is a consumer whose re-check finds the caller's work.
|
||||
pub(crate) fn wake_one_if_idle(&self) -> bool {
|
||||
fence(Ordering::SeqCst);
|
||||
if self.idle.load(Ordering::Relaxed) == 0 {
|
||||
return false;
|
||||
}
|
||||
self.wake_one()
|
||||
}
|
||||
|
||||
/// Terminal wake (replaces the AllDone wake-pipe byte): clear the mask
|
||||
/// and deliver a permit to *every* parker, parked or not. A permit set
|
||||
/// on a busy scheduler costs one spurious park return — nothing at the
|
||||
/// terminal boundary. Idempotent.
|
||||
pub(crate) fn wake_all(&self) {
|
||||
self.idle.store(0, Ordering::Release);
|
||||
for p in self.parkers.iter() {
|
||||
p.unpark();
|
||||
}
|
||||
}
|
||||
|
||||
/// Latest idle mask (RMW read — same handshake as `wake_one`). A
|
||||
/// test-only observer: production expresses the chain rule through
|
||||
/// `wake_one_if_idle` (fence + Relaxed load), not a mask read.
|
||||
#[cfg(test)]
|
||||
pub(crate) fn idle_mask(&self) -> u64 {
|
||||
self.idle.fetch_or(0, Ordering::AcqRel)
|
||||
}
|
||||
|
||||
// ----- timekeeper -----
|
||||
|
||||
/// Try to take the timekeeper role for scheduler `id` with `deadline`.
|
||||
/// MUST be called under the caller's timer serialization (the timers
|
||||
/// mutex), with `deadline` the heap minimum peeked under that same
|
||||
/// hold — this is what makes the insert-check race-free. Returns
|
||||
/// whether the role was taken (false = someone else holds it; park
|
||||
/// with no deadline).
|
||||
pub(crate) fn try_arm_timer(&self, id: usize, deadline: Instant) -> bool {
|
||||
debug_assert!(id < self.parkers.len(), "try_arm_timer: id out of range");
|
||||
if self
|
||||
.tk_holder
|
||||
.compare_exchange(NO_TIMEKEEPER, id as u64, Ordering::SeqCst, Ordering::SeqCst)
|
||||
.is_err()
|
||||
{
|
||||
return false;
|
||||
}
|
||||
// Holder-then-deadline order: an insert-check that sees the holder
|
||||
// with the deadline still NO_DEADLINE compares `new < MAX` = true
|
||||
// and over-wakes — the benign direction. (Under the mandated timer
|
||||
// serialization this interleaving cannot occur anyway.)
|
||||
self.tk_armed.store(self.deadline_nanos(deadline), Ordering::SeqCst);
|
||||
true
|
||||
}
|
||||
|
||||
/// Release the timekeeper role (the holder, on wake, before it
|
||||
/// re-peeks/fires). Callable without the timer serialization: a
|
||||
/// racing insert may wake a no-longer-holder — over-wake, benign.
|
||||
pub(crate) fn disarm_timer(&self, id: usize) {
|
||||
debug_assert_eq!(
|
||||
self.tk_holder.load(Ordering::SeqCst),
|
||||
id as u64,
|
||||
"disarm_timer by a non-holder"
|
||||
);
|
||||
// Deadline first: a concurrent insert-check then sees NO_DEADLINE
|
||||
// and skips the (now pointless) wake instead of waking a stale
|
||||
// holder id. Either order is correct; this one wastes less.
|
||||
self.tk_armed.store(NO_DEADLINE, Ordering::SeqCst);
|
||||
self.tk_holder.store(NO_TIMEKEEPER, Ordering::SeqCst);
|
||||
}
|
||||
|
||||
/// Insert-side re-arm check: if `deadline` is earlier than the armed
|
||||
/// snapshot, wake the timekeeper to re-peek. MUST be called under the
|
||||
/// same timer serialization as `try_arm_timer` (see there).
|
||||
pub(crate) fn timer_inserted(&self, deadline: Instant) {
|
||||
if self.deadline_nanos(deadline) < self.tk_armed.load(Ordering::SeqCst) {
|
||||
let holder = self.tk_holder.load(Ordering::SeqCst);
|
||||
if holder != NO_TIMEKEEPER {
|
||||
// Direct unpark, not wake_one: the wake targets the
|
||||
// timekeeper specifically (it must re-peek the heap). Its
|
||||
// idle bit stays set until it returns from park — a
|
||||
// concurrent wake_one may pick it too; over-wake, benign.
|
||||
self.parkers[holder as usize].unpark();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The armed-deadline snapshot (nanos since origin; `NO_DEADLINE` =
|
||||
/// none). Test-only introspection on the timekeeper's armed value; the
|
||||
/// busy-path due-check reads `next_deadline`, not this.
|
||||
#[cfg(test)]
|
||||
pub(crate) fn armed_deadline_nanos(&self) -> u64 {
|
||||
self.tk_armed.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
// ----- earliest-deadline snapshot (busy-path due-check) -----
|
||||
|
||||
/// Record a newly inserted timer deadline. MUST be called under the
|
||||
/// timers mutex (same serialization rule as `try_arm_timer`), which is
|
||||
/// why plain compare+store suffices for the min-maintenance. Also runs
|
||||
/// the timekeeper re-arm check (`timer_inserted`) — one call site for
|
||||
/// both consequences of an insert.
|
||||
pub(crate) fn note_deadline(&self, deadline: Instant) {
|
||||
let n = self.deadline_nanos(deadline);
|
||||
if n < self.next_deadline.load(Ordering::Relaxed) {
|
||||
self.next_deadline.store(n, Ordering::Release);
|
||||
}
|
||||
self.timer_inserted(deadline);
|
||||
}
|
||||
|
||||
/// Re-anchor the snapshot to the heap minimum (`None` = heap empty)
|
||||
/// after a `pop_due` / `clear`. MUST be called under the timers mutex.
|
||||
pub(crate) fn refresh_deadline(&self, next: Option<Instant>) {
|
||||
let n = next.map_or(NO_DEADLINE, |d| self.deadline_nanos(d));
|
||||
self.next_deadline.store(n, Ordering::Release);
|
||||
}
|
||||
|
||||
/// Busy-path due-check: is the earliest known deadline at or past
|
||||
/// `now`? One Relaxed load when no deadline is armed — the clock is
|
||||
/// read only when a timer actually exists (matching the old drain
|
||||
/// phase's is_empty guard), so the pure-compute hot path pays a load
|
||||
/// and a branch.
|
||||
pub(crate) fn deadline_due(&self) -> bool {
|
||||
let n = self.next_deadline.load(Ordering::Relaxed);
|
||||
n != NO_DEADLINE && self.deadline_nanos(Instant::now()) >= n
|
||||
}
|
||||
|
||||
/// The earliest-deadline snapshot as an `Instant` (`None` = no timer
|
||||
/// pending). Test-only: the idle path arms the timekeeper from the
|
||||
/// timer heap's own `peek_deadline` under the timers mutex.
|
||||
#[cfg(test)]
|
||||
pub(crate) fn next_deadline_instant(&self) -> Option<Instant> {
|
||||
let n = self.next_deadline.load(Ordering::Acquire);
|
||||
if n == NO_DEADLINE {
|
||||
None
|
||||
} else {
|
||||
self.origin.checked_add(std::time::Duration::from_nanos(n))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Unit tests (std build)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
#[cfg(all(test, not(loom)))]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering as O};
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
#[test]
|
||||
fn permit_before_park_returns_immediately() {
|
||||
let c = Coordinator::new(1);
|
||||
// Deliver the wake first (nobody parked: wake_one no-ops on the
|
||||
// mask, so use the timekeeper-direct path? No — permit semantics
|
||||
// are the parker's own; exercise via wake_all which permits all).
|
||||
c.wake_all();
|
||||
let t0 = Instant::now();
|
||||
let r = c.park(0, None, || false);
|
||||
assert_eq!(r, ParkResult::Woken);
|
||||
assert!(t0.elapsed() < Duration::from_millis(100), "park blocked despite permit");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn submillisecond_deadline_is_honored() {
|
||||
// Regression for the retired as_millis truncation: a 500µs deadline
|
||||
// must neither busy-return instantly forever nor round to 0/∞.
|
||||
let c = Coordinator::new(1);
|
||||
let t0 = Instant::now();
|
||||
let r = c.park(0, Some(t0 + Duration::from_micros(500)), || false);
|
||||
let dt = t0.elapsed();
|
||||
assert_eq!(r, ParkResult::TimedOut);
|
||||
assert!(dt >= Duration::from_micros(400), "woke too early: {dt:?}");
|
||||
assert!(dt < Duration::from_millis(50), "overslept: {dt:?}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn recheck_true_aborts_park_and_clears_bit() {
|
||||
let c = Coordinator::new(2);
|
||||
let r = c.park(1, None, || true);
|
||||
assert_eq!(r, ParkResult::WorkFound);
|
||||
assert_eq!(c.idle_mask(), 0, "bit not cleared after WorkFound");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn recheck_observes_own_bit_published() {
|
||||
let c = Coordinator::new(2);
|
||||
let seen = std::cell::Cell::new(0u64);
|
||||
let r = c.park(1, None, || {
|
||||
seen.set(c.idle_mask());
|
||||
true
|
||||
});
|
||||
assert_eq!(r, ParkResult::WorkFound);
|
||||
assert_eq!(seen.get() & 0b10, 0b10, "bit not published before re-check");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn wake_one_wakes_exactly_one_of_n() {
|
||||
const N: usize = 4;
|
||||
let c = Arc::new(Coordinator::new(N));
|
||||
let woken = Arc::new(AtomicUsize::new(0));
|
||||
let mut ts = Vec::new();
|
||||
for id in 0..N {
|
||||
let c = c.clone();
|
||||
let woken = woken.clone();
|
||||
ts.push(std::thread::spawn(move || {
|
||||
let r = c.park(id, None, || false);
|
||||
assert_eq!(r, ParkResult::Woken);
|
||||
woken.fetch_add(1, O::SeqCst);
|
||||
}));
|
||||
}
|
||||
// Wait until all four are published idle.
|
||||
let t0 = Instant::now();
|
||||
while c.idle_mask().count_ones() != N as u32 {
|
||||
assert!(t0.elapsed() < Duration::from_secs(5), "threads never parked");
|
||||
std::thread::yield_now();
|
||||
}
|
||||
assert!(c.wake_one());
|
||||
// Exactly one wakes; give the others a beat to (incorrectly) wake.
|
||||
let t0 = Instant::now();
|
||||
while woken.load(O::SeqCst) == 0 {
|
||||
assert!(t0.elapsed() < Duration::from_secs(5), "wake_one woke nobody");
|
||||
std::thread::yield_now();
|
||||
}
|
||||
std::thread::sleep(Duration::from_millis(100));
|
||||
assert_eq!(woken.load(O::SeqCst), 1, "wake_one woke more than one");
|
||||
assert_eq!(c.idle_mask().count_ones(), (N - 1) as u32);
|
||||
c.wake_all();
|
||||
for t in ts {
|
||||
match t.join() {
|
||||
Ok(()) => {}
|
||||
Err(p) => std::panic::resume_unwind(p),
|
||||
}
|
||||
}
|
||||
assert_eq!(woken.load(O::SeqCst), N);
|
||||
assert_eq!(c.idle_mask(), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn wake_one_prefers_highest_bit() {
|
||||
let c = Arc::new(Coordinator::new(3));
|
||||
let woken_id = Arc::new(AtomicUsize::new(usize::MAX));
|
||||
let mut ts = Vec::new();
|
||||
for id in 0..3 {
|
||||
let c = c.clone();
|
||||
let woken_id = woken_id.clone();
|
||||
ts.push(std::thread::spawn(move || {
|
||||
if c.park(id, None, || false) == ParkResult::Woken {
|
||||
let _ = woken_id.compare_exchange(usize::MAX, id, O::SeqCst, O::SeqCst);
|
||||
}
|
||||
}));
|
||||
}
|
||||
let t0 = Instant::now();
|
||||
while c.idle_mask() != 0b111 {
|
||||
assert!(t0.elapsed() < Duration::from_secs(5));
|
||||
std::thread::yield_now();
|
||||
}
|
||||
assert!(c.wake_one());
|
||||
let t0 = Instant::now();
|
||||
while woken_id.load(O::SeqCst) == usize::MAX {
|
||||
assert!(t0.elapsed() < Duration::from_secs(5));
|
||||
std::thread::yield_now();
|
||||
}
|
||||
assert_eq!(woken_id.load(O::SeqCst), 2, "LIFO-ish: highest bit first");
|
||||
c.wake_all();
|
||||
for t in ts {
|
||||
let _ = t.join();
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn wake_one_on_empty_mask_is_noop() {
|
||||
let c = Coordinator::new(2);
|
||||
assert!(!c.wake_one());
|
||||
assert_eq!(c.idle_mask(), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn timekeeper_arm_is_exclusive_and_snapshot_readable() {
|
||||
let c = Coordinator::new(2);
|
||||
let d2 = Instant::now() + Duration::from_secs(10);
|
||||
let d1 = Instant::now() + Duration::from_secs(1);
|
||||
assert_eq!(c.armed_deadline_nanos(), NO_DEADLINE);
|
||||
assert!(c.try_arm_timer(0, d2));
|
||||
assert!(!c.try_arm_timer(1, d1), "second arm must fail while held");
|
||||
assert_eq!(c.armed_deadline_nanos(), c.deadline_nanos(d2));
|
||||
c.disarm_timer(0);
|
||||
assert_eq!(c.armed_deadline_nanos(), NO_DEADLINE);
|
||||
assert!(c.try_arm_timer(1, d1), "role must be re-takeable after disarm");
|
||||
c.disarm_timer(1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn earlier_insert_wakes_timekeeper() {
|
||||
let c = Coordinator::new(2);
|
||||
let far = Instant::now() + Duration::from_secs(60);
|
||||
let near = Instant::now() + Duration::from_millis(1);
|
||||
assert!(c.try_arm_timer(0, far));
|
||||
// Holder not yet parked: the wake must land as a permit.
|
||||
c.timer_inserted(near);
|
||||
let t0 = Instant::now();
|
||||
let r = c.park(0, Some(far), || false);
|
||||
assert_eq!(r, ParkResult::Woken, "re-arm wake lost");
|
||||
assert!(t0.elapsed() < Duration::from_secs(5), "slept toward the stale deadline");
|
||||
c.disarm_timer(0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn later_insert_does_not_wake_timekeeper() {
|
||||
let c = Coordinator::new(2);
|
||||
let near = Instant::now() + Duration::from_millis(20);
|
||||
let far = Instant::now() + Duration::from_secs(60);
|
||||
assert!(c.try_arm_timer(0, near));
|
||||
c.timer_inserted(far); // later than armed: no wake
|
||||
let r = c.park(0, Some(near), || false);
|
||||
assert_eq!(r, ParkResult::TimedOut, "spurious wake for a later insert");
|
||||
c.disarm_timer(0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[should_panic(expected = "1..=64")]
|
||||
fn more_than_64_schedulers_asserts() {
|
||||
let _ = Coordinator::new(65);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn wake_one_if_idle_noop_on_empty_and_wakes_on_parked() {
|
||||
let c = Arc::new(Coordinator::new(1));
|
||||
assert!(!c.wake_one_if_idle(), "empty mask must be a no-op");
|
||||
let c2 = c.clone();
|
||||
let t = std::thread::spawn(move || {
|
||||
assert_eq!(c2.park(0, None, || false), ParkResult::Woken);
|
||||
});
|
||||
let t0 = Instant::now();
|
||||
while c.idle_mask() == 0 {
|
||||
assert!(t0.elapsed() < Duration::from_secs(5), "never parked");
|
||||
std::thread::yield_now();
|
||||
}
|
||||
assert!(c.wake_one_if_idle());
|
||||
match t.join() {
|
||||
Ok(()) => {}
|
||||
Err(p) => std::panic::resume_unwind(p),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn deadline_snapshot_min_maintenance_and_due_check() {
|
||||
let c = Coordinator::new(1);
|
||||
assert!(!c.deadline_due(), "no deadline: never due");
|
||||
assert_eq!(c.next_deadline_instant(), None);
|
||||
let far = Instant::now() + Duration::from_secs(60);
|
||||
let near = Instant::now() + Duration::from_millis(1);
|
||||
c.note_deadline(far);
|
||||
assert!(!c.deadline_due());
|
||||
c.note_deadline(near); // min wins
|
||||
assert!(c.next_deadline_instant().is_some_and(|d| d <= near));
|
||||
c.note_deadline(far); // later insert must NOT raise the snapshot
|
||||
assert!(c.next_deadline_instant().is_some_and(|d| d <= near));
|
||||
std::thread::sleep(Duration::from_millis(2));
|
||||
assert!(c.deadline_due(), "past deadline not reported due");
|
||||
c.refresh_deadline(Some(far));
|
||||
assert!(!c.deadline_due(), "refresh did not re-anchor");
|
||||
c.refresh_deadline(None);
|
||||
assert_eq!(c.next_deadline_instant(), None);
|
||||
// A deadline at/before origin encodes as 0: always due.
|
||||
c.note_deadline(Instant::now() - Duration::from_secs(1));
|
||||
assert!(c.deadline_due());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn note_deadline_earlier_wakes_timekeeper_via_snapshot_path() {
|
||||
// note_deadline must carry the timer_inserted re-arm wake too.
|
||||
let c = Coordinator::new(2);
|
||||
let far = Instant::now() + Duration::from_secs(60);
|
||||
let near = Instant::now() + Duration::from_millis(1);
|
||||
assert!(c.try_arm_timer(0, far));
|
||||
c.note_deadline(near);
|
||||
let r = c.park(0, Some(far), || false);
|
||||
assert_eq!(r, ParkResult::Woken, "re-arm wake lost through note_deadline");
|
||||
c.disarm_timer(0);
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// loom models — RUSTFLAGS="--cfg loom" cargo test --lib --release park
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
#[cfg(all(test, loom))]
|
||||
mod loom_tests {
|
||||
use super::*;
|
||||
use loom::sync::atomic::{AtomicU64 as LAtomicU64, Ordering as O};
|
||||
use loom::sync::{Arc, Mutex as LMutex};
|
||||
use loom::thread;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
/// RFC 018 loom model 1 — no lost wake.
|
||||
/// producer{publish item; wake_one} ∥ consumer{set bit; re-check; park}:
|
||||
/// the consumer always observes the item or a permit; it can never
|
||||
/// sleep past a published item (loom's deadlock detector is the
|
||||
/// assertion — a consumer parked forever fails the model).
|
||||
#[test]
|
||||
fn no_lost_wake() {
|
||||
loom::model(|| {
|
||||
let c = Arc::new(Coordinator::new(1));
|
||||
let item = Arc::new(LAtomicU64::new(0));
|
||||
|
||||
let prod = {
|
||||
let c = c.clone();
|
||||
let item = item.clone();
|
||||
thread::spawn(move || {
|
||||
// The enqueue shape: publish, then the fenced fast-path
|
||||
// wake tail (this is what the runtime's enqueue calls).
|
||||
item.store(1, O::SeqCst);
|
||||
c.wake_one_if_idle();
|
||||
})
|
||||
};
|
||||
|
||||
// Consumer: loop until the item is popped. Guarded load+CAS,
|
||||
// not a blind swap — a swap writes 0 even when empty, and
|
||||
// coherence allows that write to land after the producer's
|
||||
// store in modification order, destroying the item (a model
|
||||
// bug loom caught in an earlier draft of this test).
|
||||
loop {
|
||||
if item.load(O::SeqCst) == 1
|
||||
&& item.compare_exchange(1, 0, O::SeqCst, O::SeqCst).is_ok()
|
||||
{
|
||||
break;
|
||||
}
|
||||
let _ = c.park(0, None, || item.load(O::SeqCst) == 1);
|
||||
}
|
||||
prod.join().unwrap();
|
||||
});
|
||||
}
|
||||
|
||||
/// RFC 018 loom model 2 — chain propagation.
|
||||
/// Two items, two sleepers, ONE producer wake: the chain rule (a woken
|
||||
/// consumer that sees surplus work and a non-empty mask wakes again)
|
||||
/// must get both items consumed with no further producer action.
|
||||
#[test]
|
||||
fn chain_propagation() {
|
||||
loom::model(|| {
|
||||
let c = Arc::new(Coordinator::new(2));
|
||||
let items = Arc::new(LAtomicU64::new(0));
|
||||
let consumed = Arc::new(LAtomicU64::new(0));
|
||||
|
||||
let mut hs = Vec::new();
|
||||
for id in 0..2usize {
|
||||
let c = c.clone();
|
||||
let items = items.clone();
|
||||
let consumed = consumed.clone();
|
||||
hs.push(thread::spawn(move || loop {
|
||||
if consumed.load(O::SeqCst) == 2 {
|
||||
c.wake_all(); // release a sibling still parked
|
||||
break;
|
||||
}
|
||||
let cur = items.load(O::SeqCst);
|
||||
if cur > 0
|
||||
&& items
|
||||
.compare_exchange(cur, cur - 1, O::SeqCst, O::SeqCst)
|
||||
.is_ok()
|
||||
{
|
||||
consumed.fetch_add(1, O::SeqCst);
|
||||
// THE CHAIN RULE, exactly as production expresses it
|
||||
// (runtime.rs schedule_loop): surplus ⇒ the fenced
|
||||
// fast-path wake. Models the Relaxed-load chain path,
|
||||
// not just the RMW one.
|
||||
if items.load(O::SeqCst) > 0 {
|
||||
c.wake_one_if_idle();
|
||||
}
|
||||
continue;
|
||||
}
|
||||
let _ = c.park(id, None, || {
|
||||
items.load(O::SeqCst) > 0 || consumed.load(O::SeqCst) == 2
|
||||
});
|
||||
}));
|
||||
}
|
||||
|
||||
// Producer (main): two items, ONE wake, via the enqueue-shaped
|
||||
// fenced fast path.
|
||||
items.store(2, O::SeqCst);
|
||||
c.wake_one_if_idle();
|
||||
|
||||
for h in hs {
|
||||
h.join().unwrap();
|
||||
}
|
||||
assert_eq!(consumed.load(O::SeqCst), 2);
|
||||
assert_eq!(items.load(O::SeqCst), 0);
|
||||
});
|
||||
}
|
||||
|
||||
/// RFC 018 loom model 3 — timekeeper handoff.
|
||||
/// An earlier-deadline insert racing the parking timekeeper: the
|
||||
/// earlier deadline is always honored — either the timekeeper armed it
|
||||
/// directly (insert landed first under the timers lock) or the insert
|
||||
/// wakes the timekeeper to re-peek. A timekeeper sleeping toward the
|
||||
/// stale later deadline would deadlock the model (loom has no time).
|
||||
#[test]
|
||||
fn timekeeper_handoff() {
|
||||
loom::model(|| {
|
||||
let origin = Instant::now();
|
||||
let d_far = origin + Duration::from_secs(60);
|
||||
let d_near = origin + Duration::from_secs(1);
|
||||
|
||||
let c = Arc::new(Coordinator::new(1));
|
||||
// The timers-mutex stand-in: heap min under a lock.
|
||||
let heap_min = Arc::new(LMutex::new(d_far));
|
||||
|
||||
let tk = {
|
||||
let c = c.clone();
|
||||
let heap_min = heap_min.clone();
|
||||
thread::spawn(move || {
|
||||
// Peek + arm under the lock (the serialization rule).
|
||||
let armed = {
|
||||
let g = heap_min.lock().unwrap();
|
||||
let min = *g;
|
||||
assert!(c.try_arm_timer(0, min));
|
||||
min
|
||||
};
|
||||
if armed == d_far {
|
||||
// Insert hadn't landed: it MUST wake us. Parking
|
||||
// toward d_far with no wake = model deadlock.
|
||||
let r = c.park(0, Some(armed), || false);
|
||||
assert_eq!(r, ParkResult::Woken, "re-arm wake lost");
|
||||
}
|
||||
// Woken (or armed the near deadline directly): re-peek.
|
||||
c.disarm_timer(0);
|
||||
let g = heap_min.lock().unwrap();
|
||||
assert_eq!(*g, d_near, "earlier deadline not visible on re-peek");
|
||||
})
|
||||
};
|
||||
|
||||
// Inserter (main): publish the earlier deadline under the lock,
|
||||
// then the insert-check.
|
||||
{
|
||||
let mut g = heap_min.lock().unwrap();
|
||||
*g = d_near;
|
||||
c.timer_inserted(d_near);
|
||||
}
|
||||
|
||||
tk.join().unwrap();
|
||||
});
|
||||
}
|
||||
|
||||
/// RFC 018 loom model 4 — termination.
|
||||
/// The AllDone verdict: producer flips done and wake_all()s; consumers
|
||||
/// must never park forever past done (park's re-check + wake_all's
|
||||
/// permits close every interleaving; a stuck consumer = loom deadlock).
|
||||
#[test]
|
||||
fn termination_no_park_past_done() {
|
||||
loom::model(|| {
|
||||
let c = Arc::new(Coordinator::new(2));
|
||||
let done = Arc::new(LAtomicU64::new(0));
|
||||
|
||||
let mut hs = Vec::new();
|
||||
for id in 0..2usize {
|
||||
let c = c.clone();
|
||||
let done = done.clone();
|
||||
hs.push(thread::spawn(move || loop {
|
||||
if done.load(O::SeqCst) == 1 {
|
||||
break;
|
||||
}
|
||||
let _ = c.park(id, None, || done.load(O::SeqCst) == 1);
|
||||
}));
|
||||
}
|
||||
|
||||
done.store(1, O::SeqCst);
|
||||
c.wake_all();
|
||||
|
||||
for h in hs {
|
||||
h.join().unwrap();
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
+179
-177
@@ -86,21 +86,28 @@
|
||||
//! # Termination (counter-based)
|
||||
//!
|
||||
//! The old all-clear scanned the slot table under the big lock. Now:
|
||||
//! exit when `io_out == 0` (read *before* the queue lock, phase-1 ordering)
|
||||
//! and, under the queue lock, the queue is empty and `live_actors == 0`.
|
||||
//! `live_actors` is incremented in `spawn` before the enqueue and decremented
|
||||
//! at the very END of `finalize_actor`, strictly after every wakeup that
|
||||
//! finalize produces has been enqueued. The soundness crux: any enqueue
|
||||
//! targets a live (not-yet-finalized) actor, so `live == 0` implies no wakeup
|
||||
//! can still be in flight; combined with "spawner is itself live", observing
|
||||
//! `(queue empty, live == 0)` under the queue lock means no work can ever
|
||||
//! appear again.
|
||||
//! exit when `io_outstanding + io_fd_waiters == 0` (two Relaxed/Acquire
|
||||
//! atomic loads, read *before* the queue pop) and, under the queue lock,
|
||||
//! the queue is empty and `live_actors == 0`. `live_actors` is incremented
|
||||
//! in `spawn` before the enqueue and decremented at the very END of
|
||||
//! `finalize_actor`, strictly after every wakeup that finalize produces has
|
||||
//! been enqueued. The soundness crux: any enqueue targets a live
|
||||
//! (not-yet-finalized) actor, so `live == 0` implies no wakeup can still be
|
||||
//! in flight; combined with "spawner is itself live", observing
|
||||
//! `(queue empty, live == 0)` means no work can ever appear again.
|
||||
//!
|
||||
//! # Timer / IO drain (try-lock, one-winner)
|
||||
//! # Scheduler park/wake (RFC 018)
|
||||
//!
|
||||
//! Unchanged from phase 1: one winner per round drains due timers and IO
|
||||
//! completions from their own mutexes; wakeups go through the unpark
|
||||
//! protocol like everyone else's.
|
||||
//! Schedulers sleep on per-thread futex parkers via the coordination layer
|
||||
//! (`park.rs`), NOT on a shared wake pipe. IO backends are producers behind
|
||||
//! a two-call contract — make the actor runnable (`unpark_at`), whose
|
||||
//! `enqueue` tail wakes exactly one parked scheduler. The blocking pool and
|
||||
//! epoll thread each route their own completions (driver-enqueues); there
|
||||
//! is no shared completion queue, no drain lock, no one-winner drain phase.
|
||||
//! Timers fire two ways: a busy-path due-check every loop iteration (one
|
||||
//! Relaxed load of the earliest-deadline snapshot when no timer is armed),
|
||||
//! and the timekeeper — at most one parked scheduler holds the timer
|
||||
//! deadline, so an expiry wakes one scheduler, not a herd.
|
||||
|
||||
use crate::actor::{
|
||||
clear_current_pid, is_actor_done, reset_actor_done, set_current_actor_box,
|
||||
@@ -750,8 +757,21 @@ pub(crate) struct RuntimeInner {
|
||||
pub(crate) io: Mutex<Option<IoThread>>,
|
||||
/// Monotonic `MonitorId` source. Never reused.
|
||||
pub(crate) next_monitor_id: AtomicU64,
|
||||
/// Try-lock: exactly one scheduler thread drains timers/IO per iteration.
|
||||
drain_lock: Mutex<()>,
|
||||
/// RFC 018: the scheduler coordination layer — per-scheduler parkers,
|
||||
/// idle mask, wake protocol, timekeeper role, earliest-deadline
|
||||
/// snapshot. Arc'd because `Timers` shares it (insert-side deadline
|
||||
/// notes run under the timers mutex).
|
||||
pub(crate) coord: Arc<crate::park::Coordinator>,
|
||||
/// `block_on_io` requests in flight. Incremented by the submitter
|
||||
/// BEFORE submit (underflow-proof), decremented by the pool thread on
|
||||
/// completion. Read lock-free by the idle path's termination verdict —
|
||||
/// the per-pop `io.lock` of the drain era is gone.
|
||||
pub(crate) io_outstanding: AtomicU32,
|
||||
/// Parked fd waiters. Incremented by the registrar BEFORE
|
||||
/// `epoll_register` (rolled back on error), decremented by whoever
|
||||
/// consumes the registration (epoll thread on readiness, canceller on
|
||||
/// an unwound wait). Same lock-free verdict read as `io_outstanding`.
|
||||
pub(crate) io_fd_waiters: AtomicU32,
|
||||
/// Per-thread stats, indexed by scheduler thread slot (0..N).
|
||||
pub(crate) stats: Vec<SchedulerStats>,
|
||||
/// Global counters for RFC 000 primitives.
|
||||
@@ -802,6 +822,12 @@ impl RuntimeInner {
|
||||
let slots: Box<[Slot]> = (0..max_actors).map(|_| Slot::vacant()).collect();
|
||||
// Low indices on top of the stack so early spawns get low pids.
|
||||
let free: Vec<u32> = (0..max_actors as u32).rev().collect();
|
||||
// RFC 018: the coordination layer (asserts thread_count <= 64), and
|
||||
// the timers' hook into it — every insert under the timers mutex
|
||||
// notes its deadline (busy-path snapshot + timekeeper re-arm).
|
||||
let coord = Arc::new(crate::park::Coordinator::new(thread_count));
|
||||
let mut timers = Timers::new();
|
||||
timers.attach_coordinator(coord.clone());
|
||||
Arc::new(Self {
|
||||
run_queue: crate::run_queue::RunQueue::new(thread_count, max_actors),
|
||||
slots,
|
||||
@@ -810,10 +836,12 @@ impl RuntimeInner {
|
||||
root_bits: AtomicU64::new(u64::MAX),
|
||||
root_exited: AtomicBool::new(false),
|
||||
root_swept: AtomicBool::new(false),
|
||||
timers: Mutex::new(Timers::new()),
|
||||
timers: Mutex::new(timers),
|
||||
io: Mutex::new(None),
|
||||
next_monitor_id: AtomicU64::new(0),
|
||||
drain_lock: Mutex::new(()),
|
||||
coord,
|
||||
io_outstanding: AtomicU32::new(0),
|
||||
io_fd_waiters: AtomicU32::new(0),
|
||||
stats,
|
||||
io_parked: AtomicU32::new(0),
|
||||
sleeping: AtomicU32::new(0),
|
||||
@@ -874,6 +902,13 @@ impl RuntimeInner {
|
||||
);
|
||||
self.run_queue.push(pid);
|
||||
crate::te!(crate::trace::Event::Enqueue(pid));
|
||||
// RFC 018 enqueue wake (fixes the silent enqueue): if a scheduler
|
||||
// is parked, wake exactly one. The fast path when everyone is busy
|
||||
// is a fence + one Relaxed load of an unmodified line — the
|
||||
// pure-compute hot path pays (almost) nothing. Bias is over-wake:
|
||||
// a spurious wake costs one futex round-trip and a failed pop; a
|
||||
// missed wake would cost a stranded actor.
|
||||
self.coord.wake_one_if_idle();
|
||||
}
|
||||
|
||||
/// Make `pid` runnable if it is parked; coalesce or defer otherwise.
|
||||
@@ -1076,7 +1111,16 @@ impl Runtime {
|
||||
self.inner.live_actors.load(Ordering::Acquire), 0,
|
||||
"run() called while previous run still active"
|
||||
);
|
||||
let io_thread = match IoThread::start() {
|
||||
// RFC 018: the IO producers reach the runtime (slot table + unpark)
|
||||
// through a Weak, so no RuntimeInner → IoThread → RuntimeInner cycle
|
||||
// forms. Reset the in-flight counters BEFORE the threads can touch
|
||||
// them (a prior run left them at 0 on a clean exit; the asserts pin
|
||||
// that).
|
||||
debug_assert_eq!(self.inner.io_outstanding.load(Ordering::Acquire), 0);
|
||||
debug_assert_eq!(self.inner.io_fd_waiters.load(Ordering::Acquire), 0);
|
||||
self.inner.io_outstanding.store(0, Ordering::Release);
|
||||
self.inner.io_fd_waiters.store(0, Ordering::Release);
|
||||
let io_thread = match IoThread::start(Arc::downgrade(&self.inner)) {
|
||||
Ok(io) => io,
|
||||
Err(e) => panic!("failed to start IO thread: {e}"),
|
||||
};
|
||||
@@ -1179,6 +1223,8 @@ impl Runtime {
|
||||
}
|
||||
self.inner.io_parked.store(0, Ordering::Relaxed);
|
||||
self.inner.sleeping.store(0, Ordering::Relaxed);
|
||||
self.inner.io_outstanding.store(0, Ordering::Relaxed);
|
||||
self.inner.io_fd_waiters.store(0, Ordering::Relaxed);
|
||||
|
||||
RUNTIME.with(|r| *r.borrow_mut() = None);
|
||||
|
||||
@@ -1494,53 +1540,30 @@ fn stop_live_actors(inner: &Arc<RuntimeInner>) {
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// schedule_loop — runs on each scheduler OS thread
|
||||
// Timer firing — shared by the busy-path due-check and the timekeeper
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
crate::preempt::configure_preempt(inner.alloc_interval, inner.timeslice_cycles);
|
||||
let stats = &inner.stats[slot_idx];
|
||||
|
||||
loop {
|
||||
// ----------------------------------------------------------------
|
||||
// 1. Try to win the drain lock (timers + IO). One winner per round;
|
||||
// losers skip immediately and proceed to step 2.
|
||||
// ----------------------------------------------------------------
|
||||
if let Ok(_drain_guard) = inner.drain_lock.try_lock() {
|
||||
// Timers and IO live behind their own mutexes (phase 1), so the
|
||||
// pure-yield / pure-compute hot path never contends a global lock
|
||||
// just to discover there is nothing to drain. The clock is read
|
||||
// only when the timer heap is non-empty.
|
||||
let due = {
|
||||
let mut t = match inner.timers.lock() {
|
||||
Ok(t) => t,
|
||||
/// Pop and dispatch every due timer. `pop_due` re-anchors the
|
||||
/// earliest-deadline snapshot under the timers mutex before returning, so
|
||||
/// a caller that raced a concurrent insert simply comes back on the next
|
||||
/// due-check. Dispatch runs with the timers lock released.
|
||||
fn fire_due_timers(inner: &Arc<RuntimeInner>, try_only: bool) {
|
||||
let due = if try_only {
|
||||
// Busy path: if another scheduler is already in the timers mutex
|
||||
// (firing, inserting, or peeking) skip — the snapshot stays due
|
||||
// until someone actually pops, so the check re-fires next loop.
|
||||
match inner.timers.try_lock() {
|
||||
Ok(mut t) => t.pop_due(std::time::Instant::now()),
|
||||
Err(std::sync::TryLockError::WouldBlock) => return,
|
||||
Err(std::sync::TryLockError::Poisoned(e)) => {
|
||||
panic!("smarm: timers lock poisoned (core corrupt): {e}")
|
||||
}
|
||||
}
|
||||
} else {
|
||||
match inner.timers.lock() {
|
||||
Ok(mut t) => t.pop_due(std::time::Instant::now()),
|
||||
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
|
||||
};
|
||||
if t.is_empty() {
|
||||
Vec::new()
|
||||
} else {
|
||||
t.pop_due(std::time::Instant::now())
|
||||
}
|
||||
};
|
||||
let completions = match inner.io.lock() {
|
||||
Ok(mut io) => io
|
||||
.as_mut()
|
||||
.map(|io| {
|
||||
// Consume wake-pipe bytes ONLY here, under the drain
|
||||
// lock and strictly before draining completions.
|
||||
// Producers push their completion before writing the
|
||||
// byte, so every byte consumed here has its completion
|
||||
// visible to the drain below. Consuming bytes anywhere
|
||||
// else — in particular after an idle poll, outside the
|
||||
// lock — loses wakeups: a try_lock loser can eat the
|
||||
// byte for a completion the winner never saw, leaving
|
||||
// it stranded (and its EPOLLONESHOT fd disarmed) until
|
||||
// an unrelated timer forces another drain pass.
|
||||
crate::io::drain_wake_pipe(io.wake_fd());
|
||||
io.drain_completions()
|
||||
})
|
||||
.unwrap_or_default(),
|
||||
Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"),
|
||||
};
|
||||
for entry in due {
|
||||
match entry.reason {
|
||||
@@ -1549,9 +1572,7 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
// actor is between `timers.insert_sleep` and
|
||||
// `park_current`; RunningNotified makes the upcoming park
|
||||
// re-queue), or gone (no-op).
|
||||
crate::timer::Reason::Sleep { epoch } => {
|
||||
inner.unpark_at(entry.pid, epoch)
|
||||
}
|
||||
crate::timer::Reason::Sleep { epoch } => inner.unpark_at(entry.pid, epoch),
|
||||
crate::timer::Reason::WaitTimeout { target, epoch } => {
|
||||
// The callback may call unpark_at itself.
|
||||
target.on_timeout(entry.pid, epoch);
|
||||
@@ -1566,57 +1587,26 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
crate::timer::Reason::Send { fire } => fire(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for completion in completions {
|
||||
match completion {
|
||||
crate::io::Completion::Blocking { pid, epoch, result } => {
|
||||
match inner.io.lock() {
|
||||
Ok(mut io) => {
|
||||
if let Some(io) = io.as_mut() {
|
||||
io.outstanding = io.outstanding.saturating_sub(1);
|
||||
// ---------------------------------------------------------------------------
|
||||
// schedule_loop — runs on each scheduler OS thread
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
crate::preempt::configure_preempt(inner.alloc_interval, inner.timeslice_cycles);
|
||||
let stats = &inner.stats[slot_idx];
|
||||
|
||||
loop {
|
||||
// ----------------------------------------------------------------
|
||||
// 1. Busy-path timer due-check (RFC 018 design point (a)): under
|
||||
// saturation nobody parks, so no timekeeper exists — due timers
|
||||
// must still fire. One Relaxed load + branch when no timer is
|
||||
// armed; the clock is read only when one is.
|
||||
// ----------------------------------------------------------------
|
||||
if inner.coord.deadline_due() {
|
||||
fire_due_timers(inner, true);
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
panic!("smarm: io lock poisoned (core corrupt): {e}")
|
||||
}
|
||||
}
|
||||
// Stash the result under the cold lock, then unpark.
|
||||
// The protocol also covers the submit→park window
|
||||
// (RunningNotified), which the old code missed for
|
||||
// Blocking completions — a latent lost wakeup.
|
||||
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);
|
||||
} else {
|
||||
// Actor died (stopped) with the op in
|
||||
// flight; discard the result.
|
||||
}
|
||||
}
|
||||
inner.unpark_at(pid, epoch);
|
||||
}
|
||||
}
|
||||
crate::io::Completion::FdReady { fd, events: _ } => {
|
||||
// Resolve the parked pid under the io lock, then wake
|
||||
// through the protocol. Lock order: io before all.
|
||||
let parked = match inner.io.lock() {
|
||||
Ok(mut io) => io.as_mut().and_then(|io| {
|
||||
let entry = io.waiters.remove(&fd);
|
||||
io.epoll_deregister(fd);
|
||||
entry
|
||||
}),
|
||||
Err(e) => {
|
||||
panic!("smarm: io lock poisoned (core corrupt): {e}")
|
||||
}
|
||||
};
|
||||
if let Some((pid, epoch)) = parked {
|
||||
inner.unpark_at(pid, epoch);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} // drain_guard drops here
|
||||
|
||||
// ----------------------------------------------------------------
|
||||
// 2. Pop a runnable pid. Pop order (RFC 005): wake slot first, then
|
||||
@@ -1625,7 +1615,7 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
// ----------------------------------------------------------------
|
||||
enum Pop {
|
||||
Got(Pid),
|
||||
Idle { io_outstanding: u32, wake_fd: Option<std::os::fd::RawFd> },
|
||||
Idle,
|
||||
AllDone,
|
||||
/// Root has exited and nothing is runnable: stop the parked-forever
|
||||
/// remainder, then re-pop. Fires at most once per run.
|
||||
@@ -1650,19 +1640,12 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
crate::te!(crate::trace::Event::SlotPop(pid));
|
||||
pid
|
||||
} else {
|
||||
// Read IO liveness BEFORE the queue lock (phase-1 ordering: a
|
||||
// completion resurrects an actor only via the drain path, whose
|
||||
// enqueue would be visible under the queue lock we take next).
|
||||
let (io_out, io_fd) = {
|
||||
let io = match inner.io.lock() {
|
||||
Ok(io) => io,
|
||||
Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"),
|
||||
};
|
||||
match io.as_ref() {
|
||||
Some(io) => (io.outstanding + io.waiters.len() as u32, Some(io.wake_fd())),
|
||||
None => (0, None),
|
||||
}
|
||||
};
|
||||
// Read IO liveness BEFORE the queue pop — two atomic loads now
|
||||
// (RFC 018), not a per-pop `io.lock`: a completion resurrects
|
||||
// an actor via the producer's own unpark→enqueue, whose entry
|
||||
// would be visible to the pop below.
|
||||
let io_out = inner.io_outstanding.load(Ordering::Acquire)
|
||||
+ inner.io_fd_waiters.load(Ordering::Acquire);
|
||||
|
||||
stats.run_queue_len.store(inner.run_queue.len(), Ordering::Relaxed);
|
||||
let pop = match inner.run_queue.pop() {
|
||||
@@ -1694,7 +1677,7 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
// the idle wait below on the next pass.
|
||||
Pop::RootDrain
|
||||
} else {
|
||||
Pop::Idle { io_outstanding: io_out, wake_fd: io_fd }
|
||||
Pop::Idle
|
||||
}
|
||||
}
|
||||
};
|
||||
@@ -1709,22 +1692,15 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
Ok(mut timers) => timers.clear(),
|
||||
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
|
||||
}
|
||||
// Terminal wake: a sibling scheduler may be blocked in its
|
||||
// idle wait on a snapshot that is now terminally stale — an
|
||||
// orphaned long deadline (it would sleep it out in full) or
|
||||
// a stale `io_outstanding > 0` from a stop-cancelled waiter
|
||||
// (it would block in poll(-1) forever; cancellation produces
|
||||
// no completion, so nothing else writes the wake pipe).
|
||||
// One byte wakes every poller; each re-runs the verdict,
|
||||
// reaches AllDone itself, and re-wakes — idempotent.
|
||||
match inner.io.lock() {
|
||||
Ok(io) => {
|
||||
if let Some(io) = io.as_ref() {
|
||||
io.wake();
|
||||
}
|
||||
}
|
||||
Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"),
|
||||
}
|
||||
// Terminal wake (replaces the wake-pipe byte): a sibling
|
||||
// may be parked on a snapshot that is now terminally
|
||||
// stale — an orphaned long deadline, or a stale
|
||||
// `io_fd_waiters > 0` from a stop-cancelled waiter
|
||||
// (cancellation produces no completion, so nothing else
|
||||
// will ever wake it). `wake_all` permits every parker;
|
||||
// each sibling re-runs the verdict, reaches AllDone
|
||||
// itself, and re-wakes — idempotent.
|
||||
inner.coord.wake_all();
|
||||
return;
|
||||
}
|
||||
Pop::RootDrain => {
|
||||
@@ -1734,40 +1710,51 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
stop_live_actors(inner);
|
||||
continue;
|
||||
}
|
||||
Pop::Idle { io_outstanding, wake_fd } => {
|
||||
// Something is still in flight. Sleep on the appropriate
|
||||
// source to avoid hammering the queue mutex; retry on wake.
|
||||
let next_deadline = match inner.timers.lock() {
|
||||
Ok(timers) => timers.peek_deadline(),
|
||||
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
|
||||
Pop::Idle => {
|
||||
// Something is still in flight. Park on our own futex
|
||||
// until a producer wakes us (enqueue tail), a deadline
|
||||
// passes, or the re-check finds the world changed.
|
||||
//
|
||||
// Timekeeper (RFC 018): at most one parked scheduler
|
||||
// holds the timer deadline — the first idler to arm it
|
||||
// parks with a timeout, the rest park indefinitely, so a
|
||||
// timer expiry wakes one scheduler, not a herd. Peek and
|
||||
// arm under the timers mutex (the serialization that
|
||||
// makes the insert-side re-arm race-free).
|
||||
let tk_deadline = {
|
||||
let timers = match inner.timers.lock() {
|
||||
Ok(t) => t,
|
||||
Err(e) => {
|
||||
panic!("smarm: timers lock poisoned (core corrupt): {e}")
|
||||
}
|
||||
};
|
||||
match (next_deadline, wake_fd) {
|
||||
(Some(deadline), fd_opt) => {
|
||||
let now = std::time::Instant::now();
|
||||
if deadline > now {
|
||||
let timeout = deadline - now;
|
||||
match fd_opt {
|
||||
Some(fd) => {
|
||||
// Wake only; the byte (if any) is
|
||||
// consumed by the next drain-lock
|
||||
// winner in phase 1. Level-triggered
|
||||
// poll means an unconsumed byte makes
|
||||
// this return immediately, so a loser
|
||||
// spins briefly until the winner
|
||||
// releases — never sleeps through it.
|
||||
crate::io::poll_wake(fd, Some(timeout));
|
||||
}
|
||||
None => thread::sleep(timeout),
|
||||
}
|
||||
}
|
||||
}
|
||||
(None, Some(fd)) if io_outstanding > 0 => {
|
||||
// See above: no byte consumption outside phase 1.
|
||||
crate::io::poll_wake(fd, None);
|
||||
}
|
||||
_ => {
|
||||
thread::sleep(std::time::Duration::from_micros(100));
|
||||
}
|
||||
timers
|
||||
.peek_deadline()
|
||||
.filter(|d| inner.coord.try_arm_timer(slot_idx, *d))
|
||||
};
|
||||
// The mandatory post-publish re-check: a producer that
|
||||
// enqueued (or a verdict input that flipped) before it
|
||||
// could see our idle bit has left us the evidence.
|
||||
let _ = inner.coord.park(slot_idx, tk_deadline, || {
|
||||
!inner.run_queue.is_empty()
|
||||
|| (inner.live_actors.load(Ordering::Acquire) == 0
|
||||
&& inner.io_outstanding.load(Ordering::Acquire) == 0
|
||||
&& inner.io_fd_waiters.load(Ordering::Acquire) == 0)
|
||||
|| (inner.root_exited.load(Ordering::Acquire)
|
||||
&& !inner.root_swept.load(Ordering::Acquire))
|
||||
|| inner.coord.deadline_due()
|
||||
});
|
||||
if tk_deadline.is_some() {
|
||||
// Hand the role back BEFORE firing: pop_due can run
|
||||
// `Send` thunks that insert new timers, and the
|
||||
// insert-side re-arm check must see either no
|
||||
// timekeeper (skip) or a real parked one — never us,
|
||||
// awake and about to re-peek anyway.
|
||||
inner.coord.disarm_timer(slot_idx);
|
||||
// Woken for the deadline, for work, or to re-peek
|
||||
// after an earlier insert — fire whatever is due;
|
||||
// the next idle pass re-arms with the new minimum.
|
||||
fire_due_timers(inner, false);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
@@ -1780,6 +1767,21 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||||
// by the at-most-once-enqueued invariant nothing else can have
|
||||
// changed the state of a queued actor.
|
||||
// ----------------------------------------------------------------
|
||||
|
||||
// RFC 018 chain rule: we just took one runnable; if more remain and
|
||||
// a sibling is parked, wake exactly one so the surplus runs in
|
||||
// PARALLEL rather than serially behind us (without this the surplus
|
||||
// is not stranded — we re-pop it after resuming — but it waits out
|
||||
// our whole timeslice while an idle core sits available). Cheap: the
|
||||
// queue-length check is queue-local, and `wake_one_if_idle` is a
|
||||
// fence + one Relaxed mask load when nobody is parked. A Relaxed
|
||||
// miss here is safe — the enqueue that created the surplus already
|
||||
// issued its own wake (RFC 018 no-lost-wake); this only sharpens
|
||||
// parallelism latency.
|
||||
if !inner.run_queue.is_empty() {
|
||||
inner.coord.wake_one_if_idle();
|
||||
}
|
||||
|
||||
let slot = match inner.slot_at(pid) {
|
||||
Some(s) => s,
|
||||
None => continue, // can't happen for real pids; defensive
|
||||
|
||||
+36
-9
@@ -72,6 +72,7 @@ use crate::runtime::{
|
||||
self, RuntimeInner, YieldIntent, RUNTIME,
|
||||
};
|
||||
use crate::supervisor::Signal;
|
||||
use std::sync::atomic::Ordering;
|
||||
use std::sync::Arc;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -753,7 +754,14 @@ where
|
||||
Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"),
|
||||
};
|
||||
match io.as_mut() {
|
||||
Some(io) => io.submit(me, epoch, work),
|
||||
Some(io) => {
|
||||
// RFC 018: count the op in flight BEFORE submit — the
|
||||
// pool decrements on completion, and an increment that
|
||||
// trailed the completion would underflow. Under the io
|
||||
// lock, so ordered against the same-lock submit.
|
||||
inner.io_outstanding.fetch_add(1, Ordering::AcqRel);
|
||||
io.submit(me, epoch, work);
|
||||
}
|
||||
None => panic!("io thread not started"),
|
||||
}
|
||||
});
|
||||
@@ -813,7 +821,17 @@ fn wait_fd(fd: std::os::fd::RawFd, readable: bool, writable: bool) -> std::io::R
|
||||
Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"),
|
||||
};
|
||||
match io.as_mut() {
|
||||
Some(io) => io.epoll_register(fd, me, epoch, readable, writable),
|
||||
Some(io) => {
|
||||
// RFC 018: count the waiter BEFORE the ADD (mirror of
|
||||
// submit); roll back if the registration fails so a
|
||||
// rejected wait leaves the verdict counters clean.
|
||||
inner.io_fd_waiters.fetch_add(1, Ordering::AcqRel);
|
||||
let r = io.epoll_register(fd, me, epoch, readable, writable);
|
||||
if r.is_err() {
|
||||
inner.io_fd_waiters.fetch_sub(1, Ordering::AcqRel);
|
||||
}
|
||||
r
|
||||
}
|
||||
None => panic!("io thread not started"),
|
||||
}
|
||||
})?;
|
||||
@@ -838,9 +856,12 @@ fn wait_fd(fd: std::os::fd::RawFd, readable: bool, writable: bool) -> std::io::R
|
||||
Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"),
|
||||
};
|
||||
if let Some(io) = io.as_mut() {
|
||||
if io.waiters.get(&self.fd) == Some(&(self.me, self.epoch)) {
|
||||
io.waiters.remove(&self.fd);
|
||||
io.epoll_deregister(self.fd);
|
||||
// `cancel_waiter` removes + DELs iff still ours, all
|
||||
// under the waiters lock (the ADD/DEL serialization);
|
||||
// decrement only when we actually removed it — a
|
||||
// FdReady that consumed it already did the decrement.
|
||||
if io.cancel_waiter(self.fd, self.me, self.epoch) {
|
||||
inner.io_fd_waiters.fetch_sub(1, Ordering::AcqRel);
|
||||
}
|
||||
}
|
||||
});
|
||||
@@ -908,7 +929,14 @@ impl crate::channel::Selectable for FdArm {
|
||||
};
|
||||
match io.as_mut() {
|
||||
Some(io) => {
|
||||
io.epoll_register(self.fd, pid, epoch, self.readable, self.writable)
|
||||
inner.io_fd_waiters.fetch_add(1, Ordering::AcqRel);
|
||||
let r = io.epoll_register(
|
||||
self.fd, pid, epoch, self.readable, self.writable,
|
||||
);
|
||||
if r.is_err() {
|
||||
inner.io_fd_waiters.fetch_sub(1, Ordering::AcqRel);
|
||||
}
|
||||
r
|
||||
}
|
||||
None => panic!("io thread not started"),
|
||||
}
|
||||
@@ -936,9 +964,8 @@ impl crate::channel::Selectable for FdArm {
|
||||
Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"),
|
||||
};
|
||||
if let Some(io) = io.as_mut() {
|
||||
if io.waiters.get(&self.fd) == Some(&(pid, epoch)) {
|
||||
io.waiters.remove(&self.fd);
|
||||
io.epoll_deregister(self.fd);
|
||||
if io.cancel_waiter(self.fd, pid, epoch) {
|
||||
inner.io_fd_waiters.fetch_sub(1, Ordering::AcqRel);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
+11
-2
@@ -6,10 +6,19 @@
|
||||
//! Build the loom models with: `RUSTFLAGS="--cfg loom" cargo test --lib --release`
|
||||
|
||||
#[cfg(loom)]
|
||||
pub(crate) use loom::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
|
||||
pub(crate) use loom::sync::atomic::{fence, AtomicU64, AtomicUsize, Ordering};
|
||||
|
||||
#[cfg(not(loom))]
|
||||
pub(crate) use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
|
||||
pub(crate) use std::sync::atomic::{fence, AtomicU64, AtomicUsize, Ordering};
|
||||
|
||||
// park.rs condvar-parker (loom + non-Linux builds only; the Linux non-loom
|
||||
// build parks on a futex and never touches these — gating them identically
|
||||
// keeps the default build free of unused imports).
|
||||
#[cfg(loom)]
|
||||
pub(crate) use loom::sync::{Condvar, Mutex};
|
||||
|
||||
#[cfg(all(not(loom), not(target_os = "linux")))]
|
||||
pub(crate) use std::sync::{Condvar, Mutex};
|
||||
|
||||
/// `UnsafeCell` with loom's `with`/`with_mut` access API; pass-through cost
|
||||
/// is zero in normal builds (`#[inline]`, newtype over std's cell).
|
||||
|
||||
+37
-1
@@ -141,6 +141,14 @@ impl PartialOrd for Entry {
|
||||
|
||||
#[derive(Default)]
|
||||
pub struct Timers {
|
||||
/// RFC 018: the scheduler coordination layer. Attached once at
|
||||
/// `RuntimeInner::new`; every insert notes its deadline (min-maintained
|
||||
/// snapshot for the busy-path due-check + the timekeeper re-arm wake)
|
||||
/// and every pop/clear re-anchors the snapshot to the heap minimum.
|
||||
/// All calls happen under the timers mutex — the serialization the
|
||||
/// coordinator's timer protocol mandates. `None` only in unit tests
|
||||
/// that construct a bare `Timers`.
|
||||
coord: Option<std::sync::Arc<crate::park::Coordinator>>,
|
||||
/// Reverse-wrapped so the smallest deadline is at the top.
|
||||
heap: BinaryHeap<Reverse<Entry>>,
|
||||
/// Monotonic counter for the tiebreaker `seq` field (and the `TimerId` of a
|
||||
@@ -157,7 +165,18 @@ pub struct Timers {
|
||||
|
||||
impl Timers {
|
||||
pub fn new() -> Self {
|
||||
Self { heap: BinaryHeap::new(), next_seq: 0, armed: std::collections::HashSet::new() }
|
||||
Self {
|
||||
coord: None,
|
||||
heap: BinaryHeap::new(),
|
||||
next_seq: 0,
|
||||
armed: std::collections::HashSet::new(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Attach the scheduler coordination layer (RFC 018). Called once, at
|
||||
/// runtime construction, before any scheduler thread exists.
|
||||
pub(crate) fn attach_coordinator(&mut self, c: std::sync::Arc<crate::park::Coordinator>) {
|
||||
self.coord = Some(c);
|
||||
}
|
||||
|
||||
/// Insert a `Sleep` timer. Convenience for the common case.
|
||||
@@ -242,6 +261,13 @@ impl Timers {
|
||||
#[cfg(feature = "smarm-causal")]
|
||||
wall,
|
||||
}));
|
||||
// RFC 018: publish the (possibly new-minimum) deadline to the
|
||||
// busy-path snapshot and wake the timekeeper if it is parked
|
||||
// toward a later one. We hold the timers mutex — the mandated
|
||||
// serialization for both.
|
||||
if let Some(c) = &self.coord {
|
||||
c.note_deadline(deadline);
|
||||
}
|
||||
seq
|
||||
}
|
||||
|
||||
@@ -255,6 +281,9 @@ impl Timers {
|
||||
pub fn clear(&mut self) {
|
||||
self.heap.clear();
|
||||
self.armed.clear();
|
||||
if let Some(c) = &self.coord {
|
||||
c.refresh_deadline(None);
|
||||
}
|
||||
}
|
||||
|
||||
/// Soonest pending deadline, or `None` if the heap is empty.
|
||||
@@ -324,6 +353,13 @@ impl Timers {
|
||||
}
|
||||
out.push(entry);
|
||||
}
|
||||
// RFC 018: re-anchor the busy-path snapshot to the new heap minimum
|
||||
// (still under the timers mutex). A causal-shift re-queue above went
|
||||
// through `heap.push` directly, so this peek is the one place the
|
||||
// snapshot is guaranteed to catch up.
|
||||
if let Some(c) = &self.coord {
|
||||
c.refresh_deadline(self.peek_deadline());
|
||||
}
|
||||
out
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,69 @@
|
||||
//! RFC 018 scheduler park/wake — observable-behavior guards.
|
||||
//!
|
||||
//! These pin the two timer-latency properties the park/wake swap must
|
||||
//! preserve or introduce:
|
||||
//!
|
||||
//! - `sleep_fires_under_saturation`: due timers fire even when every
|
||||
//! scheduler is busy (nobody parked ⇒ no timekeeper) — the busy-path
|
||||
//! due-check, ratified design point (a). The old drain phase gave this
|
||||
//! for free (timers drained every loop iteration); the new design must
|
||||
//! not lose it.
|
||||
//! - `submillisecond_sleep_is_prompt`: a sub-ms sleep completes promptly.
|
||||
//! Under the old wake pipe, `poll_wake`'s `as_millis` truncation turned
|
||||
//! sub-ms deadlines into 0ms busy-polls (correct wall time, pathological
|
||||
//! CPU); under park/wake the futex timespec carries full nanosecond
|
||||
//! precision.
|
||||
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
#[test]
|
||||
fn sleep_fires_under_saturation() {
|
||||
let rt = smarm::runtime::init(smarm::runtime::Config::exact(4));
|
||||
rt.run(|| {
|
||||
let stop = Arc::new(AtomicBool::new(false));
|
||||
let mut spinners = Vec::new();
|
||||
// 8 spinners over 4 schedulers: the run queue never empties, so no
|
||||
// scheduler ever parks and no timekeeper exists. Only the busy-path
|
||||
// due-check can fire the sleeper's timer before the spinners quit.
|
||||
for _ in 0..8 {
|
||||
let stop = stop.clone();
|
||||
spinners.push(smarm::spawn(move || {
|
||||
let t0 = Instant::now();
|
||||
while !stop.load(Ordering::Relaxed) && t0.elapsed() < Duration::from_secs(5) {
|
||||
smarm::yield_now();
|
||||
}
|
||||
}));
|
||||
}
|
||||
let t0 = Instant::now();
|
||||
smarm::sleep(Duration::from_millis(10));
|
||||
let dt = t0.elapsed();
|
||||
stop.store(true, Ordering::Relaxed);
|
||||
for s in spinners {
|
||||
let _ = s.join();
|
||||
}
|
||||
assert!(
|
||||
dt < Duration::from_millis(500),
|
||||
"10ms sleep took {dt:?} under scheduler saturation — busy-path \
|
||||
timer firing is broken (timekeeper-only firing stalls under load)"
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn submillisecond_sleep_is_prompt() {
|
||||
let rt = smarm::runtime::init(smarm::runtime::Config::exact(2));
|
||||
rt.run(|| {
|
||||
// Warm one iteration, then measure.
|
||||
smarm::sleep(Duration::from_micros(500));
|
||||
let t0 = Instant::now();
|
||||
smarm::sleep(Duration::from_micros(500));
|
||||
let dt = t0.elapsed();
|
||||
assert!(dt >= Duration::from_micros(400), "woke early: {dt:?}");
|
||||
assert!(
|
||||
dt < Duration::from_millis(100),
|
||||
"500µs sleep took {dt:?} — sub-ms deadline handling is broken"
|
||||
);
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user