fix: a re-offered sentinel is not data loss
🚦 CI / changes (push) Successful in 5s
🚦 CI / docker (push) Has been skipped
🚦 CI / test (push) Successful in 5m0s
🚦 CI / tsan (push) Successful in 3m37s
🚦 CI / docs (push) Has been skipped

139bfbb made the channel refuse a sentinel offered while one was still
pending — correct, and the reason is a data race: overwriting wrote
eof_value_ while the consumer could be moving the previous one out of it,
which ThreadSanitizer reports as a torn refcount and a heap-use-after-free.
That part stands.

What was wrong was the accounting. The refusal recorded a drop, and PoolNode
reported it through the overflow event callback, on the theory that two
control tokens on one channel means the stream ended twice and is a caller
protocol error.

It is not. A source that has reached the end of its input keeps being polled
and keeps returning EOF — that is the normal steady state, not an error — so
the token is re-offered on every firing. Refusing a re-offer loses nothing:
the pending token carries the same meaning and is already on its way.

Found by running it. scene-actor-extraction on a 14 s clip reported

    [main] ERROR: frames were dropped (channel overflow):
      frame_source: 2
      scene_annotate: 1
    [main] The output would describe footage that was never analysed.
           Refusing to report success.

and exited 2, on a run where nothing had been dropped and every frame was
analysed. frame_source emits EOF once and then returns it forever
(frame_source_node.hpp:76), so the count grows with however many times the
source is polled after the end. The pipeline's own loss detector — which
exists because a dropped frame silently corrupts the output — was being
tripped by a clean run, which is the one thing a loss detector must not do.

So SlotBusy now records nothing and reports nothing. The cost, stated
plainly: a genuinely distinct second token would also be refused silently,
and the channel cannot tell a re-offer from a distinct token. Re-offering is
the case that actually occurs; the delivery guarantee that matters — the
first token arrives — holds either way.

