cheatah
Source

stdlib/memory/tests/memory_concurrency_test.cpp

1// Copyright (c) 2026 BigBrain LLC. MIT-licensed (see LICENSE).
2// Original work; see ACKNOWLEDGMENTS.md for the open-source ideas we build upon.
3// GOLD adversarial CONCURRENCY suite for the `memory` module (suite MemoryConcurrency).
4//
5// ┌───────────────────────────────────────────────────────────────────────────────────────────┐
6// │ TDD-RED until the Owner scheduling engine is implemented. These use REAL std::thread over a │
7// │ shared Owner, so they cannot even be *linked* (Owner::rread/rwrite are declared-only) — and │
8// │ an exception thrown inside a spawned thread would std::terminate the runner. They are │
9// │ therefore built ONLY behind the CMake option CHEATAH_BUILD_MEMORY_TESTS (default OFF), which │
10// │ we flip ON as we build the engine (red → green). See stdlib/memory/tests/README.md. │
11// └───────────────────────────────────────────────────────────────────────────────────────────┘
12//
13// Design principle (per the user's ask): drive the object from many threads in a way whose *final*
14// state is DETERMINISTIC even though the interleaving is not. We never do "a bunch of rotations and
15// expect the same result"; we use commuting/exclusive updates and interleaving-invariant predicates,
16// so a correct engine ALWAYS yields the exact expected answer and a broken one (lost updates, torn
17// reads, missed drains, non-exclusive writes) is caught deterministically.
19#include <atomic>
20#include <chrono>
21#include <cstdint>
22#include <cstdio>
23#include <cstdlib>
24#include <functional>
25#include <random>
26#include <string>
27#include <thread>
28#include <vector>
30#include <gtest/gtest.h>
32#include "../memory.hpp"
33#include "../../ndarray/ndarray.hpp"
35namespace mem = cheatah::memory;
36namespace nd = cheatah::ndarray;
38#if __has_include(<valgrind/valgrind.h>)
39#include <valgrind/valgrind.h>
40#define CHEATAH_HAVE_VALGRIND_H 1
41#endif
43namespace {
45/// Are we running under valgrind (helgrind, drd, memcheck)?
46///
47/// It matters because those tools serialize threads, so a test whose subject IS concurrency cannot
48/// observe its own property under them. `RUNNING_ON_VALGRIND` is valgrind's own documented client
49/// request and costs nothing when absent; the `__has_include` guard keeps the header optional so a
50/// machine without valgrind-dev still builds.
51/// @return true iff running under a valgrind tool.
52/// @complexity O(1). @alloc none.
53[[nodiscard]] bool running_under_valgrind() noexcept {
54#ifdef CHEATAH_HAVE_VALGRIND_H
55 return RUNNING_ON_VALGRIND != 0;
56#else
57 return false;
58#endif
61/**
62 * RANDOM LATENCY INJECTION — off by default, on with `CHEATAH_MEMORY_JITTER=<max_microseconds>`.
63 *
64 * WHY, given the suite already passes: a passing concurrency test proves the interleavings the
65 * scheduler HAPPENED to pick were fine. It says nothing about the ones it never picked. Two real
66 * flakes in this module were found only when something perturbed timing — one by machine load, one
67 * by ThreadSanitizer's slowdown — which means timing perturbation was doing the finding, accidentally
68 * and unrepeatably. This makes it deliberate and repeatable.
69 *
70 * REPRODUCIBILITY IS THE POINT. The seed comes from `CHEATAH_MEMORY_SEED` when set; otherwise one is
71 * drawn and PRINTED, so a failure found by a random run can be replayed exactly. A fuzzer whose
72 * failures cannot be reproduced is a rumour generator.
73 *
74 * Per-thread state, so threads do not contend on the generator and thereby add a synchronisation
75 * point that changes the very interleaving being explored.
76 */
77[[nodiscard]] unsigned jitter_max_us() {
78 static const unsigned v = [] {
79 const char* e = std::getenv("CHEATAH_MEMORY_JITTER");
80 return e ? static_cast<unsigned>(std::strtoul(e, nullptr, 10)) : 0u;
81 }();
82 return v;
85[[nodiscard]] unsigned jitter_seed() {
86 static const unsigned s = [] {
87 if (const char* e = std::getenv("CHEATAH_MEMORY_SEED"))
88 return static_cast<unsigned>(std::strtoul(e, nullptr, 10));
89 const auto drawn = static_cast<unsigned>(std::random_device{}());
90 if (jitter_max_us() != 0)
91 static_cast<void>(std::fprintf(stderr,
92 "[memory-jitter] seed=%u (replay: CHEATAH_MEMORY_SEED=%u)\n", drawn, drawn));
93 return drawn;
94 }();
95 return s;
98/// Sleep a random sub-window, or yield. No-op unless jitter is enabled.
99void jitter() {
100 const unsigned max_us = jitter_max_us();
101 if (max_us == 0) return;
102 static thread_local std::mt19937 rng{jitter_seed() ^
103 static_cast<unsigned>(
104 std::hash<std::thread::id>{}(std::this_thread::get_id()))};
105 std::uniform_int_distribution<unsigned> d(0, max_us);
106 const unsigned us = d(rng);
107 if (us == 0) std::this_thread::yield();
108 else std::this_thread::sleep_for(std::chrono::microseconds(us));
111// Spawn `n` threads running `body(i)`, join all.
112template <class F>
113void run_threads(int n, F body) {
114 std::vector<std::thread> ts;
115 ts.reserve(n);
116 for (int i = 0; i < n; ++i) ts.emplace_back([=] { body(i); });
117 for (auto& t : ts) t.join();
119constexpr int kWriters = 8;
121/// Iterations per writer: enough to expose a lost update, fast enough for the gate.
122///
123/// SCALED DOWN UNDER VALGRIND, and that is what makes a helgrind lane possible at all. Helgrind
124/// instruments every memory access and every lock operation at roughly 100x, so 8 x 20,000
125/// acquisitions is minutes per test and the whole suite does not finish inside any sane timeout —
126/// measured, not guessed. The race conditions these tests hunt are not made more likely by volume
127/// under a tool that already serializes and inspects every interleaving; volume is how we buy
128/// coverage from a NATIVE scheduler, which is a different lane. So: full count natively and under
129/// TSan, a small count under valgrind, and the reduction is announced rather than silent.
130const int kIters = [] {
131 const int n = running_under_valgrind() ? 300 : 20'000;
132 if (running_under_valgrind())
133 static_cast<void>(std::fprintf(stderr,
134 "[memory] valgrind detected: kIters reduced to %d for tractability\n", n));
135 return n;
136}();
137} // namespace
139// ── 1. Exclusive writes never lose an update: 8×N increments == 8N, exactly ───────────────
140TEST(MemoryConcurrency, ManyWritersDeterministicSum) {
141 auto o = mem::own<long long>(0);
142 run_threads(kWriters, [&](int) {
143 for (int k = 0; k < kIters; ++k) {
144 auto w = o.rwrite().acquire(); // exclusive: read-modify-write cannot interleave
145 w.write(w.read() + 1);
146 }
147 });
148 EXPECT_EQ(o.rread().acquire().read(),
149 static_cast<long long>(kWriters) * kIters); // == 400000, deterministic
152// ── 2. Owner<ndarray>: readers NEVER see a torn (non-uniform) array; final is exact ───────
153// The array starts uniform (all equal). Every writer adds 1 to *every* element under one write
154// lease, so under a correct exclusive lease the array is uniform at every quiescent point. Readers
155// assert uniformity — the value they see varies (nondeterministic), but "all elements equal" must
156// ALWAYS hold. A missed drain / non-exclusive write would let a reader observe a half-updated array.
157TEST(MemoryConcurrency, OwnerOfNdArrayStaysUniformAndSumsExactly) {
158 constexpr std::size_t N = 256;
159 auto o = mem::own(nd::basic_ndarray<long long>({N}, 0));
160 std::atomic<bool> torn{false};
161 std::atomic<bool> stop{false};
163 // 4 reader threads: continuously assert the array is uniform.
164 std::vector<std::thread> readers;
165 readers.reserve(4);
166 for (int r = 0; r < 4; ++r) {
167 readers.emplace_back([&] {
168 while (!stop.load()) {
169 auto lease = o.rread().acquire();
170 const auto& a = lease.read();
171 const long long first = a.at({0}); // const read: at() (operator[] is non-const)
172 for (std::size_t i = 1; i < a.size(); ++i)
173 if (a.at({i}) != first) { torn = true; return; }
174 }
175 });
176 }
177 // Writers: each adds 1 to every element, kIters/50 times.
178 const int rounds = kIters / 50;
179 run_threads(kWriters, [&](int) {
180 for (int k = 0; k < rounds; ++k) {
181 auto w = o.rwrite().acquire();
182 const std::size_t n = w.read().size();
183 for (std::size_t i = 0; i < n; ++i) w.write(i, w.read().at({i}) + 1); // indexed setter
184 }
185 });
186 stop = true;
187 for (auto& t : readers) t.join();
189 EXPECT_FALSE(torn.load()) << "a reader observed a partially-updated (non-uniform) array";
190 auto final = o.rread().acquire();
191 const long long expected = static_cast<long long>(kWriters) * rounds;
192 for (std::size_t i = 0; i < N; ++i)
193 ASSERT_EQ(final.read().at({i}), expected); // const read via at()
196// ── 3. Readers never see a torn write of a multi-field invariant (sum == a + b) ───────────
197TEST(MemoryConcurrency, ReadersNeverSeeATornWrite) {
198 struct Triple { long long a = 0, b = 0, sum = 0; };
199 auto o = mem::own(Triple{});
200 std::atomic<bool> torn{false}, stop{false};
202 std::thread reader([&] {
203 while (!stop.load()) {
204 auto r = o.rread().acquire();
205 const Triple& t = r.read();
206 if (t.sum != t.a + t.b) { torn = true; return; } // invariant must hold at every read
207 }
208 });
209 run_threads(kWriters, [&](int id) {
210 for (int k = 0; k < kIters; ++k) {
211 const long long a = id * 1000 + k, b = k;
212 auto w = o.rwrite().acquire();
213 w.write(Triple{a, b, a + b}); // set all three at once (one exclusive write)
214 }
215 });
216 stop = true;
217 reader.join();
218 EXPECT_FALSE(torn.load()) << "a reader observed a half-written Triple (sum != a + b)";
221// ── 4. Immediate-write preempts under load; its effect lands; normal writers still finish ─
222// Adds commute, so the final total is deterministic regardless of WHEN the immediate-write fires.
223TEST(MemoryConcurrency, ImmediateWriteLandsUnderLoad) {
224 auto o = mem::own<long long>(0);
225 std::atomic<bool> go{false};
226 std::thread emergency([&] {
227 while (!go.load()) std::this_thread::yield();
228 auto w = o.rwrite<mem::immediate>().acquire(); // jumps the queue, preempts active writer
229 w.write(w.read() + 1'000'000);
230 });
231 go = true;
232 run_threads(kWriters, [&](int) {
233 for (int k = 0; k < kIters; ++k) { auto w = o.rwrite().acquire(); w.write(w.read() + 1); }
234 });
235 emergency.join();
236 EXPECT_EQ(o.rread().acquire().read(),
237 static_cast<long long>(kWriters) * kIters + 1'000'000);
240// ── 5. A thread that reads then writes in a loop must never self-deadlock ──────────────────
241TEST(MemoryConcurrency, ReadThenWriteLoopNoDeadlock) {
242 auto o = mem::own<long long>(0);
243 run_threads(kWriters, [&](int) {
244 for (int k = 0; k < kIters / 5; ++k) {
245 { auto r = o.rread().acquire(); (void)r.read(); } // read, release
246 { auto w = o.rwrite().acquire(); w.write(w.read() + 1); } // then write — must not wait on self
247 }
248 });
249 EXPECT_EQ(o.rread().acquire().read(), static_cast<long long>(kWriters) * (kIters / 5));
252// ── 6. Renewal across relocation: writers grow a string (reallocating), readers never dangle ─
253TEST(MemoryConcurrency, RenewalAcrossRelocationNeverDangles) {
254 auto o = mem::own(std::string("a"));
255 std::atomic<bool> stop{false}, corrupt{false};
256 std::thread reader([&] {
257 while (!stop.load()) {
258 auto r = o.rread().acquire();
259 const std::string& s = r.read(); // must point at the CURRENT buffer
260 for (char c : s) if (c != 'a') { corrupt = true; return; } // every byte is 'a', never garbage
261 }
262 });
263 for (int k = 0; k < 2000; ++k) { auto w = o.rwrite().acquire(); w.write(w.read() + "a"); } // grows/relocates
264 stop = true;
265 reader.join();
266 EXPECT_FALSE(corrupt.load()) << "a renewed reader saw freed/garbage bytes after a relocation";
267 EXPECT_EQ(o.rread().acquire().read().size(), 2001u);
270// ── 7. Lease churn: hammer create/destroy so ASan/Valgrind (gate stages) can catch a leak ─
271TEST(MemoryConcurrency, LeaseChurnLeakHunt) {
272 auto o = mem::own<long long>(0);
273 run_threads(kWriters, [&](int) {
274 for (int k = 0; k < kIters; ++k) {
275 if (k & 1) { auto r = o.rread().acquire(); (void)r.read(); }
276 else { auto w = o.rwrite().acquire(); w.write(w.read() + 1); }
277 }
278 });
279 // Value check is secondary; the real assertion is "no bytes leaked" under ASan/Valgrind.
280 EXPECT_GE(o.rread().acquire().read(), 0);
283// ── 8. Concurrent readers coexist (shared), and a write still drains them ─────────────────
284// `peak > 1` is an EMERGENT property, and this test used to simply hope for it: the writer fired 200
285// back-to-back writes and set `stop`, and whether any two readers ever overlapped was left to the
286// scheduler. Under ThreadSanitizer — which the QA gate runs — thread start-up is slow enough that the
287// writer regularly finished before the readers got going at all, so `peak` stayed 0 and this failed
288// about one run in three. The property is real; it just has to be ARRANGED rather than wished for.
289//
290// Two changes, and both are about giving the property a chance instead of asserting it blindly:
291// 1. a start gate, so every reader is inside its loop before the writer begins; and
292// 2. the writer does not end the run until an overlap has actually been observed (bounded).
293//
294// The start gate is taken BEFORE any lease is acquired. That ordering is load-bearing: readers that
295// waited on each other while HOLDING read leases would block the writer's drain, and the drain is
296// what they would be waiting on — a circular wait, which is a deadlock rather than a flake.
297TEST(MemoryConcurrency, ManyReadersCoexistThenAWriteDrains) {
298 using clock = std::chrono::steady_clock;
299 constexpr int kReaders = 5;
300 constexpr auto kPatience = std::chrono::seconds(5);
302 // SKIPPED UNDER VALGRIND, and this is a tool limit rather than a weakness in the assertion.
303 // Helgrind and Memcheck run threads ONE AT A TIME — that serialization is how they get a
304 // consistent view — so two read leases can never be live simultaneously and `peak` is pinned at
305 // 1 by construction. Every other test in this file measures a FINAL state and is therefore
306 // meaningful serialized; this one measures concurrency itself, which is the one thing a
307 // serializing tool cannot show. Skipped with the reason stated, never quietly weakened to
308 // `>= 1` — an assertion that passes under the tool by no longer testing anything is worse than
309 // one that admits it did not run.
310 if (running_under_valgrind())
311 GTEST_SKIP() << "valgrind serializes threads; overlapping read leases are unobservable. "
312 "Run this lane natively or under ThreadSanitizer, which models atomics.";
314 auto o = mem::own<long long>(7);
315 std::atomic<int> peak{0}, active{0}, ready{0};
316 std::atomic<bool> stop{false};
317 run_threads(kReaders + 1, [&](int id) {
318 if (id == 0) {
319 while (ready.load() < kReaders) std::this_thread::yield(); // readers are running
320 for (int k = 0; k < 200; ++k) {
321 jitter();
322 auto w = o.rwrite().acquire();
323 jitter();
324 w.write(w.read() + 1);
325 }
326 // Hold the run open until the readers have demonstrably shared the object. Bounded, so a
327 // real regression — reads that serialize behind each other — FAILS on the assertion
328 // below instead of hanging the suite.
329 const auto deadline = clock::now() + kPatience;
330 while (peak.load() <= 1 && clock::now() < deadline)
331 std::this_thread::sleep_for(std::chrono::microseconds(100));
332 stop = true;
333 } else {
334 ++ready;
335 while (!stop.load()) {
336 jitter();
337 auto r = o.rread().acquire();
338 int a = ++active;
339 int seen = peak.load();
340 while (a > seen && !peak.compare_exchange_weak(seen, a)) {}
341 std::this_thread::yield();
342 --active;
343 EXPECT_GE(r.read(), 7); // monotonic: writes only increment
344 }
345 }
346 });
347 EXPECT_GT(peak.load(), 1) << "read leases should have coexisted (shared), not serialized";
348 EXPECT_EQ(o.rread().acquire().read(), 207);