Expand description
Concurrency primitives built and verified from first principles: bounded MPMC queues and channels, a work-stealing deque, and a work-stealing thread pool.
| what it is | |
|---|---|
LockFreeQueue | Vyukov’s per-slot sequence ring. The recommended queue. |
channel::bounded | Sender / Receiver on a LockFreeQueue, disconnecting when either side drops |
ScqQueue | Nikolaev’s SCQ: fetch-add claims, genuinely lock-free, slower on this hardware |
BlockingQueue | one mutex, two condvars: the reference implementation |
Worker / Stealer | Chase-Lev work-stealing deque |
ThreadPool, join | a work-stealing pool built from the pieces above |
For select, parallel iterators or async, use crossbeam-channel, Rayon
or an async runtime’s channels instead; the README’s “When to use
something else” section says where each one is the better choice.
Blocked threads spin briefly, then park on a futex (futex(2) on Linux and
Android, __ulock on macOS) or on std’s Condvar elsewhere, so an idle
thread does not burn a core. The only dependency is libc.
The three queues implement BoundedQueue and share one API:
push/popblock;try_push/try_popnever block;push_timeout/pop_timeoutgive up after a deadline;closerejects further pushes, wakes every waiter, and lets consumers drain the remaining items beforepopreportsPopError.
Every failed push hands the item back inside the error.
use parkring::LockFreeQueue;
let queue = LockFreeQueue::new(64);
std::thread::scope(|s| {
let consumers: Vec<_> = (0..4)
.map(|_| s.spawn(|| {
let mut popped = 0;
while queue.pop().is_ok() {
popped += 1;
}
popped
}))
.collect();
// Every producer is joined when this inner scope ends.
std::thread::scope(|p| {
for _ in 0..4 {
p.spawn(|| (0..1000).for_each(|i| queue.push(i).unwrap()));
}
});
// Consumers drain what is left, then `pop` returns `Err` and they exit.
queue.close();
let total: usize = consumers.into_iter().map(|h| h.join().unwrap()).sum();
assert_eq!(total, 4000);
});§Thread safety
Both queues are Send + Sync exactly when T: Send. Items move between
threads but are never shared, so T: Sync is not required:
fn assert_sync<T: Sync>() {}
assert_sync::<parkring::LockFreeQueue<std::rc::Rc<()>>>();fn assert_sync<T: Sync>() {}
assert_sync::<parkring::BlockingQueue<std::rc::Rc<()>>>();fn assert_sync<T: Sync>() {}
assert_sync::<parkring::ScqQueue<std::rc::Rc<()>>>();Channel handles follow the same rule, and a non-Send message type is
rejected:
fn assert_send<T: Send>() {}
assert_send::<parkring::channel::Sender<std::rc::Rc<()>>>();A deque’s Worker belongs to one thread at a time: it is Send but not
Sync. Stealer is Send + Sync.
fn assert_sync<T: Sync>() {}
assert_sync::<parkring::Worker<u32>>();fn assert_send<T: Send>() {}
assert_send::<parkring::Stealer<std::rc::Rc<()>>>();See docs/DESIGN.md in the repository for the memory-ordering argument
and how it is verified with loom and Miri.
Modules§
- channel
- Bounded multi-producer multi-consumer channels.
Structs§
- Blocking
Queue - A bounded MPMC queue guarded by one mutex, with separate
not_fullandnot_emptycondition variables. - Lock
Free Queue - A bounded MPMC queue built on Dmitry Vyukov’s per-slot sequence numbers.
- PopError
- Returned by
popwhen the queue is closed and fully drained. - Push
Error - Returned by
pushwhen the queue is closed. Carries the rejected item. - ScqQueue
64-bit - A bounded MPMC queue built on Nikolaev’s SCQ (DISC 2019): claims are
fetch_adds, so contended threads never retry them. - Stealer
- A thief’s handle to a
Worker’s deque. Cheap to clone;Send + Sync. - Thread
Pool - A work-stealing thread pool.
- Worker
- The owner’s end of a work-stealing deque.
Enums§
- PopTimeout
Error - Returned by
pop_timeout. - Push
Timeout Error - Returned by
push_timeout. - Steal
- The result of
Stealer::steal. - TryPop
Error - Returned by
try_pop. - TryPush
Error - Returned by
try_push.
Traits§
- Bounded
Queue - A bounded, closable, multi-producer multi-consumer FIFO queue.
Functions§
- join
- Runs
aandb, potentially in parallel, and returns both results.