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 .clang-tidy
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,13 @@ Checks: >
-modernize-use-trailing-return-type,
-misc-include-cleaner
WarningsAsErrors: '*'
HeaderFilterRegex: 'include/cq/.*'
# Cover every first-party directory that holds headers, not just include/cq:
# tests/queue_test_util.hpp and bench/try_operation.hpp both sit outside it and
# would otherwise receive no diagnostics at all. Dependency headers stay out
# because every dependency is SYSTEM, not because of the anchor — the anchor
# only excludes segment-substring matches like microbench/, and still matches a
# bench/ segment above the repo root.
HeaderFilterRegex: '(^|/)(include/cq|bench|tests)/'
CheckOptions:
# Complexity contributed by macro expansion (GTest asserts, etc.) is not
# the author's complexity.
Expand Down
9 changes: 8 additions & 1 deletion CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -66,11 +66,18 @@ if(CQ_BUILD_TESTS)
endif()

if(CQ_BUILD_BENCHMARKS)
# SYSTEM, like the two below: benchmark is the only dependency here that
# arrives as a plain -I (googletest self-marks its interface SYSTEM), so it is
# the only one clang-tidy's header filter can reach. That matters because the
# filter matches a path segment anywhere above the repo, so a checkout under
# e.g. /home/runner/work/bench/bench/ would otherwise lint benchmark.h under
# WarningsAsErrors.
FetchContent_Declare(
benchmark
URL https://github.com/google/benchmark/archive/refs/tags/v1.9.1.tar.gz
URL_HASH SHA256=32131c08ee31eeff2c8968d7e874f3cb648034377dfc32a4c377fa8796d84981
DOWNLOAD_EXTRACT_TIMESTAMP TRUE)
DOWNLOAD_EXTRACT_TIMESTAMP TRUE
SYSTEM)
set(BENCHMARK_ENABLE_TESTING OFF CACHE BOOL "" FORCE)
set(BENCHMARK_ENABLE_INSTALL OFF CACHE BOOL "" FORCE)

