Aleph-w 3.0
A C++ Library for Data Structures and Algorithms
Loading...
Searching...
No Matches
mpsc_queue_test.cc
Go to the documentation of this file.
1
2/*
3 Aleph_w
4
5 Data structures & Algorithms
6 https://github.com/lrleon/Aleph-w
7
8 This file is part of Aleph-w library
9
10 Copyright (c) 2002-2026 Leandro Rabindranath Leon
11
12 Permission is hereby granted, free of charge, to any person obtaining a copy
13 of this software and associated documentation files (the "Software"), to deal
14 in the Software without restriction, including without limitation the rights
15 to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
16 copies of the Software, and to permit persons to whom the Software is
17 furnished to do so, subject to the following conditions:
18
19 The above copyright notice and this permission notice shall be included in all
20 copies or substantial portions of the Software.
21
22 THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
23 IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
24 FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
25 AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
26 LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
27 OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
28 SOFTWARE.
29*/
30
36#include <gtest/gtest.h>
37
39
40#include <tpl_mpsc_queue.H>
41
42#include <algorithm>
43#include <atomic>
44#include <chrono>
45#include <memory>
46#include <stdexcept>
47#include <string>
48#include <vector>
49
50using namespace Aleph;
51using namespace Aleph::Testing;
52
53namespace
54{
55 // Thread-safe lifetime counter (concurrent constructions/destructions
56 // happen in the multi-producer tests below, unlike single-threaded
57 // lifetime probes elsewhere in the suite).
58 struct Probe
59 {
60 static std::atomic<int> live;
61 int value = 0;
62
63 Probe() noexcept { live.fetch_add(1, std::memory_order_relaxed); }
64 explicit Probe(int v) noexcept : value(v)
65 {
66 live.fetch_add(1, std::memory_order_relaxed);
67 }
68 Probe(const Probe &p) noexcept : value(p.value)
69 {
70 live.fetch_add(1, std::memory_order_relaxed);
71 }
72 Probe(Probe &&p) noexcept : value(p.value)
73 {
74 live.fetch_add(1, std::memory_order_relaxed);
75 }
76 Probe & operator=(const Probe &) = default;
77 Probe & operator=(Probe &&) noexcept = default;
78 ~Probe() { live.fetch_sub(1, std::memory_order_relaxed); }
79 };
80
81 std::atomic<int> Probe::live{0};
82
83 struct Throwing_Ctor
84 {
85 static bool should_throw;
86 int value;
87
88 explicit Throwing_Ctor(int v) : value(v)
89 {
90 if (should_throw)
91 throw std::runtime_error("Throwing_Ctor: constructor failed");
92 }
93 };
94
95 bool Throwing_Ctor::should_throw = false;
96} // namespace
97
103
105{
107 q.push(1);
108 q.push(2);
109 q.emplace(3);
111
112 int out = 0;
114 EXPECT_EQ(out, 1);
116 EXPECT_EQ(out, 2);
118 EXPECT_EQ(out, 3);
119
122}
123
125{
127 q.push(std::string("alpha"));
128 q.push(std::string("beta"));
129
130 auto a = q.try_pop();
131 ASSERT_TRUE(a.has_value());
132 EXPECT_EQ(*a, "alpha");
133
134 auto b = q.try_pop();
135 ASSERT_TRUE(b.has_value());
136 EXPECT_EQ(*b, "beta");
137
138 auto c = q.try_pop();
139 EXPECT_FALSE(c.has_value());
140}
141
143{
145 constexpr int N = 1000;
146 for (int i = 0; i < N; ++i)
147 q.push(i);
148
149 for (int i = 0; i < N; ++i)
150 {
151 int out = -1;
153 EXPECT_EQ(out, i);
154 }
156}
157
159{
161 q.push(std::make_unique<int>(10));
162 q.emplace(std::make_unique<int>(20));
163
164 std::unique_ptr<int> out;
166 ASSERT_NE(out, nullptr);
167 EXPECT_EQ(*out, 10);
168
169 auto opt = q.try_pop();
170 ASSERT_TRUE(opt.has_value());
171 ASSERT_NE(*opt, nullptr);
172 EXPECT_EQ(**opt, 20);
173}
174
176{
177 ASSERT_EQ(Probe::live.load(), 0);
178 Probe out; // outlives the inner scope below, so it can carry a live
179 // count past the queue's destructor for the final check.
180 {
182 q.emplace(1);
183 q.emplace(2);
184 q.emplace(3);
185 EXPECT_EQ(Probe::live.load(), 4); // 3 queued + out
186
188 EXPECT_EQ(out.value, 1);
189 // The popped-into `out` and the two still-queued elements are alive.
190 EXPECT_EQ(Probe::live.load(), 3);
191 }
192 // Queue destructor released the two still-queued nodes; `out` (declared
193 // above, still alive here) keeps the count at 1.
194 EXPECT_EQ(Probe::live.load(), 1);
195}
196
198{
199 ASSERT_EQ(Probe::live.load(), 0);
200 {
202 for (int i = 0; i < 50; ++i)
203 q.emplace(i);
204 EXPECT_EQ(Probe::live.load(), 50);
205 }
206 EXPECT_EQ(Probe::live.load(), 0);
207}
208
210{
212 q.emplace(1);
214
215 Throwing_Ctor::should_throw = true;
216 EXPECT_THROW(q.emplace(2), std::runtime_error);
217 Throwing_Ctor::should_throw = false;
218
219 // The successful push before the throwing one must still be there,
220 // and nothing from the failed attempt should have been enqueued.
221 Throwing_Ctor out(0);
223 EXPECT_EQ(out.value, 1);
225}
226
228{
230
232 config.producers = 1;
233 config.consumers = 1;
234 config.items_per_producer = 20000;
235 config.timeout = std::chrono::seconds(30);
236
237 const auto result = run_producer_consumer_stress(
238 config,
239 [&](size_t value) { q.push(value); },
240 [&](size_t &out) { return q.try_pop(out); });
241
242 ASSERT_FALSE(result.timed_out);
243 ASSERT_EQ(result.size(), config.items_per_producer);
244 // Single producer, single consumer: consumption order (not just the set
245 // of values) must exactly match push order.
246 for (size_t i = 0; i < result.consumed.size(); ++i)
247 EXPECT_EQ(result.consumed[i], i);
248
250}
251
253{
254 constexpr size_t producer_count = 6;
255 constexpr size_t items_per_producer = 20000;
256
258
260 config.producers = producer_count;
261 config.consumers = 1; // MpscQueue supports exactly one consumer.
262 config.items_per_producer = items_per_producer;
263 config.timeout = std::chrono::seconds(30);
264
265 // The framework encodes each pushed value as
266 // `producer_index * items_per_producer + local_sequence`, so both the
267 // exactly-once conservation check (via sorting) and the FIFO-per-producer
268 // check (via the unsorted, consumption-ordered `result.consumed`) can be
269 // derived directly from the returned values.
270 const auto result = run_producer_consumer_stress(
271 config,
272 [&](size_t value) { q.push(value); },
273 [&](size_t &out) { return q.try_pop(out); });
274
275 ASSERT_FALSE(result.timed_out);
276 const size_t total = producer_count * items_per_producer;
277 ASSERT_EQ(result.size(), total);
278
279 std::vector<long long> last_seq(producer_count, -1);
280 for (const size_t value : result.consumed)
281 {
282 const size_t producer = value / items_per_producer;
283 const size_t seq = value % items_per_producer;
284 ASSERT_LT(producer, producer_count);
285 ASSERT_GT(static_cast<long long>(seq), last_seq[producer])
286 << "producer " << producer << " delivered out of FIFO order";
287 last_seq[producer] = static_cast<long long>(seq);
288 }
289
290 auto sorted = result.consumed;
291 std::sort(sorted.begin(), sorted.end());
292 for (size_t i = 0; i < total; ++i)
293 ASSERT_EQ(sorted[i], i) << "index " << i << " missing or duplicated";
294
296}
297
299{
300 // Single-producer trace replayed against std::vector-as-model semantics:
301 // every insert corresponds to a push, and pops happen in strict FIFO
302 // order matching the model's own front-to-back consumption.
304 5000, 0xC0FFEE, 1000, {Trace_Operation_Kind::insert});
305
306 std::vector<size_t> model;
308
309 for (const auto &op : trace)
310 {
311 model.push_back(op.key);
312 subject.push(op.key);
313 }
314
315 for (size_t expected : model)
316 {
317 size_t out = 0;
318 ASSERT_TRUE(subject.try_pop(out));
320 }
321 EXPECT_TRUE(subject.is_empty());
322}
size_t size_t int32_t value
Definition ca-c-api.h:116
size_t size_t int32_t * out
Definition ca-c-api.h:120
Unbounded lock-free multi-producer/single-consumer queue.
void emplace(Args &&... args)
Construct a new element in place at the back of the queue.
void push(const T &value)
Push a copy of value onto the queue.
bool is_empty() const noexcept
Advisory check for whether the queue currently has no elements.
bool try_pop(T &out)
Attempt to pop the front element into out.
Minimal std::expected-style result type for C++20.
Reusable helpers for concurrent data-structure tests.
#define TEST(name)
#define N
Definition fib.C:294
size_t blossom_maximum_cardinality_matching(const GT &g, DynDlist< typename GT::Arc * > &matching, SA sa=SA())
Alias of compute_maximum_cardinality_general_matching().
Definition Blossom.H:466
Producer_Consumer_Stress_Result run_producer_consumer_stress(const Producer_Consumer_Stress_Config &config, Push push, TryPop try_pop)
Run a deterministic producer/consumer stress scenario.
std::vector< Trace_Operation > make_random_operation_trace(const size_t count, const uint32_t seed, const size_t key_range, const std::initializer_list< Trace_Operation_Kind > kinds={ Trace_Operation_Kind::insert, Trace_Operation_Kind::erase, Trace_Operation_Kind::contains })
Build a deterministic pseudo-random operation trace.
Main namespace for Aleph-w library functions.
Definition ah-arena.H:89
Configuration for producer/consumer stress helpers.
Unbounded lock-free multi-producer/single-consumer queue (Aleph::MpscQueue).