38#ifndef ALEPH_TESTS_CONCURRENCY_TEST_UTILS_H
39#define ALEPH_TESTS_CONCURRENCY_TEST_UTILS_H
43#include <condition_variable>
47#include <initializer_list>
50#include <system_error>
78 std::condition_variable
cv_;
93 <<
"Deterministic_Start_Gate requires at least one participant";
105 std::unique_lock lock(
mutex_);
119 std::unique_lock lock(
mutex_);
130 std::lock_guard lock(
mutex_);
148 std::lock_guard lock(
mutex_);
159 std::lock_guard lock(
mutex_);
207 if (thread.joinable())
213 catch (
const std::system_error &)
219 catch (
const std::system_error &)
254 template <
typename Worker>
258 std::vector<std::thread> threads;
267 threads.emplace_back([&, i,
worker]()
mutable
269 if (not gate.arrive_and_wait())
285 gate.wait_until_ready();
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);
306 bool timed_out =
false;
312 [[nodiscard]]
size_t size() const noexcept {
return consumed.size(); }
317 template <
typename Push>
320 if constexpr (std::is_void_v<std::invoke_result_t<Push &, size_t>>)
359 template <
typename Push,
typename TryPop>
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;
378 [&,
push, try_pop](
const size_t index)
mutable
380 if (index < config.producers)
382 const size_t p = index;
383 for (size_t i = 0; i < config.items_per_producer; ++i)
385 const size_t value = p * config.items_per_producer + i;
386 while (not detail::invoke_push(push, value))
388 if (std::chrono::steady_clock::now() >= deadline)
390 timed_out.store(true, std::memory_order_relaxed);
393 std::this_thread::yield();
399 while (consumed_count.load(std::memory_order_acquire) < total)
405 consumed_count.fetch_add(1, std::memory_order_acq_rel);
407 result.consumed[slot] = value;
411 if (std::chrono::steady_clock::now() >= deadline)
413 timed_out.store(true, std::memory_order_relaxed);
416 std::this_thread::yield();
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);
453 return kind == rhs.
kind and key == rhs.key and
value == rhs.value;
471 const size_t key_range,
472 const std::initializer_list<Trace_Operation_Kind> kinds =
474 Trace_Operation_Kind::insert,
475 Trace_Operation_Kind::erase,
476 Trace_Operation_Kind::contains
479 std::vector<Trace_Operation_Kind> allowed(kinds.begin(), kinds.end());
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;
488 std::vector<Trace_Operation> trace;
489 trace.reserve(count);
490 for (
size_t i = 0; i <
count; ++i)
493 allowed[kind_dist(rng)],
510 template <
typename SubjectStep,
typename ReferenceStep>
512 const std::vector<Trace_Operation> &trace,
513 SubjectStep subject_step,
514 ReferenceStep reference_step)
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);
Exception handling system with formatted messages for Aleph-w.
#define ah_invalid_argument_if(C)
Throws std::invalid_argument if condition holds.
bool operator==(const Time &l, const Time &r)
size_t size_t int32_t value
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().
std::condition_variable cv_
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
std::vector< std::thread > & threads_
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
Deterministic_Start_Gate & gate_
Minimal std::expected-style result type for C++20.
size_t blossom_maximum_cardinality_matching(const GT &g, DynDlist< typename GT::Arc * > &matching, SA sa=SA())
Alias of compute_maximum_cardinality_general_matching().
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.
Configuration for producer/consumer stress helpers.
size_t items_per_producer
std::chrono::milliseconds timeout
Result returned by producer/consumer stress helpers.
std::vector< size_t > consumed
size_t size() const noexcept
Return the expected number of consumed values.
One randomized operation trace entry.
Trace_Operation_Kind kind