Expand Down
16 changes: 16 additions & 0 deletions bench/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -4,3 +4,19 @@ target_link_libraries(queue_bench PRIVATE cq::cq cq_warnings benchmark::benchmar

# One-iteration smoke run; CI's Release job runs it via ctest.
add_test(NAME bench_smoke COMMAND queue_bench --benchmark_min_time=1x)

# The shallowest non-blocking run must report non-zero retry pressure, or the
# diagnostic has gone dead. Kept separate from bench_smoke so that test still
# gates the process exit code (PASS_REGULAR_EXPRESSION overrides that check).
#
# [^\n] keeps the match inside the _mean row: CTest's `.` crosses newlines, so
# without it the pattern can start on _mean and satisfy its tail from a later
# row. Asserting non-zero rather than a magnitude keeps the gate off the
# scheduler's back.
add_test(NAME bench_retry_pressure
COMMAND queue_bench
"--benchmark_filter=^MutexQueue/try_throughput/capacity:1/"
--benchmark_min_time=1000x
--benchmark_repetitions=10)
set_tests_properties(bench_retry_pressure PROPERTIES
PASS_REGULAR_EXPRESSION "_mean[^\n]*retries/op=[0-9.]*[1-9]")
148 changes: 138 additions & 10 deletions bench/queue_bench.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -11,18 +11,29 @@
// iteration, so pushes and pops balance. items/s = total ops/s (a push and
// its pop count as two). Spawn cost sits outside the timed region; add
// --benchmark_repetitions=10 for a spread.
//
// The /try_throughput sweeps use the non-blocking API instead, retrying a
// failed attempt inside the same iteration. Failures therefore cost time but
// never count as transferred work, so items/s stays the same unit as the
// blocking rows, and the pressure that caused them is reported separately as
// retries/op. Sweeping capacity is the point: it is where a lock-free ring's
// advantage over the mutex should narrow, since a shallow queue makes every
// op wait on its counterpart no matter how the waiting is implemented.

#include <cq/mpmc_queue.hpp>
#include <cq/mutex_queue.hpp>
#include <cq/spsc_queue.hpp>

#include "try_operation.hpp"

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

#include <benchmark/benchmark.h>
#include <concurrentqueue.h>
Expand Down Expand Up @@ -78,20 +89,32 @@ class MoodycamelQueue {
template <typename Queue>
std::unique_ptr<Queue> shared_queue;

// Half-full start: neither side begins blocked or spinning on the other, so
// the measurement starts in steady state. At capacity 1 the half is 0, so that
// sweep point deliberately starts empty — which is the retry pressure it is
// there to measure.
template <typename Queue>
void setup_queue(const benchmark::State& /*state*/) {
shared_queue<Queue> = std::make_unique<Queue>(kCapacity);
// Half-full start: neither side begins blocked or spinning on the other, so
// the measurement starts in steady state.
bool prefilled = true;
for (std::uint64_t i = 0; i < kCapacity / 2; ++i) {
prefilled = prefilled && shared_queue<Queue>->try_push(i);
}
if (!prefilled) {
std::abort(); // unreachable: fresh queue, stays below capacity
void install_half_full_queue(std::size_t capacity) {
shared_queue<Queue> = std::make_unique<Queue>(capacity);
for (std::uint64_t i = 0; i < capacity / 2; ++i) {
if (!shared_queue<Queue>->try_push(i)) {
std::abort(); // unreachable: fresh queue, stays below capacity
}
}
}

template <typename Queue>
void setup_queue(const benchmark::State& /*state*/) {
install_half_full_queue<Queue>(kCapacity);
}

// Same, but takes the capacity from the registered Arg so one benchmark can
// sweep it.
template <typename Queue>
void setup_queue_at_capacity(const benchmark::State& state) {
install_half_full_queue<Queue>(static_cast<std::size_t>(state.range(0)));
}

template <typename Queue>
void teardown_queue(const benchmark::State& /*state*/) {
shared_queue<Queue>.reset();
Expand Down Expand Up @@ -119,6 +142,49 @@ void BM_QueueThroughput(benchmark::State& state) {
state.SetItemsProcessed(state.iterations() * state.threads());
}

// Non-blocking counterpart of BM_QueueThroughput.
template <typename Queue>
void BM_QueueTryThroughput(benchmark::State& state) {
const bool is_producer = state.thread_index() < state.threads() / 2;
auto& queue = *shared_queue<Queue>;
std::size_t retries = 0;

if (is_producer) {
std::uint64_t item = 0;
for (auto _ : state) {
retries += cq::bench::count_failures_until_success([&] { return queue.try_push(item); });
++item;
}
} else {
std::uint64_t value = 0;
for (auto _ : state) {
retries += cq::bench::count_failures_until_success([&] { return queue.try_pop(value); });
}
}

state.SetItemsProcessed(state.iterations() * state.threads());
// Counters sum across threads, and kAvgIterations divides by the iteration
// total over *all* threads. Producers and consumers each own half of that
// total, so doubling one side's retries before the divide yields that side's
// retries per its own op — for any even producer/consumer split, verified
// bit-exact at 2 and at 8 threads. An odd split would misreport: the scale
// would have to be threads/producers on the push side and threads/consumers
// on the pop side, not a constant 2. Whether it misreports or simply hangs
// is one condition: an odd split completes only while
// (consumers - producers) * iterations-per-thread <= capacity / 2, the
// prefill being the only slack. Under the registered MinTime that product
// reaches millions, so every sweep point livelocks; an iteration-capped run
// completes only at the capacities whose prefill covers it, and there the
// counters would be silently wrong. Only even splits are registered.
using benchmark::Counter;
const auto side_retries = 2.0 * static_cast<double>(retries);
state.counters["push_retries/push"] =
Counter(is_producer ? side_retries : 0.0, Counter::kAvgIterations);
state.counters["pop_retries/pop"] =
Counter(is_producer ? 0.0 : side_retries, Counter::kAvgIterations);
state.counters["retries/op"] = Counter(static_cast<double>(retries), Counter::kAvgIterations);
}

// Uncontended single-thread round trip: the queue's raw per-op cost with no
// other thread in the picture.
template <typename Queue>
Expand Down Expand Up @@ -177,6 +243,19 @@ static_assert(kSpscThreads % 2 == 0 && kMpmcThreads % 2 == 0,
// iteration form overrides, which is what the ctest smoke run uses (1x).
constexpr double kMinTimeSeconds = 1.0;

// Capacity sweep points for the non-blocking runs: 1 and 2 force every op to
// wait on its counterpart, 8 sits at the knee, and 64 / 1024 are deep enough
// for the two sides to decouple — 1024 being the capacity the blocking rows
// use, so those rows stay comparable. Powers of two keep MpmcQueue on its mask
// fast path — except capacity 1, where the mask degenerates to 0 and it falls
// back to % anyway. Non-powers of two are supported and tested at 3; these
// values are a choice, not a requirement.
//
// const, not constexpr: ArgsProduct takes vector<vector<int64_t>>. The const
// is load-bearing beyond style — clang-tidy exempts literals in a const
// initializer, which is what retires the magic-numbers NOLINT here.
const std::vector<std::int64_t> capacity_sweep{1, 2, 8, 64, 1024};

BENCHMARK(BM_QueueThroughput<MutexQueue>)
->Setup(setup_queue<MutexQueue>)
->Teardown(teardown_queue<MutexQueue>)
Expand Down Expand Up @@ -229,23 +308,72 @@ BENCHMARK(BM_QueueThroughput<MoodycamelQueue>)
->MinTime(kMinTimeSeconds)
->Name("MoodycamelQueue/throughput");

// Non-blocking sweeps. Only the cq queues are swept: neither external adapter
// has a try_pop, and moodycamel is unbounded besides, so it has no capacity to
// vary.
//
// Any ratio read across these rows is a statement about the retry policy as
// much as about the queues: at capacity 1, Spsc/Mutex is under 2x with the
// yield and around 40x without it. See try_operation.hpp. The shallow points
// are scheduler-sensitive by construction, so they need
// --benchmark_repetitions more than the rows above do, not less.
BENCHMARK(BM_QueueTryThroughput<MutexQueue>)
->Setup(setup_queue_at_capacity<MutexQueue>)
->Teardown(teardown_queue<MutexQueue>)
->ArgName("capacity")
->ArgsProduct({capacity_sweep})
->Threads(kSpscThreads)
->UseRealTime()
->MinTime(kMinTimeSeconds)
->Name("MutexQueue/try_throughput");

BENCHMARK(BM_QueueTryThroughput<SpscQueue>)
->Setup(setup_queue_at_capacity<SpscQueue>)
->Teardown(teardown_queue<SpscQueue>)
->ArgName("capacity")
->ArgsProduct({capacity_sweep})
->Threads(kSpscThreads)
->UseRealTime()
->MinTime(kMinTimeSeconds)
->Name("SpscQueue/try_throughput");

BENCHMARK(BM_QueueTryThroughput<MpmcQueue>)
->Setup(setup_queue_at_capacity<MpmcQueue>)
->Teardown(teardown_queue<MpmcQueue>)
->ArgName("capacity")
->ArgsProduct({capacity_sweep})
->Threads(kSpscThreads)
->UseRealTime()
->MinTime(kMinTimeSeconds)
->Name("MpmcQueue/try_throughput");

// UseRealTime on the round trips too: Google Benchmark divides a rate counter
// by whichever clock the benchmark selected, so without it these rows report
// items/s per CPU-second while every threaded row above is per wall-second --
// not the same number, and not comparable, which is what the comment on
// SetItemsProcessed above claims they are.
BENCHMARK(BM_QueuePushPopSingleThread<MutexQueue>)
->UseRealTime()
->MinTime(kMinTimeSeconds)
->Name("MutexQueue/single_thread_roundtrip");

BENCHMARK(BM_QueuePushPopSingleThread<SpscQueue>)
->UseRealTime()
->MinTime(kMinTimeSeconds)
->Name("SpscQueue/single_thread_roundtrip");

BENCHMARK(BM_QueuePushPopSingleThread<MpmcQueue>)
->UseRealTime()
->MinTime(kMinTimeSeconds)
->Name("MpmcQueue/single_thread_roundtrip");

BENCHMARK(BM_QueuePushPopSingleThread<TbbBoundedQueue>)
->UseRealTime()
->MinTime(kMinTimeSeconds)
->Name("TbbBoundedQueue/single_thread_roundtrip");

BENCHMARK(BM_QueuePushPopSingleThread<MoodycamelQueue>)
->UseRealTime()
->MinTime(kMinTimeSeconds)
->Name("MoodycamelQueue/single_thread_roundtrip");

Expand Down
51 changes: 51 additions & 0 deletions bench/try_operation.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
// Retry helper shared by the non-blocking throughput benchmarks.
#ifndef CQ_BENCH_TRY_OPERATION_HPP_
#define CQ_BENCH_TRY_OPERATION_HPP_

#include <concepts>
#include <cstddef>
#include <thread>
#include <type_traits>

namespace cq::bench {

// Retry until the operation succeeds, returning how many attempts failed.
//
// The yield is a deliberate policy choice, and it is not neutral — it is worth
// naming because it moves the numbers more than the queue choice does. At
// capacity 1, where every op waits on its counterpart, removing it costs
// MutexQueue roughly a factor of ten and gains each lock-free queue a factor
// of two to three. Ratios only: the absolute rates move 10-25% with machine
// load, so quoting them would bake one afternoon's conditions into the source.
//
// The mechanism only exists for the mutex queue: try_push and try_pop take a
// blocking lock_guard, so an unyielding spinner keeps barging the lock back
// from the counterpart that would have made room. The lock-free queues have no
// lock to barge, so there the yield is pure overhead.
//
// It stays anyway: without it the mutex row degenerates to a couple of hundred
// retries per op at capacity 1, and a CI runner with two vCPUs would fare far
// worse than this 12-core box. But any ratio taken from these rows is a
// statement about this policy as much as about the queues — see the sweep's
// comment in queue_bench.cpp.
//
// By value, per the std:: algorithm convention: the operation is invoked
// repeatedly, so a forwarding reference would never actually be forwarded.
//
// Not std::predicate: that subsumes std::regular_invocable, whose contract is
// equality-preserving invocation, and this helper exists precisely to call
// something whose answer changes between calls.
template <typename Operation>
requires std::invocable<Operation&> && std::convertible_to<std::invoke_result_t<Operation&>, bool>
[[nodiscard]] std::size_t count_failures_until_success(Operation operation) {
std::size_t failures = 0;
while (!operation()) {
++failures;
std::this_thread::yield();
}
return failures;
}

} // namespace cq::bench

#endif // CQ_BENCH_TRY_OPERATION_HPP_
5 changes: 4 additions & 1 deletion tests/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,11 @@ add_executable(queue_tests
mutex_queue_test.cpp
queue_contract_test.cpp
spsc_queue_test.cpp
stress_test.cpp)
stress_test.cpp
try_operation_test.cpp)
target_link_libraries(queue_tests PRIVATE cq::cq cq_warnings cq_sanitizers GTest::gtest_main)
# The retry helper lives with the benchmarks that use it; the test reaches it.
target_include_directories(queue_tests PRIVATE ${PROJECT_SOURCE_DIR}/bench)

include(GoogleTest)
# PRE_TEST + a generous timeout: never run the (possibly TSan-instrumented)
Expand Down
23 changes: 23 additions & 0 deletions tests/try_operation_test.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
#include "try_operation.hpp"

#include <cstddef>

#include <gtest/gtest.h>

namespace cq::bench {
namespace {

TEST(TryOperation, CountsFailuresBeforeSuccess) {
std::size_t attempts = 0;

const auto failures = count_failures_until_success([&] {
++attempts;
return attempts == 3;
});

EXPECT_EQ(attempts, 3U);
EXPECT_EQ(failures, 2U);
}

} // namespace
} // namespace cq::bench
Loading