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
8 changes: 7 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,7 +76,8 @@ as a sequence of trades.
| Consumer threads | any number | **exactly one** | any number |
| Exceeding that | — | **undefined behaviour, no diagnostic** | — |
| Bounded wait (`try_push_for`) | yes | no | no |
| A waiting thread | sleeps | spins | spins |
| Bulk ops (`try_push_n` / `try_pop_n`) | no | yes | no |
| A blocked thread | sleeps | spins briefly, then sleeps | spins briefly, then sleeps |
| Built from | mutex + condition variables | atomics only | atomics only |

`MutexQueue` gives up nothing. `SpscQueue` gives up all but one thread per
Expand Down Expand Up @@ -159,6 +160,11 @@ so moving an item costs two. Higher is better.
- **`MpmcQueue` loses to the mutex at 4+4** — 13.3M against 21.6M. Lock-free
buys progress guarantees, not throughput: under contention its threads retry
while the mutex's threads sleep.
- **Sleeping when blocked costs the SPSC pair ~9%** (602M → 547M in the same
session); in exchange a blocked thread burns ~0.5% of a core instead of all
of it. Bulk transfer buys the loss back sevenfold: `try_push_n` /
`try_pop_n` publish the index once per batch and move **1.69G ops/s** at
batches of 64.

Error bars, conditions, the comparison against moodycamel and TBB, and the
reasoning: [docs/results.md](docs/results.md).
Expand Down
36 changes: 36 additions & 0 deletions bench/queue_bench.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,12 @@
#include <cq/mutex_queue.hpp>
#include <cq/spsc_queue.hpp>

#include <array>
#include <cstddef>
#include <cstdint>
#include <cstdlib>
#include <memory>
#include <span>
#include <thread>

#include <benchmark/benchmark.h>
Expand Down Expand Up @@ -136,6 +138,32 @@ using MpmcQueue = cq::MpmcQueue<std::uint64_t>;
using MutexQueue = cq::MutexQueue<std::uint64_t>;
using SpscQueue = cq::SpscQueue<std::uint64_t>;

// SPSC bulk transfer: kBulkBatch items per call, one index publish per batch.
// Each benchmark iteration moves one full batch on each side.
constexpr std::size_t kBulkBatch = 64;
void BM_SpscBulkThroughput(benchmark::State& state) {
const bool is_producer = state.thread_index() == 0;
auto& queue = *shared_queue<SpscQueue>;
std::array<std::uint64_t, kBulkBatch> buf{};
if (is_producer) {
for (auto _ : state) {
std::size_t done = 0;
while (done < kBulkBatch) {
done += queue.try_push_n(std::span{buf}.subspan(done));
}
}
} else {
for (auto _ : state) {
std::size_t done = 0;
while (done < kBulkBatch) {
done += queue.try_pop_n(std::span{buf}.subspan(done));
}
}
}
state.SetItemsProcessed(state.iterations() * static_cast<std::int64_t>(kBulkBatch) *
state.threads());
}

// SPSC: 1 producer + 1 consumer; MPMC: 4 + 4. Google Benchmark appends the
// /threads:N suffix to the reported name. SpscQueue's contract allows one
// thread per side, so it registers only the 2-thread shape.
Expand Down Expand Up @@ -166,6 +194,14 @@ BENCHMARK(BM_QueueThroughput<SpscQueue>)
->MinTime(kMinTimeSeconds)
->Name("SpscQueue/throughput");

BENCHMARK(BM_SpscBulkThroughput)
->Setup(setup_queue<SpscQueue>)
->Teardown(teardown_queue<SpscQueue>)
->Threads(kSpscThreads)
->UseRealTime()
->MinTime(kMinTimeSeconds)
->Name("SpscQueue/bulk64_throughput");

BENCHMARK(BM_QueueThroughput<MpmcQueue>)
->Setup(setup_queue<MpmcQueue>)
->Teardown(teardown_queue<MpmcQueue>)
Expand Down
33 changes: 33 additions & 0 deletions docs/results.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ class at different points in time:
| v2.1 | `SpscQueue` | Cached peer indices — same class, optimised |
| v2.5 | `MpmcQueue` | Vyukov-style bounded MPMC |
| v3 | — | The three against moodycamel and TBB |
| v3.1 | `SpscQueue`, `MpmcQueue` | Sleep when blocked; SPSC bulk ops; MPMC pow2 indexing |

## How these were measured

