fix(run_queue): make the Vyukov rings preemption-tolerant (finding 13)
A consumer OS-preempted between its dequeue_pos claim and its seq release freezes one cell; once traffic laps the ring (~cap ops ≈ 1-2ms at yield-storm throughput ≈ one scheduling quantum under load), try_push reads the stale seq and the original algorithm's 'lap behind => full' inference misfires. The old push assert then converted that liveness stall into an abort blaming a double enqueue that never happened (soak: occupancy 174-181 of cap 16384 at every failure), and the dead scheduler threads stranded actors => the observed hangs. Pristine 5504ef3 failed 14/14 under an 8-spinner soak on the 20-core box; a retry prototype passed 13/14, its one failure being a stall that outlived a fixed 1M-spin bound — pure spinning starves the descheduled consumer, so the wait must yield. Fix, following crossbeam ArrayQueue's shape (fence + opposite-counter check; cells hand off, COUNTERS give verdicts): - push: on a lap-behind cell, fence(SeqCst) + occupancy check. occ < cap => transient stall => spin-then-yield backoff and retry; occ >= cap => the REAL at-most-once-enqueued violation => panic with a truthful message and the counters. enqueue_pos is loaded before dequeue_pos so racing pops only underestimate occupancy (no spurious panic). - pop: symmetric counter check before an empty verdict; on a mid-publish producer, bounded wait then None — deliberate deviation from crossbeam's unbounded retry (spurious None is benign: RFC 018's enqueue-wake self-heals; StripedRing's probe must not hang on one stripe). - StripedRing::push probe: yield-escalating backoff after a full refused lap (was a bare spin_loop). - Hand-rolled Backoff (spin 2^n to 64, then yield_now); under loom every wait is a yield so models explore the stalled peer's progress. New loom model mpmc_lap_onto_stalled_consumer_completes reproduces the old panic in the first explored interleavings (verified FAILED against 5504ef3) and passes with the fix. Lib 60 + integration 41 + all 4 ring loom models pass. rq-striped inherits the fix (stripes are MpmcRings).
This commit is contained in:
+151
-12
@@ -51,7 +51,7 @@
|
|||||||
//! - `len()` is approximate (stats only).
|
//! - `len()` is approximate (stats only).
|
||||||
|
|
||||||
use crate::pid::Pid;
|
use crate::pid::Pid;
|
||||||
use crate::sync_shim::{AtomicUsize, Ordering, UnsafeCell};
|
use crate::sync_shim::{fence, AtomicUsize, Ordering, UnsafeCell};
|
||||||
use std::mem::MaybeUninit;
|
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
|
// rq-mutex — the baseline
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
@@ -204,15 +244,57 @@ impl MpmcRing {
|
|||||||
|
|
||||||
pub fn push(&self, pid: Pid) {
|
pub fn push(&self, pid: Pid) {
|
||||||
assert_no_preempt();
|
assert_no_preempt();
|
||||||
assert!(
|
if self.try_push(pid) {
|
||||||
self.try_push(pid),
|
return;
|
||||||
"smarm: run queue overflow — occupancy exceeded the slab bound, \
|
}
|
||||||
which the at-most-once-enqueued invariant forbids. This is a \
|
self.push_slow(pid);
|
||||||
runtime bug (double enqueue), not a capacity tuning problem."
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// 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 {
|
fn try_push(&self, pid: Pid) -> bool {
|
||||||
let mut pos = self.enqueue_pos.0.load(Ordering::Relaxed);
|
let mut pos = self.enqueue_pos.0.load(Ordering::Relaxed);
|
||||||
loop {
|
loop {
|
||||||
@@ -246,6 +328,7 @@ impl MpmcRing {
|
|||||||
|
|
||||||
pub fn pop(&self) -> Option<Pid> {
|
pub fn pop(&self) -> Option<Pid> {
|
||||||
assert_no_preempt();
|
assert_no_preempt();
|
||||||
|
let mut backoff = Backoff::new();
|
||||||
let mut pos = self.dequeue_pos.0.load(Ordering::Relaxed);
|
let mut pos = self.dequeue_pos.0.load(Ordering::Relaxed);
|
||||||
loop {
|
loop {
|
||||||
let cell = &self.buf[pos & self.mask];
|
let cell = &self.buf[pos & self.mask];
|
||||||
@@ -270,9 +353,31 @@ impl MpmcRing {
|
|||||||
Err(actual) => pos = actual,
|
Err(actual) => pos = actual,
|
||||||
}
|
}
|
||||||
} else if diff < 0 {
|
} else if diff < 0 {
|
||||||
// Empty (or the producer at this cell hasn't published yet —
|
// The cell at `pos` is unpublished: either the queue is
|
||||||
// a snapshot miss the caller's idle-retry loop absorbs).
|
// empty, or the producer that claimed it was preempted
|
||||||
return None;
|
// 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 {
|
} else {
|
||||||
pos = self.dequeue_pos.0.load(Ordering::Relaxed);
|
pos = self.dequeue_pos.0.load(Ordering::Relaxed);
|
||||||
}
|
}
|
||||||
@@ -342,6 +447,7 @@ impl StripedRing {
|
|||||||
// guarantees a free stripe exists, so the outer loop terminates.
|
// guarantees a free stripe exists, so the outer loop terminates.
|
||||||
// The retry-from-home lap handles the racy case where every stripe
|
// The retry-from-home lap handles the racy case where every stripe
|
||||||
// momentarily refused us.
|
// momentarily refused us.
|
||||||
|
let mut backoff = Backoff::new();
|
||||||
loop {
|
loop {
|
||||||
for i in 0..=self.stripe_mask {
|
for i in 0..=self.stripe_mask {
|
||||||
let s = &self.stripes[(home + i) & self.stripe_mask];
|
let s = &self.stripes[(home + i) & self.stripe_mask];
|
||||||
@@ -349,7 +455,13 @@ impl StripedRing {
|
|||||||
return;
|
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
|
/// Producer races a consumer on a single element: the consumer either
|
||||||
/// gets it or sees a clean None — never a torn/duplicated element.
|
/// gets it or sees a clean None — never a torn/duplicated element.
|
||||||
#[test]
|
#[test]
|
||||||
|
|||||||
Reference in New Issue
Block a user