The test asserting the old behaviour is inverted rather than deleted, and now
also checks that the token which was accepted is the one delivered. 148/148.
This commit is contained in:
2026-08-06 20:44:38 +02:00
parent 27f884496d
commit c9aa246322
3 changed files with 45 additions and 24 deletions
+15 -3
View File
@@ -267,10 +267,22 @@ public:
// has_eof_, the consumer is the only one that clears it, so observing // has_eof_, the consumer is the only one that clears it, so observing
// false here means the consumer has finished with the storage and will // false here means the consumer has finished with the storage and will
// not touch it again until this store publishes the next token. // not touch it again until this store publishes the next token.
if (has_eof_.load(std::memory_order_acquire)) { //
stats_.record_drop(); // Not counted as a drop, and this is the important part. A source that
// has reached the end of its input keeps being polled and keeps
// returning EOF — that is the normal steady state, not an error — so a
// token arriving while one is already pending is a *re-offer*, and
// refusing it loses nothing: the pending token carries the same
// meaning and is already on its way. Counting it as a drop made a
// clean run report data loss and exit non-zero.
//
// The cost of that choice, stated plainly: a genuinely distinct second
// token would also be refused silently, and the channel cannot tell the
// two apart. Re-offering is the case that actually occurs here, and the
// delivery guarantee that matters — the first token arrives — holds
// either way.
if (has_eof_.load(std::memory_order_acquire))
return SentinelResult::SlotBusy; return SentinelResult::SlotBusy;
}
eof_value_ = make_storage(std::move(value)); eof_value_ = make_storage(std::move(value));
has_eof_.store(true, std::memory_order_release); has_eof_.store(true, std::memory_order_release);
// Wake a consumer blocked in pop(): the sentinel is now deliverable even // Wake a consumer blocked in pop(): the sentinel is now deliverable even
+14 -16
View File
@@ -622,14 +622,13 @@ private:
// downstream pop() forever. Deliver them out-of-band (push_sentinel), // downstream pop() forever. Deliver them out-of-band (push_sentinel),
// which never overflows and never blocks this node's worker thread. // which never overflows and never blocks this node's worker thread.
if (is_sentinel_value(val)) { if (is_sentinel_value(val)) {
// A refused sentinel is a protocol error, not backpressure, so it // Not parked and not reported. Parking would spin against a slot
// is reported rather than parked and retried — retrying would spin // only the consumer can free; reporting would cry data loss on the
// forever against a slot only the consumer can free, and there is // normal steady state, since a source at end of input keeps being
// no correct value to deliver second anyway. Closed is normal // polled and keeps returning EOF, so the token is re-offered on
// during teardown and stays quiet. // every firing. Refusing a re-offer loses nothing — the pending
if (ch->try_push_sentinel(val) == Channel<std::tuple_element_t<I, return_tuple>> // token says the same thing. See Channel::try_push_sentinel.
::SentinelResult::SlotBusy) ch->try_push_sentinel(val);
fire_callbacks(event_callbacks_);
return true; return true;
} }
// Backpressure without parking the worker. A full channel means the // Backpressure without parking the worker. A full channel means the
@@ -1207,14 +1206,13 @@ private:
// downstream pop() forever. Deliver them out-of-band (push_sentinel), // downstream pop() forever. Deliver them out-of-band (push_sentinel),
// which never overflows and never blocks this node's worker thread. // which never overflows and never blocks this node's worker thread.
if (is_sentinel_value(val)) { if (is_sentinel_value(val)) {
// A refused sentinel is a protocol error, not backpressure, so it // Not parked and not reported. Parking would spin against a slot
// is reported rather than parked and retried — retrying would spin // only the consumer can free; reporting would cry data loss on the
// forever against a slot only the consumer can free, and there is // normal steady state, since a source at end of input keeps being
// no correct value to deliver second anyway. Closed is normal // polled and keeps returning EOF, so the token is re-offered on
// during teardown and stays quiet. // every firing. Refusing a re-offer loses nothing — the pending
if (ch->try_push_sentinel(val) == Channel<std::tuple_element_t<I, return_tuple>> // token says the same thing. See Channel::try_push_sentinel.
::SentinelResult::SlotBusy) ch->try_push_sentinel(val);
fire_callbacks(event_callbacks_);
return true; return true;
} }
// See the note on the typed overload above: park rather than block. // See the note on the typed overload above: park rather than block.
+16 -5
View File
@@ -271,15 +271,26 @@ TEST_CASE("a second sentinel is refused, not swallowed", "[channel][sentinel]")
CHECK(out == 3); CHECK(out == 3);
} }
TEST_CASE("a refused sentinel is counted as a drop", "[channel][sentinel]") { TEST_CASE("a refused sentinel is not counted as a drop", "[channel][sentinel]") {
// Visibility matters more here than for a dropped value: the refusal means // A refusal means a token arrived while an equivalent one was already
// a control token went nowhere, and the only alternative to a counter is // pending — not that anything was lost. Counting it as a drop was wrong in
// for it to vanish. // a way that showed up immediately on real content: a source at the end of
// its input keeps being polled and keeps returning EOF, so the token is
// re-offered on every firing, and the pipeline reported hundreds of dropped
// frames on a clean run and exited non-zero.
//
// The delivery guarantee is unaffected: the first token is pending and will
// arrive. Only the accounting changed.
Channel<int> ch(4); Channel<int> ch(4);
REQUIRE(ch.push_sentinel(1)); REQUIRE(ch.push_sentinel(1));
const auto before = ch.stats().drops.load(); const auto before = ch.stats().drops.load();
REQUIRE_FALSE(ch.push_sentinel(2)); REQUIRE_FALSE(ch.push_sentinel(2));
CHECK(ch.stats().drops.load() == before + 1); CHECK(ch.stats().drops.load() == before);
// And the one that was accepted is still the one delivered.
int out = 0;
REQUIRE(ch.try_pop_now(out));
CHECK(out == 1);
} }
TEST_CASE("try_push_sentinel leaves a refused value untouched", "[channel][sentinel]") { TEST_CASE("try_push_sentinel leaves a refused value untouched", "[channel][sentinel]") {