1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
//! Concurrency primitives built and verified from first principles: bounded
//! MPMC queues, a work-stealing deque, and a work-stealing thread pool.
//!
//! | | what it is |
//! |---|---|
//! | [`LockFreeQueue`] | Vyukov's per-slot sequence ring. The recommended queue. |
//! | [`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 a channel API, parallel iterators or `async`, use `std::sync::mpsc`,
//! crossbeam or Rayon 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:
//!
//! ```compile_fail
//! fn assert_sync<T: Sync>() {}
//! assert_sync::<parkring::LockFreeQueue<std::rc::Rc<()>>>();
//! ```
//!
//! ```compile_fail
//! fn assert_sync<T: Sync>() {}
//! assert_sync::<parkring::BlockingQueue<std::rc::Rc<()>>>();
//! ```
//!
//! ```compile_fail
//! fn assert_sync<T: Sync>() {}
//! assert_sync::<parkring::ScqQueue<std::rc::Rc<()>>>();
//! ```
//!
//! A deque's [`Worker`] belongs to one thread at a time: it is `Send` but not
//! `Sync`. [`Stealer`] is `Send + Sync`.
//!
//! ```compile_fail
//! fn assert_sync<T: Sync>() {}
//! assert_sync::<parkring::Worker<u32>>();
//! ```
//!
//! ```compile_fail
//! 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.
pub use ;
pub use ;
pub use ;
pub use ScqQueue;
pub use ;
pub use BoundedQueue;
/// Compiles and runs the README's code examples as doctests.
;