From e8f02344921c7a50ad6fda600adcae08c64f1985 Mon Sep 17 00:00:00 2001 From: "Claude (sandbox)" Date: Wed, 12 Aug 2026 19:01:07 +0000 Subject: [PATCH] feat(scheduler,runtime): non-panicking try_spawn for at-capacity load shedding MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit allocate_slot() panics on a full slab; for a load-shedding caller (an accept loop spawning one actor per connection) that panic lands in the spawning actor, which then crash-loops under Restart::Transient into the still-full slab until its restart budget is spent — and the service stops accepting entirely. Observed live (urus slowloris scaling, 2026-08-10). A full slab is a routine overload condition for such callers, not an invariant violation. - RuntimeInner::try_allocate_slot() -> Option: the non-panicking core; a single pop under the free-list lock, so the claim is atomic (claim-or-report — no check-then-spawn TOCTOU, no headroom margin). allocate_slot() is now a thin panicking wrapper over it. - scheduler::try_spawn / try_spawn_under_with -> Result: parity with spawn/spawn_under_with except a full slab returns Err(SpawnError::AtCapacity) instead of panicking. Minimal surface per the agreed strategy; the remaining _with/_addr mirrors are trivial wrappers if ever needed. - Slot-first ordering on the try path (reverse of spawn's stack-first): under overload Err is the hot path, and a rejection costs one mutex pop — no mmap/pool-pop + init + recycle per shed unit of work. A drop-guard returns the claimed slot if stack allocation panics in the claim-to-install window (would otherwise leak and trip run()'s teardown slot-leak debug_assert). - SpawnError: non_exhaustive, Display + std::error::Error. - spawn and every existing call site untouched: the panic remains the correct loud invariant check at internal/bounded spawn sites. tests/try_spawn.rs: parity when slots free; exact slab accounting at capacity (Err, no panic, repeatable); custom-shape try refuses before stack allocation; self-heal after slots free; plain spawn still panics (surfaced via JoinError payload); 4-thread race for the last slots claims exactly the free count; SpawnError impl checks. Design doc: smarm-suggestion-try-spawn.md. Downstream consumer change (canned 503 on AtCapacity in urus's accept loop) is urus scope, not smarm. (cherry picked from commit 36de4b36aeaa72b2a5f9f3797b9854652656dcf6) --- src/lib.rs | 5 +- src/runtime.rs | 21 ++++- src/scheduler.rs | 98 ++++++++++++++++++++++++ tests/try_spawn.rs | 185 +++++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 306 insertions(+), 3 deletions(-) create mode 100644 tests/try_spawn.rs diff --git a/src/lib.rs b/src/lib.rs index c108b62..7dfb3f4 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -89,8 +89,9 @@ pub use runtime::{init, Config, Runtime}; pub use scheduler::{ block_on_io, cancel_timer, request_stop, run, self_pid, send_after, send_after_named, send_after_named_wall, send_after_wall, sleep, sleep_wall, spawn, spawn_addr, spawn_addr_with, - spawn_under, spawn_under_with, spawn_with, wait_readable, wait_readable_timeout, wait_writable, - wait_writable_timeout, yield_now, FdArm, JoinError, JoinHandle, SpawnOpts, + spawn_under, spawn_under_with, spawn_with, try_spawn, try_spawn_under_with, wait_readable, + wait_readable_timeout, wait_writable, wait_writable_timeout, yield_now, FdArm, JoinError, + JoinHandle, SpawnError, SpawnOpts, }; pub use supervisor::{ChildSpec, OneForOne, Restart, Signal, Strategy}; pub use timer::TimerId; diff --git a/src/runtime.rs b/src/runtime.rs index 34d205b..be56e5e 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -1188,10 +1188,29 @@ impl RuntimeInner { 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.free.lock().pop() { + match self.try_allocate_slot() { Some(idx) => idx, None => panic!( "smarm: actor slot table exhausted — {} actors are live \ diff --git a/src/scheduler.rs b/src/scheduler.rs index 02329c9..13122e9 100644 --- a/src/scheduler.rs +++ b/src/scheduler.rs @@ -278,6 +278,37 @@ pub struct SpawnOpts { pub guard_size: Option, } +/// Why [`try_spawn`] could not start an actor. +/// +/// Marked `non_exhaustive`: today the only refusal is a full slab, but a +/// future variant (say, a shutdown-in-progress refusal) must not be a +/// breaking change for shed-path `match`es. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum SpawnError { + /// The fixed actor slab ([`Config::max_actors`] + /// (crate::runtime::Config::max_actors)) is full: every slot is claimed + /// by a live actor. This is a routine overload condition, not an + /// invariant violation — shed the unit of work (close the socket, + /// return a 503) and try again once actors have died. + AtCapacity, +} + +impl core::fmt::Display for SpawnError { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + match self { + SpawnError::AtCapacity => { + write!( + f, + "actor slab at capacity (`Config::max_actors` live actors)" + ) + } + } + } +} + +impl std::error::Error for SpawnError {} + /// Start a new actor running `f`, and return a [`JoinHandle`] for it. /// /// The new actor runs concurrently with its caller and with every other @@ -340,6 +371,73 @@ pub fn spawn_under_with( } } +/// [`spawn`] that reports a full actor slab instead of panicking. +/// +/// Behaviour parity with [`spawn`] in every case except one: when the fixed +/// slab ([`Config::max_actors`](crate::runtime::Config::max_actors)) is +/// full, this returns [`Err(SpawnError::AtCapacity)`](SpawnError::AtCapacity) +/// where `spawn` panics the calling actor. Use it at load-shedding call +/// sites — an accept loop spawning one actor per connection, a request +/// admission point — where "at capacity" is a routine overload condition to +/// handle (reject the unit of work), not an invariant violation. Internal +/// and bounded spawn sites should keep [`spawn`]: there, the panic is a +/// correct loud invariant check. +/// +/// The claim is atomic (claim-or-report): there is no +/// check-then-spawn race against other spawners for the last slot, so no +/// headroom margin is needed. +pub fn try_spawn(f: impl FnOnce() + Send + 'static) -> Result { + let parent = current_pid().unwrap_or_else(|| with_runtime(|_| crate::runtime::ROOT_PID)); + try_spawn_under_with(parent, SpawnOpts::default(), f) +} + +/// [`try_spawn`] with an explicit supervisor and per-actor stack shape +/// overrides — the full-control core the other `try_` surface is built on +/// (mirrors [`spawn_under_with`]). +pub fn try_spawn_under_with( + supervisor: Pid, + opts: SpawnOpts, + f: impl FnOnce() + Send + 'static, +) -> Result { + let supervisor = supervisor.erase(); + // Slot FIRST — deliberately the reverse of `spawn`'s stack-first order: + // under overload the Err arm is the HOT path, and a rejection must cost + // one mutex pop, not an mmap/pool-pop + init + recycle per shed unit of + // work. The claim is a single atomic pop (no TOCTOU; see + // `try_allocate_slot`). + let idx = match with_runtime(|inner| inner.try_allocate_slot()) { + Some(idx) => idx, + None => return Err(SpawnError::AtCapacity), + }; + // Between claim and install the slot is owned by this frame alone; if + // stack allocation panics in that window the slot must go back or it + // leaks for the life of the runtime (and would trip the run()-teardown + // slot-leak debug_assert). + struct ReturnOnUnwind(Option); + impl Drop for ReturnOnUnwind { + fn drop(&mut self) { + if let Some(idx) = self.0 { + with_runtime(|inner| inner.return_vacant_slot(idx)); + } + } + } + let mut claimed = ReturnOnUnwind(Some(idx)); + + let stack = with_runtime(|inner| crate::runtime::acquire_stack(inner, opts)); + let sp = init_actor_stack(stack.top(), crate::actor::trampoline); + let closure: crate::runtime::Closure = Box::new(f); + + claimed.0 = None; // install_actor takes ownership of the slot from here + let pid = with_runtime(|inner| { + crate::runtime::install_actor(inner, idx, sp, stack, supervisor, closure) + }); + + Ok(JoinHandle { + pid, + consumed: false, + }) +} + /// Spawn an actor that other actors can message directly by its [`Pid`], /// rather than only by holding on to a channel `Sender` you passed it /// yourself. diff --git a/tests/try_spawn.rs b/tests/try_spawn.rs new file mode 100644 index 0000000..859125a --- /dev/null +++ b/tests/try_spawn.rs @@ -0,0 +1,185 @@ +//! Non-panicking spawn at slab capacity (`try_spawn`). +//! +//! Covers: parity with `spawn` when slots are free; `Err(AtCapacity)` instead +//! of a panic on a full slab (the spawning actor survives — the crash-loop +//! from the motivating slowloris incident cannot start); self-heal (a freed +//! slot makes the next `try_spawn` succeed); and exact claim-or-report +//! accounting under a multi-thread race for the last slots (no TOCTOU +//! overshoot, no panic). + +use smarm::runtime::Config; +use smarm::{spawn, try_spawn, try_spawn_under_with, yield_now, SpawnError, SpawnOpts}; +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +use std::sync::Arc; + +/// A child that holds its slot until `release` flips, without parking +/// machinery: busy-yield keeps the scheduler moving and the slot occupied. +fn holder(release: Arc) -> impl FnOnce() + Send + 'static { + move || { + while !release.load(Ordering::Acquire) { + yield_now(); + } + } +} + +#[test] +fn try_spawn_is_spawn_when_slots_free() { + smarm::runtime::init(Config::exact(1)).run(|| { + let ran = Arc::new(AtomicBool::new(false)); + let flag = ran.clone(); + let h = try_spawn(move || flag.store(true, Ordering::Release)) + .expect("slots free — must behave exactly like spawn"); + h.join().unwrap(); + assert!(ran.load(Ordering::Acquire)); + }); +} + +#[test] +fn at_capacity_is_err_not_panic_and_accounting_is_exact() { + const MAX: usize = 8; + smarm::runtime::init(Config::exact(1).max_actors(MAX)).run(|| { + let release = Arc::new(AtomicBool::new(false)); + // Fill the slab from the initial actor: slots are claimed at spawn + // time, so children need not have run yet. Count until refusal. + let mut held = Vec::new(); + loop { + match try_spawn(holder(release.clone())) { + Ok(h) => held.push(h), + Err(e) => { + assert_eq!(e, SpawnError::AtCapacity); + break; + } + } + } + // Initial actor occupies one slot; the rest were spawnable. + assert_eq!(held.len(), MAX - 1, "slab accounting must be exact"); + // Still refusing (and still not panicking) on repeat. + assert!(matches!(try_spawn(|| ()), Err(SpawnError::AtCapacity))); + // The `_with` surface refuses identically — a custom shape must not + // reach stack allocation when there is no slot for it. + let opts = SpawnOpts { + stack_reserve: Some(1024 * 1024), + ..SpawnOpts::default() + }; + assert!(matches!( + try_spawn_under_with(smarm::self_pid(), opts, || ()), + Err(SpawnError::AtCapacity) + )); + + // Self-heal: free the slots, join, and the next try_spawn succeeds. + release.store(true, Ordering::Release); + for h in held { + h.join().unwrap(); + } + let h = try_spawn(|| ()).expect("slots freed — must succeed again"); + h.join().unwrap(); + }); +} + +#[test] +fn plain_spawn_still_panics_at_capacity() { + // The existing invariant-check semantics of `spawn` are untouched: at a + // full slab it panics, the panic is caught at the actor isolation + // boundary, and it surfaces as a join error — exactly as before. The + // bomb actor is spawned into the LAST slot (so the slab is full only + // once the bomb itself is live) and the panic lands inside the bomb, + // not the initial actor. + const MAX: usize = 6; + smarm::runtime::init(Config::exact(1).max_actors(MAX)).run(|| { + let release = Arc::new(AtomicBool::new(false)); + let mut held = Vec::new(); + for _ in 0..MAX - 2 { + held.push(spawn(holder(release.clone()))); + } + let armed = Arc::new(AtomicBool::new(false)); + let armed2 = armed.clone(); + let bomb = spawn(move || { + armed2.store(true, Ordering::Release); + // Slab is now full (initial + MAX−2 holders + this actor); the + // plain spawn must panic this actor. + let _ = spawn(|| ()); + unreachable!("allocate_slot must have panicked"); + }); + let err = bomb + .join() + .expect_err("bomb must die by panic, not run through"); + assert!(armed.load(Ordering::Acquire), "bomb must have actually run"); + // The panic message is a formatted String (panic! with args). + let msg = err + .payload + .downcast_ref::() + .cloned() + .unwrap_or_else(|| "".into()); + assert!( + msg.contains("slot table exhausted"), + "panic must be the slab-exhaustion invariant message, got: {msg}" + ); + release.store(true, Ordering::Release); + for h in held { + h.join().unwrap(); + } + }); +} + +#[test] +fn racing_try_spawns_claim_exactly_the_free_slots() { + // 4 scheduler threads, 4 spawner actors hammering try_spawn for a small + // pool of remaining slots. Claim-or-report must hand out exactly the + // free slots across all racers — no overshoot (TOCTOU), no panic. + const MAX: usize = 32; + const SPAWNERS: usize = 4; + smarm::runtime::init(Config::exact(4).max_actors(MAX)).run(|| { + let release = Arc::new(AtomicBool::new(false)); + let won = Arc::new(AtomicUsize::new(0)); + let done = Arc::new(AtomicUsize::new(0)); + + // Occupy some slots up front so the racers fight over a remainder. + let mut pre = Vec::new(); + for _ in 0..8 { + pre.push(spawn(holder(release.clone()))); + } + // Free slots now: MAX − 1 (initial) − 8 (pre) − SPAWNERS. + let up_for_grabs = MAX - 1 - 8 - SPAWNERS; + + let mut spawners = Vec::new(); + for _ in 0..SPAWNERS { + let release = release.clone(); + let won = won.clone(); + let done = done.clone(); + spawners.push(spawn(move || { + loop { + match try_spawn(holder(release.clone())) { + Ok(h) => { + won.fetch_add(1, Ordering::AcqRel); + drop(h); // detached; slot held by the holder + } + Err(SpawnError::AtCapacity) => break, + Err(_) => unreachable!("non_exhaustive future-proofing"), + } + } + done.fetch_add(1, Ordering::AcqRel); + })); + } + // Wait for every racer to hit AtCapacity. + while done.load(Ordering::Acquire) < SPAWNERS { + yield_now(); + } + assert_eq!(won.load(Ordering::Acquire), up_for_grabs); + + release.store(true, Ordering::Release); + for h in pre.into_iter().chain(spawners) { + h.join().unwrap(); + } + }); +} + +#[test] +fn spawn_error_is_a_real_error() { + let e = SpawnError::AtCapacity; + let msg = format!("{e}"); + assert!( + msg.contains("capacity"), + "Display should name the condition: {msg}" + ); + let _: &dyn std::error::Error = &e; +}