FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

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

Back | FazBrowse Home | New Git URL