diff --git a/.clang-tidy b/.clang-tidy index 28b0d3a..5f33529 100644 --- a/.clang-tidy +++ b/.clang-tidy @@ -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. diff --git a/CMakeLists.txt b/CMakeLists.txt index 9be0d09..fe1c200 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -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) diff --git a/bench/CMakeLists.txt b/bench/CMakeLists.txt index 27cff11..4a69f5c 100644 --- a/bench/CMakeLists.txt +++ b/bench/CMakeLists.txt @@ -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]") diff --git a/bench/queue_bench.cpp b/bench/queue_bench.cpp index 5940655..fe6d88e 100644 --- a/bench/queue_bench.cpp +++ b/bench/queue_bench.cpp @@ -11,11 +11,21 @@ // 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 #include #include +#include "try_operation.hpp" + #include #include #include @@ -23,6 +33,7 @@ #include #include #include +#include #include #include @@ -78,20 +89,32 @@ class MoodycamelQueue { template std::unique_ptr 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 -void setup_queue(const benchmark::State& /*state*/) { - shared_queue = std::make_unique(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->try_push(i); - } - if (!prefilled) { - std::abort(); // unreachable: fresh queue, stays below capacity +void install_half_full_queue(std::size_t capacity) { + shared_queue = std::make_unique(capacity); + for (std::uint64_t i = 0; i < capacity / 2; ++i) { + if (!shared_queue->try_push(i)) { + std::abort(); // unreachable: fresh queue, stays below capacity + } } } +template +void setup_queue(const benchmark::State& /*state*/) { + install_half_full_queue(kCapacity); +} + +// Same, but takes the capacity from the registered Arg so one benchmark can +// sweep it. +template +void setup_queue_at_capacity(const benchmark::State& state) { + install_half_full_queue(static_cast(state.range(0))); +} + template void teardown_queue(const benchmark::State& /*state*/) { shared_queue.reset(); @@ -119,6 +142,49 @@ void BM_QueueThroughput(benchmark::State& state) { state.SetItemsProcessed(state.iterations() * state.threads()); } +// Non-blocking counterpart of BM_QueueThroughput. +template +void BM_QueueTryThroughput(benchmark::State& state) { + const bool is_producer = state.thread_index() < state.threads() / 2; + auto& queue = *shared_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(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(retries), Counter::kAvgIterations); +} + // Uncontended single-thread round trip: the queue's raw per-op cost with no // other thread in the picture. template @@ -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>. 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 capacity_sweep{1, 2, 8, 64, 1024}; + BENCHMARK(BM_QueueThroughput) ->Setup(setup_queue) ->Teardown(teardown_queue) @@ -229,23 +308,72 @@ BENCHMARK(BM_QueueThroughput) ->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) + ->Setup(setup_queue_at_capacity) + ->Teardown(teardown_queue) + ->ArgName("capacity") + ->ArgsProduct({capacity_sweep}) + ->Threads(kSpscThreads) + ->UseRealTime() + ->MinTime(kMinTimeSeconds) + ->Name("MutexQueue/try_throughput"); + +BENCHMARK(BM_QueueTryThroughput) + ->Setup(setup_queue_at_capacity) + ->Teardown(teardown_queue) + ->ArgName("capacity") + ->ArgsProduct({capacity_sweep}) + ->Threads(kSpscThreads) + ->UseRealTime() + ->MinTime(kMinTimeSeconds) + ->Name("SpscQueue/try_throughput"); + +BENCHMARK(BM_QueueTryThroughput) + ->Setup(setup_queue_at_capacity) + ->Teardown(teardown_queue) + ->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) + ->UseRealTime() ->MinTime(kMinTimeSeconds) ->Name("MutexQueue/single_thread_roundtrip"); BENCHMARK(BM_QueuePushPopSingleThread) + ->UseRealTime() ->MinTime(kMinTimeSeconds) ->Name("SpscQueue/single_thread_roundtrip"); BENCHMARK(BM_QueuePushPopSingleThread) + ->UseRealTime() ->MinTime(kMinTimeSeconds) ->Name("MpmcQueue/single_thread_roundtrip"); BENCHMARK(BM_QueuePushPopSingleThread) + ->UseRealTime() ->MinTime(kMinTimeSeconds) ->Name("TbbBoundedQueue/single_thread_roundtrip"); BENCHMARK(BM_QueuePushPopSingleThread) + ->UseRealTime() ->MinTime(kMinTimeSeconds) ->Name("MoodycamelQueue/single_thread_roundtrip"); diff --git a/bench/try_operation.hpp b/bench/try_operation.hpp new file mode 100644 index 0000000..c962357 --- /dev/null +++ b/bench/try_operation.hpp @@ -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 +#include +#include +#include + +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 + requires std::invocable && std::convertible_to, 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_ diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 560eb4c..c0a4137 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -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) diff --git a/tests/try_operation_test.cpp b/tests/try_operation_test.cpp new file mode 100644 index 0000000..b795a48 --- /dev/null +++ b/tests/try_operation_test.cpp @@ -0,0 +1,23 @@ +#include "try_operation.hpp" + +#include + +#include + +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