parkring
Concurrency primitives in Rust: a bounded MPMC ring, SCQ, a Chase-Lev deque, and a work-stealing pool, each checked with loom and Miri.
Results
| faster than a blocking queue at 16 threads | 2.7× |
| vs 1,090 µs blocking, 64-slot buffer | 396 µs |
| primitives, each model-checked with loom and run under Miri | 4 |
Problem
A blocking queue serialises every producer and consumer through one mutex, so throughput falls as thread count rises. The goal was a fixed-capacity multi-producer multi-consumer queue with no global lock whose performance still holds at 16 threads, built from atomics rather than a library.
Approach
A fixed-size ring buffer where every slot carries an atomic sequence number. Producers and consumers claim positions with compare-and-swap on the tail and head and use the slot's sequence to tell whether it is ready, so no thread waits on another. Cache-line padding keeps head and tail apart, AcqRel ordering is used only where the algorithm needs it, power-of-two capacity turns modulo into a mask, and exponential backoff stops a CAS storm under contention. Waiting threads spin briefly and then park on a futex, so an idle consumer costs no CPU. The crate grew into parkring: the same discipline applied to Nikolaev's SCQ, a Chase-Lev work-stealing deque, and a thread pool, with loom model checks and Miri runs in CI.
Latency under contention
16 threads · 64-slot buffer, lower is better.
What's inside
BoundedRing<T>: Multi-producer, multi-consumer ring with per-slot sequence stamps.Scq<T>: Scalable circular queue for higher contention.ChaseLev<T>: Work-stealing deque: owner pushes and pops, thieves steal.Pool: Work-stealing thread pool with futex-based parking for idle workers.
The core loop
pub fn try_push(&self, item: T) -> Result<(), TryPushError<T>> {
let (mut backoff, mut tail) = (Backoff::new(), self.tail.load(Relaxed));
loop {
let slot = self.slot(tail);
let diff = pos_diff(slot.sequence.load(Acquire), tail);
if diff == 0 { // empty for this lap: claim it
match self.tail.compare_exchange_weak(tail, pos_add(tail, 1), AcqRel, Relaxed) {
Ok(_) => {
unsafe { slot.write(item) };
slot.sequence.store(pos_add(tail, 1), Release);
self.consumers.notify_one();
return Ok(());
}
Err(actual) => { tail = actual; backoff.spin(); }
}
} else if diff < 0 { /* full, or a consumer mid-recycle */ }
else { backoff.snooze(); tail = self.tail.load(Relaxed); }
}
}
src/queue/lockfree.rs · try_push, condensed
Verification
Every primitive runs under loom, which explores thread interleavings exhaustively, and under Miri, which catches undefined behaviour in the unsafe code paths.
What I'd do next
- Publish the crossbeam ArrayQueue and Rayon comparisons as a tracked benchmark table across 2–64 threads.
- Add an async-aware wait strategy so the pool can back a Tokio executor.