feat(scheduler,runtime): non-panicking try_spawn for at-capacity load shedding
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<u32>: 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<JoinHandle, SpawnError>: 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)
This commit is contained in:
committed by
Claude (sandbox)
parent
95306c7f60
commit
e8f0234492
@@ -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<AtomicBool>) -> 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::<String>()
|
||||
.cloned()
|
||||
.unwrap_or_else(|| "<non-string payload>".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;
|
||||
}
|
||||
Reference in New Issue
Block a user