Aleph-w 3.0
A C++ Library for Data Structures and Algorithms
Loading...
Searching...
No Matches
concurrency_test_utils.H
Go to the documentation of this file.
1/*
2 Aleph_w
3
4 Data structures & Algorithms
5 version 2.0.0b
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
38#ifndef ALEPH_TESTS_CONCURRENCY_TEST_UTILS_H
39#define ALEPH_TESTS_CONCURRENCY_TEST_UTILS_H
40
41#include <atomic>
42#include <chrono>
43#include <condition_variable>
44#include <cstddef>
45#include <cstdint>
46#include <exception>
47#include <initializer_list>
48#include <mutex>
49#include <random>
50#include <system_error>
51#include <thread>
52#include <type_traits>
53#include <utility>
54#include <vector>
55
56#include <ah-errors.H>
57
59{
72 {
73 const size_t expected_ = 0;
74 size_t arrived_ = 0;
75 bool released_ = false;
76 bool cancelled_ = false;
77 mutable std::mutex mutex_;
78 std::condition_variable cv_;
79
80 public:
89 explicit Deterministic_Start_Gate(const size_t expected)
91 {
93 << "Deterministic_Start_Gate requires at least one participant";
94 }
95
104 {
105 std::unique_lock lock(mutex_);
106 ++arrived_;
107 cv_.notify_all();
108 cv_.wait(lock, [this] { return released_ or cancelled_; });
109 return released_;
110 }
111
118 {
119 std::unique_lock lock(mutex_);
120 cv_.wait(lock, [this] { return arrived_ == expected_ or cancelled_; });
121 }
122
128 void release()
129 {
130 std::lock_guard lock(mutex_);
131 released_ = true;
132 cv_.notify_all();
133 }
134
147 {
148 std::lock_guard lock(mutex_);
149 cancelled_ = true;
150 cv_.notify_all();
151 }
152
157 [[nodiscard]] size_t arrived() const
158 {
159 std::lock_guard lock(mutex_);
160 return arrived_;
161 }
162 };
163
164 namespace detail
165 {
183 {
184 std::vector<std::thread> &threads_;
186 bool launch_completed_ = false;
187
188 public:
189 Joining_Thread_Guard(std::vector<std::thread> &threads,
191 : threads_(threads), gate_(gate)
192 {}
193
196
201
203 {
205 gate_.cancel();
206 for (auto &thread : threads_)
207 if (thread.joinable())
208 {
209 try
210 {
211 thread.join();
212 }
213 catch (const std::system_error &)
214 {
215 try
216 {
217 thread.detach();
218 }
219 catch (const std::system_error &)
220 {
221 // Nothing actionable left from a destructor.
222 }
223 }
224 }
225 }
226 };
227 } // namespace detail
228
254 template <typename Worker>
256 {
258 std::vector<std::thread> threads;
259 threads.reserve(thread_count);
260 std::mutex exception_mutex;
261 std::exception_ptr first_exception;
262
263 {
265
266 for (size_t i = 0; i < thread_count; ++i)
267 threads.emplace_back([&, i, worker]() mutable
268 {
269 if (not gate.arrive_and_wait())
270 return;
271 try
272 {
273 worker(i);
274 }
275 catch (...)
276 {
277 std::lock_guard lock(exception_mutex);
279 first_exception = std::current_exception();
280 }
281 });
282
283 guard.dismiss(); // every thread launched; let them reach the gate
284 // on their own instead of racing to cancel() it.
285 gate.wait_until_ready();
286 gate.release();
287 } // guard destructor joins every thread (see Joining_Thread_Guard).
288
289 if (first_exception)
290 std::rethrow_exception(first_exception);
291 }
292
295 {
296 size_t producers = 1;
297 size_t consumers = 1;
298 size_t items_per_producer = 1024;
299 std::chrono::milliseconds timeout = std::chrono::seconds(10);
300 };
301
304 {
305 std::vector<size_t> consumed;
306 bool timed_out = false;
307
312 [[nodiscard]] size_t size() const noexcept { return consumed.size(); }
313 };
314
315 namespace detail
316 {
317 template <typename Push>
318 bool invoke_push(Push &push, const size_t value)
319 {
320 if constexpr (std::is_void_v<std::invoke_result_t<Push &, size_t>>)
321 {
322 push(value);
323 return true;
324 }
325 else
326 {
327 return static_cast<bool>(push(value));
328 }
329 }
330 }
331
359 template <typename Push, typename TryPop>
362 Push push,
363 TryPop try_pop)
364 {
365 const size_t total = config.producers * config.items_per_producer;
367 result.consumed.resize(total);
368
369 std::atomic<size_t> consumed_count{0};
370 std::atomic<bool> timed_out{false};
371 const auto deadline = std::chrono::steady_clock::now() + config.timeout;
372
373 // Workers `[0, producers)` produce; `[producers, producers+consumers)`
374 // consume. `push`/`try_pop` are captured by value here once, then
375 // copied again per-thread by run_workers() itself (each thread's copy
376 // of this worker closure carries its own copy of both callables).
377 run_workers(config.producers + config.consumers,
378 [&, push, try_pop](const size_t index) mutable
379 {
380 if (index < config.producers)
381 {
382 const size_t p = index;
383 for (size_t i = 0; i < config.items_per_producer; ++i)
384 {
385 const size_t value = p * config.items_per_producer + i;
386 while (not detail::invoke_push(push, value))
387 {
388 if (std::chrono::steady_clock::now() >= deadline)
389 {
390 timed_out.store(true, std::memory_order_relaxed);
391 return;
392 }
393 std::this_thread::yield();
394 }
395 }
396 }
397 else
398 {
399 while (consumed_count.load(std::memory_order_acquire) < total)
400 {
401 size_t value = 0;
402 if (try_pop(value))
403 {
404 const size_t slot =
405 consumed_count.fetch_add(1, std::memory_order_acq_rel);
406 if (slot < total)
407 result.consumed[slot] = value;
408 }
409 else
410 {
411 if (std::chrono::steady_clock::now() >= deadline)
412 {
413 timed_out.store(true, std::memory_order_relaxed);
414 return;
415 }
416 std::this_thread::yield();
417 }
418 }
419 }
420 });
421
422 const size_t n = consumed_count.load(std::memory_order_acquire);
423 if (n < result.consumed.size())
424 result.consumed.resize(n);
425 result.timed_out = timed_out.load(std::memory_order_relaxed);
426 return result;
427 }
428
431 {
432 push,
433 pop,
434 insert,
435 erase,
437 };
438
441 {
442 Trace_Operation_Kind kind = Trace_Operation_Kind::contains;
443 size_t key = 0;
444 size_t value = 0;
445
451 [[nodiscard]] bool operator == (const Trace_Operation &rhs) const noexcept
452 {
453 return kind == rhs.kind and key == rhs.key and value == rhs.value;
454 }
455 };
456
468 inline std::vector<Trace_Operation> make_random_operation_trace(
469 const size_t count,
470 const uint32_t seed,
471 const size_t key_range,
472 const std::initializer_list<Trace_Operation_Kind> kinds =
473 {
474 Trace_Operation_Kind::insert,
475 Trace_Operation_Kind::erase,
476 Trace_Operation_Kind::contains
477 })
478 {
479 std::vector<Trace_Operation_Kind> allowed(kinds.begin(), kinds.end());
480 ah_invalid_argument_if(key_range == 0) << "key_range must be > 0";
481 ah_invalid_argument_if(allowed.empty()) << "kinds must not be empty";
482
483 std::mt19937 rng(seed);
484 std::uniform_int_distribution<size_t> kind_dist(0, allowed.size() - 1);
485 std::uniform_int_distribution<size_t> key_dist(0, key_range - 1);
486 std::uniform_int_distribution<size_t> value_dist;
487
488 std::vector<Trace_Operation> trace;
489 trace.reserve(count);
490 for (size_t i = 0; i < count; ++i)
491 trace.push_back(
492 {
493 allowed[kind_dist(rng)],
494 key_dist(rng),
495 value_dist(rng)
496 });
497
498 return trace;
499 }
500
510 template <typename SubjectStep, typename ReferenceStep>
512 const std::vector<Trace_Operation> &trace,
513 SubjectStep subject_step,
514 ReferenceStep reference_step)
515 {
516 std::vector<size_t> mismatches;
517 for (size_t i = 0; i < trace.size(); ++i)
518 if (not (subject_step(trace[i]) == reference_step(trace[i])))
519 mismatches.push_back(i);
520 return mismatches;
521 }
522} // namespace Aleph::Testing
523
524#endif // ALEPH_TESTS_CONCURRENCY_TEST_UTILS_H
Exception handling system with formatted messages for Aleph-w.
#define ah_invalid_argument_if(C)
Throws std::invalid_argument if condition holds.
Definition ah-errors.H:644
bool operator==(const Time &l, const Time &r)
Definition ah-time.H:133
size_t size_t int32_t value
Definition ca-c-api.h:116
Single-use deterministic start gate for concurrent tests.
void wait_until_ready()
Wait until all expected workers are ready, or the gate is cancelled.
void cancel() noexcept
Abort the gate, waking every worker currently blocked in arrive_and_wait() or wait_until_ready().
size_t arrived() const
Return the number of arrived workers.
bool arrive_and_wait()
Mark the caller ready and block until release() or cancel().
Deterministic_Start_Gate(const size_t expected)
Construct a gate for the given number of workers.
void release()
Release all waiting workers.
RAII guard that leaves no thread joinable when it goes out of scope, cancelling gate first so any thr...
Joining_Thread_Guard & operator=(const Joining_Thread_Guard &)=delete
Joining_Thread_Guard(const Joining_Thread_Guard &)=delete
void dismiss() noexcept
Mark the launch loop as having completed without throwing.
Joining_Thread_Guard(std::vector< std::thread > &threads, Deterministic_Start_Gate &gate) noexcept
Minimal std::expected-style result type for C++20.
static mt19937 rng
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
bool invoke_push(Push &push, const size_t value)
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.
void run_workers(const size_t thread_count, Worker worker)
Launch thread_count workers with a deterministic simultaneous start, join every one of them,...
Trace_Operation_Kind
Operation kind for randomized reference traces.
std::vector< size_t > replay_trace_and_collect_mismatches(const std::vector< Trace_Operation > &trace, SubjectStep subject_step, ReferenceStep reference_step)
Replay a trace against a subject and a reference implementation.
bool contains(const std::string_view &str, const std::string_view &substr)
Check if substr appears inside str.
Itor::difference_type count(const Itor &beg, const Itor &end, const T &value)
Count elements equal to a value.
Definition ahAlgo.H:127
Configuration for producer/consumer stress helpers.
Result returned by producer/consumer stress helpers.
size_t size() const noexcept
Return the expected number of consumed values.
One randomized operation trace entry.
ValueArg< size_t > seed
Definition testHash.C:53