diff --git a/src/run_queue.rs b/src/run_queue.rs index 6814e8b..a2d1187 100644 --- a/src/run_queue.rs +++ b/src/run_queue.rs @@ -51,7 +51,7 @@ //! - `len()` is approximate (stats only). use crate::pid::Pid; -use crate::sync_shim::{AtomicUsize, Ordering, UnsafeCell}; +use crate::sync_shim::{fence, AtomicUsize, Ordering, UnsafeCell}; use std::mem::MaybeUninit; // --------------------------------------------------------------------------- @@ -92,6 +92,46 @@ fn assert_no_preempt() { ); } +/// Escalating wait for transient ring stalls: a peer preempted by the OS +/// inside its ~100ns claim→publish window (finding 13). Spin first (the +/// common stall is a peer that is merely slow, gone within a few hundred +/// cycles), then donate the timeslice — under oversubscription pure +/// spinning STARVES the descheduled peer of the CPU it needs to publish +/// (measured in the finding-13 soak: 10⁶ pure spins can outlast the very +/// stall they prolong). OS-level yielding is orthogonal to +/// `assert_no_preempt`, which guards smarm signal preemption only. +struct Backoff(u32); + +impl Backoff { + const SPIN_LIMIT: u32 = 6; + + fn new() -> Self { + Self(0) + } + + /// Total waits so far — lets bounded callers cap the yield phase. + fn steps(&self) -> u32 { + self.0 + } + + #[inline] + fn wait(&mut self) { + #[cfg(loom)] + // Loom has no notion of spinning time; every wait is a scheduling + // point so the model explores the stalled peer's progress. + loom::thread::yield_now(); + #[cfg(not(loom))] + if self.0 <= Self::SPIN_LIMIT { + for _ in 0..1u32 << self.0 { + std::hint::spin_loop(); + } + } else { + std::thread::yield_now(); + } + self.0 += 1; + } +} + // --------------------------------------------------------------------------- // rq-mutex — the baseline // --------------------------------------------------------------------------- @@ -204,15 +244,57 @@ impl MpmcRing { pub fn push(&self, pid: Pid) { assert_no_preempt(); - assert!( - self.try_push(pid), - "smarm: run queue overflow — occupancy exceeded the slab bound, \ - which the at-most-once-enqueued invariant forbids. This is a \ - runtime bug (double enqueue), not a capacity tuning problem." - ); + if self.try_push(pid) { + return; + } + self.push_slow(pid); } - /// One full claim attempt; `false` only if the ring is full. + /// Cold path: the cell at `enqueue_pos` is still a lap behind. Two + /// worlds are indistinguishable at the cell (finding 13): a consumer + /// preempted between its `dequeue_pos` claim and its seq release while + /// the ring lapped onto that cell (transient), or a genuine occupancy + /// overflow (a runtime bug). The COUNTERS discriminate — the same shape + /// as crossbeam `ArrayQueue::push`'s fence + opposite-counter check: + /// occupancy < capacity ⇒ transient ⇒ wait for the stalled peer; + /// occupancy ≥ capacity ⇒ the at-most-once-enqueued invariant really is + /// broken ⇒ panic (a legal push starts from occupancy ≤ max_actors − 1). + #[cold] + fn push_slow(&self, pid: Pid) { + let mut backoff = Backoff::new(); + loop { + // Order the counter reads after the failed cell read. + // `enqueue_pos` is loaded BEFORE `dequeue_pos`, so a pop racing + // us can only make the computed occupancy an UNDERestimate — + // conservative in the safe direction (never a spurious panic; + // a genuine violation is persistent and caught next lap). + fence(Ordering::SeqCst); + let enq = self.enqueue_pos.0.load(Ordering::SeqCst); + let deq = self.dequeue_pos.0.load(Ordering::SeqCst); + let occ = enq.wrapping_sub(deq); + assert!( + occ <= self.mask, + "smarm: run queue occupancy {} reached capacity {} — a pid \ + was enqueued more than once (or more pids exist than \ + max_actors); the at-most-once-enqueued invariant is broken \ + (enq={} deq={})", + occ, + self.mask + 1, + enq, + deq + ); + backoff.wait(); + if self.try_push(pid) { + return; + } + } + } + + /// One full claim attempt; `false` when the cell at `enqueue_pos` is + /// still a lap behind — which means EITHER genuinely full OR a + /// lap-stalled consumer (finding 13). The cell cannot tell the two + /// apart; callers must disambiguate via the counters (`push_slow`) or + /// tolerate refusal (`StripedRing`'s probe). fn try_push(&self, pid: Pid) -> bool { let mut pos = self.enqueue_pos.0.load(Ordering::Relaxed); loop { @@ -246,6 +328,7 @@ impl MpmcRing { pub fn pop(&self) -> Option { assert_no_preempt(); + let mut backoff = Backoff::new(); let mut pos = self.dequeue_pos.0.load(Ordering::Relaxed); loop { let cell = &self.buf[pos & self.mask]; @@ -270,9 +353,31 @@ impl MpmcRing { Err(actual) => pos = actual, } } else if diff < 0 { - // Empty (or the producer at this cell hasn't published yet — - // a snapshot miss the caller's idle-retry loop absorbs). - return None; + // The cell at `pos` is unpublished: either the queue is + // empty, or the producer that claimed it was preempted + // inside its claim→publish window (finding 13). The cell + // cannot tell the two apart; the counters can. + fence(Ordering::SeqCst); + let deq = self.dequeue_pos.0.load(Ordering::SeqCst); + if deq != pos { + // Stale head — no verdict; re-probe at the real head. + pos = deq; + continue; + } + if self.enqueue_pos.0.load(Ordering::SeqCst) == pos { + return None; // counters agree: genuinely empty + } + // Producer mid-publish. Wait briefly, then report None + // anyway — a DELIBERATE bounded deviation from crossbeam's + // unbounded retry: a spurious None is correctness-benign + // here (the stalled push completes and its RFC 018 + // enqueue-wake re-wakes a parked scheduler), a parked + // scheduler beats a yielding one, and StripedRing's pop + // probe must not hang on one stripe. + if backoff.steps() > Backoff::SPIN_LIMIT + 8 { + return None; + } + backoff.wait(); } else { pos = self.dequeue_pos.0.load(Ordering::Relaxed); } @@ -342,6 +447,7 @@ impl StripedRing { // guarantees a free stripe exists, so the outer loop terminates. // The retry-from-home lap handles the racy case where every stripe // momentarily refused us. + let mut backoff = Backoff::new(); loop { for i in 0..=self.stripe_mask { let s = &self.stripes[(home + i) & self.stripe_mask]; @@ -349,7 +455,13 @@ impl StripedRing { return; } } - std::hint::spin_loop(); + // Every stripe refused this lap: either transiently full (the + // headroom argument above guarantees a genuinely free stripe + // exists) or the probes landed on lap-stalled cells (finding + // 13 — try_push cannot tell the two apart). Waiting is correct + // either way; the backoff escalates to an OS yield so stalled + // consumers get the CPU they need to release their cells. + backoff.wait(); } } @@ -560,6 +672,33 @@ mod loom_tests { }); } + /// Finding 13 regression: a consumer stalled between its dequeue_pos + /// claim and its seq release must NOT make a lapping producer conclude + /// "full" (the old code asserted here — occupancy never exceeded the + /// bound; the cell just hadn't been recycled). Capacity 2: fill, pop + /// once on each thread, then push a third element — in the + /// interleavings where a pop's release is still pending, that push + /// laps onto the stalled cell and must wait, not die. + #[test] + fn mpmc_lap_onto_stalled_consumer_completes() { + loom::model(|| { + let q = Arc::new(MpmcRing::with_capacity(2)); + q.push(pid(0)); + q.push(pid(1)); + let q2 = q.clone(); + let h = thread::spawn(move || q2.pop().expect("ring has two elements")); + let a = q.pop().expect("ring has two elements"); + // enqueue_pos = 2 → cell 0: laps onto the other thread's cell + // whenever its release is delayed. + q.push(pid(2)); + let b = h.join().unwrap(); + let c = q.pop().expect("the lapping push must have landed"); + let mut got = vec![a.index(), b.index(), c.index()]; + got.sort_unstable(); + assert_eq!(got, vec![0, 1, 2]); + }); + } + /// Producer races a consumer on a single element: the consumer either /// gets it or sees a clean None — never a torn/duplicated element. #[test]