Expand Down Expand Up @@ -141,4 +142,36 @@ shared ring. The queue that beat it gave both of those up. Lock-free bought
progress guarantees and large wins where contention is structurally limited —
not throughput under real MPMC contention.

## v3.1 — sleep when blocked, bulk transfer, pow2 indexing

Three features closing the gap to the industrial queues, measured in one
session (mutex control reproduced within noise; load ~5).

**Blocked ops now sleep instead of spinning.** Blocking `push`/`pop` yield
briefly, then sleep with doubling backoff (`cq/backoff.hpp`): a consumer
blocked for 2 s used to burn **1975 ms** of CPU, now **10 ms** (~0.5% of a
core), with wakeup latency bounded at 1 ms. Two rejected designs earned their
place in this writeup: a futex gate the publisher must check cost **70%** of
the SPSC pair throughput (the check after a release store forces the store
buffer to drain), and merely inlining the sleep machinery into `push`/`pop`
cost **4×** on the uncontended round trip by pushing the functions past the
inliner's budget — hence the `noinline` cold `Backoff::wait()`. What remains:
the pair throughput runs ~9% below pure spinning (547M vs 602M ops/s),
because a stalled side now sleeps through the stall instead of burning a core
to catch the first free slot. The round trip is unchanged at 1.80G ops/s
(this session's control level; the v2.1 session above recorded 1.56G —
cross-session levels drift with machine load, ratios within a session are
the signal).

**SPSC bulk transfer.** `try_push_n`/`try_pop_n` move up to a span's worth
of items with one index publish per call — the batched-publish lever the
v2.1 section named.
At batches of 64 the pair moves **1.69G ± 0.09G ops/s**, 3.1× the single-op
pair, while crossing cores.

