diff --git a/include/kpn/channel.hpp b/include/kpn/channel.hpp index 4466fd8..99ebe51 100644 --- a/include/kpn/channel.hpp +++ b/include/kpn/channel.hpp @@ -180,7 +180,16 @@ public: if (h == t) { // Ring drained — deliver any pending out-of-band sentinel (EOF) // now, so it always arrives after the data pushed before it. - { T s; if (take_sentinel(s)) return s; } + // + // Re-confirm emptiness against a fresh tail_ first: the snapshot + // at the top of the loop may be stale (the producer can push more + // values *and* the sentinel in the window since), and the sentinel + // must never jump ahead of ring values pushed before it. The spin + // and post-spin takes below already reload tail_ on the line above + // them; this is the one take that used the loop-top snapshot. + if (h == tail_.load(std::memory_order_acquire)) { + T s; if (take_sentinel(s)) return s; + } if (!accepting_.load(std::memory_order_acquire)) throw ChannelClosedError{}; diff --git a/tests/test_channel_stress.cpp b/tests/test_channel_stress.cpp index bd113d0..eb58704 100644 --- a/tests/test_channel_stress.cpp +++ b/tests/test_channel_stress.cpp @@ -180,28 +180,24 @@ TEST_CASE("SPSC: push_callback fires on each empty->non-empty transition", REQUIRE(callbacks.load() <= N); } -// What the out-of-band sentinel guarantees under contention — and what it does -// not. push_sentinel() publishes has_eof_ (release) after the producer's N ring +// Ordering contract of the out-of-band sentinel under contention. +// +// push_sentinel() publishes has_eof_ (release) after the producer's N ring // pushes; a consumer that observes has_eof_ (acquire) therefore also observes -// every value pushed before it. What these tests assert: +// every value pushed before it. Both pop() and try_pop_now() only surface the +// sentinel once the ring is *freshly* observed empty, so the sentinel is the +// strictly last item received — it never jumps ahead of a ring value pushed +// before it. These tests treat the sentinel as a hard "last message" barrier +// (the consumer stops draining the moment it sees it) and assert that all N +// values arrived, in a contiguous 0..N-1 sequence, before it. // -// * Losslessness — every value 0..N-1 is delivered exactly once (contiguous, -// no gaps, no duplicates) and the sentinel is delivered exactly once. This -// is the invariant that must hold on every run; a broken acquire/release -// pairing would surface as a lost/duplicated value or (under TSan) a data -// race on has_eof_/eof_value_. -// -// What they deliberately do NOT assert is that the sentinel is the *strictly -// last* item popped. pop() checks emptiness (h == t) using a tail_ snapshot -// taken at the top of its loop; the producer can push more values *and* the -// sentinel in the window before take_sentinel() runs, so the consumer may -// surface the sentinel with a few real values still queued behind it. Those -// values are not lost — a consumer that keeps draining still receives them — -// but "sentinel arrives dead last" is not a property the channel promises, so -// asserting it would be flaky. We track how many values trailed the sentinel -// for visibility without failing on it. +// Regression guard: an earlier version of pop() checked emptiness against a +// stale tail_ snapshot from the top of its loop, so under load the sentinel +// could surface with a few real values still queued — breaking in_order / +// values==N here. Under TSan these also cover the has_eof_/eof_value_ +// acquire/release handshake and the spin/futex wakeup on push_sentinel(). -TEST_CASE("SPSC: sentinel and all values survive contention (blocking pop)", +TEST_CASE("SPSC: sentinel is strictly last, after every value (blocking pop)", "[channel][stress]") { constexpr int N = 20'000; constexpr int SENTINEL = -1; @@ -221,34 +217,32 @@ TEST_CASE("SPSC: sentinel and all values survive contention (blocking pop)", ch.push_sentinel(SENTINEL); // must-deliver, never overflows/blocks }); - std::vector seen(N, false); - int values = 0; - int sentinels = 0; - bool duplicate = false; - // Drain until the sentinel AND all N values have been received; the - // sentinel may arrive before the last few values (see note above). - while (values < N || sentinels == 0) { + int expected = 0; + bool in_order = true; + bool saw_sentinel = false; + // Treat the sentinel as EOF: stop draining the instant it appears. + for (;;) { int v = ch.pop(); - if (v == SENTINEL) { ++sentinels; continue; } - if (seen[v]) duplicate = true; else seen[v] = true; - ++values; + if (v == SENTINEL) { saw_sentinel = true; break; } + if (v != expected) in_order = false; + ++expected; } producer.join(); - REQUIRE_FALSE(duplicate); - REQUIRE(values == N); // every value delivered exactly once - REQUIRE(sentinels == 1); // sentinel delivered exactly once + REQUIRE(saw_sentinel); + REQUIRE(in_order); + REQUIRE(expected == N); // all N values received before the sentinel REQUIRE(ch.size() == 0); REQUIRE(ch.approx_size() == 0); } } -TEST_CASE("SPSC: sentinel and all values survive contention (try_pop_now)", +TEST_CASE("SPSC: sentinel is strictly last, after every value (try_pop_now)", "[channel][stress]") { // The pool-node consume path is try_pop_now(), not pop(): it must surface - // the out-of-band sentinel once the ring is observed empty. The consumer - // spins with no sleeps, racing the producer at full tilt across the - // empty-ring boundary where take_sentinel() is reached. + // the out-of-band sentinel only once the ring is freshly observed empty. + // The consumer spins with no sleeps, racing the producer at full tilt + // across the empty-ring boundary where take_sentinel() is reached. constexpr int N = 20'000; constexpr int SENTINEL = -1; @@ -265,23 +259,22 @@ TEST_CASE("SPSC: sentinel and all values survive contention (try_pop_now)", ch.push_sentinel(SENTINEL); }); - std::vector seen(N, false); - int values = 0; - int sentinels = 0; - bool duplicate = false; + int expected = 0; + bool in_order = true; + bool saw_sentinel = false; int v; - while (values < N || sentinels == 0) { + for (;;) { if (!ch.try_pop_now(v)) { std::this_thread::yield(); continue; } - if (v == SENTINEL) { ++sentinels; continue; } - if (seen[v]) duplicate = true; else seen[v] = true; - ++values; + if (v == SENTINEL) { saw_sentinel = true; break; } + if (v != expected) in_order = false; + ++expected; } producer.join(); - REQUIRE_FALSE(duplicate); - REQUIRE(values == N); - REQUIRE(sentinels == 1); - // Sentinel held no ring slot; once drained the channel is fully empty. + REQUIRE(saw_sentinel); + REQUIRE(in_order); + REQUIRE(expected == N); + // Sentinel held no ring slot; once taken the channel is fully empty. REQUIRE(ch.size() == 0); REQUIRE(ch.approx_size() == 0); }