Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 10 additions & 4 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -32,16 +32,22 @@ jobs:
- name: Test
run: ctest --test-dir build-tsan --output-on-failure --no-tests=error

bench-build:
name: Benchmarks (Release, smoke-run)
release:
name: Tests + benchmark smoke-run (Release)
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Configure
run: cmake -B build-rel -DCMAKE_BUILD_TYPE=Release -DCQ_BUILD_TESTS=OFF
# Tests stay ON. Without this job they only ever run Debug+TSan, which
# compiles no NDEBUG path, applies no optimiser, and runs the checksum
# stress test roughly an order of magnitude slower — a very different
# and much narrower slice of the interleaving space. Two of the three
# queues here are lock-free, which is exactly the code where -O0 and
# -O3 diverge.
run: cmake -B build-rel -DCMAKE_BUILD_TYPE=Release
- name: Build
run: cmake --build build-rel
- name: Smoke-run benchmarks
- name: Test and smoke-run benchmarks
run: ctest --test-dir build-rel --output-on-failure --no-tests=error

lint:
Expand Down
6 changes: 4 additions & 2 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -36,8 +36,10 @@ target_link_libraries(cq INTERFACE Threads::Threads)
add_library(cq_warnings INTERFACE)
target_compile_options(cq_warnings INTERFACE
$<$<CXX_COMPILER_ID:GNU,Clang,AppleClang>:-Wall -Wextra -Wpedantic -Wconversion>
# Validates Doxygen comments against signatures (clang-only).
$<$<CXX_COMPILER_ID:Clang,AppleClang>:-Wdocumentation>
# Validates Doxygen comments against signatures (clang-only). An error, not a
# warning, because STYLE.md promises the check is enforced — as a warning it
# exits 0 and a stale @param ships green.
$<$<CXX_COMPILER_ID:Clang,AppleClang>:-Wdocumentation -Werror=documentation>
$<$<CXX_COMPILER_ID:MSVC>:/W4>)

# Applied only to our own targets: instrumenting the FetchContent deps just
Expand Down
63 changes: 60 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,8 +49,8 @@ entirely, which is only possible by giving something else up.

## The contract

All three promise the same things. Choose between them on threading, not on
behaviour.
All three promise the same things, with one exception noted below. Choose
between them on threading, not on behaviour.

