2446 lines
112 KiB
Rust
2446 lines
112 KiB
Rust
//! Multi-scheduler runtime: configuration, initialisation, and the shared
|
||
//! state that all scheduler OS threads operate against.
|
||
//!
|
||
//! # Architecture (post slot-table split, ROADMAP_v0.5 phase 2)
|
||
//!
|
||
//! ```text
|
||
//! init(Config) → Runtime (Arc<RuntimeInner>)
|
||
//!
|
||
//! RuntimeInner {
|
||
//! slots: Box<[Slot]> ← FIXED slab, max_actors entries, lock-free lookup
|
||
//! free: RawMutex<Vec<u32>> ← vacant slot indices
|
||
//! run_queue: RunQueue ← compile-time selected (src/run_queue.rs)
|
||
//! timers: Mutex<Timers>
|
||
//! io: Mutex<Option<IoThread>>
|
||
//! live_actors: AtomicU32 ← spawned-but-not-finalized count (termination)
|
||
//! stats: Vec<SchedulerStats> ← one per thread, lockless atomics (RFC 000)
|
||
//! }
|
||
//!
|
||
//! Slot {
|
||
//! word: AtomicU64 ← (gen << 32) | (epoch << 8) | state — THE state machine
|
||
//! sp: AtomicUsize ← saved stack pointer
|
||
//! stop_ptr:AtomicPtr<...> ← into the actor's Arc<AtomicBool>
|
||
//! closure: AtomicPtr<...> ← first-resume closure, swap-to-take
|
||
//! cold: RawMutex<SlotCold> ← lifecycle collections (waiters/monitors/links/…)
|
||
//! }
|
||
//! ```
|
||
//!
|
||
//! # The per-slot state machine
|
||
//!
|
||
//! Scheduling state lives in one atomic word per slot packing
|
||
//! `(generation, park-epoch, state)`, where state is one of:
|
||
//!
|
||
//! ```text
|
||
//! Vacant ─spawn→ Queued ─pop→ Running ─yield→ Queued
|
||
//! ↑ │ │
|
||
//! │ park │ unpark while running
|
||
//! │ ↓ ↓
|
||
//! unpark ←──── Parked RunningNotified ─park→ Queued
|
||
//! (re-queued immediately)
|
||
//! Running|RunningNotified ─actor returns→ Done ─reclaim→ Vacant(gen+1)
|
||
//! ```
|
||
//!
|
||
//! Every transition is a CAS on the packed word, so:
|
||
//!
|
||
//! - The generation check is **atomic with the transition** — a stale `Pid`
|
||
//! can never act on a recycled slot (no ABA, no spurious unparks).
|
||
//! - The park-epoch (middle 24 bits) is the actor's *wait identity*: opened
|
||
//! by `begin_wait` before any registration, consumed (bumped) by every
|
||
//! successful wake. Registration-based wakers carry `(pid, epoch)` and use
|
||
//! `unpark_at`, so a waker holding a registration from an already-woken
|
||
//! wait — a `select` loser arm, a satisfied wait's timer — fails the epoch
|
||
//! check and no-ops instead of faulting a later one-shot park. The only
|
||
//! wildcard wake is `request_stop`, which is terminal. Full rules in
|
||
//! slot_state.rs.
|
||
//! - `RunningNotified` replaces the old `pending_unpark` bool: an unpark that
|
||
//! races the prep-to-park window is a *state*, resolved by the scheduler's
|
||
//! park-return CAS, not a flag read under a lock. This also closes a latent
|
||
//! lost-wakeup in the old Blocking-IO completion path, which set the result
|
||
//! for a still-Running actor without flagging it.
|
||
//! - **A pid is in the run queue at most once**: the only pushes are paired
|
||
//! 1:1 with successful transitions *into* `Queued`, and only the scheduler
|
||
//! transitions `Queued → Running` (paired 1:1 with pops).
|
||
//!
|
||
//! Memory ordering: all word CASes are `AcqRel` (failure `Acquire`), plain
|
||
//! word stores are `Release`, loads are `Acquire`. The chain that matters:
|
||
//! the park path stores `sp` (Relaxed) *before* its Release transition; any
|
||
//! later Acquire transition/load of the word therefore observes that `sp`.
|
||
//! RFC 019's `hwm` (and the shrink that reads it) piggybacks this exact
|
||
//! pattern in the same pre-Release window and adds no edges.
|
||
//! The run-queue mutex independently provides the same edges today; the
|
||
//! word's own ordering is what phase 3's lock-free queue will rely on.
|
||
//!
|
||
//! # Locks and ordering
|
||
//!
|
||
//! - The run queue is its own module (`run_queue.rs`), selected at compile
|
||
//! time (`rq-mutex` / `rq-mpmc` / `rq-striped`). Queue ops require
|
||
//! preemption disabled (debug-asserted there); when the mutex variant is in
|
||
//! play it is the innermost lock — nothing else is acquired under it.
|
||
//! - Per-slot `cold` locks (`RawMutex`, non-poisoning, guard enters
|
||
//! `NoPreempt`) guard the lifecycle collections. **Leaf rule: never hold
|
||
//! two cold locks at once** — `finalize_actor`'s link cascade and `link()`
|
||
//! lock peers one at a time (correctness arguments at the call sites).
|
||
//! Holding a cold lock while pushing to the run queue is permitted.
|
||
//! - Lock order overall: `io` → (slot `cold` | `free` | `stack_pool`,
|
||
//! mutually leaf) → run queue (innermost). `timers` is independent (never
|
||
//! nested with any of the above on either side).
|
||
//!
|
||
//! # Termination (counter-based)
|
||
//!
|
||
//! The old all-clear scanned the slot table under the big lock. Now:
|
||
//! 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.
|
||
//!
|
||
//! # Scheduler park/wake (RFC 018)
|
||
//!
|
||
//! 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, set_current_pid,
|
||
take_last_outcome, Actor, Outcome,
|
||
};
|
||
use crate::channel::Sender;
|
||
use crate::context::switch_to_actor;
|
||
use crate::io::IoThread;
|
||
use crate::monitor::{Down, DownReason, MonitorId};
|
||
use crate::pid::Pid;
|
||
use crate::preempt::PREEMPTION_ENABLED;
|
||
use crate::raw_mutex::RawMutex;
|
||
use crate::slot_state::{StateWord, Status, Unpark};
|
||
use crate::supervisor::Signal;
|
||
use crate::timer::Timers;
|
||
|
||
use std::sync::atomic::{AtomicBool, AtomicPtr, AtomicU32, AtomicU64, AtomicUsize, Ordering};
|
||
use std::sync::{Arc, Mutex, Weak};
|
||
use std::thread;
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Config
|
||
// ---------------------------------------------------------------------------
|
||
|
||
/// Default capacity of the actor slot table. Slots are ~256 bytes, so the
|
||
/// default costs ~4 MiB, allocated once at `init`. See [`Config::max_actors`].
|
||
pub const DEFAULT_MAX_ACTORS: usize = 16_384;
|
||
|
||
/// Runtime configuration.
|
||
///
|
||
/// ```
|
||
/// use smarm::runtime::Config;
|
||
///
|
||
/// // Use all available CPUs (default):
|
||
/// let _all = Config::default();
|
||
///
|
||
/// // Exactly 4 scheduler threads:
|
||
/// let _four = Config::exact(4);
|
||
///
|
||
/// // Between 2 and 8, clamped to available parallelism:
|
||
/// let _clamped = Config::new(2, 8, None);
|
||
/// ```
|
||
#[derive(Clone, Debug)]
|
||
pub struct Config {
|
||
min: usize,
|
||
max: usize,
|
||
exact: Option<usize>,
|
||
alloc_interval: u32,
|
||
timeslice_cycles: u64,
|
||
stack_pool_cap: usize,
|
||
stack_reserve: usize,
|
||
stack_guard: usize,
|
||
max_actors: usize,
|
||
wake_slot: bool,
|
||
node_id: crate::pg::NodeId,
|
||
incarnation: crate::pg::Incarnation,
|
||
}
|
||
|
||
impl Config {
|
||
/// Exact thread count; takes precedence over min/max.
|
||
pub fn exact(n: usize) -> Self {
|
||
assert!(n >= 1, "scheduler thread count must be ≥ 1");
|
||
Self {
|
||
min: n,
|
||
max: n,
|
||
exact: Some(n),
|
||
alloc_interval: crate::preempt::DEFAULT_ALLOC_INTERVAL,
|
||
timeslice_cycles: crate::preempt::DEFAULT_TIMESLICE_CYCLES,
|
||
stack_pool_cap: n * 4,
|
||
stack_reserve: DEFAULT_STACK_RESERVE,
|
||
stack_guard: DEFAULT_STACK_GUARD,
|
||
max_actors: DEFAULT_MAX_ACTORS,
|
||
wake_slot: true,
|
||
node_id: crate::pg::DEFAULT_NODE_ID,
|
||
incarnation: crate::pg::DEFAULT_INCARNATION,
|
||
}
|
||
}
|
||
|
||
/// Bounded range. Thread count = clamp(available_parallelism, min, max).
|
||
pub fn new(min: usize, max: usize, exact: Option<usize>) -> Self {
|
||
assert!(min >= 1, "min must be ≥ 1");
|
||
assert!(max >= min, "max must be ≥ min");
|
||
if let Some(e) = exact {
|
||
assert!(e >= 1, "exact must be ≥ 1");
|
||
}
|
||
Self {
|
||
min,
|
||
max,
|
||
exact,
|
||
alloc_interval: crate::preempt::DEFAULT_ALLOC_INTERVAL,
|
||
timeslice_cycles: crate::preempt::DEFAULT_TIMESLICE_CYCLES,
|
||
stack_pool_cap: max * 4,
|
||
stack_reserve: DEFAULT_STACK_RESERVE,
|
||
stack_guard: DEFAULT_STACK_GUARD,
|
||
max_actors: DEFAULT_MAX_ACTORS,
|
||
wake_slot: true,
|
||
node_id: crate::pg::DEFAULT_NODE_ID,
|
||
incarnation: crate::pg::DEFAULT_INCARNATION,
|
||
}
|
||
}
|
||
|
||
/// How many allocations (or `smarm::check!()` calls) between RDTSC checks.
|
||
/// Lower = more responsive preemption, higher = less overhead.
|
||
/// Default: 128.
|
||
pub fn alloc_interval(mut self, n: u32) -> Self {
|
||
assert!(n >= 1, "alloc_interval must be ≥ 1");
|
||
self.alloc_interval = n;
|
||
self
|
||
}
|
||
|
||
/// How many TSC cycles constitute one timeslice.
|
||
/// Default: 300_000 (≈ 100µs on a 3 GHz CPU).
|
||
pub fn timeslice_cycles(mut self, n: u64) -> Self {
|
||
assert!(n >= 1, "timeslice_cycles must be ≥ 1");
|
||
self.timeslice_cycles = n;
|
||
self
|
||
}
|
||
|
||
/// Maximum number of stacks kept in the pool for reuse across spawns.
|
||
/// A larger cap reduces `mmap`/`munmap` syscalls at the cost of idle memory.
|
||
/// Default: `thread_count * 4`.
|
||
pub fn stack_pool_cap(mut self, n: usize) -> Self {
|
||
self.stack_pool_cap = n;
|
||
self
|
||
}
|
||
|
||
/// Default per-actor stack reserve (RFC 019). A *virtual* reservation —
|
||
/// anonymous mmap is demand-paged, so RSS follows touched pages, not
|
||
/// this number — but overflowing it hits the guard and dies. Page-rounded.
|
||
/// Per-actor override: `SpawnOpts::stack_reserve`.
|
||
/// Default: [`DEFAULT_STACK_RESERVE`] (64 KiB) — the million-cheap-actors
|
||
/// story is unchanged; big stacks are opt-in.
|
||
pub fn stack_reserve(mut self, n: usize) -> Self {
|
||
assert!(n > 0, "stack_reserve must be non-zero");
|
||
self.stack_reserve = n;
|
||
self
|
||
}
|
||
|
||
/// Default PROT_NONE guard below each stack (RFC 019). Address space
|
||
/// only. Page-rounded. Rust overflow is caught by any single page
|
||
/// (probestack touches pages in order); the wide default exists for
|
||
/// unprobed FFI frames, which can step over a small guard in one
|
||
/// `sub rsp`. Per-actor override: `SpawnOpts::guard_size`.
|
||
/// Default: [`DEFAULT_STACK_GUARD`] (1 MiB — the kernel's
|
||
/// `stack_guard_gap` convention; see its doc for why width is free).
|
||
pub fn stack_guard(mut self, n: usize) -> Self {
|
||
assert!(n > 0, "stack_guard must be non-zero");
|
||
self.stack_guard = n;
|
||
self
|
||
}
|
||
|
||
/// Capacity of the actor slot table — the maximum number of
|
||
/// **simultaneously live** actors (total spawned over a run is unbounded;
|
||
/// slots are recycled). The table is a fixed slab allocated once at
|
||
/// `init`: slots never move, which is what makes lock-free slot lookup
|
||
/// sound. Exhausting it is a loud panic naming this knob.
|
||
/// Default: [`DEFAULT_MAX_ACTORS`] (16_384, ~4 MiB).
|
||
pub fn max_actors(mut self, n: usize) -> Self {
|
||
assert!(n >= 1, "max_actors must be ≥ 1");
|
||
assert!(
|
||
n < u32::MAX as usize,
|
||
"max_actors must fit a u32 slot index (ROOT_PID reserves u32::MAX)"
|
||
);
|
||
self.max_actors = n;
|
||
self
|
||
}
|
||
|
||
/// Enable the per-scheduler wake slot (RFC 005): a thread-local,
|
||
/// capacity-one wake cache checked before the shared run queue. A wake
|
||
/// performed from actor context parks the woken pid in the waking
|
||
/// thread's slot; it is resumed next on that core and inherits the
|
||
/// remainder of the waker's timeslice. Scheduler-context wakes
|
||
/// (timer/IO drain) and spawns always go to the shared queue.
|
||
/// Default: `true` (accepted 2026-08-18, history.md finding 17/18).
|
||
pub fn wake_slot(mut self, on: bool) -> Self {
|
||
self.wake_slot = on;
|
||
self
|
||
}
|
||
|
||
/// This runtime's node identity (RFC 012). Defaults to a fixed single-node
|
||
/// value; clustering (RFC 010) will supply a real one. Threaded through pg
|
||
/// storage/eviction so the process-group public API never changes to
|
||
/// acquire it.
|
||
pub fn node_id(mut self, id: impl Into<crate::pg::NodeId>) -> Self {
|
||
self.node_id = id.into();
|
||
self
|
||
}
|
||
|
||
/// This runtime's incarnation epoch (RFC 012) — the BEAM `Creation` analogue
|
||
/// that separates a crashed node from its restart. Defaults to a fixed
|
||
/// single-node value; constant for the life of a run.
|
||
pub fn incarnation(mut self, inc: impl Into<crate::pg::Incarnation>) -> Self {
|
||
self.incarnation = inc.into();
|
||
self
|
||
}
|
||
|
||
/// The number of scheduler threads this config resolves to.
|
||
pub fn resolved_thread_count(&self) -> usize {
|
||
if let Some(e) = self.exact {
|
||
return e;
|
||
}
|
||
let avail = thread::available_parallelism()
|
||
.map(|n| n.get())
|
||
.unwrap_or(1);
|
||
avail.clamp(self.min, self.max)
|
||
}
|
||
}
|
||
|
||
impl Default for Config {
|
||
fn default() -> Self {
|
||
let avail = thread::available_parallelism()
|
||
.map(|n| n.get())
|
||
.unwrap_or(1);
|
||
Self {
|
||
min: 1,
|
||
max: avail,
|
||
exact: None,
|
||
alloc_interval: crate::preempt::DEFAULT_ALLOC_INTERVAL,
|
||
timeslice_cycles: crate::preempt::DEFAULT_TIMESLICE_CYCLES,
|
||
stack_pool_cap: avail * 4,
|
||
stack_reserve: DEFAULT_STACK_RESERVE,
|
||
stack_guard: DEFAULT_STACK_GUARD,
|
||
max_actors: DEFAULT_MAX_ACTORS,
|
||
wake_slot: true,
|
||
node_id: crate::pg::DEFAULT_NODE_ID,
|
||
incarnation: crate::pg::DEFAULT_INCARNATION,
|
||
}
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Per-thread stats (RFC 000 Layer 1 primitives)
|
||
// ---------------------------------------------------------------------------
|
||
|
||
/// Lockless per-scheduler-thread counters. Written only by the owning thread;
|
||
/// readable from any thread (introspection actor, tests).
|
||
#[repr(align(64))]
|
||
pub struct SchedulerStats {
|
||
/// PID index of the actor currently on-CPU, or `u32::MAX` when idle.
|
||
pub current_pid_index: AtomicU32,
|
||
/// Snapshot of run queue length maintained on every push/pop.
|
||
pub run_queue_len: AtomicU64,
|
||
/// RFC 005: wakes resumed from this thread's wake slot.
|
||
pub slot_hits: AtomicU64,
|
||
/// RFC 005: slot occupants displaced to the shared queue by a newer wake.
|
||
pub slot_displacements: AtomicU64,
|
||
// --- wake-path diagnostics (target 5). Relaxed counters, cheap. ---
|
||
/// unpark hit Parked from actor context → slot_push.
|
||
pub unpark_slot: AtomicU64,
|
||
/// unpark hit Parked from scheduler/foreign context → shared enqueue.
|
||
pub unpark_queue: AtomicU64,
|
||
/// unpark hit Running → RunningNotified (peer had not parked yet).
|
||
pub unpark_notified: AtomicU64,
|
||
/// park_return found the flag consumed → shared re-enqueue.
|
||
pub park_flag_consumed: AtomicU64,
|
||
/// Yield-intent re-enqueues.
|
||
pub yield_requeues: AtomicU64,
|
||
/// `enqueue` tail wake actually delivered a futex permit.
|
||
pub enqueue_wakes: AtomicU64,
|
||
/// Chain-rule wake actually delivered a permit.
|
||
pub chain_wakes: AtomicU64,
|
||
/// Times this scheduler entered Pop::Idle (about to futex-park).
|
||
pub idle_parks: AtomicU64,
|
||
/// Idle parks whose recheck aborted (WorkFound) — no futex.
|
||
pub idle_recheck_hits: AtomicU64,
|
||
}
|
||
|
||
impl SchedulerStats {
|
||
fn new() -> Self {
|
||
Self {
|
||
current_pid_index: AtomicU32::new(u32::MAX),
|
||
run_queue_len: AtomicU64::new(0),
|
||
slot_hits: AtomicU64::new(0),
|
||
slot_displacements: AtomicU64::new(0),
|
||
unpark_slot: AtomicU64::new(0),
|
||
unpark_queue: AtomicU64::new(0),
|
||
unpark_notified: AtomicU64::new(0),
|
||
park_flag_consumed: AtomicU64::new(0),
|
||
yield_requeues: AtomicU64::new(0),
|
||
enqueue_wakes: AtomicU64::new(0),
|
||
chain_wakes: AtomicU64::new(0),
|
||
idle_parks: AtomicU64::new(0),
|
||
idle_recheck_hits: AtomicU64::new(0),
|
||
}
|
||
}
|
||
|
||
fn reset(&self) {
|
||
for c in [
|
||
&self.slot_hits,
|
||
&self.slot_displacements,
|
||
&self.unpark_slot,
|
||
&self.unpark_queue,
|
||
&self.unpark_notified,
|
||
&self.park_flag_consumed,
|
||
&self.yield_requeues,
|
||
&self.enqueue_wakes,
|
||
&self.chain_wakes,
|
||
&self.idle_parks,
|
||
&self.idle_recheck_hits,
|
||
] {
|
||
c.store(0, Ordering::Relaxed);
|
||
}
|
||
}
|
||
}
|
||
|
||
/// Bump a per-thread diagnostic counter on the calling scheduler thread.
|
||
macro_rules! diag {
|
||
($inner:expr, $field:ident) => {
|
||
SCHED_SLOT.with(|s| $inner.stats[s.get()].$field.fetch_add(1, Ordering::Relaxed))
|
||
};
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Runtime stats snapshot (for tests / introspection)
|
||
// ---------------------------------------------------------------------------
|
||
|
||
pub struct RuntimeStats {
|
||
pub(crate) inner: Arc<RuntimeInner>,
|
||
}
|
||
|
||
impl RuntimeStats {
|
||
/// Sum of run queue lengths across all scheduler threads.
|
||
pub fn total_run_queue_len(&self) -> u64 {
|
||
self.inner
|
||
.stats
|
||
.iter()
|
||
.map(|s| s.run_queue_len.load(Ordering::Relaxed))
|
||
.sum()
|
||
}
|
||
|
||
/// Number of scheduler threads.
|
||
pub fn scheduler_count(&self) -> usize {
|
||
self.inner.stats.len()
|
||
}
|
||
|
||
/// Actors currently parked on IO.
|
||
pub fn io_parked_count(&self) -> u32 {
|
||
self.inner.io_parked.load(Ordering::Relaxed)
|
||
}
|
||
|
||
/// Actors currently sleeping on a timer.
|
||
pub fn sleeping_count(&self) -> u32 {
|
||
self.inner.sleeping.load(Ordering::Relaxed)
|
||
}
|
||
|
||
/// RFC 005: total wakes resumed from a wake slot, summed across
|
||
/// scheduler threads. Counters are reset at the start of each `run()`,
|
||
/// so after a run this reads that run's total.
|
||
pub fn slot_hits(&self) -> u64 {
|
||
self.inner
|
||
.stats
|
||
.iter()
|
||
.map(|s| s.slot_hits.load(Ordering::Relaxed))
|
||
.sum()
|
||
}
|
||
|
||
/// RFC 005: total slot occupants displaced to the shared queue, summed
|
||
/// across scheduler threads. Reset at the start of each `run()`.
|
||
pub fn slot_displacements(&self) -> u64 {
|
||
self.inner
|
||
.stats
|
||
.iter()
|
||
.map(|s| s.slot_displacements.load(Ordering::Relaxed))
|
||
.sum()
|
||
}
|
||
|
||
/// Wake-path diagnostics summed across threads, as `name=value` pairs
|
||
/// (target 5 instrumentation). Reset at the start of each `run()`.
|
||
pub fn wake_diag(&self) -> String {
|
||
let sum = |f: fn(&SchedulerStats) -> &AtomicU64| -> u64 {
|
||
self.inner.stats.iter().map(|s| f(s).load(Ordering::Relaxed)).sum()
|
||
};
|
||
format!(
|
||
"slot_hits={} displaced={} unpark_slot={} unpark_queue={} unpark_notified={} \
|
||
flag_consumed={} yield_requeues={} enqueue_wakes={} chain_wakes={} \
|
||
idle_parks={} idle_recheck_hits={}",
|
||
sum(|s| &s.slot_hits),
|
||
sum(|s| &s.slot_displacements),
|
||
sum(|s| &s.unpark_slot),
|
||
sum(|s| &s.unpark_queue),
|
||
sum(|s| &s.unpark_notified),
|
||
sum(|s| &s.park_flag_consumed),
|
||
sum(|s| &s.yield_requeues),
|
||
sum(|s| &s.enqueue_wakes),
|
||
sum(|s| &s.chain_wakes),
|
||
sum(|s| &s.idle_parks),
|
||
sum(|s| &s.idle_recheck_hits),
|
||
)
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Slot — packed state word + hot atomics + cold lifecycle data
|
||
// ---------------------------------------------------------------------------
|
||
|
||
/// Default usable stack reserve per actor (RFC 019). See [`Config::stack_reserve`].
|
||
pub const DEFAULT_STACK_RESERVE: usize = 64 * 1024;
|
||
|
||
/// Default PROT_NONE guard below each actor stack (RFC 019). Raised from one
|
||
/// page so unprobed C frames cannot leap it. See [`Config::stack_guard`].
|
||
///
|
||
/// 1 MiB, following the kernel's own answer to the same problem: after Stack
|
||
/// Clash (2017) the main-thread guard gap became `stack_guard_gap` = 256
|
||
/// pages, because 4 KiB was jumpable by one honest `sub rsp` and no small
|
||
/// constant was defensible. Guard pages are PROT_NONE: virtual address space
|
||
/// only — zero RSS, zero page-table entries, no overcommit charge — so the
|
||
/// wide default is free at any actor count (1 M actors ≈ 1 TiB of VA against
|
||
/// a 128 TiB budget). A frame that jumps even this lands in the tier-2
|
||
/// overshoot window of the SIGSEGV diagnostic (`signal.rs`) instead of
|
||
/// silence.
|
||
pub const DEFAULT_STACK_GUARD: usize = 1024 * 1024;
|
||
|
||
/// RFC 019 §3: minimum releasable span (`sp − hwm` at park) before the
|
||
/// park-path shrink spends a syscall. A constant, not a `Config` field
|
||
/// (ratified): nobody tunes this well and the measured stakes are low — a
|
||
/// threshold-sized `MADV_FREE` costs ~3 µs against a ~100 ns park, paid
|
||
/// only on spike-recovery parks, which are rare by construction and *were*
|
||
/// the spike. Steady-state actors never reach the syscall: their check is
|
||
/// two Relaxed loads and a compare on a line the context-save just wrote.
|
||
pub const SHRINK_THRESHOLD: usize = 256 * 1024;
|
||
|
||
/// RFC 019 §3: parks between shrinks of one actor. Guards a few-µs cost, so
|
||
/// it can be coarse (parks, not wall time); the kernel's
|
||
/// reclaim-under-pressure-only handling of `MADV_FREE` is the real release
|
||
/// hysteresis — re-touched-before-pressure pages cost a 0.24 µs/page
|
||
/// cancel-write and no fault. A constant, not `Config` (ratified, same
|
||
/// rationale as [`SHRINK_THRESHOLD`]).
|
||
pub const SHRINK_COOLDOWN: u32 = 64;
|
||
|
||
/// RFC 019 §6: the entry-end span (highest addresses — the frames the next
|
||
/// actor faults first) a recycled stack keeps resident; everything below it
|
||
/// is `MADV_DONTNEED`ed before the stack re-enters the pool. Ratified as a
|
||
/// constant, not Config, alongside the shrink knobs; the 64 KiB value was a
|
||
/// flagged Claude-solo call at ratification — it equals the default reserve,
|
||
/// so with an unraised Config the zap is a no-op and only Configs that raise
|
||
/// the default reserve pay it.
|
||
pub const RECYCLE_RETAIN: usize = 64 * 1024;
|
||
|
||
pub(crate) type Closure = Box<dyn FnOnce() + Send>;
|
||
|
||
/// Lifecycle data, mutated only under the slot's cold [`RawMutex`]. Everything
|
||
/// here is touched O(1) times per actor lifetime (spawn / join / monitor /
|
||
/// link / finalize), never on the yield/park/unpark hot path.
|
||
pub(crate) struct SlotCold {
|
||
pub(crate) actor: Option<Actor>,
|
||
/// Parked joiners as `(pid, park-epoch)`; finalize wakes each via the
|
||
/// epoch-matched unpark.
|
||
pub(crate) waiters: Vec<(Pid, u32)>,
|
||
pub(crate) outcome: Option<Outcome>,
|
||
/// The slot's most recent *watchable-tenancy* death: `(generation,
|
||
/// reason)`, stamped by `finalize_actor` — but only for a tenancy whose
|
||
/// `watchable` bit was set — and deliberately never cleared: a new
|
||
/// tenant's install leaves it standing (it describes the previous
|
||
/// tenancy), and only the next *watchable* death overwrites it.
|
||
/// Anonymous green-thread churn must not evict it: the free list is
|
||
/// LIFO, so the just-freed slot is the first recycled, and an
|
||
/// unconditional stamp made a watchable tenancy's record the
|
||
/// shortest-lived data in the runtime. Read generation-matched via
|
||
/// [`terminal_reason`](crate::monitor::terminal_reason), so a watch that
|
||
/// raced its target's death can recover the real down reason instead of
|
||
/// a blanket `NoProc` (bridge soak signatures 4 and 5).
|
||
pub(crate) terminal: Option<(u32, DownReason)>,
|
||
/// Stamp eligibility for `terminal` above: someone could plausibly hold
|
||
/// a watch on this tenancy. Two set-sites, both while the tenancy is
|
||
/// live: `register_with` *before* the binding lands (no successfully
|
||
/// registered actor can die unflagged; a failed register's overshoot is
|
||
/// harmless), and [`mark_watchable`](crate::monitor::mark_watchable) —
|
||
/// the bridge calls it wherever a pid is encoded across the boundary,
|
||
/// because BEAM can only watch pids it holds and can only hold pids
|
||
/// that crossed. Reset at reclaim.
|
||
pub(crate) watchable: bool,
|
||
pub(crate) supervisor_channel: Option<Sender<Signal>>,
|
||
/// Watchers registered via `monitor()`, each tagged with its
|
||
/// `MonitorId` so `demonitor` can remove exactly one. Each receives one
|
||
/// `Down` when this actor terminates (drained in `finalize_actor`).
|
||
/// Distinct from `supervisor_channel`, which is the parent's single funnel.
|
||
pub(crate) monitors: Vec<(MonitorId, Sender<Down>)>,
|
||
/// Bidirectional links (roadmap #3). Each entry is a peer whose abnormal
|
||
/// death propagates to this actor (and vice versa). Entries may be
|
||
/// momentarily or persistently stale (peer already dead) — every walk
|
||
/// re-verifies the peer's word, so stale entries are benign no-ops.
|
||
pub(crate) links: Vec<Pid>,
|
||
pub(crate) outstanding_handles: u32,
|
||
pub(crate) pending_io_result: Option<crate::io::IoResult>,
|
||
}
|
||
|
||
/// One actor slot. Hot scheduling state is atomic; cold lifecycle state is
|
||
/// behind `cold`. Slots live in a fixed slab and never move.
|
||
///
|
||
/// `align(128)` keeps two adjacent slots' hot words off each other's
|
||
/// cache-line pair (x86 prefetches lines in pairs), avoiding false sharing
|
||
/// between unrelated actors.
|
||
#[repr(align(128))]
|
||
pub(crate) struct Slot {
|
||
/// `(generation << 32) | state` — the state machine, factored into
|
||
/// `slot_state.rs` (loom-modeled there; every transition self-asserts).
|
||
word: StateWord,
|
||
/// Saved stack pointer. Written by the owning scheduler thread before the
|
||
/// Release transition out of Running; read after the Acquire transition
|
||
/// Queued→Running. Relaxed is sufficient — ordering rides on `word`.
|
||
sp: AtomicUsize,
|
||
/// RFC 019: sampled stack high-water — the minimum `sp` ever stored above,
|
||
/// i.e. the deepest excursion *observed at a switch point*. Advisory:
|
||
/// correctness never depends on it; its one job is "is a shrink worth a
|
||
/// syscall?". Declared adjacent to `sp` so the min-update dirties the
|
||
/// line the context-save just wrote. Same single-writer Relaxed
|
||
/// discipline as `sp`; reset to the fresh `sp` at install.
|
||
hwm: AtomicUsize,
|
||
/// RFC 019: parks since the last shrink (or install). Counted on every
|
||
/// pass through the Park arm by the owning scheduler thread; the shrink
|
||
/// fires only once this clears [`SHRINK_COOLDOWN`] *and* the releasable
|
||
/// span clears [`SHRINK_THRESHOLD`]. Single-writer Relaxed.
|
||
parks_since_shrink: AtomicU32,
|
||
/// RFC 019: shrinks performed on this incarnation (introspection lands
|
||
/// with the RFC's introspect surface; the counter exists from birth so
|
||
/// tests can rely on install resetting it). Single-writer Relaxed.
|
||
shrink_count: AtomicU32,
|
||
/// RFC 019 §7 — stack geometry for the SIGSEGV classifier, readable
|
||
/// without the cold lock (the `Stack` itself lives under it). Written in
|
||
/// `install_actor` before the Release publish; consulted by the handler
|
||
/// only while `preempt::CURRENT_SLOT` points here, i.e. while this actor
|
||
/// is on-CPU, so the values are never stale where they are read. 0 =
|
||
/// never installed. Usable top of the stack.
|
||
pub(crate) diag_stack_top: AtomicUsize,
|
||
/// See `diag_stack_top`: the reserve (usable) size.
|
||
pub(crate) diag_stack_reserve: AtomicUsize,
|
||
/// See `diag_stack_top`: the guard size.
|
||
pub(crate) diag_stack_guard: AtomicUsize,
|
||
/// See `diag_stack_top`: `(idx << 32) | generation`, for the message.
|
||
pub(crate) diag_pid: AtomicU64,
|
||
/// Pointer into the actor's `Arc<AtomicBool>` stop flag. Set at spawn,
|
||
/// nulled at finalize. The box outlives every read: it is only ever read
|
||
/// on the resume path while the actor cannot be finalized (it is on-CPU).
|
||
stop_ptr: AtomicPtr<AtomicBool>,
|
||
/// First-resume closure, double-boxed so it fits an `AtomicPtr`
|
||
/// (`Box<Closure>` is a thin pointer). Swap-to-take; null when absent.
|
||
closure: AtomicPtr<Closure>,
|
||
/// RFC 016 Chunk 2 — per-actor timeslice overrun tally. Single-writer: only
|
||
/// the on-CPU actor's scheduler thread increments it (at the slice-expiry
|
||
/// site in `preempt.rs`, reached via the stashed slot pointer), so the
|
||
/// writes are plain Relaxed load+store with no atomic-RMW traffic; the
|
||
/// snapshot reads it Relaxed from any thread. Lives in the hot region rather
|
||
/// than `SlotCold` so the increment needs no lock; reset across reuse like
|
||
/// every other slot field (`vacant` / `reclaim_slot` / `install_actor`).
|
||
overruns: AtomicU64,
|
||
/// RFC 016 Chunk 2 — messages this actor has received (dequeued). Same
|
||
/// single-writer discipline as `overruns`: only the receiving actor, on its
|
||
/// own thread, increments it on the receive path (D4/D5), Relaxed load+store
|
||
/// with no RMW; snapshot reads Relaxed.
|
||
messages_received: AtomicU64,
|
||
/// RFC 016 Chunk 2 — cumulative on-CPU cycles this incarnation has consumed
|
||
/// (the cycle-accurate "budget used", an analogue of OTP reductions).
|
||
/// Written only by the scheduler thread that ran the actor, once per resume,
|
||
/// and only when the `budget-accounting` feature is on (it costs two RDTSC
|
||
/// per resume); stays 0 otherwise. Same single-writer Relaxed discipline.
|
||
budget_cycles: AtomicU64,
|
||
/// RFC 007 (`smarm-causal`) — id of the causal-profiling site this actor
|
||
/// is currently inside (0 = none). Lives in the slot, not a thread-local,
|
||
/// so it survives preemption and cross-scheduler migration. Written only
|
||
/// by the actor itself (guard enter/exit on its own thread), read by that
|
||
/// same thread in `maybe_preempt` — single-writer Relaxed, like `overruns`.
|
||
/// Exists regardless of the feature; stays 0 without it.
|
||
causal_site: AtomicU32,
|
||
/// RFC 007 (`smarm-causal`) — virtual-speedup delay cycles this actor has
|
||
/// absorbed (or been credited). Compared against the global ledger in
|
||
/// `causal::check`; fast-forwarded on resume-from-park so blocked time
|
||
/// absorbs delay for free (Coz's blocked-thread rule). Same single-writer
|
||
/// discipline.
|
||
causal_delay: AtomicU64,
|
||
/// RFC 007 (`smarm-causal`) — true iff this actor's last deschedule was a
|
||
/// real park (a successful `park_return`, not a slice yield and not the
|
||
/// instant-wake re-queue). Consumed by `causal::on_resume`: only a wake
|
||
/// from genuine blocking forgives outstanding virtual delay; a merely
|
||
/// preempted (runnable) actor stays in debt and must pay by spinning.
|
||
/// Starts false: a fresh spawn is born current (`reset_counters` sets
|
||
/// `causal_delay` to the global ledger), and anything injected while it
|
||
/// sits spawn-queued is owed, not waived. Written by the owning
|
||
/// scheduler thread at the deschedule /
|
||
/// resume boundary only — same single-writer discipline as `causal_delay`.
|
||
causal_parked: AtomicBool,
|
||
/// RFC 007 — TSC at this actor's last *runnable* (yield) deschedule
|
||
/// while it sat in the live experiment's target site; 0 = no gap
|
||
/// pending. Paired with `causal_desched_epoch`; written on the
|
||
/// deschedule path, consumed at the next resume — same single-writer
|
||
/// discipline as `causal_delay`.
|
||
causal_desched_tsc: AtomicU64,
|
||
/// RFC 007 — experiment epoch live at that deschedule (see
|
||
/// `causal::EXPERIMENT_EPOCH`): the resume counts the gap only into
|
||
/// the same window, so one straddling `end()` — or a later `begin()`
|
||
/// with an identical site+pct word — is dropped instead of leaking a
|
||
/// cooldown across windows.
|
||
causal_desched_epoch: AtomicU64,
|
||
/// Cold lifecycle data. See [`SlotCold`].
|
||
pub(crate) cold: RawMutex<SlotCold>,
|
||
}
|
||
|
||
impl Slot {
|
||
fn vacant() -> Self {
|
||
Self {
|
||
word: StateWord::new(),
|
||
sp: AtomicUsize::new(0),
|
||
hwm: AtomicUsize::new(0),
|
||
parks_since_shrink: AtomicU32::new(0),
|
||
shrink_count: AtomicU32::new(0),
|
||
diag_stack_top: AtomicUsize::new(0),
|
||
diag_stack_reserve: AtomicUsize::new(0),
|
||
diag_stack_guard: AtomicUsize::new(0),
|
||
diag_pid: AtomicU64::new(0),
|
||
stop_ptr: AtomicPtr::new(std::ptr::null_mut()),
|
||
closure: AtomicPtr::new(std::ptr::null_mut()),
|
||
overruns: AtomicU64::new(0),
|
||
messages_received: AtomicU64::new(0),
|
||
budget_cycles: AtomicU64::new(0),
|
||
causal_site: AtomicU32::new(0),
|
||
causal_delay: AtomicU64::new(0),
|
||
// Overwritten at every occupancy by `reset_counters` (born
|
||
// current); false so a hypothetical reset-skipping path cannot
|
||
// waive the whole global backlog at first resume.
|
||
causal_parked: AtomicBool::new(false),
|
||
causal_desched_tsc: AtomicU64::new(0),
|
||
causal_desched_epoch: AtomicU64::new(0),
|
||
cold: RawMutex::new(SlotCold {
|
||
actor: None,
|
||
waiters: Vec::new(),
|
||
outcome: None,
|
||
terminal: None,
|
||
watchable: false,
|
||
supervisor_channel: None,
|
||
monitors: Vec::new(),
|
||
links: Vec::new(),
|
||
outstanding_handles: 0,
|
||
pending_io_result: None,
|
||
}),
|
||
}
|
||
}
|
||
|
||
/// Current generation (of whatever occupies the slot — pair with a
|
||
/// status check or a CAS before acting on it).
|
||
#[inline]
|
||
pub(crate) fn generation(&self) -> u32 {
|
||
self.word.generation()
|
||
}
|
||
|
||
/// Raw packed state word, for introspection's lock-free classify
|
||
/// (`introspect.rs`). The coarse `status_for` only distinguishes
|
||
/// Live/Done/Stale; the snapshot needs the fine scheduling state.
|
||
#[inline]
|
||
pub(crate) fn state_word(&self) -> u64 {
|
||
self.word.load()
|
||
}
|
||
|
||
/// Tally one timeslice overrun (RFC 016 Chunk 2). Single-writer: only the
|
||
/// on-CPU actor's own thread calls this, at the slice-expiry site, so a
|
||
/// Relaxed load+store is sufficient and avoids the cache-line lock of an
|
||
/// atomic RMW.
|
||
#[inline]
|
||
pub(crate) fn record_overrun(&self) {
|
||
let v = self.overruns.load(Ordering::Relaxed);
|
||
self.overruns.store(v.wrapping_add(1), Ordering::Relaxed);
|
||
}
|
||
|
||
/// Read the overrun tally (Relaxed; the snapshot reads cross-thread).
|
||
#[inline]
|
||
/// RFC 019 §8 — the stack introspection tuple, all lock-free:
|
||
/// `(reserve, guard, top, hwm, parks_since_shrink, shrink_count)`.
|
||
/// Geometry from the c6 diag atomics (install-time, gen-coherent under
|
||
/// `read_slot`'s gen check exactly like the other counters); `hwm` is the
|
||
/// §2 sampled high-water (lowest saved sp). All zeros before first
|
||
/// install.
|
||
pub(crate) fn stack_introspect(&self) -> (usize, usize, usize, usize, u32, u32) {
|
||
(
|
||
self.diag_stack_reserve.load(Ordering::Relaxed),
|
||
self.diag_stack_guard.load(Ordering::Relaxed),
|
||
self.diag_stack_top.load(Ordering::Relaxed),
|
||
self.hwm.load(Ordering::Relaxed),
|
||
self.parks_since_shrink.load(Ordering::Relaxed),
|
||
self.shrink_count.load(Ordering::Relaxed),
|
||
)
|
||
}
|
||
|
||
pub(crate) fn overruns(&self) -> u64 {
|
||
self.overruns.load(Ordering::Relaxed)
|
||
}
|
||
|
||
/// Tally one received (dequeued) message. Same single-writer Relaxed
|
||
/// discipline as `record_overrun`; called by the receiving actor on its own
|
||
/// thread, so no RMW.
|
||
#[inline]
|
||
pub(crate) fn record_message(&self) {
|
||
let v = self.messages_received.load(Ordering::Relaxed);
|
||
self.messages_received
|
||
.store(v.wrapping_add(1), Ordering::Relaxed);
|
||
}
|
||
|
||
/// Read the received-message tally (Relaxed; cross-thread snapshot read).
|
||
#[inline]
|
||
pub(crate) fn messages_received(&self) -> u64 {
|
||
self.messages_received.load(Ordering::Relaxed)
|
||
}
|
||
|
||
/// Accumulate on-CPU cycles consumed in one resume (RFC 016 Chunk 2,
|
||
/// `budget-accounting`). Single-writer (the scheduler thread that ran the
|
||
/// actor), Relaxed load+store. Approximate by design: the figure is
|
||
/// `now − slice-start`, and a wake-slot resume inherits the waker's slice,
|
||
/// so a handed-off actor is charged a little of the chain's time — noise
|
||
/// that averages out across runs, traded for one RDTSC instead of two.
|
||
#[cfg(feature = "budget-accounting")]
|
||
#[inline]
|
||
pub(crate) fn add_budget(&self, cycles: u64) {
|
||
let v = self.budget_cycles.load(Ordering::Relaxed);
|
||
self.budget_cycles
|
||
.store(v.wrapping_add(cycles), Ordering::Relaxed);
|
||
}
|
||
|
||
/// Read the accumulated budget cycles (Relaxed). Always 0 unless the
|
||
/// `budget-accounting` feature is enabled.
|
||
#[inline]
|
||
pub(crate) fn budget_cycles(&self) -> u64 {
|
||
self.budget_cycles.load(Ordering::Relaxed)
|
||
}
|
||
|
||
/// Zero the per-actor introspection counters. Called at every point a slot
|
||
/// is recycled or freshly occupied (`reclaim_slot`, `install_actor`) so a
|
||
/// reused slot never carries a previous incarnation's counts — the standing
|
||
/// slot-lifecycle reset invariant (RFC 016 D7).
|
||
#[inline]
|
||
pub(crate) fn reset_counters(&self) {
|
||
self.overruns.store(0, Ordering::Relaxed);
|
||
self.messages_received.store(0, Ordering::Relaxed);
|
||
self.budget_cycles.store(0, Ordering::Relaxed);
|
||
self.causal_site.store(0, Ordering::Relaxed);
|
||
// Born current (Coz's new-thread rule): a fresh incarnation neither
|
||
// owes the process's accumulated delay history nor books it as park
|
||
// forgiveness. The old init (delay 0, parked true) waived the whole
|
||
// monotone backlog once per spawn via the first resume — found live
|
||
// under close-mode conn churn (~95k spawns/s): millions of phantom
|
||
// forgiven ms per 700ms window, even in 0% cells. Parked starts
|
||
// false: delay injected while spawn-queued is *owed* (the newborn is
|
||
// runnable, not blocked) and paid at its first check — the semantics
|
||
// `audit_zero_pct_window_absorbs_leftover_debt` pins.
|
||
#[cfg(feature = "smarm-causal")]
|
||
self.causal_delay
|
||
.store(crate::causal::global_delay_cycles(), Ordering::Relaxed);
|
||
#[cfg(not(feature = "smarm-causal"))]
|
||
self.causal_delay.store(0, Ordering::Relaxed);
|
||
self.causal_parked.store(false, Ordering::Relaxed);
|
||
self.causal_desched_tsc.store(0, Ordering::Relaxed);
|
||
self.causal_desched_epoch.store(0, Ordering::Relaxed);
|
||
}
|
||
|
||
/// RFC 007 — mark that this actor's deschedule was a genuine park.
|
||
/// Called from the scheduler's `YieldIntent::Park` branch on a successful
|
||
/// `park_return` only.
|
||
#[cfg_attr(not(feature = "smarm-causal"), allow(dead_code))]
|
||
#[inline]
|
||
pub(crate) fn set_causal_parked(&self) {
|
||
self.causal_parked.store(true, Ordering::Relaxed);
|
||
}
|
||
|
||
/// RFC 007 — consume the parked marker at resume: returns whether the
|
||
/// last deschedule was a real park, and clears it so the next resume
|
||
/// defaults to "was runnable" unless the park branch says otherwise.
|
||
#[cfg_attr(not(feature = "smarm-causal"), allow(dead_code))]
|
||
#[inline]
|
||
pub(crate) fn take_causal_parked(&self) -> bool {
|
||
self.causal_parked.swap(false, Ordering::Relaxed)
|
||
}
|
||
|
||
/// RFC 007 — stash the runnable-deschedule instant and the experiment
|
||
/// epoch it happened under (offcpu-gap audit; see `causal::on_resume`).
|
||
#[cfg_attr(not(feature = "smarm-causal"), allow(dead_code))]
|
||
#[inline]
|
||
pub(crate) fn set_causal_desched(&self, tsc: u64, epoch: u64) {
|
||
self.causal_desched_tsc.store(tsc, Ordering::Relaxed);
|
||
self.causal_desched_epoch.store(epoch, Ordering::Relaxed);
|
||
}
|
||
|
||
/// RFC 007 — consume the stash: `(tsc, epoch)`; `(0, _)` = none pending.
|
||
#[cfg_attr(not(feature = "smarm-causal"), allow(dead_code))]
|
||
#[inline]
|
||
pub(crate) fn take_causal_desched(&self) -> (u64, u64) {
|
||
let tsc = self.causal_desched_tsc.swap(0, Ordering::Relaxed);
|
||
let epoch = self.causal_desched_epoch.load(Ordering::Relaxed);
|
||
(tsc, epoch)
|
||
}
|
||
|
||
/// RFC 007 — current causal site id (0 = none). Single-writer: only the
|
||
/// on-CPU actor's thread writes, via the site-guard enter/exit.
|
||
#[cfg_attr(not(feature = "smarm-causal"), allow(dead_code))]
|
||
#[inline]
|
||
pub(crate) fn causal_site(&self) -> u32 {
|
||
self.causal_site.load(Ordering::Relaxed)
|
||
}
|
||
|
||
/// RFC 007 — write the current causal site id (guard enter/exit).
|
||
#[cfg_attr(not(feature = "smarm-causal"), allow(dead_code))]
|
||
#[inline]
|
||
pub(crate) fn set_causal_site(&self, id: u32) {
|
||
self.causal_site.store(id, Ordering::Relaxed);
|
||
}
|
||
|
||
/// RFC 007 — absorbed/credited virtual-delay cycles.
|
||
#[cfg_attr(not(feature = "smarm-causal"), allow(dead_code))]
|
||
#[inline]
|
||
pub(crate) fn causal_delay(&self) -> u64 {
|
||
self.causal_delay.load(Ordering::Relaxed)
|
||
}
|
||
|
||
/// RFC 007 — set the absorbed-delay ledger (spin-absorb, credit, or the
|
||
/// resume-path fast-forward). Single-writer per the resume protocol: the
|
||
/// actor's own thread while on-CPU, the resuming scheduler thread at the
|
||
/// resume boundary — never both at once.
|
||
#[cfg_attr(not(feature = "smarm-causal"), allow(dead_code))]
|
||
#[inline]
|
||
pub(crate) fn set_causal_delay(&self, v: u64) {
|
||
self.causal_delay.store(v, Ordering::Relaxed);
|
||
}
|
||
|
||
/// A pid's-eye snapshot of the slot. Cold paths re-read this under the
|
||
/// cold lock (generation can't change while it is held).
|
||
#[inline]
|
||
pub(crate) fn status_for(&self, pid: Pid) -> Status {
|
||
self.word.status_for(pid.generation())
|
||
}
|
||
|
||
/// Does the slot currently hold the actor `pid` names, in a non-terminal
|
||
/// state? (Snapshot — callers that mutate must re-verify under `cold` or
|
||
/// CAS on the word.)
|
||
#[inline]
|
||
pub(crate) fn is_live_for(&self, pid: Pid) -> bool {
|
||
self.status_for(pid) == Status::Live
|
||
}
|
||
|
||
fn store_closure(&self, c: Closure) {
|
||
let raw = Box::into_raw(Box::new(c));
|
||
let prev = self.closure.swap(raw, Ordering::Release);
|
||
debug_assert!(prev.is_null(), "slot already had a pending closure");
|
||
}
|
||
|
||
fn take_closure(&self) -> Option<Closure> {
|
||
// Fast path: every resume after the first (the overwhelming case)
|
||
// finds null. A plain load suffices to prove it — `store_closure`
|
||
// runs only before `publish_queued`, whose Release/Acquire pairing
|
||
// with the claimer's `try_claim` orders it before this call, so no
|
||
// writer can race the load within an occupancy. This keeps the
|
||
// locked RMW (full barrier, ~20+ cycles) off the per-resume path.
|
||
if self.closure.load(Ordering::Relaxed).is_null() {
|
||
return None;
|
||
}
|
||
let raw = self.closure.swap(std::ptr::null_mut(), Ordering::Acquire);
|
||
if raw.is_null() {
|
||
None
|
||
} else {
|
||
// SAFETY: non-null values in `closure` are exclusively
|
||
// `Box::into_raw(Box<Closure>)` from `store_closure`, and the
|
||
// swap above made us the unique owner.
|
||
Some(*unsafe { Box::from_raw(raw) })
|
||
}
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// RuntimeInner — the shared core behind an Arc
|
||
// ---------------------------------------------------------------------------
|
||
|
||
pub(crate) struct RuntimeInner {
|
||
/// The run queue, compile-time selected (see `run_queue.rs` for the
|
||
/// contract: ops require preemption disabled, push is infallible,
|
||
/// pop-None is a snapshot).
|
||
pub(crate) run_queue: crate::run_queue::RunQueue,
|
||
/// The fixed actor slot table. Allocated once; slots never move.
|
||
pub(crate) slots: Box<[Slot]>,
|
||
/// Vacant slot indices. RawMutex leaf; never held with a cold lock.
|
||
pub(crate) free: RawMutex<Vec<u32>>,
|
||
/// Spawned-but-not-finalized actor count; the termination criterion.
|
||
/// Incremented in `spawn` before the enqueue; decremented at the very end
|
||
/// of `finalize_actor`, after every wakeup finalize produces.
|
||
pub(crate) live_actors: AtomicU32,
|
||
/// Packed `(index << 32 | generation)` of the run's root (initial) actor,
|
||
/// or `u64::MAX` (the ROOT_PID sentinel) before one is set. When this actor
|
||
/// finalizes, `finalize_actor` runs the root-exit shutdown (see
|
||
/// `shutdown_forest_roots`). Set once per `run()`, right after the initial
|
||
/// spawn.
|
||
pub(crate) root_bits: AtomicU64,
|
||
/// Timer heap. Independent lock: never nested with any other.
|
||
pub(crate) timers: Mutex<Timers>,
|
||
/// IO subsystem. `None` between runs. Lock order: io before everything.
|
||
pub(crate) io: Mutex<Option<IoThread>>,
|
||
/// Monotonic `MonitorId` source. Never reused.
|
||
pub(crate) next_monitor_id: AtomicU64,
|
||
/// 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.
|
||
pub(crate) io_parked: AtomicU32,
|
||
pub(crate) sleeping: AtomicU32,
|
||
/// Preemption knobs, written into each scheduler thread's locals on startup.
|
||
pub(crate) alloc_interval: u32,
|
||
pub(crate) timeslice_cycles: u64,
|
||
/// RFC 005: whether actor-context wakes route through the per-scheduler
|
||
/// wake slot. Read-only after init; one predictable branch per wake.
|
||
pub(crate) wake_slot: bool,
|
||
/// The name <-> pid registry (bidirectional). RawMutex Leaf: never held
|
||
/// with any other lock; liveness checks under it read only the atomic
|
||
/// slot word.
|
||
pub(crate) registry: RawMutex<crate::registry::Registry>,
|
||
/// Runtime identity (RFC 012). Read-only after init — one node, one fixed
|
||
/// incarnation until clustering (RFC 010) supplies real values. Carried
|
||
/// like `wake_slot`; pg fills these into every `Member` so the public
|
||
/// surface stays Pid-shaped.
|
||
pub(crate) node_id: crate::pg::NodeId,
|
||
pub(crate) incarnation: crate::pg::Incarnation,
|
||
/// Process groups: `name -> multiset<Member>` (RFC 012). RawMutex Leaf,
|
||
/// exactly like `registry`: never held with any other lock; liveness
|
||
/// checks under it read only the atomic slot word, and the eviction path
|
||
/// keeps it off the send path.
|
||
pub(crate) process_groups: RawMutex<crate::pg::ProcessGroups>,
|
||
/// Recycled stacks waiting to be reused by the next spawn.
|
||
pub(crate) stack_pool: RawMutex<Vec<crate::stack::Stack>>,
|
||
/// Maximum number of stacks to retain in the pool.
|
||
pub(crate) stack_pool_cap: usize,
|
||
/// Default stack shape (RFC 019), pre-page-rounded so it compares exactly
|
||
/// against `Stack::shape()`. Only stacks of exactly this shape are pooled.
|
||
pub(crate) stack_reserve: usize,
|
||
pub(crate) stack_guard: usize,
|
||
}
|
||
|
||
impl RuntimeInner {
|
||
// Private constructor taking the parsed Config fields one-for-one; a params
|
||
// struct would only move the same 10 values across the call boundary.
|
||
#[allow(clippy::too_many_arguments)]
|
||
fn new(
|
||
thread_count: usize,
|
||
alloc_interval: u32,
|
||
timeslice_cycles: u64,
|
||
stack_pool_cap: usize,
|
||
stack_reserve: usize,
|
||
stack_guard: usize,
|
||
max_actors: usize,
|
||
wake_slot: bool,
|
||
node_id: crate::pg::NodeId,
|
||
incarnation: crate::pg::Incarnation,
|
||
) -> Arc<Self> {
|
||
let stats = (0..thread_count).map(|_| SchedulerStats::new()).collect();
|
||
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,
|
||
free: RawMutex::new(free),
|
||
live_actors: AtomicU32::new(0),
|
||
root_bits: AtomicU64::new(u64::MAX),
|
||
timers: Mutex::new(timers),
|
||
io: Mutex::new(None),
|
||
next_monitor_id: AtomicU64::new(0),
|
||
coord,
|
||
io_outstanding: AtomicU32::new(0),
|
||
io_fd_waiters: AtomicU32::new(0),
|
||
stats,
|
||
io_parked: AtomicU32::new(0),
|
||
sleeping: AtomicU32::new(0),
|
||
alloc_interval,
|
||
timeslice_cycles,
|
||
wake_slot,
|
||
registry: RawMutex::new(crate::registry::Registry::new()),
|
||
node_id,
|
||
incarnation,
|
||
process_groups: RawMutex::new(crate::pg::ProcessGroups::new()),
|
||
stack_pool: RawMutex::new(Vec::new()),
|
||
stack_pool_cap,
|
||
stack_reserve: crate::stack::round_to_pages(stack_reserve),
|
||
stack_guard: crate::stack::round_to_pages(stack_guard),
|
||
})
|
||
}
|
||
|
||
/// Slot lookup by index only — bounds-checked, NOT generation-checked.
|
||
/// `ROOT_PID` (index `u32::MAX`) is out of bounds by construction and
|
||
/// resolves to `None`. Callers verify the generation atomically: either
|
||
/// inside a CAS on the word, or by re-reading the word under the cold lock.
|
||
#[inline]
|
||
pub(crate) fn slot_at(&self, pid: Pid) -> Option<&Slot> {
|
||
self.slots.get(pid.index() as usize)
|
||
}
|
||
|
||
/// Record `pid` as this run's root actor. Called once per `run()`, right
|
||
/// after the initial spawn and before any scheduler thread starts, so no
|
||
/// finalize can observe the count before the root is set.
|
||
#[inline]
|
||
pub(crate) fn set_root(&self, pid: Pid) {
|
||
self.root_bits.store(Self::pack(pid), Ordering::Relaxed);
|
||
}
|
||
|
||
/// Is `pid` (index + generation) this run's root actor?
|
||
#[inline]
|
||
pub(crate) fn is_root(&self, pid: Pid) -> bool {
|
||
self.root_bits.load(Ordering::Relaxed) == Self::pack(pid)
|
||
}
|
||
|
||
#[inline]
|
||
fn pack(pid: Pid) -> u64 {
|
||
((pid.index() as u64) << 32) | pid.generation() as u64
|
||
}
|
||
|
||
/// Push to the run queue. Callers must have just transitioned the pid
|
||
/// into `Queued` (spawn's publish, the unpark protocol, or the
|
||
/// scheduler's yield/notified-park return paths).
|
||
pub(crate) fn enqueue(&self, pid: Pid) {
|
||
// Every push pairs 1:1 with a transition INTO Queued, and nothing can
|
||
// move the word off Queued until this very entry is popped — so the
|
||
// word must read EXACTLY (gen, Queued) here. This is the at-most-once-
|
||
// enqueued invariant the bounded rings' capacity proof leans on.
|
||
debug_assert!(
|
||
self.slot_at(pid).map(|s| s.word.load()).is_some_and(|w| {
|
||
crate::slot_state::word_gen(w) == pid.generation()
|
||
&& crate::slot_state::word_state(w) == crate::slot_state::ST_QUEUED
|
||
}),
|
||
"enqueue of a pid not in (gen, Queued)"
|
||
);
|
||
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.
|
||
if self.coord.wake_one_if_idle() {
|
||
diag!(self, enqueue_wakes);
|
||
}
|
||
}
|
||
|
||
/// Make `pid` runnable if it is parked; coalesce or defer otherwise.
|
||
/// The runtime-internal core of `scheduler::unpark`. WILDCARD wake:
|
||
/// consumes the epoch but does not check it — reserved for terminal
|
||
/// wakes (`request_stop`); see slot_state.rs.
|
||
pub(crate) fn unpark(&self, pid: Pid) {
|
||
self.unpark_inner(pid, None);
|
||
}
|
||
|
||
/// Epoch-matched wake: lands only if `pid`'s current wait is still the
|
||
/// one the waker registered for. The form every registration-based waker
|
||
/// (channel senders, mutex grants, wait-timers, …) must use.
|
||
pub(crate) fn unpark_at(&self, pid: Pid, epoch: u32) {
|
||
self.unpark_inner(pid, Some(epoch));
|
||
}
|
||
|
||
/// Open a new wait for `pid` (the calling actor itself): bump its
|
||
/// park-epoch and return it. See slot_state.rs for the rules.
|
||
#[must_use]
|
||
pub(crate) fn begin_wait(&self, pid: Pid) -> u32 {
|
||
let slot = match self.slot_at(pid) {
|
||
Some(slot) => slot,
|
||
None => panic!("begin_wait: own slot vanished: {:?}", pid),
|
||
};
|
||
slot.word.begin_wait(pid.generation())
|
||
}
|
||
|
||
/// Retire the calling actor's current wait without parking on it: bump
|
||
/// the epoch (invalidating every in-flight registration-based wake),
|
||
/// then eat a notification that already landed. The caller MUST
|
||
/// re-check its stop flag afterwards — see `StateWord::clear_notify`.
|
||
pub(crate) fn retire_wait(&self, pid: Pid) {
|
||
let slot = match self.slot_at(pid) {
|
||
Some(slot) => slot,
|
||
None => panic!("retire_wait: own slot vanished: {:?}", pid),
|
||
};
|
||
let _ = slot.word.begin_wait(pid.generation());
|
||
slot.word.clear_notify(pid.generation());
|
||
}
|
||
|
||
fn unpark_inner(&self, pid: Pid, want: Option<u32>) {
|
||
if let Some(slot) = self.slot_at(pid) {
|
||
match slot.word.unpark(pid.generation(), want) {
|
||
Unpark::Enqueue => {
|
||
crate::te!(crate::trace::Event::UnparkDirect(pid));
|
||
// RFC 005: a wake from ACTOR context is slot-eligible —
|
||
// the woken actor's message bytes are hot in this core's
|
||
// cache. Scheduler-context wakes (timer/IO drain, where
|
||
// current_pid is None) have no locality to exploit,
|
||
// arrive in bursts that would thrash the slot, and
|
||
// concentrate on the drain winner by construction:
|
||
// shared queue, always. Spawns never come through here
|
||
// (install_actor enqueues directly), so they bypass the
|
||
// slot by construction.
|
||
if self.wake_slot && crate::actor::current_pid().is_some() {
|
||
diag!(self, unpark_slot);
|
||
self.slot_push(pid);
|
||
} else {
|
||
diag!(self, unpark_queue);
|
||
self.enqueue(pid);
|
||
}
|
||
}
|
||
Unpark::Notified => {
|
||
diag!(self, unpark_notified);
|
||
crate::te!(crate::trace::Event::UnparkDeferred(pid));
|
||
}
|
||
Unpark::Noop => {}
|
||
}
|
||
}
|
||
}
|
||
|
||
/// RFC 005: park `pid` in this thread's wake slot instead of the shared
|
||
/// queue. Replaces `enqueue` at the tail of the wake protocol's
|
||
/// `Parked → Queued` CAS, so the caller has just transitioned the pid
|
||
/// into `Queued` — same precondition, same invariant, different home.
|
||
/// Displacement: the NEW wake takes the slot (newest is hottest; the old
|
||
/// occupant was about to lose its locality window anyway) and the old
|
||
/// occupant is pushed to the shared queue.
|
||
fn slot_push(&self, pid: Pid) {
|
||
debug_assert!(
|
||
!crate::preempt::PREEMPTION_ENABLED.with(|c| c.get()),
|
||
"slot_push with preemption enabled — a switch mid-op could \
|
||
migrate the actor and split the slot access across threads"
|
||
);
|
||
debug_assert!(
|
||
self.slot_at(pid).map(|s| s.word.load()).is_some_and(|w| {
|
||
crate::slot_state::word_gen(w) == pid.generation()
|
||
&& crate::slot_state::word_state(w) == crate::slot_state::ST_QUEUED
|
||
}),
|
||
"slot_push of a pid not in (gen, Queued)"
|
||
);
|
||
let displaced = WAKE_SLOT.with(|s| s.replace(Some(pid)));
|
||
crate::te!(crate::trace::Event::SlotPush(pid));
|
||
if let Some(old) = displaced {
|
||
SCHED_SLOT.with(|s| {
|
||
self.stats[s.get()]
|
||
.slot_displacements
|
||
.fetch_add(1, Ordering::Relaxed)
|
||
});
|
||
self.enqueue(old);
|
||
}
|
||
}
|
||
|
||
/// Allocate the next process-unique `MonitorId`. Lock-free; monitors are a
|
||
/// cold path but there is no reason to serialize id minting under any lock.
|
||
pub(crate) fn alloc_monitor_id(&self) -> MonitorId {
|
||
MonitorId(self.next_monitor_id.fetch_add(1, Ordering::Relaxed) + 1)
|
||
}
|
||
|
||
/// Pop a vacant slot index, or `None` when the slab is full. The claim
|
||
/// is atomic — a single pop under the free-list lock — so callers get
|
||
/// claim-or-report semantics with no check-then-spawn TOCTOU: whoever
|
||
/// gets `Some` owns that slot, full stop.
|
||
pub(crate) fn try_allocate_slot(&self) -> Option<u32> {
|
||
self.free.lock().pop()
|
||
}
|
||
|
||
/// Return a slot claimed by [`try_allocate_slot`](Self::try_allocate_slot)
|
||
/// that never had an actor installed into it (e.g. stack allocation
|
||
/// panicked between claim and install). NOT for dead actors — their
|
||
/// slots go back through `reclaim_slot`, which handles generation bump,
|
||
/// waiter/monitor/link teardown, and stack recycling.
|
||
pub(crate) fn return_vacant_slot(&self, idx: u32) {
|
||
self.free.lock().push(idx);
|
||
}
|
||
|
||
/// Pop a vacant slot index, or die loudly. The fixed slab is a deliberate
|
||
/// v0.5 simplification (ROADMAP: "Deferred"); the panic names the fix.
|
||
/// Callers that can shed load instead use [`try_allocate_slot`]
|
||
/// (Self::try_allocate_slot) via `scheduler::try_spawn`.
|
||
pub(crate) fn allocate_slot(&self) -> u32 {
|
||
match self.try_allocate_slot() {
|
||
Some(idx) => idx,
|
||
None => panic!(
|
||
"smarm: actor slot table exhausted — {} actors are live \
|
||
simultaneously, which is the configured maximum. \
|
||
Fix: raise the cap at runtime init, e.g. \
|
||
`smarm::init(Config::default().max_actors({}))`. \
|
||
(Slots are ~256 bytes each; the table is allocated up-front.)",
|
||
self.slots.len(),
|
||
self.slots.len() * 2
|
||
),
|
||
}
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Runtime — the public handle
|
||
// ---------------------------------------------------------------------------
|
||
|
||
pub struct Runtime {
|
||
inner: Arc<RuntimeInner>,
|
||
thread_count: usize,
|
||
}
|
||
|
||
/// Initialise the runtime with the given config. Returns a reusable handle.
|
||
pub fn init(config: Config) -> Runtime {
|
||
// RFC 019 §7: one process-global SIGSEGV handler, installed before any
|
||
// scheduler thread (and so before any classifiable fault) can exist.
|
||
crate::signal::install_once();
|
||
let n = config.resolved_thread_count();
|
||
Runtime {
|
||
inner: RuntimeInner::new(
|
||
n,
|
||
config.alloc_interval,
|
||
config.timeslice_cycles,
|
||
config.stack_pool_cap,
|
||
config.stack_reserve,
|
||
config.stack_guard,
|
||
config.max_actors,
|
||
config.wake_slot,
|
||
config.node_id,
|
||
config.incarnation,
|
||
),
|
||
thread_count: n,
|
||
}
|
||
}
|
||
|
||
impl Runtime {
|
||
/// Run `f` as the initial actor, block until all actors finish.
|
||
/// Can be called multiple times sequentially on the same `Runtime`.
|
||
///
|
||
/// A panic in the initial actor propagates out of `run` (re-raised
|
||
/// after teardown, so the `Runtime` stays reusable even under a
|
||
/// caller's `catch_unwind`). Panics in other actors surface through
|
||
/// joins, links, monitors, and supervisors as usual.
|
||
pub fn run(&self, f: impl FnOnce() + Send + 'static) {
|
||
// Install smarm's panic hook on first call. The default Rust hook is
|
||
// not reentrant — concurrent actor panics can trigger a double-panic
|
||
// abort when the backtrace printer takes an internal lock that is
|
||
// already held. smarm catches every actor panic via `catch_unwind` in
|
||
// the trampoline, so panics never need to reach the hook for runtime
|
||
// correctness; the hook fires only as a side-effect of unwinding before
|
||
// `catch_unwind` catches it.
|
||
//
|
||
// We install once and leave it installed: the previous hook is chained
|
||
// so that panics outside actor context (e.g. in the test harness
|
||
// itself) are still reported normally.
|
||
static HOOK_INSTALLED: std::sync::OnceLock<()> = std::sync::OnceLock::new();
|
||
HOOK_INSTALLED.get_or_init(|| {
|
||
let prev = std::panic::take_hook();
|
||
std::panic::set_hook(Box::new(move |info| {
|
||
// If we are currently executing inside an actor trampoline the
|
||
// panic will be caught by `catch_unwind` momentarily. Suppress
|
||
// the hook output to avoid interleaved noise and reentrancy.
|
||
// Outside actor context, delegate to the previous hook so that
|
||
// genuine runtime panics are still reported.
|
||
if crate::actor::current_pid().is_some() {
|
||
// Inside an actor — catch_unwind handles it; stay silent.
|
||
} else {
|
||
prev(info);
|
||
}
|
||
}));
|
||
});
|
||
|
||
// Open the trace store for this run (no-op without smarm-trace).
|
||
#[cfg(feature = "smarm-trace")]
|
||
crate::trace::open();
|
||
|
||
// Re-initialise shared state for this run.
|
||
assert_eq!(
|
||
self.inner.run_queue.len(),
|
||
0,
|
||
"run() called while previous run still active"
|
||
);
|
||
debug_assert_eq!(
|
||
self.inner.live_actors.load(Ordering::Acquire),
|
||
0,
|
||
"run() called while previous run still active"
|
||
);
|
||
// 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}"),
|
||
};
|
||
match self.inner.io.lock() {
|
||
Ok(mut io) => *io = Some(io_thread),
|
||
Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"),
|
||
}
|
||
|
||
// RFC 005: slot counters reset at the START of a run (not the end),
|
||
// so `stats()` read after `run()` returns reports that run's totals.
|
||
for stat in &self.inner.stats {
|
||
stat.reset();
|
||
}
|
||
|
||
// Spawn the initial actor through the public spawn path (which
|
||
// requires a running runtime in the thread-local).
|
||
RUNTIME.with(|r| *r.borrow_mut() = Some(self.inner.clone()));
|
||
let initial_handle = crate::scheduler::spawn(f);
|
||
// The initial actor is the run's root: its exit means "the program is
|
||
// done" — every remaining top-level actor is asked to shut down (see
|
||
// `finalize_actor` / `shutdown_forest_roots`).
|
||
self.inner.set_root(initial_handle.pid());
|
||
|
||
// Launch N-1 extra scheduler threads, named `smarm-sched-{slot}` so
|
||
// they are identifiable in `/proc/<pid>/task/*/comm`, stack dumps and
|
||
// debuggers. The calling thread is thread 0 and keeps its caller-given
|
||
// name (an embedder typically names it when spawning `run`).
|
||
let mut os_threads = Vec::new();
|
||
for slot in 1..self.thread_count {
|
||
let inner = self.inner.clone();
|
||
let t = thread::Builder::new()
|
||
.name(format!("smarm-sched-{slot}"))
|
||
.spawn(move || {
|
||
RUNTIME.with(|r| *r.borrow_mut() = Some(inner.clone()));
|
||
SCHED_SLOT.with(|s| s.set(slot));
|
||
schedule_loop(&inner, slot);
|
||
RUNTIME.with(|r| *r.borrow_mut() = None);
|
||
});
|
||
// `thread::spawn` (the previous form) also panics when the OS
|
||
// refuses a thread, so this keeps the failure semantics.
|
||
let t = match t {
|
||
Ok(t) => t,
|
||
Err(e) => panic!("failed to spawn smarm scheduler thread: {e}"),
|
||
};
|
||
os_threads.push(t);
|
||
}
|
||
|
||
// Thread 0 runs the loop on the calling thread.
|
||
SCHED_SLOT.with(|s| s.set(0));
|
||
schedule_loop(&self.inner, 0);
|
||
|
||
// Wait for all other scheduler threads.
|
||
for t in os_threads {
|
||
let _ = t.join();
|
||
}
|
||
|
||
// The root's outcome, read before its handle drops (the outstanding
|
||
// handle pins the slot: no reclaim, no generation change under us).
|
||
// A root panic must escape `run()` — propagated at the tail below,
|
||
// after teardown. Dropping it unread here made every assert inside
|
||
// `run` silently vacuous (found live: a failing-first test passed).
|
||
let root_outcome = {
|
||
let slot = match self.inner.slot_at(initial_handle.pid()) {
|
||
Some(slot) => slot,
|
||
None => panic!("run(): root pid out of range"),
|
||
};
|
||
debug_assert!(
|
||
slot.status_for(initial_handle.pid()) == Status::Done,
|
||
"root not Done at run() teardown"
|
||
);
|
||
slot.cold.lock().outcome.take()
|
||
};
|
||
|
||
// Drop initial handle (decrements outstanding_handles count).
|
||
drop(initial_handle);
|
||
|
||
// Tear down IO and clean up for the next run() call.
|
||
match self.inner.io.lock() {
|
||
Ok(mut io) => drop(io.take()), // joins IO threads
|
||
Err(e) => panic!("smarm: io lock poisoned (core corrupt): {e}"),
|
||
}
|
||
match self.inner.timers.lock() {
|
||
Ok(mut timers) => timers.clear(),
|
||
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
|
||
}
|
||
self.inner.next_monitor_id.store(0, Ordering::Relaxed);
|
||
// Every slot must have come back: any leak here is a runtime bug
|
||
// (a JoinHandle held across run() is decremented just above).
|
||
debug_assert_eq!(
|
||
self.inner.free.lock().len(),
|
||
self.inner.slots.len(),
|
||
"slot leak across run()"
|
||
);
|
||
// Reset per-thread stats.
|
||
for stat in &self.inner.stats {
|
||
stat.current_pid_index.store(u32::MAX, Ordering::Relaxed);
|
||
stat.run_queue_len.store(0, Ordering::Relaxed);
|
||
}
|
||
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);
|
||
|
||
// Flush trace to disk (no-op without smarm-trace).
|
||
#[cfg(feature = "smarm-trace")]
|
||
crate::trace::flush();
|
||
|
||
// Propagate a root panic — last, after the full teardown above, so
|
||
// a `catch_unwind` around `run()` leaves the `Runtime` reusable.
|
||
// `Exit` and `Stopped` return normally: a cooperatively stopped
|
||
// root is not an error. `resume_unwind` re-raises the original
|
||
// payload; the throw-site hook output was suppressed in-actor, so
|
||
// the caller sees the payload message without the origin file:line.
|
||
if let Some(Outcome::Panic(payload)) = root_outcome {
|
||
// Surface the message: the throw-site hook output was suppressed
|
||
// in-actor, so without this a harness shows a bare FAILED.
|
||
let msg: Option<&str> = payload
|
||
.downcast_ref::<&'static str>()
|
||
.copied()
|
||
.or_else(|| payload.downcast_ref::<String>().map(String::as_str));
|
||
eprintln!(
|
||
"smarm: root actor panicked: {}",
|
||
msg.unwrap_or("<non-string panic payload>")
|
||
);
|
||
std::panic::resume_unwind(payload);
|
||
}
|
||
}
|
||
|
||
/// Snapshot of runtime statistics for introspection / tests.
|
||
pub fn stats(&self) -> RuntimeStats {
|
||
RuntimeStats {
|
||
inner: self.inner.clone(),
|
||
}
|
||
}
|
||
|
||
/// A `Send + Sync` handle to this runtime, usable from any thread —
|
||
/// including threads that are not smarm schedulers (an OS-signal handler
|
||
/// thread, an external event source). Grab it *before* [`run`](Self::run)
|
||
/// and hand it to e.g. a signal thread; that thread can then
|
||
/// [`request_stop`](RuntimeHandle::request_stop) the runtime's top
|
||
/// supervisor to drive an ordered shutdown from outside the runtime.
|
||
///
|
||
/// The in-runtime primitives ([`scheduler::request_stop`](crate::request_stop)
|
||
/// and friends) reach the runtime through a thread-local that is unset on
|
||
/// any non-scheduler thread, so they are silent no-ops off-runtime; this
|
||
/// handle carries its own reference and closes that gap.
|
||
pub fn handle(&self) -> RuntimeHandle {
|
||
RuntimeHandle {
|
||
inner: Arc::downgrade(&self.inner),
|
||
}
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// RuntimeHandle — off-runtime wake/stop
|
||
// ---------------------------------------------------------------------------
|
||
|
||
/// A `Send + Sync` handle to a [`Runtime`], obtained from
|
||
/// [`Runtime::handle`]. Lets a thread that is *not* a smarm scheduler thread
|
||
/// drive a cooperative stop into the runtime — the off-runtime counterpart to
|
||
/// [`scheduler::request_stop`](crate::request_stop).
|
||
///
|
||
/// Holds a [`Weak`] to the runtime, for the same reason the IO backend does
|
||
/// (RFC 018): a lingering handle can never keep the runtime's slot table alive
|
||
/// and can never block [`Runtime::run`] from finishing. Once the `Runtime` is
|
||
/// dropped every method is a harmless no-op — the same end state as calling
|
||
/// `request_stop` on an actor that has already exited.
|
||
#[derive(Clone)]
|
||
pub struct RuntimeHandle {
|
||
inner: Weak<RuntimeInner>,
|
||
}
|
||
|
||
impl RuntimeHandle {
|
||
/// Ask `pid` to stop cooperatively, from any thread. The off-runtime
|
||
/// equivalent of [`scheduler::request_stop`](crate::request_stop): it sets
|
||
/// the target's stop flag and wakes it, so a parked actor unwinds at its
|
||
/// next checkpoint exactly as it would for an in-runtime stop. A no-op if
|
||
/// the runtime has been dropped, or if the actor has already exited.
|
||
pub fn request_stop<A>(&self, pid: Pid<A>) {
|
||
let pid = pid.erase();
|
||
// Upgrade the Weak per call, like the IO backend does (io.rs): a live
|
||
// runtime yields the inner and we drive the same stop the in-runtime
|
||
// path would; a dropped runtime makes this a no-op.
|
||
if let Some(inner) = self.inner.upgrade() {
|
||
crate::scheduler::request_stop_inner(&inner, pid);
|
||
}
|
||
}
|
||
|
||
/// Ask `pid` to shut down gracefully, from any thread. The off-runtime
|
||
/// equivalent of [`scheduler::request_shutdown`](crate::request_shutdown);
|
||
/// the delivered [`ExitSignal`](crate::ExitSignal) carries `from ==
|
||
/// ROOT_PID`, since no actor made the request. A no-op if the runtime has
|
||
/// been dropped, or if the actor has already exited.
|
||
pub fn request_shutdown<A>(&self, pid: Pid<A>) {
|
||
let pid = pid.erase();
|
||
if let Some(inner) = self.inner.upgrade() {
|
||
crate::scheduler::request_shutdown_inner(&inner, pid, ROOT_PID);
|
||
}
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Thread-locals
|
||
// ---------------------------------------------------------------------------
|
||
|
||
use std::cell::{Cell, RefCell};
|
||
|
||
thread_local! {
|
||
/// The RuntimeInner for the current run(). Set by run() on the calling
|
||
/// thread and by each spawned scheduler thread.
|
||
pub(crate) static RUNTIME: RefCell<Option<Arc<RuntimeInner>>> =
|
||
const { RefCell::new(None) };
|
||
|
||
/// This scheduler thread's index into RuntimeInner::stats.
|
||
static SCHED_SLOT: Cell<usize> = const { Cell::new(0) };
|
||
|
||
/// RFC 005: the per-scheduler wake slot — a capacity-one wake cache
|
||
/// checked before the shared run queue. Holds a pid in state
|
||
/// `(gen, Queued)` exactly as a shared-queue entry would; the
|
||
/// at-most-once-enqueued invariant reads "in (slot ⊕ shared queue) at
|
||
/// most once". All access is from the owning thread with preemption
|
||
/// disabled (the existing queue-op contract), so plain Cell ops suffice:
|
||
/// no atomics, nothing to steal, nothing to model. Empty whenever the
|
||
/// thread reaches the idle or termination path (pop order drains it
|
||
/// first), so it never holds a pid across the end of a run.
|
||
static WAKE_SLOT: Cell<Option<Pid>> = const { Cell::new(None) };
|
||
|
||
/// What the actor wants when it yields back to the scheduler.
|
||
static YIELD_INTENT: Cell<YieldIntent> = const { Cell::new(YieldIntent::Yield) };
|
||
}
|
||
|
||
#[derive(Copy, Clone)]
|
||
pub(crate) enum YieldIntent {
|
||
Yield,
|
||
Park,
|
||
}
|
||
|
||
pub(crate) fn set_yield_intent(i: YieldIntent) {
|
||
YIELD_INTENT.with(|c| c.set(i));
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Sentinel root PID
|
||
// ---------------------------------------------------------------------------
|
||
|
||
/// Index `u32::MAX` is out of bounds for any slab (Config asserts
|
||
/// `max_actors < u32::MAX`), so every slot lookup on ROOT_PID resolves to
|
||
/// `None` — the root "actor" silently absorbs supervisor signals.
|
||
pub const ROOT_PID: Pid = Pid::new(u32::MAX, u32::MAX);
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Spawn-side slot installation
|
||
// ---------------------------------------------------------------------------
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Stack shrink — RFC 019 §3 (park path only)
|
||
// ---------------------------------------------------------------------------
|
||
|
||
/// The per-park shrink check. Called from the `YieldIntent::Park` arm inside
|
||
/// the owned window (see the assert-comment at the call site). Fast path —
|
||
/// no spike since the last shrink — is two Relaxed loads, a compare, and the
|
||
/// park counter bump, all on the slot line the context-save just wrote.
|
||
///
|
||
/// On a shrink: `MADV_FREE` the whole pages of `[hwm, sp − redzone)` (the
|
||
/// inward-rounded range from [`crate::stack::shrink_range`]), then reset
|
||
/// `hwm = sp` and the park counter. MADV_FREE only *marks*: the kernel
|
||
/// reclaims under pressure, skips re-dirtied pages, and refaults zero pages
|
||
/// for writes after reclaim — so an over-eager mark costs a cancel-write,
|
||
/// never data.
|
||
fn maybe_shrink_stack(slot: &Slot) {
|
||
let parks = slot
|
||
.parks_since_shrink
|
||
.load(Ordering::Relaxed)
|
||
.saturating_add(1);
|
||
slot.parks_since_shrink.store(parks, Ordering::Relaxed);
|
||
|
||
let sp = slot.sp.load(Ordering::Relaxed);
|
||
let hwm = slot.hwm.load(Ordering::Relaxed);
|
||
if sp.wrapping_sub(hwm) < SHRINK_THRESHOLD || sp < hwm {
|
||
return; // common case: nothing worth a syscall
|
||
}
|
||
if parks < SHRINK_COOLDOWN {
|
||
return;
|
||
}
|
||
let page = crate::stack::page_size();
|
||
if let Some((addr, len)) = crate::stack::shrink_range(hwm, sp, page) {
|
||
// Advisory: on the (kernel-config) chance MADV_FREE is unsupported,
|
||
// failing silently degrades to "never shrinks", which is correct.
|
||
unsafe {
|
||
libc::madvise(addr as *mut libc::c_void, len, libc::MADV_FREE);
|
||
}
|
||
slot.hwm.store(sp, Ordering::Relaxed);
|
||
slot.parks_since_shrink.store(0, Ordering::Relaxed);
|
||
slot.shrink_count.store(
|
||
slot.shrink_count.load(Ordering::Relaxed).saturating_add(1),
|
||
Ordering::Relaxed,
|
||
);
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Stack acquisition / recycling — RFC 019 pool rule
|
||
// ---------------------------------------------------------------------------
|
||
|
||
/// Get a stack of the shape `opts` requests (`None` fields ⇒ the runtime
|
||
/// defaults).
|
||
///
|
||
/// Pool rule (RFC 019 §1): the pool is a uniform `Vec<Stack>` of
|
||
/// default-shaped stacks and stays that way. Default-shaped requests try the
|
||
/// pool first; custom shapes always mmap fresh (and `recycle_stack` never
|
||
/// admits them, so a pooled stack is default-shaped by induction). The pool
|
||
/// lock is dropped before any mmap: no syscall ever stalls another spawner.
|
||
pub(crate) fn acquire_stack(
|
||
inner: &RuntimeInner,
|
||
opts: crate::scheduler::SpawnOpts,
|
||
) -> crate::stack::Stack {
|
||
let reserve = opts.stack_reserve.unwrap_or(inner.stack_reserve);
|
||
let guard = opts.guard_size.unwrap_or(inner.stack_guard);
|
||
let default_shaped = crate::stack::round_to_pages(reserve) == inner.stack_reserve
|
||
&& crate::stack::round_to_pages(guard) == inner.stack_guard;
|
||
if default_shaped {
|
||
if let Some(stack) = inner.stack_pool.lock().pop() {
|
||
return stack;
|
||
}
|
||
}
|
||
match crate::stack::Stack::new(reserve, guard) {
|
||
Ok(stack) => stack,
|
||
Err(e) => panic!("stack allocation failed: {e}"),
|
||
}
|
||
}
|
||
|
||
/// Return a dead actor's stack: pooled if default-shaped and under cap,
|
||
/// otherwise dropped here → munmap (custom shapes and cap overflow alike).
|
||
pub(crate) fn recycle_stack(inner: &RuntimeInner, stack: crate::stack::Stack) {
|
||
if stack.shape() == (inner.stack_reserve, inner.stack_guard) {
|
||
// RFC 019 §6: zap the dead spike before pooling, BEFORE taking the
|
||
// pool lock — acquire_stack's invariant is that no syscall ever
|
||
// stalls another spawner under it. On the rare cap-overflow the zap
|
||
// is wasted work ahead of the munmap; harmless, and cheaper than a
|
||
// second lock round-trip to find out.
|
||
stack.recycle_zap(RECYCLE_RETAIN);
|
||
let mut pool = inner.stack_pool.lock();
|
||
if pool.len() < inner.stack_pool_cap {
|
||
pool.push(stack);
|
||
}
|
||
// else: fall through — drop → munmap.
|
||
}
|
||
// Custom-shaped (or cap overflow): `stack` drops here → munmap.
|
||
}
|
||
|
||
/// Install a freshly spawned actor into the slot `idx` (which must have come
|
||
/// from `allocate_slot`) and publish it as Queued. Returns the new `Pid`.
|
||
/// Called by `scheduler::spawn_under`; lives here next to its inverse
|
||
/// (`reclaim_slot`) so the lifecycle is in one file.
|
||
pub(crate) fn install_actor(
|
||
inner: &RuntimeInner,
|
||
idx: u32,
|
||
sp: usize,
|
||
stack: crate::stack::Stack,
|
||
supervisor: Pid,
|
||
closure: Closure,
|
||
monitor: Option<(crate::monitor::MonitorId, crate::channel::Sender<crate::monitor::Down>)>,
|
||
) -> Pid {
|
||
let slot = &inner.slots[idx as usize];
|
||
let gen = slot.generation(); // stable: we own the vacant slot via the free list
|
||
let pid = Pid::new(idx, gen);
|
||
|
||
let stop = Arc::new(AtomicBool::new(false));
|
||
// RFC 019 §7: geometry for the SIGSEGV classifier, captured before the
|
||
// Stack moves under the cold lock. Ordered before readers by the
|
||
// publish below.
|
||
let (diag_reserve, diag_guard) = stack.shape();
|
||
let diag_top = stack.top() as usize;
|
||
slot.stop_ptr
|
||
.store(Arc::as_ptr(&stop) as *mut _, Ordering::Release);
|
||
{
|
||
let mut cold = slot.cold.lock();
|
||
debug_assert!(cold.actor.is_none(), "install over live actor");
|
||
debug_assert!(cold.waiters.is_empty() && cold.monitors.is_empty() && cold.links.is_empty());
|
||
cold.actor = Some(Actor {
|
||
pid,
|
||
stack,
|
||
supervisor,
|
||
stop,
|
||
trap: None,
|
||
});
|
||
cold.outstanding_handles = 1;
|
||
cold.outcome = None;
|
||
cold.pending_io_result = None;
|
||
// `spawn_monitor`: register before publish, so no scheduler can run
|
||
// (and finalize) the child before its monitor exists. Same slot as a
|
||
// `monitor()` registration; the send-from-finalize path is unchanged.
|
||
if let Some(m) = monitor {
|
||
cold.monitors.push(m);
|
||
}
|
||
}
|
||
slot.sp.store(sp, Ordering::Relaxed);
|
||
// RFC 019: a fresh incarnation starts with its high-water at the fresh
|
||
// top-of-stack `sp` and its shrink bookkeeping zeroed.
|
||
slot.hwm.store(sp, Ordering::Relaxed);
|
||
slot.parks_since_shrink.store(0, Ordering::Relaxed);
|
||
slot.shrink_count.store(0, Ordering::Relaxed);
|
||
slot.diag_stack_top.store(diag_top, Ordering::Relaxed);
|
||
slot.diag_stack_reserve
|
||
.store(diag_reserve, Ordering::Relaxed);
|
||
slot.diag_stack_guard.store(diag_guard, Ordering::Relaxed);
|
||
slot.diag_pid
|
||
.store(((idx as u64) << 32) | gen as u64, Ordering::Relaxed);
|
||
slot.store_closure(closure);
|
||
slot.reset_counters();
|
||
inner.live_actors.fetch_add(1, Ordering::Relaxed);
|
||
|
||
// Publish: only now can pops, unparks, or stops find the actor. The
|
||
// Release store orders everything above before any Acquire reader.
|
||
slot.word.publish_queued(gen);
|
||
inner.enqueue(pid);
|
||
crate::te!(crate::trace::Event::Spawn {
|
||
parent: supervisor,
|
||
child: pid
|
||
});
|
||
pid
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Slot reclamation
|
||
// ---------------------------------------------------------------------------
|
||
|
||
/// Reclaim `pid`'s slot if (still) eligible: generation matches, state is
|
||
/// Done, and no handles are outstanding. Safe to call from racing sites
|
||
/// (finalize tail vs. JoinHandle drop): the first caller bumps the
|
||
/// generation under the cold lock, the loser sees the mismatch and no-ops.
|
||
///
|
||
/// Channel senders extracted from the slot are dropped *after* the cold lock
|
||
/// is released — a last-sender drop can unpark a receiver, which takes the
|
||
/// run-queue mutex; legal under a cold lock, but pointless to nest.
|
||
pub(crate) fn reclaim_slot(inner: &RuntimeInner, pid: Pid) {
|
||
let Some(slot) = inner.slot_at(pid) else {
|
||
return;
|
||
};
|
||
let dropped_outside;
|
||
{
|
||
let mut cold = slot.cold.lock();
|
||
if slot.status_for(pid) != Status::Done || cold.outstanding_handles != 0 {
|
||
return; // already reclaimed, or not yet eligible
|
||
}
|
||
debug_assert!(
|
||
cold.actor.is_none(),
|
||
"reclaiming a slot that still owns an actor"
|
||
);
|
||
dropped_outside = (
|
||
cold.outcome.take(),
|
||
cold.supervisor_channel.take(),
|
||
cold.pending_io_result.take(),
|
||
slot.take_closure(), // an actor stopped before first resume
|
||
);
|
||
cold.waiters.clear();
|
||
cold.monitors.clear();
|
||
cold.links.clear();
|
||
cold.watchable = false;
|
||
slot.reset_counters();
|
||
slot.stop_ptr.store(std::ptr::null_mut(), Ordering::Release);
|
||
// The generation bump IS the reclaim: every stale pid is dead from
|
||
// this store onwards (unpark protocol, pops, cold-path re-verifies).
|
||
slot.word.reclaim(pid.generation());
|
||
}
|
||
drop(dropped_outside);
|
||
inner.free.lock().push(pid.index());
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// finalize_actor
|
||
// ---------------------------------------------------------------------------
|
||
|
||
fn finalize_actor(inner: &Arc<RuntimeInner>, pid: Pid, outcome: Outcome) {
|
||
let (joiner_outcome, sup_signal, down_reason) = match outcome {
|
||
Outcome::Exit => (Outcome::Exit, Signal::Exit(pid), DownReason::Exit),
|
||
Outcome::Panic(payload) => (
|
||
Outcome::Panic(payload),
|
||
Signal::Panic(pid, Box::new(()) as Box<dyn std::any::Any + Send>),
|
||
DownReason::Panic,
|
||
),
|
||
// Cooperative cancellation: kept distinct from a normal Exit so a
|
||
// supervisor's await logic (roadmap #2) can tell "I stopped it" apart
|
||
// from "it finished on its own".
|
||
Outcome::Stopped => (Outcome::Stopped, Signal::Stopped(pid), DownReason::Stopped),
|
||
};
|
||
|
||
let slot = match inner.slot_at(pid) {
|
||
Some(slot) => slot,
|
||
None => panic!("finalize_actor: pid out of range: {:?}", pid),
|
||
};
|
||
let (waiters, monitors, links, actor) = {
|
||
let mut cold = slot.cold.lock();
|
||
let actor = match cold.actor.take() {
|
||
Some(actor) => actor,
|
||
None => panic!("finalize_actor: actor vanished"),
|
||
};
|
||
cold.outcome = Some(joiner_outcome);
|
||
// Terminal record (soak sig 4): stamped before the generation ever
|
||
// bumps, under the cold lock, so a reader that resolved this pid can
|
||
// recover the reason after the slot moves on — but only for a
|
||
// tenancy that ever held a name. The free list is LIFO, so the slot
|
||
// this death frees is the very next one recycled; if every green
|
||
// thread's exit stamped too, the churn behind any real workload
|
||
// would evict a watchable tenancy's record in well under the race
|
||
// window this exists to cover. Overwritten only by the slot's next
|
||
// *watchable* death.
|
||
if cold.watchable {
|
||
cold.terminal = Some((pid.generation(), down_reason));
|
||
}
|
||
slot.stop_ptr.store(std::ptr::null_mut(), Ordering::Release);
|
||
// Done is published under the cold lock, so join's
|
||
// check-Done-or-register-waiter (also under it) can never miss: it
|
||
// either sees Done and takes the outcome, or its waiter registration
|
||
// happens before our take() below and is woken further down.
|
||
// (set_done self-asserts the Running|Notified precondition + gen.)
|
||
slot.word.set_done(pid.generation());
|
||
(
|
||
std::mem::take(&mut cold.waiters),
|
||
std::mem::take(&mut cold.monitors),
|
||
std::mem::take(&mut cold.links),
|
||
actor,
|
||
)
|
||
};
|
||
|
||
// Recycle the stack outside the cold lock; drop the rest of the Actor
|
||
// (the trap sender can unpark its receiver — keep that outside too).
|
||
let supervisor_pid = actor.supervisor;
|
||
let Actor { stack, .. } = actor;
|
||
recycle_stack(inner, stack);
|
||
|
||
// Deliver to supervisor. ROOT_PID resolves to no slot → silently absorbed.
|
||
let sender = inner.slot_at(supervisor_pid).and_then(|sup| {
|
||
let cold = sup.cold.lock();
|
||
if sup.generation() == supervisor_pid.generation() {
|
||
cold.supervisor_channel.clone()
|
||
} else {
|
||
None
|
||
}
|
||
});
|
||
if let Some(sender) = sender {
|
||
let _ = sender.send(sup_signal);
|
||
}
|
||
|
||
// Notify monitors. Sent outside any slot lock: `send` may unpark a parked
|
||
// receiver, which takes the run-queue mutex.
|
||
for (_, m) in monitors {
|
||
let _ = m.send(Down {
|
||
pid,
|
||
reason: down_reason,
|
||
});
|
||
}
|
||
|
||
// Walk linked peers ONE AT A TIME (cold locks are leaves). For every
|
||
// peer: remove the back-link to this (dying) actor; on abnormal death,
|
||
// also fetch its trap sender and deliver after unlocking.
|
||
//
|
||
// Acyclicity: the back-link removal happens under the peer's cold lock
|
||
// *before* any stop is delivered to it, so when the peer later dies its
|
||
// own cascade no longer contains us. Two peers finalizing concurrently
|
||
// each find the other already Done (set above, before any cascade) and
|
||
// skip — no ping-pong, no deadlock (never two cold locks held).
|
||
let abnormal = matches!(down_reason, DownReason::Panic | DownReason::Stopped);
|
||
for peer in links {
|
||
let trap = match inner.slot_at(peer) {
|
||
Some(ps) => {
|
||
let mut cold = ps.cold.lock();
|
||
if ps.status_for(peer) == Status::Live {
|
||
cold.links.retain(|p| *p != pid);
|
||
if abnormal {
|
||
Some(cold.actor.as_ref().and_then(|a| a.trap.clone()))
|
||
} else {
|
||
None // normal exit never propagates
|
||
}
|
||
} else {
|
||
None // peer already gone; nothing to do
|
||
}
|
||
}
|
||
None => None,
|
||
};
|
||
match trap {
|
||
Some(Some(tx)) => {
|
||
let _ = tx.send(crate::link::ExitSignal {
|
||
from: pid,
|
||
reason: down_reason,
|
||
});
|
||
}
|
||
Some(None) => crate::scheduler::request_stop(peer),
|
||
None => {}
|
||
}
|
||
}
|
||
|
||
// Unpark joiners (epoch-matched: each registered under the cold lock).
|
||
for (joiner, epoch) in waiters {
|
||
inner.unpark_at(joiner, epoch);
|
||
}
|
||
|
||
// Reclaim if no outstanding handles (re-verified inside).
|
||
reclaim_slot(inner, pid);
|
||
|
||
// Root exit = the program is done. Ask every top-level survivor to shut
|
||
// down, right here, before the live-count decrement below: any wake this
|
||
// produces is then ordered before `live_actors` can be observed at its
|
||
// decremented value, same as every other wakeup finalize issues.
|
||
if inner.is_root(pid) {
|
||
shutdown_forest_roots(inner, pid);
|
||
}
|
||
|
||
// The decrement is LAST: every wakeup this finalize produced (joiners,
|
||
// monitor/trap sends, stop cascades) is enqueued before `live_actors`
|
||
// can be observed at its decremented value. See the termination note in
|
||
// the module docs.
|
||
let prev = inner.live_actors.fetch_sub(1, Ordering::Release);
|
||
debug_assert!(prev >= 1, "live_actors underflow — double finalize");
|
||
}
|
||
|
||
/// The root-exit shutdown. Delivers [`request_shutdown`](crate::request_shutdown)
|
||
/// to every **forest root**: each live actor whose recorded parent
|
||
/// (`Actor::supervisor` — the spawner for a plain `spawn`, the supervisor for
|
||
/// `spawn_under`) is the run itself (`ROOT_PID`) or is no longer live. Actors
|
||
/// under a live parent are not addressed — that parent is responsible for
|
||
/// them: a supervisor traps and runs its ordered, policy-driven shutdown; a
|
||
/// bare parent that dies takes non-trapping children with it via the next
|
||
/// pass of this same rule only if it dies *now*, so a parent that outlives
|
||
/// this scan and later dies leaves its subtree to itself (Erlang semantics: an
|
||
/// unlinked spawn is nobody's child).
|
||
///
|
||
/// Semantics per target follow `request_shutdown`: a trapping actor receives
|
||
/// `ExitSignal { from: root, reason: Shutdown }` and may finish work — drain,
|
||
/// keep its timers ticking, then stop itself; a non-trapping one is stopped
|
||
/// outright. There is no second, forcing sweep: an actor that traps and never
|
||
/// stops keeps the run alive by design (put it under a supervisor with a
|
||
/// `Shutdown::Timeout` policy if that is not wanted). Runs once, on the root's
|
||
/// finalize path, so it races only against actors that are still running —
|
||
/// each `request_shutdown_inner` re-verifies its target under the cold lock,
|
||
/// so a slot that dies or is reused mid-scan is a no-op.
|
||
fn shutdown_forest_roots(inner: &Arc<RuntimeInner>, root: Pid) {
|
||
for idx in 0..inner.slots.len() as u32 {
|
||
let slot = &inner.slots[idx as usize];
|
||
let pid = Pid::new(idx, slot.generation());
|
||
if pid == root {
|
||
continue;
|
||
}
|
||
// Read the parent under the cold lock (generation-verified); act
|
||
// outside it — `request_shutdown_inner` sends and may unpark.
|
||
let parent = {
|
||
let cold = slot.cold.lock();
|
||
if slot.generation() != pid.generation() {
|
||
continue;
|
||
}
|
||
match cold.actor.as_ref() {
|
||
Some(a) => a.supervisor,
|
||
None => continue,
|
||
}
|
||
};
|
||
let parent_live = inner
|
||
.slot_at(parent)
|
||
.is_some_and(|ps| ps.is_live_for(parent));
|
||
if !parent_live {
|
||
// `_probe` so the trace can say what each leftover was: under
|
||
// `smarm-trace` every swept actor is a `root_sweep` line — the
|
||
// visibility that makes "a forgotten actor costs only a slot, never
|
||
// a hung run" a checkable claim rather than a hope.
|
||
let found = crate::scheduler::request_shutdown_inner_probe(inner, pid, root);
|
||
#[cfg_attr(not(feature = "smarm-trace"), allow(unused_variables))]
|
||
if let Some(trapping) = found {
|
||
crate::te!(crate::trace::Event::RootSweep {
|
||
target: pid,
|
||
trapping
|
||
});
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// Timer firing — shared by the busy-path due-check and the timekeeper
|
||
// ---------------------------------------------------------------------------
|
||
|
||
/// 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}"),
|
||
}
|
||
};
|
||
for entry in due {
|
||
match entry.reason {
|
||
// A sleep expiry is just an unpark: the protocol handles
|
||
// every interleaving — Parked (re-queue), Running (the
|
||
// 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::WaitTimeout { target, epoch } => {
|
||
// The callback may call unpark_at itself.
|
||
target.on_timeout(entry.pid, epoch);
|
||
}
|
||
// A `send_after` deadline: run the captured delivery thunk.
|
||
// It resolves the destination through the registry and
|
||
// sends now (a send can unpark a receiver) — same as any
|
||
// other in-loop unpark. The timers lock is already
|
||
// released; lock order Leaf -> Channel is preserved by the
|
||
// send itself. `pop_due` only returns still-armed Sends, so
|
||
// a cancelled one never reaches here.
|
||
crate::timer::Reason::Send { fire } => fire(),
|
||
}
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------------
|
||
// schedule_loop — runs on each scheduler OS thread
|
||
// ---------------------------------------------------------------------------
|
||
|
||
fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
||
// RFC 019 §7: a guard hit leaves no stack to handle the signal on.
|
||
crate::signal::register_altstack();
|
||
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);
|
||
}
|
||
|
||
// ----------------------------------------------------------------
|
||
// 2. Pop a runnable pid. Pop order (RFC 005): wake slot first, then
|
||
// shared queue. The queue mutex covers ONLY the pop; the slot's
|
||
// own atomics carry everything needed to resume.
|
||
// ----------------------------------------------------------------
|
||
enum Pop {
|
||
Got(Pid),
|
||
Idle,
|
||
AllDone,
|
||
}
|
||
|
||
// 2a. RFC 005: drain this thread's wake slot before touching the
|
||
// shared queue. Two consequences fall out of slot-first order:
|
||
// the idle path below is only reachable with an empty slot, and so
|
||
// is AllDone — an occupied slot on ANOTHER thread holds a Queued
|
||
// (hence live) actor, so `live_actors > 0` and termination cannot
|
||
// fire; the counter-first argument is untouched.
|
||
let slot_pid = if inner.wake_slot {
|
||
WAKE_SLOT.with(|s| s.take())
|
||
} else {
|
||
None
|
||
};
|
||
let from_slot = slot_pid.is_some();
|
||
|
||
let pid = if let Some(pid) = slot_pid {
|
||
stats.slot_hits.fetch_add(1, Ordering::Relaxed);
|
||
crate::te!(crate::trace::Event::SlotPop(pid));
|
||
pid
|
||
} else {
|
||
// 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() {
|
||
Some(pid) => Pop::Got(pid),
|
||
None => {
|
||
// Termination does not lean on pop-None being a fence (with
|
||
// the ring queues it is only a snapshot). The argument is
|
||
// counter-first: every queue entry's target stays `Queued` —
|
||
// hence un-finalized, hence counted live — until that very
|
||
// entry is popped. So `live == 0` (Acquire, pairing with
|
||
// finalize's Release decrement, which strictly follows all
|
||
// wakeup enqueues) by itself implies no entry is in, or can
|
||
// ever again enter, the queue: enqueues only target live
|
||
// actors, and a spawner is itself live. The pop-None above
|
||
// is then just the cheap fast-path filter; io_out was read
|
||
// before it per the phase-1 ordering. `live == 0` is also
|
||
// final — no spawn can resurrect the count — so every
|
||
// scheduler thread independently reaches this same verdict.
|
||
let live = inner.live_actors.load(Ordering::Acquire);
|
||
if live == 0 && io_out == 0 {
|
||
Pop::AllDone
|
||
} else {
|
||
Pop::Idle
|
||
}
|
||
}
|
||
};
|
||
|
||
match pop {
|
||
Pop::Got(pid) => pid,
|
||
Pop::AllDone => {
|
||
// Remaining timer entries are orphaned (no live actor can be
|
||
// woken by them — e.g. a sleeper cancelled out of its sleep);
|
||
// they must not keep the runtime alive. Drop them on the way out.
|
||
match inner.timers.lock() {
|
||
Ok(mut timers) => timers.clear(),
|
||
Err(e) => panic!("smarm: timers 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::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}")
|
||
}
|
||
};
|
||
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.
|
||
stats.idle_parks.fetch_add(1, Ordering::Relaxed);
|
||
let pr = 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.coord.deadline_due()
|
||
});
|
||
if matches!(pr, crate::park::ParkResult::WorkFound) {
|
||
stats.idle_recheck_hits.fetch_add(1, Ordering::Relaxed);
|
||
}
|
||
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;
|
||
}
|
||
}
|
||
};
|
||
|
||
// ----------------------------------------------------------------
|
||
// 3. Claim and resume the actor: CAS Queued → Running. A failure
|
||
// means the pid is stale (slot recycled — generation mismatch);
|
||
// 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() {
|
||
stats.chain_wakes.fetch_add(1, Ordering::Relaxed);
|
||
}
|
||
|
||
let slot = match inner.slot_at(pid) {
|
||
Some(s) => s,
|
||
None => continue, // can't happen for real pids; defensive
|
||
};
|
||
if !slot.word.try_claim(pid.generation()) {
|
||
continue; // stale pid: retry immediately (never the idle path)
|
||
}
|
||
crate::te!(crate::trace::Event::Dequeue(pid));
|
||
|
||
let sp = slot.sp.load(Ordering::Relaxed);
|
||
let stop_flag = slot.stop_ptr.load(Ordering::Relaxed);
|
||
// First resume: move the closure into the trampoline's thread-local.
|
||
if let Some(b) = slot.take_closure() {
|
||
set_current_actor_box(b);
|
||
}
|
||
|
||
// Update per-thread stats: record who's on-CPU.
|
||
stats
|
||
.current_pid_index
|
||
.store(pid.index(), Ordering::Relaxed);
|
||
|
||
set_current_pid(pid);
|
||
crate::preempt::set_current_stop(stop_flag);
|
||
crate::preempt::set_current_slot(slot as *const Slot);
|
||
reset_actor_done();
|
||
YIELD_INTENT.with(|c| c.set(YieldIntent::Yield));
|
||
// RFC 005 timeslice inheritance: a slot-popped actor does NOT get a
|
||
// fresh slice — it inherits the waker's remaining one (this thread's
|
||
// TIMESLICE_START/ALLOC_COUNT carry over from the waker's run, with
|
||
// only scheduler bookkeeping in between). A chain of slot handoffs
|
||
// is therefore collectively bounded by one slice, after which
|
||
// preemption fires and the preempt-yield re-enqueue goes to the
|
||
// SHARED queue (a yield is not a wake — never slot-eligible). The
|
||
// shared queue is thus consulted at least once per slice per
|
||
// scheduler: the one-slice starvation bound, zero new counters.
|
||
if !from_slot {
|
||
crate::preempt::reset_timeslice();
|
||
}
|
||
PREEMPTION_ENABLED.with(|c| c.set(true));
|
||
// RFC 007: delay accrued while this actor was off-CPU is absorbed for
|
||
// free (Coz's blocked-thread rule) — fast-forward its ledger and arm
|
||
// the per-thread sample clock before it runs.
|
||
#[cfg(feature = "smarm-causal")]
|
||
crate::causal::on_resume(slot);
|
||
|
||
crate::te!(crate::trace::Event::Resume(pid));
|
||
let saved_sp = unsafe { switch_to_actor(sp) };
|
||
|
||
PREEMPTION_ENABLED.with(|c| c.set(false));
|
||
// RFC 016 Chunk 2: charge the cycles this resume consumed to the actor
|
||
// (approximate; reuses the slice-start timestamp — one RDTSC). Read
|
||
// before the next resume re-arms TIMESLICE_START.
|
||
#[cfg(feature = "budget-accounting")]
|
||
slot.add_budget(crate::preempt::elapsed_slice_cycles());
|
||
stats.current_pid_index.store(u32::MAX, Ordering::Relaxed);
|
||
clear_current_pid();
|
||
crate::preempt::clear_current_stop();
|
||
crate::preempt::clear_current_slot();
|
||
|
||
let intent = YIELD_INTENT.with(|c| c.get());
|
||
slot.sp.store(saved_sp, Ordering::Relaxed);
|
||
// RFC 019 §2: sampled high-water — one branch + at most one store
|
||
// into the line the store above just dirtied. Relaxed and advisory;
|
||
// it piggybacks the existing Relaxed-store-before-Release pattern
|
||
// (mod docs, "Memory ordering") and adds no edges.
|
||
if saved_sp < slot.hwm.load(Ordering::Relaxed) {
|
||
slot.hwm.store(saved_sp, Ordering::Relaxed);
|
||
}
|
||
|
||
if is_actor_done() {
|
||
crate::te!(crate::trace::Event::Done(pid));
|
||
let outcome = take_last_outcome().unwrap_or(Outcome::Exit);
|
||
finalize_actor(inner, pid, outcome);
|
||
} else {
|
||
let gen = pid.generation();
|
||
match intent {
|
||
YieldIntent::Yield => {
|
||
// RFC 007 audit: measure the sample tail this yield drops
|
||
// (near-zero for slice-expiry yields, which sample at the
|
||
// same checkpoint; fat for explicit yield_now in-site).
|
||
#[cfg(feature = "smarm-causal")]
|
||
crate::causal::on_deschedule(slot, false);
|
||
// Running OR RunningNotified → Queued; a notification
|
||
// arriving mid-run coalesces into the re-queue.
|
||
crate::te!(crate::trace::Event::Yield(pid));
|
||
slot.word.yield_return(gen);
|
||
diag!(inner, yield_requeues);
|
||
inner.enqueue(pid);
|
||
}
|
||
YieldIntent::Park => {
|
||
// RFC 019 §3 shrink window (correctness obligation 1):
|
||
// this site sits after the `sp` store above and before
|
||
// the `park_return` Release transition below publishes
|
||
// Parked — the scheduler is on its own stack and the
|
||
// actor is saved but not yet stealable, so the madvise
|
||
// races nothing (belt). MADV_FREE's cancel-on-write is
|
||
// the suspenders: even a racing writer could lose
|
||
// nothing written after the mark, and everything below
|
||
// live `sp` is dead by definition. Runs on BOTH arms of
|
||
// the park_return race — a consumed unpark flag means a
|
||
// wasted-but-harmless madvise on a rare window.
|
||
//
|
||
// This is the ONLY shrink site: the preempt/yield path
|
||
// deliberately never checks (§4's bounded leak under
|
||
// saturation — syscalls must not fire when scheduler
|
||
// cycles are scarcest).
|
||
maybe_shrink_stack(slot);
|
||
if slot.word.park_return(gen) {
|
||
// RFC 007 audit: an in-site park drops its sample
|
||
// tail (nothing flushes it; on_resume re-arms).
|
||
#[cfg(feature = "smarm-causal")]
|
||
crate::causal::on_deschedule(slot, true);
|
||
// RFC 007: a real park — the eventual wake forgives
|
||
// delay accrued while blocked (the instant-wake
|
||
// re-queue below does NOT: that actor never blocked).
|
||
#[cfg(feature = "smarm-causal")]
|
||
slot.set_causal_parked();
|
||
crate::te!(crate::trace::Event::Park(pid));
|
||
} else {
|
||
// An unpark landed in the prep-to-park window; the
|
||
// word is back to Queued — re-queue instead of
|
||
// parking. The lost-wakeup window, closed: this actor
|
||
// never blocked, so its dropped tail counts as a
|
||
// yield in the RFC 007 audit.
|
||
#[cfg(feature = "smarm-causal")]
|
||
crate::causal::on_deschedule(slot, false);
|
||
crate::te!(crate::trace::Event::UnparkFlagConsumed(pid));
|
||
diag!(inner, park_flag_consumed);
|
||
inner.enqueue(pid);
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|