Skip to main content

Crate parkring

Crate parkring 

Source
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
LockFreeQueueVyukov’s per-slot sequence ring. The recommended queue.
channel::boundedSender / Receiver on a LockFreeQueue, disconnecting when either side drops
ScqQueueNikolaev’s SCQ: fetch-add claims, genuinely lock-free, slower on this hardware
BlockingQueueone mutex, two condvars: the reference implementation
Worker / StealerChase-Lev work-stealing deque
ThreadPool, joina 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 / pop block;
  • try_push / try_pop never block;
  • push_timeout / pop_timeout give up after a deadline;
  • close rejects further pushes, wakes every waiter, and lets consumers drain the remaining items before pop reports PopError.

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§

BlockingQueue
A bounded MPMC queue guarded by one mutex, with separate not_full and not_empty condition variables.
LockFreeQueue
A bounded MPMC queue built on Dmitry Vyukov’s per-slot sequence numbers.
PopError
Returned by pop when the queue is closed and fully drained.
PushError
Returned by push when the queue is closed. Carries the rejected item.
ScqQueue64-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.
ThreadPool
A work-stealing thread pool.
Worker
The owner’s end of a work-stealing deque.

Enums§

PopTimeoutError
Returned by pop_timeout.
PushTimeoutError
Returned by push_timeout.
Steal
The result of Stealer::steal.
TryPopError
Returned by try_pop.
TryPushError
Returned by try_push.

Traits§

BoundedQueue
A bounded, closable, multi-producer multi-consumer FIFO queue.

Functions§

join
Runs a and b, potentially in parallel, and returns both results.