parkring
Bounded channels and queues, a work-stealing deque and a work-stealing thread pool for Rust. Each one is implemented from its paper, checked with the loom model checker, Miri, fuzzing and ThreadSanitizer, and benchmarked against the crate you would otherwise use.
A thread that has to wait (for an empty queue, a full channel, or work in the
pool) spins briefly and then sleeps on a futex (futex(2) on Linux,
__ulock on macOS, Condvar elsewhere), so an idle thread uses no CPU. The
only dependency is libc.
- Install
- Quick start
- Using parkring: channels, queues, shutdown and timeouts, the work-stealing deque, the thread pool
- Choosing a type
- Performance
- How it is verified
Install
[]
= "1"
Rust 1.85 or newer. Works on any target with std; see
Stability for the platforms CI tests.
Quick start
A bounded channel between worker threads. When every sender is dropped, the receivers finish what was sent and their loops end:
use thread;
use channel;
let = bounded;
// Four producers, each with its own clone of the sender.
let producers: =
.map
.collect;
drop; // only the producers' clones keep the channel open now
// `iter()` ends once all senders are gone and the channel is empty.
let total: u64 = rx.iter.sum;
assert_eq!;
for p in producers
Using parkring
Channels
channel::bounded(n) returns a Sender and a Receiver. Clone either one to
send or receive from more threads; each message goes to exactly one receiver.
| method | when the channel is full / empty | after the other side is gone |
|---|---|---|
send(msg) / recv() |
sleeps until there is room / a message | Err; send hands the message back |
try_send(msg) / try_recv() |
returns Full / Empty at once |
returns Disconnected |
send_timeout / recv_timeout |
waits up to the timeout, then Timeout |
returns Disconnected |
iter() / try_iter() |
blocks / stops | ends after the last message |
Disconnection is driven by dropping handles, so there is no close to call:
- drop the last
Sender: receivers still get every message that was sent, thenrecvreturnsErr(RecvError); - drop the last
Receiver:sendfails immediately and returns your message insideSendError.
use ;
let = bounded;
tx.try_send.unwrap;
tx.try_send.unwrap;
assert_eq!; // full: nothing waits
assert_eq!;
drop; // the only receiver
assert_eq!; // you get "d" back
A channel holds at least the capacity you ask for; capacity() tells you
the real size. It is bounded only, and there is no select; for those, see
when to use something else.
Queues
The queues are the layer below the channel: one value that all threads share
by reference, with an explicit close. Use them when you want that shared
object (in an Arc, a static, or borrowed by scoped threads) rather than
sender and receiver handles.
use Arc;
use thread;
use LockFreeQueue;
let queue = new;
let consumer = ;
for i in 0..10_000
queue.close;
assert_eq!;
There are three queues with the same API:
LockFreeQueue: the one to use. Dmitry Vyukov's bounded ring: no locks on the fast path, and waiting threads park.ScqQueue: Nikolaev's SCQ, with a stronger progress guarantee (some thread always completes its operation, even if others are suspended). It is several times slower here; use it only if you need that guarantee. 64-bit targets only.BlockingQueue: one mutex and two condition variables. The simple reference version.
All three implement the BoundedQueue trait, so code can be written once for
any of them (fn run<Q: BoundedQueue<Job>>(queue: &Q)), or take a
&dyn BoundedQueue<Job>.
| method | queue full / empty | queue closed |
|---|---|---|
push(item) / pop() |
sleeps | push: Err(PushError(item)). pop: drains what is left, then Err(PopError) |
try_push / try_pop |
Err(TryPushError::Full(item)) / Err(TryPopError::Empty) |
Closed(item) / Closed once drained |
push_timeout / pop_timeout |
waits up to the timeout, then Timeout |
Closed |
close() |
wakes every waiting thread; returns true the first time |
Every error from a push gives the item back (into_inner()), so nothing is
ever dropped silently.
Shutdown and timeouts
Timeouts let a producer notice that consumers have fallen behind, and let a
consumer do other work while the queue is idle. close (or dropping the last
sender) stops new work without losing what was already accepted:
use Duration;
use ;
let jobs = new;
jobs.push.unwrap;
jobs.push.unwrap;
// Full: give up after 10 ms instead of blocking, and keep the job.
match jobs.push_timeout
jobs.close;
// Items accepted before `close` are still delivered...
assert_eq!;
assert_eq!;
// ...then the queue reports that it is closed and empty.
assert_eq!;
examples/shutdown.rs shows a full producer and
consumer doing this with load shedding.
The work-stealing deque
A Worker is a deque owned by one thread: it pushes and pops at one end,
newest first. Any number of Stealer handles take from the other end, oldest
first. This is the building block of work-stealing schedulers: each thread
works through its own tasks and steals from the others when it runs out.
use ;
let worker = new;
let stealer = worker.stealer; // `Clone`; send it to other threads
worker.push;
worker.push;
worker.push;
assert_eq!; // owner: newest first
assert_eq!; // thief: oldest first
steal can return Steal::Retry when it loses a race with another thread;
the deque may still have work, so try again. A Worker can be moved to
another thread but not shared (Send, not Sync); a Stealer can be both.
examples/scheduler.rs is a small scheduler built
this way.
The thread pool
ThreadPool runs closures on a fixed set of worker threads that steal work
from each other. Inside the pool, join(a, b) runs a and b, possibly in
parallel: b is offered to idle workers, and if nobody takes it, the current
thread runs it too. Recursive join spreads divide-and-conquer work across
the pool without any queues of your own.
use ;
let pool = new;
let values: = .collect;
let total = pool.install; // run on the pool and wait
assert_eq!;
install(f)runsfon the pool and returns its result. A panic insidefis passed back to the caller.join(a, b)outside any pool simply runsaand thenb.spawn(f)runsfin the background without waiting. Dropping the pool waits for spawned jobs. A panic in a spawned job aborts the process, since there is nobody to return it to.
examples/quicksort.rs sorts in parallel with
join.
Examples
Choosing a type
| you want | use |
|---|---|
| to pass messages between threads, with shutdown when one side goes away | channel::bounded |
| one shared bounded queue whose waiting threads sleep instead of spinning | LockFreeQueue |
| a queue with a lock-free progress guarantee, and you accept lower throughput | ScqQueue (64-bit targets) |
| the simplest correct queue, for reference or low traffic | BlockingQueue |
| per-thread task deques for your own scheduler | Worker / Stealer |
fork-join parallelism (join, install, spawn) |
ThreadPool |
When to use something else
parkring is small and heavily verified, but the established crates are better in several places, and the benchmarks show where:
- Many producers or consumers at maximum throughput: crossbeam's
ArrayQueueis 10 to 17% faster from 2 + 2 threads up, and much faster with many producers feeding one consumer. select, zero-capacity (rendezvous) or unbounded channels: usecrossbeam-channelorflume.- Parallel iterators, or fine-grained
joinat scale: use Rayon. It matches parkring's pool on coarse work and is faster on very fine-grained joins at 8 threads. asynccode: parkring blocks threads; it has noasyncAPI yet.
Performance
Measured on an Apple M4; the method, every table and the charts are in docs/BENCHMARKS.md. In short:
- Faster than crossbeam with one producer and one consumer (91 against 80 million items/s), and faster than crossbeam-deque with one thief draining (85 against 70).
- Level with Rayon on compute-bound fork-join work (
fib(32)on 8 threads: 1.36 ms each). - Slower than crossbeam under contention (10 to 17% from 2 + 2 threads,
about 70% with 8 producers and 1 consumer).
ScqQueueis 5 to 8 times slower thanLockFreeQueue. - Idle threads sleep: a waiting consumer wakes in about 9 µs and uses
under 2% of a core, like
std::sync::mpsc. A queue that spins instead (crossbeam'sArrayQueuehas no blocking API) wakes in 0.3 µs but uses a whole core.
What the verification found
The tests were written to fail on real bugs, and they did. Each item links to the write-up.
- The first version (0.1) failed spuriously in
try_push, lost items at non-power-of-two capacities, and spun forever when idle (DESIGN.md §9). - A capacity-1 overwrite in the Vyukov ring, found by drop accounting and independently by proptest, which shrank it to capacity 1 (DESIGN.md §5).
- A lost wakeup loom could not verify, because loom treats
SeqCstaccesses asAcqRel. The parking protocol was re-derived on read-modify-writes and release sequences, which loom does model, and the hot path got cheaper (DESIGN.md §4). - A hole in the SCQ paper's threshold bound: with more threads than capacity, an item could be stranded forever. A stress test hung 11 times in 40; the fix and the reasoning are in SCQ.md §4.
- The classic Chase-Lev double take: remove either
SeqCstfence and loom producesan element was taken twice: [0, 1, 1]. CI builds each mutant and requires that failure (DEQUE.md §3). - Two aliasing violations Miri caught and loom could not: retiring a
deque buffer through
Box::from_rawretags memory a thief may still be reading (DEQUE.md §6), and a latch's&selfargument stayed protected while the waiting thread freed it (POOL.md).
Verification
| covers | |
|---|---|
| loom | every interleaving (up to a preemption bound) of the parking protocol, close and channel-disconnect races, both queues' claims, the deque's pop/steal races, and the pool's sleep/wake; four deliberately broken builds must fail |
| Miri | the unsafe code in every component: uninitialised reads, double drops, leaks, aliasing and data races, under both stacked and tree borrows, with both parkers (Linux's futex through Miri's emulation of it), and 16 schedules per test for the channel and the Vyukov queue |
| fuzzing | the queues, the deque (including index wraparound) and the channel against sequential models; one minute per target on every pull request, fifteen minutes nightly |
| ThreadSanitizer | the queue and channel tests with both parkers (TSan cannot model the deque's standalone fences; loom and Miri cover those) |
| proptest | each queue against a VecDeque model (capacities 1–17), the deque against a VecDeque with wrapping indices |
| concurrency tests | exactly-once delivery and per-producer FIFO across many shapes; per-thief ordering for the deque; repeated runs to flush out rare schedules |
getrusage |
parked queues and an idle pool use about 35–55 µs of CPU over 300 ms |
RUSTFLAGS="--cfg loom"
RUSTFLAGS="--cfg loom"
&&
CI runs all of it on Linux, macOS and Windows, and on 32-bit and AArch64
Linux, plus the MSRV, docs, a FreeBSD check, both parkers under loom, the loom
mutants, cargo-deny, semver checks and the public API snapshot. The
unsafe code and its invariants are listed in UNSAFE.md.
Stability
parkring follows semver. From 1.0, these are promises:
- The public API is exactly what
public-api.txtlists, including which types areSend,Sync,Unpinand unwind-safe. CI fails if it changes without the file being updated, and a change that breaks it needs a major release. BoundedQueueis sealed. Only parkring's queues implement it, so methods can be added in minor releases.- Error enums and
Stealare exhaustive, likestd::sync::mpsc's and crossbeam's: you can match every variant. A new kind of failure would get a new type, not a new variant. - Capacity: a queue or channel holds at least the capacity you ask for;
capacity()reports the real number. How much it rounds up is not part of the contract. stdfeature: on by default and currently required. It exists so that a futureno_stdmode can be added without breaking anyone.- MSRV: Rust 1.85. Raising it is not a breaking change, but only happens in a minor release, never a patch, and is noted in the changelog. parkring supports at least the last four stable Rust releases.
- Platforms: tested on x86-64 Linux and Windows, AArch64 Linux and
macOS, and 32-bit i686 Linux;
ScqQueueexists only on 64-bit targets. Other targets withstduse the portableMutex+Condvarparker. - Items marked
#[doc(hidden)]or behind features whose names start with__are not public API.
Documentation
- DESIGN.md: the Vyukov queue, parking, and closing.
- SCQ.md: the fetch-add queue, its departures from the paper, and why it is slower here.
- DEQUE.md: the work-stealing deque and its memory orderings.
- POOL.md: the thread pool.
- BENCHMARKS.md: methodology and every measurement.
- UNSAFE.md: every
unsafeblock, the invariant it relies on, and which tool checks it. - CONTRIBUTING.md and SECURITY.md: how to run the checks, and how to report a soundness bug privately.
Project history
This crate began as bounded_mpmc_queue 0.1: a mutex queue and a Vyukov
queue. Version 0.2 audited that first version, fixed its bugs and added
parking, shutdown and the verification suite. Version 0.3
renamed it to parkring and added futex parking, SCQ, the work-stealing deque
and the pool. The git history shows each step.
Layout
src/
queue/lockfree.rs LockFreeQueue queue/scq/ ScqQueue
queue/blocking.rs BlockingQueue deque/ Worker, Stealer
pool/ ThreadPool, join sync/futex/ futex backends
sync/wait_queue/ parking sync/primitives.rs std/loom shim
tests/ loom, proptest, drop accounting, regressions, CPU checks
crates/parkring-bench/ benchmarks and chart generation (unpublished)
License
Licensed under the MIT license.