| | |
|---|---|
Expand All @@ -61,11 +61,67 @@ behaviour.
| **Shutdown** | `close()` refuses new items and wakes every waiter |
| **Draining** | Items already queued still come out after `close()` |
| **Loss** | Nothing accepted is dropped, duplicated, or reordered |
| **Element type** | Any `T` that is `DefaultConstructible` and `MoveAssignable` |
| **Element type** | Any `T` that is `DefaultConstructible` and `MoveAssignable` — plus `noexcept` move assignment for `MpmcQueue` and for `SpscQueue`'s bulk ops, both of which `static_assert` it |
| **Lifetime** | The queue must outlive every thread using it |
| **Non-blocking loops** | `try_push`/`try_pop` return `false` for "not now" and "never again" alike; `closed()` is what lets a retry loop terminate — [see below](#non-blocking-loops) |
| **A throwing move** | The one place the three differ — see below |

Every operation reports whether it succeeded, and every one is `[[nodiscard]]`.

**If `T`'s move assignment can throw**, the three part company. `MutexQueue`
and `SpscQueue` keep their indices intact — nothing lost or duplicated — but a
failed pop leaves both the destination and the still-queued element
valid-but-unspecified. `SpscQueue`'s bulk ops publish one index per batch and
would lose or duplicate elements, so they reject a throwing `T` at compile
time. `MpmcQueue` cannot survive a throw at all — it strands a claimed ticket
and the queue stalls on it forever — so it rejects such a `T` at compile time.

### Non-blocking loops

`try_push`/`try_pop` return `false` for "not now" and for "never again" alike,
so a loop that only tests the return value never terminates after `close()`.
`closed()` is the discriminator, and both sides have a trap.

**Producer.** `try_push` is a by-value sink, so a failed attempt has already
consumed an rvalue argument. Re-materialise it every pass:

```cpp
while (!q.try_push(make_value())) { // NOT try_push(std::move(v))
if (q.closed()) break;
}
```

A move-only value that cannot be re-created has no correct `try_push` loop at
all — after one failed attempt it is gone. Blocking `push()` narrows the window
to "closed" but does not close it: it too drops the value it was given.

**Consumer.** Observing `closed()` is not the same as the queue being empty: a
producer may have pushed between the failed `try_pop` and the check. Re-attempt
once after observing it, and take whatever that attempt gives — which is what
`SpscQueue::pop`/`MpmcQueue::pop` do verbatim (`if (closed()) return
try_pop(out);`). `MutexQueue::pop` reaches the same result differently, by
resolving "closed and drained" under a single lock:

```cpp
for (T item;;) {
if (!q.try_pop(item)) {
if (!q.closed()) continue; // not now — keep trying
if (!q.try_pop(item)) break; // closed and drained
}
use(item);
}
```

The re-attempt is load-bearing on all three. `try_pop`, `closed()` and the
second `try_pop` are three separate steps — breaking on `closed()` alone
strands whatever arrived between the first two, and breaking only when the
re-attempt *fails* discards the element it just retrieved. `MutexQueue` is not
exempt: the lock hand-off after a failed `try_pop` is a wide window, not a
narrow one. What it lacks is a different hazard — a push racing `close()` —
which is why its `close()` has no stop-producers-first precondition while the
other two do. For those two, the re-attempt is also what orders the consumer
after the producer's last push.

## The three queues

They differ on one axis — **how many threads may touch each side** — and read
Expand All @@ -78,6 +134,7 @@ as a sequence of trades.
| Exceeding that | — | **undefined behaviour, no diagnostic** | — |
| Bounded wait (`try_push_for`) | yes | no | no |
| Bulk ops (`try_push_n` / `try_pop_n`) | no | yes | no |
| Throwing move assignment | survivable | survivable (single-element ops; bulk ops reject it) | rejected at compile time |
| A blocked thread | sleeps | spins briefly, then sleeps | spins briefly, then sleeps |
| Built from | mutex + condition variables | atomics only | atomics only |

Expand Down
4 changes: 2 additions & 2 deletions STYLE.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,8 @@ checked in CI).
- **Internal code** (tests, benchmarks, private members, function bodies):
plain `//` prose. Explain *why*, not *what*.
- Tag hygiene is compiler-enforced: clang builds compile with
`-Wdocumentation`, which rejects `@param` names that do not match the
signature. Keep comments in sync with code or the build fails.
`-Wdocumentation -Werror=documentation`, which rejects `@param` names that
do not match the signature. Keep comments in sync with code or the build fails.
- `TODO(username): description` for known follow-ups.

## Layout
Expand Down
37 changes: 24 additions & 13 deletions include/cq/mpmc_queue.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@

#include <atomic>
#include <cstddef>
#include <type_traits>
#include <vector>