**Power-of-two indexing for MpmcQueue.** When capacity is a power of two the
ticket→slot map is one AND instead of a division. On the M2 the difference is
within noise (the divider is not this queue's bottleneck); kept because it is
strictly cheaper and capacity stays arbitrary.

_Roadmap complete through v3. Stretch (a thread pool) remains._
58 changes: 58 additions & 0 deletions include/cq/backoff.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
#ifndef CQ_BACKOFF_HPP_
#define CQ_BACKOFF_HPP_

#include <algorithm>
#include <chrono>
#include <thread>

namespace cq {

// Blocked-side wait policy: yield briefly, then sleep with doubling backoff.
// The sleeping side pays for its own wakeup latency (bounded by kMaxSleep);
// the publish path pays nothing at all — both alternatives measured worse. A
// futex gate the publisher must check cost 70% of the SPSC pair throughput,
// and even inlining this wait's sleep machinery into push()/pop() cost 4x on
// the uncontended round trip, which is why wait() is noinline.
/// wait() calls that yield before the first sleep.
inline constexpr int kSpinsBeforeSleep = 64;
/// First sleep duration; doubles on each subsequent sleep.
inline constexpr auto kInitialSleep = std::chrono::microseconds{4};
/// Sleep cap — the bound on wakeup latency once a waiter is asleep.
inline constexpr auto kMaxSleep = std::chrono::microseconds{1000};

#if defined(__GNUC__) || defined(__clang__)
#define CQ_NOINLINE [[gnu::noinline]]
#else
#define CQ_NOINLINE
#endif

/// One blocked wait: construct fresh, call wait() each time the predicate
/// still fails. Yields for the first kSpinsBeforeSleep calls, then sleeps
/// with doubling duration up to kMaxSleep.
class Backoff {
public:
/// One step of the schedule: a yield while spins remain, then a sleep
/// whose duration doubles up to kMaxSleep. noinline keeps the cold sleep
/// machinery out of the caller's inlining budget (see the file comment).
CQ_NOINLINE void wait();

private:
int spins_ = 0;
std::chrono::microseconds delay_ = kInitialSleep;
};

#undef CQ_NOINLINE

inline void Backoff::wait() {
if (spins_ < kSpinsBeforeSleep) {
++spins_;
std::this_thread::yield();
return;
}
std::this_thread::sleep_for(delay_);
delay_ = std::min(delay_ * 2, kMaxSleep);
}

} // namespace cq

#endif // CQ_BACKOFF_HPP_
15 changes: 11 additions & 4 deletions include/cq/mpmc_queue.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -5,15 +5,17 @@
#include <cstddef>
#include <vector>

#include "cq/backoff.hpp"
#include "cq/cache_line.hpp"

namespace cq {

/// v2.5: bounded FIFO ring for any number of producer and consumer threads,
/// synchronized with atomics only — Dmitry Vyukov's bounded MPMC design,
/// where every slot carries its own sequence counter and the two position
/// counters only hand out tickets. The blocking push()/pop() spin with
/// std::this_thread::yield() instead of sleeping.
/// counters only hand out tickets. The blocking push()/pop() spin briefly,
/// then sleep with doubling timed backoff — near-zero CPU while blocked,
/// wakeup within kMaxSleep.
///
/// Thread-safety: after construction, all member functions may be called
/// concurrently from any number of threads. closed() and size() return
Expand All @@ -34,7 +36,7 @@ namespace cq {
/// constructed up front) and MoveAssignable.
template <typename T>
// The "excessive padding" the analyzer flags is deliberate: each position
// counter gets a private cache line (see kCacheLineSize).
// counter gets a private cache line (see cq/cache_line.hpp).
// NOLINTNEXTLINE(clang-analyzer-optin.performance.Padding)
class MpmcQueue {
public:
Expand Down Expand Up @@ -109,8 +111,13 @@ class MpmcQueue {

[[nodiscard]] static std::vector<Slot> make_slots(std::size_t capacity);

// ticket -> slot position: a masked AND when capacity is a power of two,
// a division otherwise.
[[nodiscard]] std::size_t slot_index(std::size_t ticket) const noexcept;

std::vector<Slot> slots_;
// Tickets only ever increase; a slot is addressed by ticket % slots_.size().
std::size_t mask_; // slots_.size() - 1 if it is a power of two, else 0
// Tickets only ever increase; slot_index() maps them into the ring.
// The counters are on separate lines: producers hammer one, consumers the
// other, and the slots' own sequence counters carry the actual handoff.
alignas(kCacheLineSize) std::atomic<std::size_t> enqueue_pos_ = 0;
Expand Down
25 changes: 18 additions & 7 deletions include/cq/mpmc_queue.ipp
Original file line number Diff line number Diff line change
Expand Up @@ -5,15 +5,16 @@
#define CQ_MPMC_QUEUE_IPP_

#include <atomic>
#include <bit>
#include <cstddef>
#include <stdexcept>
#include <thread>
#include <utility>

namespace cq {

template <typename T>
MpmcQueue<T>::MpmcQueue(std::size_t capacity) : slots_(make_slots(capacity)) {
MpmcQueue<T>::MpmcQueue(std::size_t capacity)
: slots_(make_slots(capacity)), mask_(std::has_single_bit(capacity) ? capacity - 1 : 0) {
for (std::size_t i = 0; i < slots_.size(); ++i) {
// Lap zero: every slot starts free for the producer holding ticket i.
// Relaxed is enough — no other thread can touch the queue yet.
Expand Down Expand Up @@ -67,15 +68,16 @@ std::vector<typename MpmcQueue<T>::Slot> MpmcQueue<T>::make_slots(std::size_t ca
template <typename T>
bool MpmcQueue<T>::push(T value) {
// The closed check comes first so a close is honored even when the ring
// has room.
// has room. Yield briefly, then sleep with backoff instead of burning CPU.
Backoff backoff;
while (true) {
if (closed_.load(std::memory_order_relaxed)) {
return false;
}
if (try_enqueue(value)) {
return true;
}
std::this_thread::yield();
backoff.wait();
}
}

Expand All @@ -89,6 +91,7 @@ bool MpmcQueue<T>::try_push(T value) {

template <typename T>
bool MpmcQueue<T>::pop(T& out) {
Backoff backoff;
while (true) {
if (try_dequeue(out)) {
return true;
Expand All @@ -99,7 +102,7 @@ bool MpmcQueue<T>::pop(T& out) {
if (closed_.load(std::memory_order_acquire)) {
return try_dequeue(out);
}
std::this_thread::yield();
backoff.wait();
}
}

Expand All @@ -112,7 +115,7 @@ template <typename T>
bool MpmcQueue<T>::try_enqueue(T& value) {
auto ticket = enqueue_pos_.load(std::memory_order_relaxed);
while (true) {
auto& slot = slots_[ticket % slots_.size()];
auto& slot = slots_[slot_index(ticket)];
const auto seq = slot.sequence.load(std::memory_order_acquire);
const auto dif = static_cast<std::ptrdiff_t>(seq - (2 * ticket));
if (dif == 0) { // Slot is free for this ticket — race to claim it
Expand All @@ -134,7 +137,7 @@ template <typename T>
bool MpmcQueue<T>::try_dequeue(T& out) {
auto ticket = dequeue_pos_.load(std::memory_order_relaxed);
while (true) {
auto& slot = slots_[ticket % slots_.size()];
auto& slot = slots_[slot_index(ticket)];
const auto seq = slot.sequence.load(std::memory_order_acquire);
const auto dif = static_cast<std::ptrdiff_t>(seq - ((2 * ticket) + 1));
if (dif == 0) { // Slot holds data for this ticket — race to claim it
Expand All @@ -157,6 +160,14 @@ void MpmcQueue<T>::close() noexcept {
closed_.store(true, std::memory_order_release);
}

template <typename T>
std::size_t MpmcQueue<T>::slot_index(std::size_t ticket) const noexcept {
// One AND when the mask is nonzero (capacity a power of two; capacity 1
// computes 1 - 1 == 0 and falls through to the equally free % 1), a
// division otherwise.
return mask_ != 0 ? (ticket & mask_) : (ticket % slots_.size());
}

template <typename T>
bool MpmcQueue<T>::closed() const noexcept {
return closed_.load(std::memory_order_acquire);
Expand Down
35 changes: 26 additions & 9 deletions include/cq/spsc_queue.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -3,18 +3,21 @@

#include <atomic>
#include <cstddef>
#include <span>
#include <vector>

#include "cq/backoff.hpp"
#include "cq/cache_line.hpp"

namespace cq {

/// v2: bounded FIFO ring for exactly one producer thread and one consumer
/// thread, synchronized with atomics only — no mutex, no condition variables.
/// The blocking push()/pop() spin with std::this_thread::yield() instead of
/// sleeping. Each side caches the peer's index (added in v2.1), so a steady
/// stream of pushes and pops decides "there is room" / "there is data"
/// without touching the other core's cache line.
/// The blocking push()/pop() spin briefly, then sleep with doubling timed
/// backoff — near-zero CPU while blocked, wakeup within kMaxSleep. Each side
/// caches the peer's index (added in v2.1), so a steady stream of pushes and
/// pops decides "there is room" / "there is data" without touching the other
/// core's cache line.
///
/// Thread-safety: at most one thread may call the producer side (push,
/// try_push) and at most one thread the consumer side (pop, try_pop),
Expand All @@ -34,7 +37,7 @@ namespace cq {
/// constructed up front) and MoveAssignable.
template <typename T>
// The "excessive padding" the analyzer flags is deliberate: head_ and tail_
// each get a private cache line (see kCacheLineSize above).
// each get a private cache line (see cq/cache_line.hpp).
// NOLINTNEXTLINE(clang-analyzer-optin.performance.Padding)
class SpscQueue {
public:
Expand All @@ -51,7 +54,8 @@ class SpscQueue {
SpscQueue(SpscQueue&&) = delete;
SpscQueue& operator=(SpscQueue&&) = delete;

/// Enqueues a value, spinning while the queue is full. Producer side.
/// 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).
[[nodiscard]] bool push(T value);
Expand All @@ -61,8 +65,8 @@ class SpscQueue {
/// @return false if the queue is full or closed.
[[nodiscard]] bool try_push(T value);

/// Dequeues into out, spinning while the queue is empty and open.
/// Consumer side.
/// Dequeues into out, waiting while the queue is empty and open (brief
/// spin, then a timed sleep). Consumer side.
/// @param[out] out Receives the dequeued element on success.
/// @return false once the queue is closed and drained.
[[nodiscard]] bool pop(T& out);
Expand All @@ -73,9 +77,22 @@ class SpscQueue {
/// @return false if the queue is empty.
[[nodiscard]] bool try_pop(T& out);

/// Moves up to items.size() items from the front of items into the ring
/// with a single index publish — the batched-publish optimization.
/// Producer side.
/// @param items Source span; the first k elements are moved from.
/// @return k, the number enqueued (0 if the ring is full or closed).
[[nodiscard]] std::size_t try_push_n(std::span<T> items);

/// Moves up to out.size() items into the front of out with a single index
/// publish. Consumer side.
/// @param[out] out Destination span; the first k elements are written.
/// @return k, the number dequeued (0 if the ring is empty).
[[nodiscard]] std::size_t try_pop_n(std::span<T> out);

/// Closes the queue: push() and try_push() refuse new values, pop() drains
/// what remains and then returns false. Idempotent; callable from any
/// thread. Spinning push()/pop() calls return once they observe the close.
/// thread. Waiting push()/pop() calls wake and return.
///
/// To guarantee the consumer drains every item, stop the producer before
/// calling close(): a push racing with close() may enqueue an item after
Expand Down
Loading
Loading