//! 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 { //! slots: Box<[Slot]> ← FIXED slab, max_actors entries, lock-free lookup //! free: RawMutex> ← vacant slot indices //! run_queue: RunQueue ← compile-time selected (src/run_queue.rs) //! timers: Mutex //! io: Mutex> //! live_actors: AtomicU32 ← spawned-but-not-finalized count (termination) //! stats: Vec ← 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 //! closure: AtomicPtr<...> ← first-resume closure, swap-to-take //! cold: RawMutex ← 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 c = Config::default(); /// /// // Exactly 4 scheduler threads: /// let c = Config::exact(4); /// /// // Between 2 and 8, clamped to available parallelism: /// let c = Config::new(2, 8, None); /// ``` #[derive(Clone, Debug)] pub struct Config { min: usize, max: usize, exact: Option, 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: false, 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) -> 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: false, 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: `false` (off until the slot shootout accepts it). 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) -> 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) -> 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: false, 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, } 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), } } } // --------------------------------------------------------------------------- // Runtime stats snapshot (for tests / introspection) // --------------------------------------------------------------------------- pub struct RuntimeStats { pub(crate) inner: Arc, } 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() } } // --------------------------------------------------------------------------- // 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; /// 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, /// Parked joiners as `(pid, park-epoch)`; finalize wakes each via the /// epoch-matched unpark. pub(crate) waiters: Vec<(Pid, u32)>, pub(crate) outcome: Option, /// 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>, /// 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)>, /// 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, pub(crate) outstanding_handles: u32, pub(crate) pending_io_result: Option, } /// 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` 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, /// First-resume closure, double-boxed so it fits an `AtomicPtr` /// (`Box` is a thin pointer). Swap-to-take; null when absent. closure: AtomicPtr, /// 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, } 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 { 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)` 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>, /// 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, /// IO subsystem. `None` between runs. Lock order: io before everything. pub(crate) io: Mutex>, /// 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, /// `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, /// 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, /// 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` (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, /// Recycled stacks waiting to be reused by the next spawn. pub(crate) stack_pool: RawMutex>, /// 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 { 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 = (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. self.coord.wake_one_if_idle(); } /// 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) { 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() { self.slot_push(pid); } else { self.enqueue(pid); } } 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 { 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, 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.slot_hits.store(0, Ordering::Relaxed); stat.slot_displacements.store(0, Ordering::Relaxed); } // 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//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::().map(String::as_str)); eprintln!( "smarm: root actor panicked: {}", msg.unwrap_or("") ); 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, } 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(&self, pid: Pid) { 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(&self, pid: Pid) { 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>> = const { RefCell::new(None) }; /// This scheduler thread's index into RuntimeInner::stats. static SCHED_SLOT: Cell = 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> = const { Cell::new(None) }; /// What the actor wants when it yields back to the scheduler. static YIELD_INTENT: Cell = 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` 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, ) -> 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; } 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, 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), 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, 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, 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, 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. let _ = inner.coord.park(slot_idx, tk_deadline, || { !inner.run_queue.is_empty() || (inner.live_actors.load(Ordering::Acquire) == 0 && inner.io_outstanding.load(Ordering::Acquire) == 0 && inner.io_fd_waiters.load(Ordering::Acquire) == 0) || inner.coord.deadline_due() }); if tk_deadline.is_some() { // Hand the role back BEFORE firing: pop_due can run // `Send` thunks that insert new timers, and the // insert-side re-arm check must see either no // timekeeper (skip) or a real parked one — never us, // awake and about to re-peek anyway. inner.coord.disarm_timer(slot_idx); // Woken for the deadline, for work, or to re-peek // after an earlier insert — fire whatever is due; // the next idle pass re-arms with the new minimum. fire_due_timers(inner, false); } continue; } } }; // ---------------------------------------------------------------- // 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(); } 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); 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)); inner.enqueue(pid); } } } } } }