use crate::loom_compat::{AtomicU32, AtomicU64, AtomicUsize, Mutex, Ordering, UnsafeCell};
pub(crate) trait FlagArray {
fn flag(&self, slot: usize) -> &AtomicU32;
}
struct SlotCell<P>(UnsafeCell<Option<P>>);
unsafe impl<P: Send> Sync for SlotCell<P> {}
pub(crate) struct SlotTable<P, F: FlagArray> {
flags: F,
active: Vec<AtomicU64>,
payload: Vec<SlotCell<P>>,
free: Mutex<Vec<usize>>,
n_armed: AtomicUsize,
num_slots: usize,
}
impl<P: Send, F: FlagArray> SlotTable<P, F> {
pub(crate) fn new(num_slots: usize, flags: F) -> Self {
let num_words = num_slots.div_ceil(64);
SlotTable {
flags,
active: (0..num_words).map(|_| AtomicU64::new(0)).collect(),
payload: (0..num_slots)
.map(|_| SlotCell(UnsafeCell::new(None)))
.collect(),
free: Mutex::new((0..num_slots).rev().collect()),
n_armed: AtomicUsize::new(0),
num_slots,
}
}
#[allow(dead_code)]
pub(crate) fn is_idle(&self) -> bool {
self.n_armed.load(Ordering::Acquire) == 0
}
fn lock_free(&self) -> impl std::ops::DerefMut<Target = Vec<usize>> + '_ {
self.free
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub(crate) fn claim(&self) -> Option<usize> {
self.lock_free().pop()
}
pub(crate) fn release(&self, slot: usize) {
self.lock_free().push(slot);
}
pub(crate) fn reset_flag(&self, slot: usize) {
self.flags.flag(slot).store(0, Ordering::Release);
}
pub(crate) fn publish(&self, slot: usize, waker: P) -> bool {
self.payload[slot]
.0
.with_mut(|p| unsafe { *p = Some(waker) });
let was_idle = self.n_armed.fetch_add(1, Ordering::AcqRel) == 0;
self.active[slot / 64].fetch_or(1u64 << (slot % 64), Ordering::Release);
was_idle
}
pub(crate) fn scan_once(&self, woken: &mut Vec<P>) -> bool {
let mut any_active = false;
let mut retired: Vec<usize> = Vec::new();
for word in 0..self.active.len() {
let mut bits = self.active[word].load(Ordering::Acquire);
if bits == 0 {
continue;
}
any_active = true;
while bits != 0 {
let bit = bits.trailing_zeros() as usize;
bits &= bits - 1;
let slot = word * 64 + bit;
if slot >= self.num_slots {
break;
}
if self.flags.flag(slot).load(Ordering::Acquire) == 1 {
self.active[word].fetch_and(!(1u64 << bit), Ordering::AcqRel);
let taken = self.payload[slot].0.with_mut(|p| unsafe { (*p).take() });
retired.push(slot);
if let Some(waker) = taken {
woken.push(waker);
}
}
}
}
if !retired.is_empty() {
let n = retired.len();
self.lock_free().append(&mut retired);
let prev = self.n_armed.fetch_sub(n, Ordering::AcqRel);
debug_assert!(prev >= n, "n_armed underflow: {prev} < {n}");
}
any_active
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::loom_compat::Arc;
struct MockFlags(Vec<AtomicU32>);
impl FlagArray for MockFlags {
fn flag(&self, slot: usize) -> &AtomicU32 {
&self.0[slot]
}
}
fn mock_flags(n: usize) -> MockFlags {
MockFlags((0..n).map(|_| AtomicU32::new(0)).collect())
}
#[cfg(loom)]
#[test]
fn loom_single_slot_handoff() {
loom::model(|| {
let table = Arc::new(SlotTable::new(1, mock_flags(1)));
let slot = table.claim().expect("slot");
table.reset_flag(slot);
let scanner = {
let table = table.clone();
loom::thread::spawn(move || {
let mut woken = Vec::new();
while woken.is_empty() {
table.scan_once(&mut woken);
loom::thread::yield_now();
}
woken
})
};
let gpu = {
let table = table.clone();
loom::thread::spawn(move || {
table.flags.flag(slot).store(1, Ordering::Release);
})
};
let was_idle = table.publish(slot, 7usize);
assert!(
was_idle,
"first publish from an idle table must report idle→active"
);
gpu.join().unwrap();
let woken = scanner.join().unwrap();
assert_eq!(
woken,
vec![7usize],
"payload delivered exactly once, value intact"
);
assert_eq!(table.lock_free().len(), 1, "slot recycled");
assert!(table.is_idle(), "n_armed returns to 0 after retirement");
});
}
#[cfg(not(loom))]
#[test]
fn stress_many_producers_one_scanner() {
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering as O};
use std::sync::mpsc;
#[cfg(miri)]
const SLOTS: usize = 8;
#[cfg(miri)]
const THREADS: usize = 2;
#[cfg(miri)]
const PER_THREAD: usize = 15;
#[cfg(not(miri))]
const SLOTS: usize = 256;
#[cfg(not(miri))]
const THREADS: usize = 8;
#[cfg(not(miri))]
const PER_THREAD: usize = 500;
const TOTAL: usize = THREADS * PER_THREAD;
let table = Arc::new(SlotTable::<usize, MockFlags>::new(SLOTS, mock_flags(SLOTS)));
let done = Arc::new(AtomicBool::new(false));
let (arm_tx, arm_rx) = mpsc::channel::<usize>();
let gpu = {
let table = table.clone();
std::thread::spawn(move || {
while let Ok(slot) = arm_rx.recv() {
table.flags.flag(slot).store(1, O::Release);
}
})
};
let seen = Arc::new((0..TOTAL).map(|_| AtomicUsize::new(0)).collect::<Vec<_>>());
let scanner = {
let (table, done, seen) = (table.clone(), done.clone(), seen.clone());
std::thread::spawn(move || {
let mut woken = Vec::new();
let mut total = 0usize;
while total < TOTAL {
table.scan_once(&mut woken);
for token in woken.drain(..) {
seen[token].fetch_add(1, O::SeqCst);
total += 1;
}
if !done.load(O::Relaxed) {
std::hint::spin_loop();
}
}
})
};
let workers: Vec<_> = (0..THREADS)
.map(|t| {
let (table, arm_tx) = (table.clone(), arm_tx.clone());
std::thread::spawn(move || {
for i in 0..PER_THREAD {
let token = t * PER_THREAD + i;
let slot = loop {
if let Some(s) = table.claim() {
break s;
}
std::hint::spin_loop();
};
table.reset_flag(slot);
table.publish(slot, token);
arm_tx.send(slot).unwrap();
}
})
})
.collect();
for w in workers {
w.join().unwrap();
}
drop(arm_tx);
gpu.join().unwrap();
done.store(true, O::Relaxed);
scanner.join().unwrap();
for (token, count) in seen.iter().enumerate() {
assert_eq!(
count.load(O::SeqCst),
1,
"token {token} not delivered exactly once"
);
}
assert_eq!(table.lock_free().len(), SLOTS, "slots leaked");
}
#[cfg(not(loom))]
#[test]
fn claim_signals_exhaustion_and_recycle_restores_capacity() {
const SLOTS: usize = 4;
let table = SlotTable::<usize, MockFlags>::new(SLOTS, mock_flags(SLOTS));
let mut claimed: Vec<usize> = (0..SLOTS)
.map(|_| table.claim().expect("slot available"))
.collect();
claimed.sort_unstable();
claimed.dedup();
assert_eq!(claimed.len(), SLOTS, "claims must be distinct");
assert_eq!(table.claim(), None, "exhausted pool must signal None");
let slot = claimed[0];
table.reset_flag(slot);
table.publish(slot, 99);
table.flags.flag(slot).store(1, Ordering::Release);
let mut woken = Vec::new();
table.scan_once(&mut woken);
assert_eq!(woken, vec![99], "completed slot delivered");
assert!(table.claim().is_some(), "recycled slot must be claimable");
assert_eq!(table.claim(), None, "only the recycled slot returned");
}
}
#[cfg(all(test, not(loom), not(miri)))]
mod ab_bench {
use super::*;
use crate::loom_compat::Arc;
use std::sync::atomic::{AtomicUsize, Ordering as O};
use std::thread;
use std::time::Instant;
struct Flags(Vec<AtomicU32>);
impl FlagArray for Flags {
fn flag(&self, s: usize) -> &AtomicU32 {
&self.0[s]
}
}
#[derive(Clone, Copy, PartialEq)]
enum Unpark {
Always,
EmptyWake,
}
#[derive(Clone, Copy, PartialEq)]
enum Idle {
Park,
NeverSpin,
}
struct Result {
ops_per_sec: f64,
unparks: usize,
scan_passes: usize,
}
fn run(unpark: Unpark, idle: Idle, saturated: bool) -> Result {
const SLOTS: usize = 1024;
const THREADS: usize = 8;
const OPS_PER: usize = 150_000;
let total = THREADS * OPS_PER;
let table = Arc::new(SlotTable::<usize, Flags>::new(
SLOTS,
Flags((0..SLOTS).map(|_| AtomicU32::new(0)).collect()),
));
let n_armed = Arc::new(AtomicUsize::new(0));
let unparks = Arc::new(AtomicUsize::new(0));
let scan_passes = Arc::new(AtomicUsize::new(0));
let done = Arc::new(AtomicUsize::new(0));
let scanner = {
let (table, n_armed, scan_passes, done) = (
table.clone(),
n_armed.clone(),
scan_passes.clone(),
done.clone(),
);
thread::Builder::new()
.name("ab-scanner".into())
.spawn(move || {
let mut woken = Vec::new();
loop {
table.scan_once(&mut woken);
scan_passes.fetch_add(1, O::Relaxed);
let retired = woken.len();
if retired > 0 {
n_armed.fetch_sub(retired, O::AcqRel);
let d = done.fetch_add(retired, O::Relaxed) + retired;
woken.clear();
if d >= total {
break;
}
continue;
}
match idle {
Idle::Park => {
if n_armed.load(O::Acquire) == 0 {
thread::park();
}
}
Idle::NeverSpin => std::hint::spin_loop(),
}
if done.load(O::Relaxed) >= total {
break;
}
}
})
.unwrap()
};
let scanner_thread = scanner.thread().clone();
let start = Instant::now();
let workers: Vec<_> = (0..THREADS)
.map(|t| {
let (table, n_armed, unparks, scanner_thread) = (
table.clone(),
n_armed.clone(),
unparks.clone(),
scanner_thread.clone(),
);
thread::spawn(move || {
for i in 0..OPS_PER {
let slot = loop {
if let Some(s) = table.claim() {
break s;
}
std::hint::spin_loop();
};
table.reset_flag(slot);
table.publish(slot, t * OPS_PER + i);
let was_idle = n_armed.fetch_add(1, O::AcqRel) == 0;
table.flags.flag(slot).store(1, O::Release);
let do_unpark = match unpark {
Unpark::Always => true,
Unpark::EmptyWake => was_idle,
};
if do_unpark {
unparks.fetch_add(1, O::Relaxed);
scanner_thread.unpark();
}
if !saturated && i % 256 == 0 {
thread::yield_now();
}
}
})
})
.collect();
for w in workers {
w.join().unwrap();
}
while done.load(O::Relaxed) < total {
scanner_thread.unpark();
std::hint::spin_loop();
}
scanner.join().unwrap();
let secs = start.elapsed().as_secs_f64();
Result {
ops_per_sec: total as f64 / secs,
unparks: unparks.load(O::Relaxed),
scan_passes: scan_passes.load(O::Relaxed),
}
}
#[test]
#[ignore = "perf A/B; run explicitly with --release --nocapture --ignored"]
fn ab_reactor_variants() {
let configs = [
(Unpark::Always, Idle::Park, "always+park"),
(Unpark::EmptyWake, Idle::Park, "empty→wake+park"),
(Unpark::Always, Idle::NeverSpin, "always+never-park"),
(Unpark::EmptyWake, Idle::NeverSpin, "empty→wake+never-park"),
];
for (saturated, label) in [(true, "SATURATED"), (false, "BURSTY")] {
eprintln!("\n## {label} (8 threads x 150k ops)\n");
eprintln!("| config | Mops/s | unparks | scan passes |");
eprintln!("|---|---|---|---|");
for (u, i, name) in configs {
let mut best = run(u, i, saturated);
for _ in 0..2 {
let r = run(u, i, saturated);
if r.ops_per_sec > best.ops_per_sec {
best = r;
}
}
eprintln!(
"| {name} | {:.2} | {} | {} |",
best.ops_per_sec / 1e6,
best.unparks,
best.scan_passes
);
}
}
}
}