From 6802328e97245a623c5cae6be9566bad3fd37791 Mon Sep 17 00:00:00 2001 From: Duncan Tourolle Date: Thu, 6 Aug 2026 20:50:34 +0200 Subject: [PATCH] fix: a push must wake its consumer, even when the ring looked non-empty push(), try_push() and push_blocking() fired push_callback_ only on the empty->non-empty edge, and computed that edge from a head_ sampled before the item was published. A PoolNode consumer decides whether to run again from the level (count_ready -> approx_size), so a pop landing in that window left both sides standing down: producer (push) consumer (PoolNode firing) ------------------------ ---------------------------- samples t=782, h=781 -> was_empty = false, no wake pops idx 781, head_ = 782 count_ready(): head_==tail_==782 -> not ready, gate released to Idle tail_.store(783) The item is in the ring, the node is idle, and no wake is outstanding. The failure is absorbing: every later push then sees a non-empty ring, so the edge never fires again and the node sleeps while its backlog grows. Observed as a hang in bench_pipeline at (chain, depth=4, work_us=10, shared pool): all pool workers asleep in worker_loop, the reader blocked in pop(), and 218 items stranded in one channel with head_ stopped at exactly the index where the edge was dropped. Re-reading head_ after the tail_ store does not fix this. That is the store-buffer pattern, and under acquire/release both sides may legally read stale; forbidding it needs seq_cst on the producer's tail_ store and head_ load *and* on the consumer's head_ store and tail_ load, a fence on both hot paths. Firing unconditionally is correct by construction: the callback runs after the publishing store, so a consumer that observes the level at all observes the item. The redundant wakes are cheap -- on_input_ready re-checks the level and SubmitGate::claim() collapses a wake arriving mid-firing into the firing already in flight. The stress test named this exact hazard and could not detect it: it asserted only 1 <= callbacks <= N, which a *missed* callback satisfies. It now requires one callback per successful push, and fails at 1325/10000 against the old code. Co-Authored-By: Claude Opus 5 --- include/kpn/channel.hpp | 48 ++++++++++++++++++++++++++++++----- tests/test_channel_stress.cpp | 22 ++++++++++------ 2 files changed, 55 insertions(+), 15 deletions(-) diff --git a/include/kpn/channel.hpp b/include/kpn/channel.hpp index e7420ff..d586c1e 100644 --- a/include/kpn/channel.hpp +++ b/include/kpn/channel.hpp @@ -136,7 +136,6 @@ public: throw ChannelOverflowError(capacity_); } - const bool was_empty = (t == h); buf_[t & ring_mask_] = make_storage(std::move(value)); tail_.store(t + 1, std::memory_order_release); stats_.record_push(t - h + 1, data_bytes); @@ -144,7 +143,8 @@ public: wake_.fetch_add(1, std::memory_order_release); wake_.notify_one(); - if (was_empty && push_callback_) + // Level-triggered, not edge-triggered — see set_push_callback. + if (push_callback_) push_callback_(); } @@ -187,13 +187,13 @@ public: if (t - h >= capacity_) return PushResult::Full; const std::size_t data_bytes = ChannelDataSize::bytes(value); - const bool was_empty = (t == h); buf_[t & ring_mask_] = make_storage(std::move(value)); tail_.store(t + 1, std::memory_order_release); stats_.record_push(t - h + 1, data_bytes); wake_.fetch_add(1, std::memory_order_release); wake_.notify_one(); - if (was_empty && push_callback_) push_callback_(); + // Level-triggered, not edge-triggered — see set_push_callback. + if (push_callback_) push_callback_(); return PushResult::Taken; } @@ -212,13 +212,13 @@ public: const std::size_t h = head_.load(std::memory_order_acquire); if (t - h < capacity_) { // space available → normal push const std::size_t data_bytes = ChannelDataSize::bytes(value); - const bool was_empty = (t == h); buf_[t & ring_mask_] = make_storage(std::move(value)); tail_.store(t + 1, std::memory_order_release); stats_.record_push(t - h + 1, data_bytes); wake_.fetch_add(1, std::memory_order_release); wake_.notify_one(); - if (was_empty && push_callback_) push_callback_(); + // Level-triggered, not edge-triggered — see set_push_callback. + if (push_callback_) push_callback_(); return true; } // full: yield briefly and retry (consumer will drain) @@ -388,7 +388,41 @@ public: wake_.notify_all(); } - // Register a callback fired when the queue transitions empty→non-empty. + // Register a callback fired after every successful push. + // + // It fires on every push, not on the empty→non-empty transition, and that + // is a correctness requirement rather than a simplification. + // + // The edge version tested `was_empty = (t == h)` using an `h` sampled + // *before* the item was published. A PoolNode consumer decides whether to + // run again from the level (count_ready → approx_size), so the two sides + // could each read the other as stale and both stand down: + // + // producer (push) consumer (PoolNode firing) + // ------------------------ ---------------------------- + // samples t=782, h=781 + // -> was_empty = false, no wake + // pops idx 781, head_ = 782 + // count_ready(): head_==tail_==782 + // -> not ready, gate released to Idle + // tail_.store(783) + // + // The item is in the ring, the node is idle, and no wake is outstanding. + // Worse, the failure is absorbing: every later push now sees a non-empty + // ring, so `was_empty` is false forever and the callback never fires again. + // The node sleeps while its backlog grows and its consumer waits on it. + // + // Re-reading head_ after the tail_ store does not fix it. That is the + // store-buffer pattern, and under acquire/release both sides may legally + // read stale; forbidding it needs seq_cst on the producer's tail_ store and + // head_ load *and* on the consumer's head_ store and tail_ load — a fence + // on both hot paths. Firing unconditionally is correct by construction: + // the callback runs after the publishing store, so a consumer that observes + // the level at all observes the item. + // + // The redundant wakes are cheap. on_input_ready re-checks the level, and + // SubmitGate::claim() collapses a wake arriving during a firing into the + // firing already in flight, so the cost is one CAS, not one extra run. void set_push_callback(std::function cb) { push_callback_ = std::move(cb); } diff --git a/tests/test_channel_stress.cpp b/tests/test_channel_stress.cpp index 9f851e3..58ac3a0 100644 --- a/tests/test_channel_stress.cpp +++ b/tests/test_channel_stress.cpp @@ -153,11 +153,18 @@ TEST_CASE("SPSC: producer racing a disable() never throws and never hangs", } } -TEST_CASE("SPSC: push_callback fires on each empty->non-empty transition", +TEST_CASE("SPSC: push_callback fires for every push, never missed", "[channel][stress]") { - // The empty->non-empty callback ([channel.hpp] was_empty branch) is read by - // the consumer-side notification path. Run it under contention to make sure - // the was_empty detection isn't torn by a concurrent pop(). + // Regression: this callback is the *only* thing that wakes a PoolNode, and + // it used to fire only on the empty->non-empty edge, computed from a head_ + // sampled before the item was published. A concurrent pop() could drain the + // ring to empty in that window, so neither side saw the other: the item sat + // in the ring with the consumer idle, and because the trigger was an edge it + // never recovered. See set_push_callback in channel.hpp. + // + // The old version of this test asserted only `1 <= callbacks <= N`, which a + // *missed* callback satisfies — it named the hazard and could not detect it. + // One callback per successful push is the contract, so assert exactly that. Channel ch(/*capacity=*/4, /*spin_count=*/4); std::atomic callbacks{0}; ch.set_push_callback([&] { callbacks.fetch_add(1, std::memory_order_relaxed); }); @@ -175,10 +182,9 @@ TEST_CASE("SPSC: push_callback fires on each empty->non-empty transition", for (int i = 0; i < N; ++i) (void)ch.pop(); producer.join(); - // At least one transition, at most one per item; mainly we assert the run - // completed without TSan flagging a race on push_callback_/was_empty. - REQUIRE(callbacks.load() >= 1); - REQUIRE(callbacks.load() <= N); + // Exactly one callback per successful push. Fewer means a wake was dropped, + // which is the bug; more would mean a spurious wake was manufactured. + REQUIRE(callbacks.load() == N); } // Ordering contract of the out-of-band sentinel under contention.