#include "cq/backoff.hpp"
Expand All @@ -20,25 +21,34 @@ namespace cq {
/// Thread-safety: after construction, all member functions may be called
/// concurrently from any number of threads. closed() and size() return
/// advisory snapshots — drive control flow off the push/pop return values
/// instead.
/// instead, with one exception: try_push()/try_pop() return false for "not
/// now" and for "never again" alike, so a non-blocking retry loop needs
/// closed() to terminate.
///
/// Lifetime: the queue must outlive every thread using it — call close() and
/// join all producers/consumers before destruction. Destroying the queue
/// while a thread is spinning in push()/pop() is undefined behavior.
///
/// Exceptions: if T's move assignment throws while a slot is claimed, that
/// slot's sequence is never re-published and the queue degrades (later
/// operations on the slot spin); unlike the locked queue there is no way to
/// return a claimed ticket. Use element types whose move assignment cannot
/// throw.
/// Exceptions: a throwing move assignment would strand a claimed ticket whose
/// sequence is never re-published — the consumer can never advance past it and
/// the producer stalls once the ring wraps onto it — with no way to return the
/// ticket. A static_assert therefore requires a noexcept move assignment.
///
/// Arguments: push operations take T by value. An rvalue argument is moved
/// from at the call — even when the push fails, in which case the value is
/// discarded. An lvalue argument is copied and left intact. A try_push()
/// retry loop must therefore re-materialise its argument every pass.
///
/// @tparam T Element type. Must be DefaultConstructible (ring slots are
/// constructed up front) and MoveAssignable.
/// constructed up front) and nothrow-MoveAssignable (see Exceptions).
template <typename T>
// The "excessive padding" the analyzer flags is deliberate: each position
// counter gets a private cache line (see cq/cache_line.hpp).
// NOLINTNEXTLINE(clang-analyzer-optin.performance.Padding)
class MpmcQueue {
static_assert(std::is_nothrow_move_assignable_v<T>,
"MpmcQueue requires a noexcept move assignment");

public:
/// @param capacity Fixed number of slots; never resized.
/// @throws std::invalid_argument if capacity is 0.
Expand All @@ -52,13 +62,14 @@ class MpmcQueue {
MpmcQueue& operator=(MpmcQueue&&) = delete;

/// Enqueues a value, spinning while the queue is full.
/// @param value Element to enqueue; consumed even when the push fails.
/// @return false if the queue is closed (the value is dropped).
/// @param value Element to enqueue; see the class note on by-value arguments.
/// @return false if the queue is closed; the value is discarded.
[[nodiscard]] bool push(T value);

/// Enqueues a value without blocking.
/// @param value Element to enqueue; consumed even when the push fails.
/// @return false if the queue is full or closed.
/// @param value Element to enqueue; see the class note on by-value arguments.
/// @return false if the queue is full or closed. closed() tells them apart;
/// see the README's "Non-blocking loops" for the retry idiom.
[[nodiscard]] bool try_push(T value);

/// Dequeues into out, spinning while the queue is empty and open.
Expand All @@ -67,8 +78,8 @@ class MpmcQueue {
[[nodiscard]] bool pop(T& out);

/// Dequeues into out without blocking.
/// @param[out] out Receives the dequeued element on success; untouched on
/// failure.
/// @param[out] out Receives the dequeued element on success; untouched on a
/// false return; disturbed if T's move assignment throws (see Exceptions).
/// @return false if the queue is empty (including transiently, while a
/// producer has claimed the next slot but not yet published it).
[[nodiscard]] bool try_pop(T& out);
Expand Down
36 changes: 23 additions & 13 deletions include/cq/mutex_queue.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -19,15 +19,24 @@ namespace cq {
/// Thread-safety: after construction, all member functions may be called
/// concurrently from any number of producer and consumer threads. closed()
/// and size() return advisory snapshots — drive control flow off the
/// push/pop return values instead.
/// push/pop return values instead, with one exception: try_push()/try_pop()
/// return false for "not now" and for "never again" alike, so a non-blocking
/// retry loop needs closed() to terminate.
///
/// Lifetime: the queue must outlive every thread using it — call close()
/// and join all producers/consumers before destruction. Destroying the
/// queue while a thread is blocked in push()/pop() is undefined behavior.
///
/// Exceptions: if T's move assignment throws, the failing push()/try_push()
/// enqueues nothing and the failing pop()/try_pop() leaves the element
/// queued — the queue itself stays consistent.
/// Exceptions: if T's move assignment throws, the indices do not move — nothing
/// is lost or duplicated — but element values are not protected: a failed
/// push() enqueues nothing, and a failed pop() leaves both out and the
/// still-queued element valid-but-unspecified. Unreachable for a noexcept
/// move assignment.
///
/// Arguments: push operations take T by value. An rvalue argument is moved
/// from at the call — even when the push fails, in which case the value is
/// discarded. An lvalue argument is copied and left intact. A try_push()
/// retry loop must therefore re-materialise its argument every pass.
///
/// @tparam T Element type. Must be DefaultConstructible (ring slots are
/// constructed up front) and MoveAssignable.
Expand All @@ -46,13 +55,14 @@ class MutexQueue {
MutexQueue& operator=(MutexQueue&&) = delete;

/// Enqueues a value, blocking while the queue is full.
/// @param value Element to enqueue; consumed even when the push fails.
/// @return false if the queue is closed (the value is dropped).
/// @param value Element to enqueue; see the class note on by-value arguments.
/// @return false if the queue is closed; the value is discarded.
[[nodiscard]] bool push(T value);

/// Enqueues a value without blocking.
/// @param value Element to enqueue; consumed even when the push fails.
/// @return false if the queue is full or closed.
/// @param value Element to enqueue; see the class note on by-value arguments.
/// @return false if the queue is full or closed. closed() tells them apart;
/// see the README's "Non-blocking loops" for the retry idiom.
[[nodiscard]] bool try_push(T value);

/// Dequeues into out, blocking while the queue is empty and open.
Expand All @@ -61,8 +71,8 @@ class MutexQueue {
[[nodiscard]] bool pop(T& out);

/// Dequeues into out without blocking.
/// @param[out] out Receives the dequeued element on success; untouched on
/// failure.
/// @param[out] out Receives the dequeued element on success; untouched on a
/// false return; disturbed if T's move assignment throws (see Exceptions).
/// @return false if the queue is empty.
[[nodiscard]] bool try_pop(T& out);

Expand All @@ -71,7 +81,7 @@ class MutexQueue {
/// indefinitely, and try_push(), which does not wait at all.
/// @tparam Rep Arithmetic type of the timeout's tick count.
/// @tparam Period std::ratio giving the timeout's tick period.
/// @param value Element to enqueue; consumed even when the push fails.
/// @param value Element to enqueue; see the class note on by-value arguments.
/// @param timeout Longest time to wait. A non-positive timeout makes this
/// equivalent to try_push().
/// @return false if the timeout elapsed with the queue still full, or if
Expand All @@ -85,8 +95,8 @@ class MutexQueue {
/// and drains, or timeout elapses.
/// @tparam Rep Arithmetic type of the timeout's tick count.
/// @tparam Period std::ratio giving the timeout's tick period.
/// @param[out] out Receives the dequeued element on success; untouched on
/// failure.
/// @param[out] out Receives the dequeued element on success; untouched on a
/// false return; disturbed if T's move assignment throws (see Exceptions).
/// @param timeout Longest time to wait. A non-positive timeout makes this
/// equivalent to try_pop().
/// @return false if the timeout elapsed with the queue still empty, or
Expand Down
36 changes: 25 additions & 11 deletions include/cq/spsc_queue.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -23,18 +23,31 @@ namespace cq {
/// try_push) and at most one thread the consumer side (pop, try_pop),
/// concurrently with each other. close(), closed(), size(), and capacity()
/// may be called from any thread. closed() and size() return advisory
/// snapshots — drive control flow off the push/pop return values instead.
/// snapshots — drive control flow off the push/pop return values instead,
/// with one exception: try_push()/try_pop() return false for "not now" and
/// for "never again" alike, so a non-blocking retry loop needs closed() to
/// terminate.
///
/// Lifetime: the queue must outlive both threads using it — call close() and
/// join the producer/consumer before destruction. Destroying the queue while
/// a thread is spinning in push()/pop() is undefined behavior.
///
/// Exceptions: if T's move assignment throws, the failing push()/try_push()
/// enqueues nothing and the failing pop()/try_pop() leaves the element
/// queued — the queue itself stays consistent.
/// Exceptions: for single-element operations, if T's move assignment throws
/// the indices do not move — nothing is lost or duplicated — but element
/// values are not protected: a failed push() enqueues nothing, and a failed
/// pop() leaves both out and the still-queued element valid-but-unspecified.
/// Unreachable for a noexcept move assignment. try_push_n()/try_pop_n()
/// publish one index per batch, so a throw part-way through would lose or
/// duplicate elements; they static_assert a noexcept move assignment instead.
///
/// Arguments: push operations take T by value. An rvalue argument is moved
/// from at the call — even when the push fails, in which case the value is
/// discarded. An lvalue argument is copied and left intact. A try_push()
/// retry loop must therefore re-materialise its argument every pass.
///
/// @tparam T Element type. Must be DefaultConstructible (ring slots are
/// constructed up front) and MoveAssignable.
/// constructed up front) and MoveAssignable; nothrow for the bulk ops (see
/// Exceptions).
template <typename T>
// The "excessive padding" the analyzer flags is deliberate: head_ and tail_
// each get a private cache line (see cq/cache_line.hpp).
Expand All @@ -56,13 +69,14 @@ class SpscQueue {

/// Enqueues a value, waiting while the queue is full (brief spin, then a
/// timed sleep — near-zero CPU while blocked). Producer side.
/// @param value Element to enqueue; consumed even when the push fails.
/// @return false if the queue is closed (the value is dropped).
/// @param value Element to enqueue; see the class note on by-value arguments.
/// @return false if the queue is closed; the value is discarded.
[[nodiscard]] bool push(T value);

/// Enqueues a value without blocking. Producer side.
/// @param value Element to enqueue; consumed even when the push fails.
/// @return false if the queue is full or closed.
/// @param value Element to enqueue; see the class note on by-value arguments.
/// @return false if the queue is full or closed. closed() tells them apart;
/// see the README's "Non-blocking loops" for the retry idiom.
[[nodiscard]] bool try_push(T value);

/// Dequeues into out, waiting while the queue is empty and open (brief
Expand All @@ -72,8 +86,8 @@ class SpscQueue {
[[nodiscard]] bool pop(T& out);

/// Dequeues into out without blocking. Consumer side.
/// @param[out] out Receives the dequeued element on success; untouched on
/// failure.
/// @param[out] out Receives the dequeued element on success; untouched on a
/// false return; disturbed if T's move assignment throws (see Exceptions).
/// @return false if the queue is empty.
[[nodiscard]] bool try_pop(T& out);

Expand Down
7 changes: 7 additions & 0 deletions include/cq/spsc_queue.ipp
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
#include <limits>
#include <span>
#include <stdexcept>
#include <type_traits>
#include <utility>

namespace cq {
Expand Down Expand Up @@ -138,6 +139,10 @@ void SpscQueue<T>::dequeue(std::size_t head, T& out) {

template <typename T>
std::size_t SpscQueue<T>::try_push_n(std::span<T> items) {
// In the body, not on the class: single-element use of a throwing T must
// still compile.
static_assert(std::is_nothrow_move_assignable_v<T>,
"SpscQueue bulk transfer requires a noexcept move assignment");
if (items.empty() || closed_.load(std::memory_order_relaxed)) {
return 0;
}
Expand All @@ -162,6 +167,8 @@ std::size_t SpscQueue<T>::try_push_n(std::span<T> items) {

template <typename T>
std::size_t SpscQueue<T>::try_pop_n(std::span<T> out) {
static_assert(std::is_nothrow_move_assignable_v<T>,
"SpscQueue bulk transfer requires a noexcept move assignment");
if (out.empty()) {
return 0;
}
Expand Down
Loading
Loading