diff --git a/README.md b/README.md index 1b56c33..9338499 100644 --- a/README.md +++ b/README.md @@ -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 @@ -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). diff --git a/bench/queue_bench.cpp b/bench/queue_bench.cpp index 81b41de..5940655 100644 --- a/bench/queue_bench.cpp +++ b/bench/queue_bench.cpp @@ -16,10 +16,12 @@ #include #include +#include #include #include #include #include +#include #include #include @@ -136,6 +138,32 @@ using MpmcQueue = cq::MpmcQueue; using MutexQueue = cq::MutexQueue; using SpscQueue = cq::SpscQueue; +// 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; + std::array 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(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. @@ -166,6 +194,14 @@ BENCHMARK(BM_QueueThroughput) ->MinTime(kMinTimeSeconds) ->Name("SpscQueue/throughput"); +BENCHMARK(BM_SpscBulkThroughput) + ->Setup(setup_queue) + ->Teardown(teardown_queue) + ->Threads(kSpscThreads) + ->UseRealTime() + ->MinTime(kMinTimeSeconds) + ->Name("SpscQueue/bulk64_throughput"); + BENCHMARK(BM_QueueThroughput) ->Setup(setup_queue) ->Teardown(teardown_queue) diff --git a/docs/results.md b/docs/results.md index a668776..500435a 100644 --- a/docs/results.md +++ b/docs/results.md @@ -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 @@ -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._ diff --git a/include/cq/backoff.hpp b/include/cq/backoff.hpp new file mode 100644 index 0000000..1faa0b3 --- /dev/null +++ b/include/cq/backoff.hpp @@ -0,0 +1,58 @@ +#ifndef CQ_BACKOFF_HPP_ +#define CQ_BACKOFF_HPP_ + +#include +#include +#include + +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_ diff --git a/include/cq/mpmc_queue.hpp b/include/cq/mpmc_queue.hpp index 9c9da97..a90ce89 100644 --- a/include/cq/mpmc_queue.hpp +++ b/include/cq/mpmc_queue.hpp @@ -5,6 +5,7 @@ #include #include +#include "cq/backoff.hpp" #include "cq/cache_line.hpp" namespace cq { @@ -12,8 +13,9 @@ 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 @@ -34,7 +36,7 @@ namespace cq { /// constructed up front) and MoveAssignable. template // 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: @@ -109,8 +111,13 @@ class MpmcQueue { [[nodiscard]] static std::vector 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 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 enqueue_pos_ = 0; diff --git a/include/cq/mpmc_queue.ipp b/include/cq/mpmc_queue.ipp index 02faef6..c642b23 100644 --- a/include/cq/mpmc_queue.ipp +++ b/include/cq/mpmc_queue.ipp @@ -5,15 +5,16 @@ #define CQ_MPMC_QUEUE_IPP_ #include +#include #include #include -#include #include namespace cq { template -MpmcQueue::MpmcQueue(std::size_t capacity) : slots_(make_slots(capacity)) { +MpmcQueue::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. @@ -67,7 +68,8 @@ std::vector::Slot> MpmcQueue::make_slots(std::size_t ca template bool MpmcQueue::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; @@ -75,7 +77,7 @@ bool MpmcQueue::push(T value) { if (try_enqueue(value)) { return true; } - std::this_thread::yield(); + backoff.wait(); } } @@ -89,6 +91,7 @@ bool MpmcQueue::try_push(T value) { template bool MpmcQueue::pop(T& out) { + Backoff backoff; while (true) { if (try_dequeue(out)) { return true; @@ -99,7 +102,7 @@ bool MpmcQueue::pop(T& out) { if (closed_.load(std::memory_order_acquire)) { return try_dequeue(out); } - std::this_thread::yield(); + backoff.wait(); } } @@ -112,7 +115,7 @@ template bool MpmcQueue::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(seq - (2 * ticket)); if (dif == 0) { // Slot is free for this ticket — race to claim it @@ -134,7 +137,7 @@ template bool MpmcQueue::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(seq - ((2 * ticket) + 1)); if (dif == 0) { // Slot holds data for this ticket — race to claim it @@ -157,6 +160,14 @@ void MpmcQueue::close() noexcept { closed_.store(true, std::memory_order_release); } +template +std::size_t MpmcQueue::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 bool MpmcQueue::closed() const noexcept { return closed_.load(std::memory_order_acquire); diff --git a/include/cq/spsc_queue.hpp b/include/cq/spsc_queue.hpp index 88478e8..52d7ae5 100644 --- a/include/cq/spsc_queue.hpp +++ b/include/cq/spsc_queue.hpp @@ -3,18 +3,21 @@ #include #include +#include #include +#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), @@ -34,7 +37,7 @@ namespace cq { /// constructed up front) and MoveAssignable. template // 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: @@ -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); @@ -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); @@ -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 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 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 diff --git a/include/cq/spsc_queue.ipp b/include/cq/spsc_queue.ipp index fc61963..fa1d1f0 100644 --- a/include/cq/spsc_queue.ipp +++ b/include/cq/spsc_queue.ipp @@ -4,11 +4,12 @@ #ifndef CQ_SPSC_QUEUE_IPP_ #define CQ_SPSC_QUEUE_IPP_ +#include #include #include #include +#include #include -#include #include namespace cq { @@ -36,6 +37,10 @@ std::size_t SpscQueue::ring_slots(std::size_t capacity) { // visible to the producer before it overwrites the slot. Each side loads its // *own* index relaxed — it is the only writer of that index. // +// A blocked side needs no wakeup from the publisher: it sleeps with timed +// backoff and re-polls (see cq/backoff.hpp), so the publish path carries no +// waiter bookkeeping at all. +// // close() is a release store; the consumer's acquire load of closed_ in pop() // therefore also makes every push that preceded the close visible, which is // what lets pop() decide "closed and drained" with one final try_pop. @@ -64,8 +69,9 @@ template bool SpscQueue::push(T value) { const auto tail = tail_.load(std::memory_order_relaxed); const auto slot_after = next(tail); - // Wait for a free slot, bailing out if the queue closes first. The closed - // check stays first so a close is honored even when the ring has room. + // The closed check stays first so a close is honored even when there is + // room. + Backoff backoff; while (true) { if (closed_.load(std::memory_order_relaxed)) { return false; @@ -73,7 +79,7 @@ bool SpscQueue::push(T value) { if (has_room(slot_after)) { break; } - std::this_thread::yield(); + backoff.wait(); } enqueue(tail, std::move(value)); return true; @@ -94,6 +100,7 @@ bool SpscQueue::try_push(T value) { template bool SpscQueue::pop(T& out) { + Backoff backoff; while (true) { if (try_pop(out)) { return true; @@ -103,7 +110,7 @@ bool SpscQueue::pop(T& out) { if (closed_.load(std::memory_order_acquire)) { return try_pop(out); } - std::this_thread::yield(); + backoff.wait(); } } @@ -129,6 +136,51 @@ void SpscQueue::dequeue(std::size_t head, T& out) { head_.store(next(head), std::memory_order_release); } +template +std::size_t SpscQueue::try_push_n(std::span items) { + if (items.empty() || closed_.load(std::memory_order_relaxed)) { + return 0; + } + auto tail = tail_.load(std::memory_order_relaxed); + // Free slots as seen through the cache; refresh once if it says none. + const auto free_slots = [&] { + return (head_cache_ + buffer_.size() - 1 - tail) % buffer_.size(); + }; + if (free_slots() == 0) { + head_cache_ = head_.load(std::memory_order_acquire); + } + const auto count = std::min(items.size(), free_slots()); + for (std::size_t i = 0; i < count; ++i) { + buffer_[tail] = std::move(items[i]); + tail = next(tail); + } + if (count != 0) { + tail_.store(tail, std::memory_order_release); // one publish for the batch + } + return count; +} + +template +std::size_t SpscQueue::try_pop_n(std::span out) { + if (out.empty()) { + return 0; + } + auto head = head_.load(std::memory_order_relaxed); + const auto available = [&] { return (tail_cache_ + buffer_.size() - head) % buffer_.size(); }; + if (available() == 0) { + tail_cache_ = tail_.load(std::memory_order_acquire); + } + const auto count = std::min(out.size(), available()); + for (std::size_t i = 0; i < count; ++i) { + out[i] = std::move(buffer_[head]); + head = next(head); + } + if (count != 0) { + head_.store(head, std::memory_order_release); // one publish for the batch + } + return count; +} + template bool SpscQueue::has_room(std::size_t slot_after) noexcept { if (slot_after != head_cache_) { diff --git a/tests/spsc_queue_test.cpp b/tests/spsc_queue_test.cpp index f027d33..6e4f27c 100644 --- a/tests/spsc_queue_test.cpp +++ b/tests/spsc_queue_test.cpp @@ -3,8 +3,10 @@ #include +#include #include #include +#include #include #include @@ -56,5 +58,83 @@ TEST(SpscQueue, RepeatedFillAndDrainCyclesPreserveFifoOrder) { } } +// --- bulk operations ------------------------------------------------------ +// +// try_push_n / try_pop_n move up to n items with one index publish per call. + +TEST(SpscQueue, TryPushNStopsAtCapacity) { + SpscQueue q(4); + ASSERT_TRUE(q.try_push(0)); + ASSERT_TRUE(q.try_push(1)); + std::array items = {2, 3, 4, 4}; // the last item won't fit anyway + EXPECT_EQ(q.try_push_n(items), 2U); // only two slots left + EXPECT_EQ(q.size(), 4U); + int out = 0; + for (int expected = 0; expected < 4; ++expected) { + ASSERT_TRUE(q.try_pop(out)); + EXPECT_EQ(out, expected); + } +} + +TEST(SpscQueue, TryPopNStopsAtAvailable) { + SpscQueue q(4); + ASSERT_TRUE(q.try_push(7)); + ASSERT_TRUE(q.try_push(8)); + std::array out = {}; + EXPECT_EQ(q.try_pop_n(out), 2U); + EXPECT_EQ(out[0], 7); + EXPECT_EQ(out[1], 8); + EXPECT_EQ(q.try_pop_n(out), 0U); // empty +} + +TEST(SpscQueue, BulkAndSingleOpsInterleavePreservingFifo) { + SpscQueue q(3); + std::array items = {1, 2}; + EXPECT_EQ(q.try_push_n(items), 2U); + ASSERT_TRUE(q.try_push(3)); + std::array out = {}; + EXPECT_EQ(q.try_pop_n(out), 2U); + EXPECT_EQ(out[0], 1); + EXPECT_EQ(out[1], 2); + int single = 0; + ASSERT_TRUE(q.try_pop(single)); + EXPECT_EQ(single, 3); +} + +TEST(SpscQueue, BulkOpsWrapTheRing) { + constexpr int kCycles = 10; + SpscQueue q(3); + std::array out = {}; + for (int cycle = 0; cycle < kCycles; ++cycle) { + std::array items = {cycle * 3, (cycle * 3) + 1, (cycle * 3) + 2}; + ASSERT_EQ(q.try_push_n(items), 3U); + ASSERT_EQ(q.try_pop_n(out), 3U); + EXPECT_EQ(out[0], cycle * 3); + EXPECT_EQ(out[2], (cycle * 3) + 2); + } +} + +TEST(SpscQueue, TryPushNAfterCloseReturnsZero) { + SpscQueue q(4); + q.close(); + std::array items = {1, 2}; + EXPECT_EQ(q.try_push_n(items), 0U); +} + +// NOLINTBEGIN(clang-analyzer-cplusplus.NewDeleteLeaks): the analyzer cannot +// see the ownership handoff through the ring on libstdc++ and reports a +// false leak; the asserts below prove both pointers arrive intact. +TEST(SpscQueue, BulkOpsSupportMoveOnlyTypes) { + SpscQueue> q(4); + std::array, 2> in = {std::make_unique(1), std::make_unique(2)}; + EXPECT_EQ(q.try_push_n(in), 2U); + EXPECT_EQ(in[0], nullptr) << "pushed items must be moved from"; + std::array, 2> out; + EXPECT_EQ(q.try_pop_n(out), 2U); + EXPECT_EQ(*out[0], 1); + EXPECT_EQ(*out[1], 2); +} +// NOLINTEND(clang-analyzer-cplusplus.NewDeleteLeaks) + } // namespace } // namespace cq