GitHub Viewer
#define DOCTEST_CONFIG_IMPLEMENT_WITH_MAIN
#include
#include
#include
#include
#include
#include
#include
#include
// ============================================================================
// Shared helpers
// ============================================================================
static inline void tiny_jitter(std::mt19937& rng) {
std::uniform_int_distribution pick(0, 9);
int x = pick(rng);
if (x < 4) {
std::this_thread::yield();
} else if (x < 8) {
volatile int sink = 0;
int spins = 20 + (pick(rng) * 30);
for (int i = 0; i < spins; ++i) sink += i;
(void)sink;
} else {
std::this_thread::sleep_for(std::chrono::nanoseconds(200));
}
}
enum class NotificationType { ONE, N, ALL };
// ============================================================================
// no_missing_notify_all
// Each round: all N threads prepare_wait(), pause in the window between
// prepare and commit until all threads are ready, then commit_wait().
// Main calls notify_all() in that window. All commit_wait() calls must
// return without blocking (no lost wakeup).
// End check: completed == R*N, num_waiters() == 0.
// ============================================================================
template
void no_missing_notify_all(size_t N) {
T notifier(N);
REQUIRE(notifier.size() == N);
size_t R = 20 * (N + 1);
if (N >= 31) R = 1 * (N + 1);
std::atomic prepared(0);
std::atomic completed(0);
std::atomic round(0);
std::atomic stop(false);
std::vector threads;
threads.reserve(N);
for (size_t i = 0; i < N; ++i) {
threads.emplace_back([&, i]() {
size_t local_round = 0;
while (!stop.load(std::memory_order_relaxed)) {
while (round.load(std::memory_order_acquire) == local_round &&
!stop.load(std::memory_order_relaxed)) {
std::this_thread::yield();
}
if (stop.load(std::memory_order_relaxed)) break;
notifier.prepare_wait(i);
prepared.fetch_add(1, std::memory_order_relaxed);
while (prepared.load(std::memory_order_acquire) < (local_round + 1) * N &&
!stop.load(std::memory_order_relaxed)) {
std::this_thread::yield();
}
notifier.commit_wait(i);
completed.fetch_add(1, std::memory_order_relaxed);
local_round++;
}
});
}
for (size_t r = 0; r < R; ++r) {
round.fetch_add(1, std::memory_order_release);
while (prepared.load(std::memory_order_acquire) != (r + 1) * N) {
std::this_thread::yield();
}
notifier.notify_all();
}
stop.store(true, std::memory_order_release);
notifier.notify_all();
for (auto& t : threads) t.join();
REQUIRE(completed.load() == R * N);
REQUIRE(notifier.num_waiters() == 0);
}
// ============================================================================
// no_missing_notify_one
// Each round: all N threads prepare_wait() then commit_wait().
// Main calls notify_one() and verifies at least 1 thread unblocks,
// then notify_all() drains the rest. num_waiters() == 0 at end.
// ============================================================================
template
void no_missing_notify_one(size_t N) {
T notifier(N);
REQUIRE(notifier.size() == N);
size_t R = 20 * (N + 1);
if (N >= 31) R = 1 * (N + 1);
std::atomic round(0);
std::atomic prepared(0);
std::atomic committed(0);
std::atomic stop(false);
std::vector threads;
threads.reserve(N);
for (size_t i = 0; i < N; ++i) {
threads.emplace_back([&, i]() {
size_t local_round = 0;
while (!stop.load(std::memory_order_relaxed)) {
while (round.load(std::memory_order_acquire) = 31) R = 1 * (N + 1);
std::atomic round(0);
std::atomic prepared(0);
std::atomic committed(0);
std::atomic stop(false);
std::vector threads;
threads.reserve(N);
for (size_t i = 0; i < N; ++i) {
threads.emplace_back([&, i]() {
size_t local_round = 0;
while (!stop.load(std::memory_order_relaxed)) {
while (round.load(std::memory_order_acquire) = 31) R = 1 * (N + 1);
std::atomic round(0);
std::atomic prepared(0);
std::atomic committed(0);
std::atomic stop(false);
std::vector threads;
threads.reserve(N);
for (size_t i = 0; i < N; ++i) {
threads.emplace_back([&, i]() {
size_t local_round = 0;
while (!stop.load(std::memory_order_relaxed)) {
while (round.load(std::memory_order_acquire) = 31) R = 1 * (N + 1);
std::atomic round(0);
std::atomic prepared(0);
std::atomic canceled(0);
std::atomic committed(0);
std::atomic stop(false);
std::atomic go(false);
std::vector has_work(N);
for (size_t i = 0; i < N; ++i) has_work[i].store(false, std::memory_order_relaxed);
std::vector threads;
threads.reserve(N);
for (size_t i = 0; i < N; ++i) {
threads.emplace_back([&, i]() {
size_t local_round = 0;
while (!stop.load(std::memory_order_relaxed)) {
while (round.load(std::memory_order_acquire) 0) {
std::uniform_int_distribution pickX(0, expected_commits);
size_t X = pickX(rng);
size_t first = expected_commits - X;
notifier.notify_n(first);
for (size_t j = 0; j < X; ++j) notifier.notify_one();
}
// Ensure the round always completes.
notifier.notify_all();
while ((canceled.load(std::memory_order_acquire) +
committed.load(std::memory_order_acquire)) != N) {
std::this_thread::yield();
}
REQUIRE(notifier.num_waiters() == 0);
}
stop.store(true, std::memory_order_release);
go.store(true, std::memory_order_release);
notifier.notify_all();
for (auto& t : threads) t.join();
REQUIRE(notifier.num_waiters() == 0);
}
// ============================================================================
// no_missing_notifications (compile-time switch: ONE / N / ALL)
// M concurrent notifier threads continuously hammer one notify variant.
// N worker threads follow prepare -> cancel_or_commit protocol.
// Round ends when every worker has either canceled or returned from commit.
// num_waiters() == 0 after every round.
// ============================================================================
template
void no_missing_notifications(size_t N, size_t M = 4, uint32_t seed = 12345) {
T notifier(N);
REQUIRE(notifier.size() == N);
size_t R = 20 * (N + 1);
if (N >= 31) R = 1 * (N + 1);
std::atomic round(0);
std::atomic prepared(0);
std::atomic canceled(0);
std::atomic committed(0);
std::atomic stop(false);
std::atomic go(false);
std::vector has_work(N);
for (size_t i = 0; i < N; ++i) has_work[i].store(false, std::memory_order_relaxed);
std::vector workers;
workers.reserve(N);
for (size_t i = 0; i < N; ++i) {
workers.emplace_back([&, i]() {
size_t local_round = 0;
while (!stop.load(std::memory_order_relaxed)) {
while (round.load(std::memory_order_acquire) 0 (slow path exercised).
// - fast_path > 0 (predicate check before commit works).
// ============================================================================
template
void fuzz_stress_notifier(size_t N, size_t M_notifiers, size_t rounds, uint32_t seed) {
T notifier(N);
REQUIRE(notifier.size() == N);
std::atomic stop{false};
std::atomic signal_epoch{0};
std::atomic prepares{0};
std::atomic cancels{0};
std::atomic commits_entered{0};
std::atomic commits_returned{0};
std::atomic fast_path{0};
std::vector workers;
workers.reserve(N);
for (size_t i = 0; i < N; ++i) {
workers.emplace_back([&, i] {
std::mt19937 rng(seed ^ (0x9e3779b9u + (uint32_t)i * 101u));
uint64_t local = signal_epoch.load(std::memory_order_relaxed);
std::uniform_int_distribution coin(0, 99);
for (size_t it = 0; it < rounds && !stop.load(std::memory_order_relaxed); ++it) {
// Partial participation: sometimes do "work" without waiting.
if (coin(rng) < 15) { tiny_jitter(rng); continue; }
uint64_t cur = signal_epoch.load(std::memory_order_acquire);
if (cur != local) {
local = cur;
fast_path.fetch_add(1, std::memory_order_relaxed);
tiny_jitter(rng);
continue;
}
// Two-phase wait protocol.
tiny_jitter(rng);
notifier.prepare_wait(i);
prepares.fetch_add(1, std::memory_order_relaxed);
tiny_jitter(rng);
cur = signal_epoch.load(std::memory_order_acquire);
if (cur != local) {
notifier.cancel_wait(i);
cancels.fetch_add(1, std::memory_order_relaxed);
local = cur;
tiny_jitter(rng);
continue;
}
commits_entered.fetch_add(1, std::memory_order_relaxed);
notifier.commit_wait(i);
commits_returned.fetch_add(1, std::memory_order_relaxed);
local = signal_epoch.load(std::memory_order_acquire);
tiny_jitter(rng);
}
});
}
std::vector notifiers;
notifiers.reserve(M_notifiers);
for (size_t t = 0; t < M_notifiers; ++t) {
notifiers.emplace_back([&, t] {
std::mt19937 rng(seed + (uint32_t)(777u + t * 17u));
std::uniform_int_distribution which(0, 99);
while (!stop.load(std::memory_order_relaxed)) {
tiny_jitter(rng);
int w = which(rng);
// Occasionally send "empty" notifies (no predicate change) — must be harmless.
if (w < 15) {
int kind = which(rng) % 3;
if (kind == 0) notifier.notify_one();
else if (kind == 1) notifier.notify_n((size_t)(which(rng) % (int)(N + 1)));
else notifier.notify_all();
continue;
}
// Normal: advance predicate then notify (condition-variable style).
signal_epoch.fetch_add(1, std::memory_order_release);
if (w < 60) notifier.notify_one();
else if (w < 85) notifier.notify_n((size_t)(which(rng) % (int)(N + 1)));
else notifier.notify_all();
}
});
}
for (auto& w : workers) w.join();
stop.store(true, std::memory_order_release);
notifier.notify_all();
for (auto& n : notifiers) n.join();
REQUIRE(notifier.num_waiters() == 0);
REQUIRE(commits_returned.load() == commits_entered.load());
REQUIRE(prepares.load() > 0);
REQUIRE(fast_path.load() > 0);
}
// ============================================================================
// notify_n_releases_committed
// Forces ALL N threads to become committed waiters (num_waiters() == N),
// then verifies notify_n(k) alone releases at least min(k, N) of them.
// notify_all() drains the rest so the round completes.
// This is a targeted regression for notify_n boundary semantics and for
// spurious early-exit from commit_wait due to an epoch field bug.
// ============================================================================
template
void notify_n_releases_committed(size_t N, size_t k, size_t rounds, uint32_t seed) {
(void)seed;
T notifier(N);
REQUIRE(notifier.size() == N);
std::atomic round(0);
std::atomic prepared(0);
std::atomic committed_done(0);
std::atomic stop(false);
std::vector workers;
workers.reserve(N);
for (size_t i = 0; i < N; ++i) {
workers.emplace_back([&, i] {
size_t local_round = 0;
while (!stop.load(std::memory_order_relaxed)) {
while (round.load(std::memory_order_acquire) == local_round &&
!stop.load(std::memory_order_relaxed)) {
std::this_thread::yield();
}
if (stop.load(std::memory_order_relaxed)) break;
// No cancel path: we want all threads to become committed waiters.
notifier.prepare_wait(i);
prepared.fetch_add(1, std::memory_order_release);
notifier.commit_wait(i);
committed_done.fetch_add(1, std::memory_order_release);
local_round++;
}
});
}
for (size_t r = 1; r = target);
// Drain remaining waiters.
notifier.notify_all();
auto deadline2 = std::chrono::steady_clock::now() + std::chrono::seconds(10);
while (committed_done.load(std::memory_order_acquire) != N &&
std::chrono::steady_clock::now() < deadline2) {
notifier.notify_all();
std::this_thread::yield();
}
REQUIRE(committed_done.load(std::memory_order_acquire) == N);
REQUIRE(notifier.num_waiters() == 0);
}
stop.store(true, std::memory_order_release);
notifier.notify_all();
for (auto& t : workers) t.join();
REQUIRE(notifier.num_waiters() == 0);
}
// ============================================================================
// stress_test_notifier
// Validates:
// 1. notify_n(k) wakes at least k threads if k