Author SHA1 Message Date
dtourolleandClaude Opus 5 ac36159f37 fix: apply lossless push in PoolObjectNode too
PoolNode and PoolObjectNode each have their own push_one_out(). The previous
commit only reached PoolNode's copy, so ObjectNode-based pipelines — the common
case — still dropped on overflow with lossless enabled.

Verified end to end: a 2300-frame source now reports exactly 2300 frames at
every stage, where it previously delivered 774 to the second node.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-25 23:12:51 +02:00
dtourolleandClaude Opus 5 4e81752838 feat: opt-in lossless (blocking) output for nodes and fanouts
Channels drop on overflow by default. That is the right behaviour for live
sources, where a stale item is worth less than a fresh one, but it makes the
library unusable for offline batch work: items vanish with no diagnostic, and
any downstream analysis that assumes a fixed sample rate is silently invalid.

Auto-inserted fanouts were the harder half of this. They are created inside
make_network(), so user code cannot reach them to configure, and they drop
per-output inside a swallowed catch — so a pipeline whose own nodes were all
configured lossless could still lose items with nothing reported anywhere. In
the pipeline this came from, capture's 2505 frames arrived at the detector as
774 while the overflow counter read zero.

Adds:
  - INode::set_lossless_output(bool), defaulted to a no-op so node types with
    no output channels ignore it
  - PoolNode / PoolObjectNode: route push_one_out() through push_blocking()
  - FanoutNode: same, plus set_lossy_output(i) to opt a single branch back out
  - StaticNetwork::set_lossless(), which reaches user nodes and fanouts alike
  - StaticNetwork::drain(), a public wrapper over the existing private
    drain_all_channels(), so callers can flush in-flight work before stop()

Default behaviour is unchanged; every path is off unless explicitly enabled.

Blocking output is only safe when every consumer eventually drains. A branch
that can stall indefinitely — a display node nobody is servicing — will apply
backpressure to the whole pipeline, which is what set_lossy_output() is for.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-25 23:03:56 +02:00
9 changed files with 117 additions and 654 deletions
-1
View File
@@ -27,4 +27,3 @@ Thumbs.db
# Claude Code local settings # Claude Code local settings
.claude/settings.local.json .claude/settings.local.json
include/kpn/ort_cache/
+1 -44
View File
@@ -136,42 +136,6 @@ public:
push_callback_(); push_callback_();
} }
/// Called when a pop frees a slot in a previously-full ring.
///
/// The mirror of `set_push_callback`, and it exists for the same reason:
/// a producer must be able to *park* rather than spin. Without it the only
/// lossless option is `push_blocking`, which sleeps inside the caller's
/// thread — and when that thread is a scheduler worker, parking it starves
/// every node pinned to it (see the hold-and-wait note on push_sentinel).
void set_space_callback(std::function<void()> cb) { space_callback_ = std::move(cb); }
/// True when a push would currently succeed. Used to close the lost-wakeup
/// race: a producer that parks must re-check after clearing its queued flag,
/// because a space_callback fired in between would otherwise be swallowed.
bool has_space() const {
return tail_.load(std::memory_order_relaxed) -
head_.load(std::memory_order_acquire) < capacity_;
}
/// Non-blocking, lossless push. Returns false when the ring is full, having
/// changed nothing — the caller keeps the value and retries when woken.
bool try_push(T& value) {
if (!accepting_.load(std::memory_order_acquire)) { stats_.record_drop(); return true; }
const std::size_t t = tail_.load(std::memory_order_relaxed);
const std::size_t h = head_.load(std::memory_order_acquire);
if (t - h >= capacity_) return false;
const std::size_t data_bytes = ChannelDataSize<T>::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_();
return true;
}
// Lossless push with BACKPRESSURE: if the ring is full, wait for the consumer to // Lossless push with BACKPRESSURE: if the ring is full, wait for the consumer to
// drain instead of dropping (the throwing push()) — the producer just runs slower. // drain instead of dropping (the throwing push()) — the producer just runs slower.
// Use when every value must be delivered (e.g. replaying a dump for scoring, where // Use when every value must be delivered (e.g. replaying a dump for scoring, where
@@ -284,8 +248,6 @@ public:
throw ChannelClosedError{}; throw ChannelClosedError{};
T value = extract(std::move(buf_[h & ring_mask_])); T value = extract(std::move(buf_[h & ring_mask_]));
head_.store(h + 1, std::memory_order_release); head_.store(h + 1, std::memory_order_release);
// A slot just freed: wake any producer parked on this channel.
if (t - h >= capacity_ && space_callback_) space_callback_();
stats_.record_pop(); stats_.record_pop();
return value; return value;
} }
@@ -307,15 +269,11 @@ public:
// so pool nodes — which pop only via this path — still receive the token. // so pool nodes — which pop only via this path — still receive the token.
bool try_pop_now(T& out) { bool try_pop_now(T& out) {
const std::size_t h = head_.load(std::memory_order_relaxed); const std::size_t h = head_.load(std::memory_order_relaxed);
const std::size_t t = tail_.load(std::memory_order_acquire); if (h == tail_.load(std::memory_order_acquire))
if (h == t)
return take_sentinel(out); return take_sentinel(out);
out = extract(std::move(buf_[h & ring_mask_])); out = extract(std::move(buf_[h & ring_mask_]));
head_.store(h + 1, std::memory_order_release); head_.store(h + 1, std::memory_order_release);
stats_.record_pop(); stats_.record_pop();
// Pool nodes pop only through here, so this is where a parked producer
// gets woken: the ring was full, and it no longer is.
if (t - h >= capacity_ && space_callback_) space_callback_();
return true; return true;
} }
@@ -405,7 +363,6 @@ private:
std::size_t ring_mask_; std::size_t ring_mask_;
std::unique_ptr<storage_type[]> buf_; std::unique_ptr<storage_type[]> buf_;
std::function<void()> push_callback_; std::function<void()> push_callback_;
std::function<void()> space_callback_;
ChannelStats stats_; ChannelStats stats_;
// Out-of-band sentinel (EOF): stored outside the ring so its delivery never // Out-of-band sentinel (EOF): stored outside the ring so its delivery never
+25 -2
View File
@@ -115,6 +115,16 @@ public:
out_channels_[I] = ch; out_channels_[I] = ch;
} }
// Lossless fanout: block until each consumer drains rather than dropping.
void set_lossless_output(bool on) override { lossless_ = on; }
// Opt a single output back out of blocking. Needed when one branch may
// stall indefinitely — a display tap nobody is servicing, say — since
// blocking on it would apply backpressure to every other branch too.
void set_lossy_output(std::size_t i, bool lossy = true) {
if (i < N) lossy_out_[i] = lossy;
}
private: private:
void run_loop() { void run_loop() {
while (!stop_flag_.load(std::memory_order_relaxed)) { while (!stop_flag_.load(std::memory_order_relaxed)) {
@@ -126,8 +136,19 @@ private:
for (std::size_t i = 0; i < N; ++i) { for (std::size_t i = 0; i < N; ++i) {
if (out_channels_[i]) { if (out_channels_[i]) {
try { out_channels_[i]->push(val); } // Lossless: block until this consumer drains. Note the
catch (const ChannelOverflowError&) {} // drop for this output independently // branches differ in more than blocking — the dropping
// path discards per-output independently and silently,
// so a slow consumer on one branch costs frames on that
// branch only, with no diagnostic. That is the right
// default for display taps but hides frame loss from
// analysis branches.
if (lossless_ && !lossy_out_[i])
out_channels_[i]->push_blocking(val);
else {
try { out_channels_[i]->push(val); }
catch (const ChannelOverflowError&) {} // drop independently
}
} }
} }
@@ -142,6 +163,8 @@ private:
std::string name_; std::string name_;
std::size_t fifo_capacity_; std::size_t fifo_capacity_;
bool lossless_{false};
std::array<bool, N> lossy_out_{}; // per-output opt-out of blocking
std::shared_ptr<Channel<T>> input_ch_; std::shared_ptr<Channel<T>> input_ch_;
std::array<Channel<T>*, N> out_channels_{}; std::array<Channel<T>*, N> out_channels_{};
std::atomic<bool> stop_flag_{false}; std::atomic<bool> stop_flag_{false};
+4 -8
View File
@@ -34,17 +34,13 @@ struct INode {
virtual void set_network_overflow_callback(NodeEventCallback) {} virtual void set_network_overflow_callback(NodeEventCallback) {}
virtual void set_network_closed_callback(NodeEventCallback) {} virtual void set_network_closed_callback(NodeEventCallback) {}
// Network-level error listener. Consulted when a node's function throws
// and no per-node handler resolved it. Without this the exception is
// discarded and the failure is only visible as a Closed event, which says
// a node stopped but not why — the difference between a diagnosis and a
// guess. Same contract as NodeErrorHandler: true to continue, false to
// stop the node.
virtual void set_network_error_callback(NodeErrorHandler) {}
// halt(): alias for stop() — immediate, discards in-flight work. // halt(): alias for stop() — immediate, discards in-flight work.
virtual void halt() { stop(); } virtual void halt() { stop(); }
// Opt into lossless (blocking) output for nodes that support it. Default
// is a no-op so node types with no output channels ignore it.
virtual void set_lossless_output(bool) {}
// shutdown(): graceful drain before stopping. Base implementation falls // shutdown(): graceful drain before stopping. Base implementation falls
// back to stop(). Network and StaticNetwork override with topo-ordered drain. // back to stop(). Network and StaticNetwork override with topo-ordered drain.
virtual void shutdown() { stop(); } virtual void shutdown() { stop(); }
+62 -391
View File
@@ -15,7 +15,6 @@
#include <iostream> #include <iostream>
#include <memory> #include <memory>
#include <optional> #include <optional>
#include <variant>
#include <stdexcept> #include <stdexcept>
#include <thread> #include <thread>
#include <tuple> #include <tuple>
@@ -132,9 +131,16 @@ public:
void set_name(std::string name) override { name_ = std::move(name); } void set_name(std::string name) override { name_ = std::move(name); }
void set_error_handler(NodeErrorHandler h) { error_handler_ = std::move(h); } void set_error_handler(NodeErrorHandler h) { error_handler_ = std::move(h); }
void set_network_error_callback(NodeErrorHandler h) override { net_error_handler_ = std::move(h); }
void set_max_exec_time(std::chrono::milliseconds t) { max_exec_time_ = t; } void set_max_exec_time(std::chrono::milliseconds t) { max_exec_time_ = t; }
// Lossless output: block the producer until the consumer drains rather than
// dropping on a full channel. Default is drop, which suits live sources
// where a stale frame is worth less than a fresh one. Enable for offline
// batch runs, where a dropped item leaves a gap that downstream analysis
// cannot recover. Must be set before start().
void set_lossless_output(bool on) override { lossless_ = on; }
void set_lossless(bool on) { set_lossless_output(on); }
void set_overflow_callback(NodeEventCallback cb) { event_callbacks_[0] = std::move(cb); } void set_overflow_callback(NodeEventCallback cb) { event_callbacks_[0] = std::move(cb); }
void set_network_overflow_callback(NodeEventCallback cb) override { event_callbacks_[1] = std::move(cb); } void set_network_overflow_callback(NodeEventCallback cb) override { event_callbacks_[1] = std::move(cb); }
void set_closed_callback(NodeEventCallback cb) { closed_callbacks_[0] = std::move(cb); } void set_closed_callback(NodeEventCallback cb) { closed_callbacks_[0] = std::move(cb); }
@@ -236,8 +242,6 @@ private:
template<std::size_t... Is> template<std::size_t... Is>
void register_callbacks(std::index_sequence<Is...>) { void register_callbacks(std::index_sequence<Is...>) {
// A parked producer is re-submitted when its output drains.
register_space_callbacks(std::make_index_sequence<output_count>{});
(std::get<Is>(input_channels_)->set_push_callback( (std::get<Is>(input_channels_)->set_push_callback(
[this] { on_input_ready(); }), ...); [this] { on_input_ready(); }), ...);
} }
@@ -247,26 +251,10 @@ private:
for (auto& cb : cbs) if (cb) cb(ts); for (auto& cb : cbs) if (cb) cb(ts);
} }
template<std::size_t... Os>
bool outputs_have_space(std::index_sequence<Os...>) const {
return (... && (!std::get<Os>(output_channels_) ||
std::get<Os>(output_channels_)->has_space()));
}
template<std::size_t... Os>
void register_space_callbacks(std::index_sequence<Os...>) {
((std::get<Os>(output_channels_)
? (void)std::get<Os>(output_channels_)->set_space_callback(
[this] { try_submit(0.5f); })
: (void)0), ...);
}
void self_stop() { void self_stop() {
disable_inputs(std::make_index_sequence<input_count>{}); disable_inputs(std::make_index_sequence<input_count>{});
disable_outputs(std::make_index_sequence<output_count>{}); disable_outputs(std::make_index_sequence<output_count>{});
stats_.exec_start_us.store(0, std::memory_order_relaxed); stats_.exec_start_us.store(0, std::memory_order_relaxed);
// Plain store, not release_and_recheck(): this node is stopping, and
// honouring a pending wake here would resubmit a dead node.
queued_.store(false, std::memory_order_release); queued_.store(false, std::memory_order_release);
stop_flag_.store(true, std::memory_order_relaxed); stop_flag_.store(true, std::memory_order_relaxed);
} }
@@ -300,51 +288,11 @@ private:
return ((std::get<Is>(input_channels_)->approx_size() > 0 ? 1u : 0u) + ...); return ((std::get<Is>(input_channels_)->approx_size() > 0 ? 1u : 0u) + ...);
} }
/// Priority in [0,1], higher runs sooner.
///
/// Input fill alone answers "how much work is waiting for me". That is
/// only half the question: a node whose *outputs* are already full cannot
/// deliver anything: running it produces a value with nowhere to go, so it
/// immediately parks and the slot is wasted. Meanwhile the node that would
/// have drained that full channel waits behind it.
///
/// So occupancy of the outputs is deducted from occupancy of the inputs.
/// The scheduler then naturally favours whoever is furthest downstream of
/// a bottleneck — the node whose inputs are backed up but whose outputs
/// have room is exactly the one whose execution frees the most capacity —
/// and defers producers that would only deepen a queue that is already
/// full.
///
/// Mapped as 0.5·(1 + in - out) rather than clamping (in - out) at zero:
/// both terms are mean fills in [0,1], so the difference is in [-1,1], and
/// the affine map keeps the whole range distinguishable instead of
/// collapsing every output-saturated node onto the same value. 0.5 remains
/// the neutral point, matching the default used for source nodes.
float compute_priority() { float compute_priority() {
if constexpr (input_count == 0) return 0.5f; if constexpr (input_count == 0) return 0.5f;
float in = 0.0f; float sum = 0.0f;
sum_fill(in, std::make_index_sequence<input_count>{}); sum_fill(sum, std::make_index_sequence<input_count>{});
in /= static_cast<float>(input_count); return sum / static_cast<float>(input_count);
if constexpr (output_count == 0) return in;
float out = 0.0f;
sum_output_fill(out, std::make_index_sequence<output_count>{});
out /= static_cast<float>(output_count);
const float p = 0.5f * (1.0f + in - out);
return p < 0.0f ? 0.0f : (p > 1.0f ? 1.0f : p);
}
/// Mean fill of the output channels, same normalisation as sum_fill.
/// An unconnected output holds nothing back, so it contributes 0.
template<std::size_t... Os>
void sum_output_fill(float& sum, std::index_sequence<Os...>) {
((sum += (std::get<Os>(output_channels_) &&
std::get<Os>(output_channels_)->capacity() > 0)
? float(std::get<Os>(output_channels_)->approx_size())
/ float(std::get<Os>(output_channels_)->capacity())
: 0.0f), ...);
} }
template<std::size_t... Is> template<std::size_t... Is>
@@ -355,32 +303,10 @@ private:
: 0.5f), ...); : 0.5f), ...);
} }
/// Submit unless already queued. A wake that arrives while this node is
/// queued or running is *recorded*, never dropped.
///
/// Wakes are edge-triggered: a channel fires its space callback on the
/// transition, once. If that lands while queued_ is up, the CAS below fails
/// and — before wake_pending_ — the wake was gone. A node could then park a
/// value, release its worker, and sleep forever holding output its consumer
/// was waiting for, with every worker idle in cond_wait and nothing left to
/// re-trigger it. Recording the drop turns the signal level-triggered: the
/// invariant is that a node never sleeps with a wake outstanding, enforced
/// by release_and_recheck() at every point that releases the node.
void try_submit(float priority) { void try_submit(float priority) {
bool expected = false; bool expected = false;
if (queued_.compare_exchange_strong(expected, true, std::memory_order_acq_rel)) if (queued_.compare_exchange_strong(expected, true, std::memory_order_acq_rel))
scheduler_->submit([this] { fire_once(); }, priority); scheduler_->submit([this] { fire_once(); }, priority);
else
wake_pending_.store(true, std::memory_order_release);
}
/// Clear queued_, then honour any wake that was dropped while it was up.
/// Every path that finishes or parks a firing must release the node through
/// here rather than storing queued_ directly.
void release_and_recheck(float priority = 0.5f) {
queued_.store(false, std::memory_order_release);
if (wake_pending_.exchange(false, std::memory_order_acq_rel))
try_submit(priority);
} }
// ── Execution ───────────────────────────────────────────────────────────── // ── Execution ─────────────────────────────────────────────────────────────
@@ -397,77 +323,6 @@ private:
t0.time_since_epoch()).count(); t0.time_since_epoch()).count();
stats_.exec_start_us.store(now_us, std::memory_order_relaxed); stats_.exec_start_us.store(now_us, std::memory_order_relaxed);
// Parked from a previous firing: retry that value before touching the
// inputs. Returning here releases the worker — the channel's space
// callback re-submits this node when the consumer drains a slot.
if constexpr (!std::is_void_v<return_raw>) {
if (pending_) {
push_outputs(std::move(*pending_), std::make_index_sequence<output_count>{});
release_and_recheck();
if (pending_) {
// Close the lost-wakeup race: a space_callback that fired
// between the failed push and clearing queued_ was
// swallowed, and nothing else will wake this node. Re-check
// now that the flag is down.
if (outputs_have_space(std::make_index_sequence<output_count>{}))
try_submit(0.5f);
return; // parked
}
// Drained: resume normal firing, resubmitting exactly the way
// the normal tail below does. An unconditional try_submit here
// would fire a node whose inputs are empty, and pop_inputs
// reports an empty channel as ChannelClosedError — which this
// node treats as "upstream finished" and self-stops on. That
// is a live node killing itself purely because it was woken by
// *output* space rather than by input arrival.
if constexpr (input_count == 0) try_submit(0.5f);
else on_input_ready();
return;
}
}
// Parked from a previous firing: retry that value before touching the
// inputs. Returning here releases the worker — the channel's space
// callback re-submits this node once the consumer drains a slot.
if constexpr (!std::is_void_v<return_raw>) {
if (pending_) {
push_outputs(std::move(*pending_), std::make_index_sequence<output_count>{});
release_and_recheck();
if (pending_) {
// Close the lost-wakeup race: a space_callback that fired
// between the failed push and clearing queued_ was
// swallowed, and nothing else will wake this node. Re-check
// now that the flag is down.
if (outputs_have_space(std::make_index_sequence<output_count>{}))
try_submit(0.5f);
return; // parked
}
// Drained: resume normal firing, resubmitting exactly the way
// the normal tail below does. An unconditional try_submit here
// would fire a node whose inputs are empty, and pop_inputs
// reports an empty channel as ChannelClosedError — which this
// node treats as "upstream finished" and self-stops on. That
// is a live node killing itself purely because it was woken by
// *output* space rather than by input arrival.
if constexpr (input_count == 0) try_submit(0.5f);
else on_input_ready();
return;
}
}
// Woken by output space rather than by input arrival, with nothing
// parked left to flush: there is no work to do. Falling through would
// read an empty channel, and pop_one reports empty as
// ChannelClosedError — self-stopping a live node. Release the worker;
// on_input_ready() resubmits when data actually lands.
if constexpr (input_count > 0) {
if (count_ready(std::make_index_sequence<input_count>{}) != input_count) {
release_and_recheck();
on_input_ready(); // data may have arrived while we checked
return;
}
}
try { try {
auto args = pop_inputs(std::make_index_sequence<input_count>{}); auto args = pop_inputs(std::make_index_sequence<input_count>{});
auto t1 = clock_t::now(); auto t1 = clock_t::now();
@@ -493,11 +348,7 @@ private:
} catch (const ChannelOverflowError&) { } catch (const ChannelOverflowError&) {
fire_callbacks(event_callbacks_); fire_callbacks(event_callbacks_);
} catch (...) { } catch (...) {
auto eptr = std::current_exception(); if (error_handler_ && error_handler_(name_, std::current_exception())) {
const bool handled =
(error_handler_ && error_handler_(name_, eptr)) ||
(net_error_handler_ && net_error_handler_(name_, eptr));
if (handled) {
// continue — fall through to resubmit check // continue — fall through to resubmit check
} else { } else {
fire_callbacks(closed_callbacks_); fire_callbacks(closed_callbacks_);
@@ -507,27 +358,10 @@ private:
} }
stats_.exec_start_us.store(0, std::memory_order_relaxed); stats_.exec_start_us.store(0, std::memory_order_relaxed);
release_and_recheck(); queued_.store(false, std::memory_order_release);
if (stop_flag_.load(std::memory_order_relaxed)) return; if (stop_flag_.load(std::memory_order_relaxed)) return;
// Parked by the push above. Same situation as the retry path at the top
// of fire_once — and the same lost-wakeup race, which that path closes
// and this one did not. A space callback that fired while queued_ was
// still up got swallowed by try_submit's CAS, and the resubmit below
// cannot cover it: this firing consumed its input, so inputs are empty
// and on_input_ready() will not resubmit. The node would then hold its
// value forever while its consumer waits for exactly that value and its
// producer parks on an input channel that never drains. Re-check now
// that the flag is down.
if constexpr (!std::is_void_v<return_raw>) {
if (pending_) {
if (outputs_have_space(std::make_index_sequence<output_count>{}))
try_submit(0.5f);
return; // parked
}
}
// Source nodes always resubmit; others resubmit only if inputs are ready. // Source nodes always resubmit; others resubmit only if inputs are ready.
if constexpr (input_count == 0) { if constexpr (input_count == 0) {
try_submit(0.5f); try_submit(0.5f);
@@ -559,41 +393,36 @@ private:
} }
template<std::size_t... Is> template<std::size_t... Is>
/// Pushes what it can and parks the rest. `done_` marks the elements that
/// were taken, so a retry never re-pushes one — a duplicate would be as
/// wrong as a drop, just harder to notice.
void push_outputs(return_tuple&& result, std::index_sequence<Is...>) { void push_outputs(return_tuple&& result, std::index_sequence<Is...>) {
bool all = true; (push_one_out<Is>(std::get<Is>(std::move(result))), ...);
((pending_done_[Is] = pending_done_[Is] ||
push_one_out<Is>(std::get<Is>(std::move(result))),
all = all && pending_done_[Is]), ...);
if (all) { pending_.reset(); pending_done_.fill(false); }
else pending_ = std::move(result);
} }
/// Returns false when the ring was full and the value was NOT taken; the
/// caller must keep it and retry after the channel signals space.
template<std::size_t I> template<std::size_t I>
bool push_one_out(std::tuple_element_t<I, return_tuple>&& val) { void push_one_out(std::tuple_element_t<I, return_tuple>&& val) {
auto* ch = std::get<I>(output_channels_); auto* ch = std::get<I>(output_channels_);
if (!ch) return true; if (!ch) return;
// Sentinels (EOF) must never be dropped: a lost token wedges every // Sentinels (EOF) must never be dropped: a lost token wedges every
// 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)) {
ch->push_sentinel(std::move(val)); ch->push_sentinel(std::move(val));
return true; return;
}
// Lossless mode: block until the consumer drains instead of dropping.
// Dropping is the right default for live sources (a stale frame is
// worth less than a fresh one), but for offline batch work every sample
// matters — dropped frames leave a non-uniformly sampled series, which
// silently invalidates any fixed-rate spectral analysis downstream.
if (lossless_) {
ch->push_blocking(std::move(val));
return;
}
try {
ch->push(std::move(val));
} catch (const ChannelOverflowError&) {
throw ChannelOverflowError(ch->capacity(),
"pool node '" + name_ + "' " + output_port_label<I>());
} }
// Backpressure without parking the worker. A full channel means the
// consumer is behind; the value is kept and this node stops running
// until the channel signals space (set_space_callback re-submits it).
//
// Blocking here instead would sleep inside a scheduler worker, and
// nodes are pinned to workers — park enough of them and nothing is left
// to run the consumer that would drain the channel. That is the
// hold-and-wait deadlock channel.hpp warns about for sentinels; it
// applies to data pushes just as much.
return ch->try_push(val);
} }
template<std::size_t I> template<std::size_t I>
@@ -610,25 +439,14 @@ private:
std::shared_ptr<IScheduler> scheduler_; std::shared_ptr<IScheduler> scheduler_;
std::string name_; std::string name_;
bool lossless_{false};
std::size_t fifo_capacity_; std::size_t fifo_capacity_;
input_channels_t input_channels_; input_channels_t input_channels_;
output_channels_t output_channels_{}; output_channels_t output_channels_{};
std::atomic<bool> stop_flag_{true}; std::atomic<bool> stop_flag_{true};
std::atomic<bool> queued_{false}; std::atomic<bool> queued_{false};
/// A wake that arrived while queued_ was up. See try_submit.
std::atomic<bool> wake_pending_{false};
/// The hidden one-slot output buffer (see push_outputs). Holding the value
/// here is what lets a node stop running without dropping it or occupying a
/// scheduler worker. One slot suffices because at most one fire_once() runs
/// per node at a time — with concurrent firing it would have to hold a whole
/// FIFO's worth.
std::conditional_t<std::is_void_v<return_raw>, std::monostate,
std::optional<return_tuple>> pending_{};
std::array<bool, (output_count ? output_count : 1)> pending_done_{};
NodeStats stats_; NodeStats stats_;
NodeErrorHandler error_handler_; NodeErrorHandler error_handler_;
NodeErrorHandler net_error_handler_;
std::chrono::milliseconds max_exec_time_{0}; std::chrono::milliseconds max_exec_time_{0};
std::array<NodeEventCallback, 2> event_callbacks_{}; // [0]=user [1]=network std::array<NodeEventCallback, 2> event_callbacks_{}; // [0]=user [1]=network
std::array<NodeEventCallback, 2> closed_callbacks_{}; std::array<NodeEventCallback, 2> closed_callbacks_{};
@@ -696,9 +514,16 @@ public:
bool running() const override { return !stop_flag_.load(std::memory_order_relaxed); } bool running() const override { return !stop_flag_.load(std::memory_order_relaxed); }
void set_name(std::string name) override { name_ = std::move(name); } void set_name(std::string name) override { name_ = std::move(name); }
void set_error_handler(NodeErrorHandler h) { error_handler_ = std::move(h); } void set_error_handler(NodeErrorHandler h) { error_handler_ = std::move(h); }
void set_network_error_callback(NodeErrorHandler h) override { net_error_handler_ = std::move(h); }
void set_max_exec_time(std::chrono::milliseconds t) { max_exec_time_ = t; } void set_max_exec_time(std::chrono::milliseconds t) { max_exec_time_ = t; }
// Lossless output: block the producer until the consumer drains rather than
// dropping on a full channel. Default is drop, which suits live sources
// where a stale frame is worth less than a fresh one. Enable for offline
// batch runs, where a dropped item leaves a gap that downstream analysis
// cannot recover. Must be set before start().
void set_lossless_output(bool on) override { lossless_ = on; }
void set_lossless(bool on) { set_lossless_output(on); }
void set_overflow_callback(NodeEventCallback cb) { event_callbacks_[0] = std::move(cb); } void set_overflow_callback(NodeEventCallback cb) { event_callbacks_[0] = std::move(cb); }
void set_network_overflow_callback(NodeEventCallback cb) override { event_callbacks_[1] = std::move(cb); } void set_network_overflow_callback(NodeEventCallback cb) override { event_callbacks_[1] = std::move(cb); }
void set_closed_callback(NodeEventCallback cb) { closed_callbacks_[0] = std::move(cb); } void set_closed_callback(NodeEventCallback cb) { closed_callbacks_[0] = std::move(cb); }
@@ -768,8 +593,6 @@ private:
} }
template<std::size_t... Is> template<std::size_t... Is>
void register_callbacks(std::index_sequence<Is...>) { void register_callbacks(std::index_sequence<Is...>) {
// A parked producer is re-submitted when its output drains.
register_space_callbacks(std::make_index_sequence<output_count>{});
(std::get<Is>(input_channels_)->set_push_callback([this] { on_input_ready(); }), ...); (std::get<Is>(input_channels_)->set_push_callback([this] { on_input_ready(); }), ...);
} }
@@ -778,26 +601,10 @@ private:
for (auto& cb : cbs) if (cb) cb(ts); for (auto& cb : cbs) if (cb) cb(ts);
} }
template<std::size_t... Os>
bool outputs_have_space(std::index_sequence<Os...>) const {
return (... && (!std::get<Os>(output_channels_) ||
std::get<Os>(output_channels_)->has_space()));
}
template<std::size_t... Os>
void register_space_callbacks(std::index_sequence<Os...>) {
((std::get<Os>(output_channels_)
? (void)std::get<Os>(output_channels_)->set_space_callback(
[this] { try_submit(0.5f); })
: (void)0), ...);
}
void self_stop() { void self_stop() {
disable_inputs(std::make_index_sequence<input_count>{}); disable_inputs(std::make_index_sequence<input_count>{});
disable_outputs(std::make_index_sequence<output_count>{}); disable_outputs(std::make_index_sequence<output_count>{});
stats_.exec_start_us.store(0, std::memory_order_relaxed); stats_.exec_start_us.store(0, std::memory_order_relaxed);
// Plain store, not release_and_recheck(): this node is stopping, and
// honouring a pending wake here would resubmit a dead node.
queued_.store(false, std::memory_order_release); queued_.store(false, std::memory_order_release);
stop_flag_.store(true, std::memory_order_relaxed); stop_flag_.store(true, std::memory_order_relaxed);
} }
@@ -825,51 +632,11 @@ private:
return ((std::get<Is>(input_channels_)->approx_size() > 0 ? 1u : 0u) + ...); return ((std::get<Is>(input_channels_)->approx_size() > 0 ? 1u : 0u) + ...);
} }
/// Priority in [0,1], higher runs sooner.
///
/// Input fill alone answers "how much work is waiting for me". That is
/// only half the question: a node whose *outputs* are already full cannot
/// deliver anything: running it produces a value with nowhere to go, so it
/// immediately parks and the slot is wasted. Meanwhile the node that would
/// have drained that full channel waits behind it.
///
/// So occupancy of the outputs is deducted from occupancy of the inputs.
/// The scheduler then naturally favours whoever is furthest downstream of
/// a bottleneck — the node whose inputs are backed up but whose outputs
/// have room is exactly the one whose execution frees the most capacity —
/// and defers producers that would only deepen a queue that is already
/// full.
///
/// Mapped as 0.5·(1 + in - out) rather than clamping (in - out) at zero:
/// both terms are mean fills in [0,1], so the difference is in [-1,1], and
/// the affine map keeps the whole range distinguishable instead of
/// collapsing every output-saturated node onto the same value. 0.5 remains
/// the neutral point, matching the default used for source nodes.
float compute_priority() { float compute_priority() {
if constexpr (input_count == 0) return 0.5f; if constexpr (input_count == 0) return 0.5f;
float in = 0.0f; float sum = 0.0f;
sum_fill(in, std::make_index_sequence<input_count>{}); sum_fill(sum, std::make_index_sequence<input_count>{});
in /= static_cast<float>(input_count); return sum / static_cast<float>(input_count);
if constexpr (output_count == 0) return in;
float out = 0.0f;
sum_output_fill(out, std::make_index_sequence<output_count>{});
out /= static_cast<float>(output_count);
const float p = 0.5f * (1.0f + in - out);
return p < 0.0f ? 0.0f : (p > 1.0f ? 1.0f : p);
}
/// Mean fill of the output channels, same normalisation as sum_fill.
/// An unconnected output holds nothing back, so it contributes 0.
template<std::size_t... Os>
void sum_output_fill(float& sum, std::index_sequence<Os...>) {
((sum += (std::get<Os>(output_channels_) &&
std::get<Os>(output_channels_)->capacity() > 0)
? float(std::get<Os>(output_channels_)->approx_size())
/ float(std::get<Os>(output_channels_)->capacity())
: 0.0f), ...);
} }
template<std::size_t... Is> template<std::size_t... Is>
void sum_fill(float& sum, std::index_sequence<Is...>) { void sum_fill(float& sum, std::index_sequence<Is...>) {
@@ -879,32 +646,10 @@ private:
: 0.5f), ...); : 0.5f), ...);
} }
/// Submit unless already queued. A wake that arrives while this node is
/// queued or running is *recorded*, never dropped.
///
/// Wakes are edge-triggered: a channel fires its space callback on the
/// transition, once. If that lands while queued_ is up, the CAS below fails
/// and — before wake_pending_ — the wake was gone. A node could then park a
/// value, release its worker, and sleep forever holding output its consumer
/// was waiting for, with every worker idle in cond_wait and nothing left to
/// re-trigger it. Recording the drop turns the signal level-triggered: the
/// invariant is that a node never sleeps with a wake outstanding, enforced
/// by release_and_recheck() at every point that releases the node.
void try_submit(float priority) { void try_submit(float priority) {
bool expected = false; bool expected = false;
if (queued_.compare_exchange_strong(expected, true, std::memory_order_acq_rel)) if (queued_.compare_exchange_strong(expected, true, std::memory_order_acq_rel))
scheduler_->submit([this] { fire_once(); }, priority); scheduler_->submit([this] { fire_once(); }, priority);
else
wake_pending_.store(true, std::memory_order_release);
}
/// Clear queued_, then honour any wake that was dropped while it was up.
/// Every path that finishes or parks a firing must release the node through
/// here rather than storing queued_ directly.
void release_and_recheck(float priority = 0.5f) {
queued_.store(false, std::memory_order_release);
if (wake_pending_.exchange(false, std::memory_order_acq_rel))
try_submit(priority);
} }
void fire_once() { void fire_once() {
@@ -917,46 +662,6 @@ private:
t0.time_since_epoch()).count(); t0.time_since_epoch()).count();
stats_.exec_start_us.store(now_us, std::memory_order_relaxed); stats_.exec_start_us.store(now_us, std::memory_order_relaxed);
// Parked from a previous firing: retry that value before touching the
// inputs. Returning here releases the worker — the channel's space
// callback re-submits this node once the consumer drains a slot.
if constexpr (!std::is_void_v<return_raw>) {
if (pending_) {
push_outputs(std::move(*pending_), std::make_index_sequence<output_count>{});
release_and_recheck();
if (pending_) {
// Close the lost-wakeup race: a space_callback that fired
// between the failed push and clearing queued_ was
// swallowed, and nothing else will wake this node. Re-check
// now that the flag is down.
if (outputs_have_space(std::make_index_sequence<output_count>{}))
try_submit(0.5f);
return; // parked
}
// Drained: resume normal firing, resubmitting exactly the way
// the normal tail below does. An unconditional try_submit here
// would fire a node whose inputs are empty, and pop_inputs
// reports an empty channel as ChannelClosedError — which this
// node treats as "upstream finished" and self-stops on. That
// is a live node killing itself purely because it was woken by
// *output* space rather than by input arrival.
if constexpr (input_count == 0) try_submit(0.5f);
else on_input_ready();
return;
}
}
// See the equivalent guard in PoolNode::fire_once: a space-callback
// wake with nothing parked must release the worker, not fall through
// into pop_inputs on an empty channel.
if constexpr (input_count > 0) {
if (count_ready(std::make_index_sequence<input_count>{}) != input_count) {
release_and_recheck();
on_input_ready();
return;
}
}
try { try {
auto args = pop_inputs(std::make_index_sequence<input_count>{}); auto args = pop_inputs(std::make_index_sequence<input_count>{});
auto t1 = clock_t::now(); auto t1 = clock_t::now();
@@ -980,11 +685,7 @@ private:
} catch (const ChannelOverflowError&) { } catch (const ChannelOverflowError&) {
fire_callbacks(event_callbacks_); fire_callbacks(event_callbacks_);
} catch (...) { } catch (...) {
auto eptr = std::current_exception(); if (error_handler_ && error_handler_(name_, std::current_exception())) {
const bool handled =
(error_handler_ && error_handler_(name_, eptr)) ||
(net_error_handler_ && net_error_handler_(name_, eptr));
if (handled) {
} else { } else {
fire_callbacks(closed_callbacks_); fire_callbacks(closed_callbacks_);
self_stop(); self_stop();
@@ -993,26 +694,8 @@ private:
} }
stats_.exec_start_us.store(0, std::memory_order_relaxed); stats_.exec_start_us.store(0, std::memory_order_relaxed);
release_and_recheck(); queued_.store(false, std::memory_order_release);
if (stop_flag_.load(std::memory_order_relaxed)) return; if (stop_flag_.load(std::memory_order_relaxed)) return;
// Parked by the push above. Same situation as the retry path at the top
// of fire_once — and the same lost-wakeup race, which that path closes
// and this one did not. A space callback that fired while queued_ was
// still up got swallowed by try_submit's CAS, and the resubmit below
// cannot cover it: this firing consumed its input, so inputs are empty
// and on_input_ready() will not resubmit. The node would then hold its
// value forever while its consumer waits for exactly that value and its
// producer parks on an input channel that never drains. Re-check now
// that the flag is down.
if constexpr (!std::is_void_v<return_raw>) {
if (pending_) {
if (outputs_have_space(std::make_index_sequence<output_count>{}))
try_submit(0.5f);
return; // parked
}
}
if constexpr (input_count == 0) try_submit(0.5f); if constexpr (input_count == 0) try_submit(0.5f);
else on_input_ready(); else on_input_ready();
} }
@@ -1035,56 +718,44 @@ private:
} }
template<std::size_t... Is> template<std::size_t... Is>
/// Pushes what it can and parks the rest. `pending_done_` marks the elements
/// already taken, so a retry never re-pushes one — a duplicate is as wrong as
/// a drop and harder to notice.
void push_outputs(return_tuple&& result, std::index_sequence<Is...>) { void push_outputs(return_tuple&& result, std::index_sequence<Is...>) {
bool all = true; (push_one_out<Is>(std::get<Is>(std::move(result))), ...);
((pending_done_[Is] = pending_done_[Is] ||
push_one_out<Is>(std::get<Is>(std::move(result))),
all = all && pending_done_[Is]), ...);
if (all) { pending_.reset(); pending_done_.fill(false); }
else pending_ = std::move(result);
} }
/// Returns false when the ring was full and the value was NOT taken; the
/// caller must keep it and retry after the channel signals space.
template<std::size_t I> template<std::size_t I>
bool push_one_out(std::tuple_element_t<I, return_tuple>&& val) { void push_one_out(std::tuple_element_t<I, return_tuple>&& val) {
auto* ch = std::get<I>(output_channels_); auto* ch = std::get<I>(output_channels_);
if (!ch) return true; if (!ch) return;
// Sentinels (EOF) must never be dropped: a lost token wedges every // Sentinels (EOF) must never be dropped: a lost token wedges every
// 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)) {
ch->push_sentinel(std::move(val)); ch->push_sentinel(std::move(val));
return true; return;
}
// Lossless mode: block until the consumer drains instead of dropping.
if (lossless_) {
ch->push_blocking(std::move(val));
return;
}
try {
ch->push(std::move(val));
} catch (const ChannelOverflowError&) {
throw ChannelOverflowError(ch->capacity(),
"pool node '" + name_ + "'");
} }
// See the note on the typed overload above: park rather than block.
return ch->try_push(val);
} }
Obj& obj_; Obj& obj_;
std::shared_ptr<IScheduler> scheduler_; std::shared_ptr<IScheduler> scheduler_;
std::string name_; std::string name_;
bool lossless_{false};
std::size_t fifo_capacity_; std::size_t fifo_capacity_;
input_channels_t input_channels_; input_channels_t input_channels_;
output_channels_t output_channels_{}; output_channels_t output_channels_{};
std::atomic<bool> stop_flag_{true}; std::atomic<bool> stop_flag_{true};
std::atomic<bool> queued_{false}; std::atomic<bool> queued_{false};
/// A wake that arrived while queued_ was up. See try_submit.
std::atomic<bool> wake_pending_{false};
/// The hidden one-slot output buffer (see push_outputs). Holding the value
/// here is what lets a node stop running without dropping it or occupying a
/// scheduler worker. One slot suffices because at most one fire_once() runs
/// per node at a time — with concurrent firing it would have to hold a whole
/// FIFO's worth.
std::conditional_t<std::is_void_v<return_raw>, std::monostate,
std::optional<return_tuple>> pending_{};
std::array<bool, (output_count ? output_count : 1)> pending_done_{};
NodeStats stats_; NodeStats stats_;
NodeErrorHandler error_handler_; NodeErrorHandler error_handler_;
NodeErrorHandler net_error_handler_;
std::chrono::milliseconds max_exec_time_{0}; std::chrono::milliseconds max_exec_time_{0};
std::array<NodeEventCallback, 2> event_callbacks_{}; // [0]=user [1]=network std::array<NodeEventCallback, 2> event_callbacks_{}; // [0]=user [1]=network
std::array<NodeEventCallback, 2> closed_callbacks_{}; std::array<NodeEventCallback, 2> closed_callbacks_{};
+18 -13
View File
@@ -122,10 +122,6 @@ public:
[this, n](auto ts) { event_handler_(n, NodeEvent::Closed, ts); }); [this, n](auto ts) { event_handler_(n, NodeEvent::Closed, ts); });
} }
} }
if (error_handler_) {
for (auto* node : user_nodes_topo_)
node->set_network_error_callback(error_handler_);
}
for (auto* n : user_nodes_topo_) n->start(); for (auto* n : user_nodes_topo_) n->start();
for (auto* n : fanout_nodes_ptr_) n->start(); for (auto* n : fanout_nodes_ptr_) n->start();
#ifdef KPN_WEB_DEBUG #ifdef KPN_WEB_DEBUG
@@ -185,14 +181,6 @@ public:
void set_event_handler(EventHandler h) { event_handler_ = std::move(h); } void set_event_handler(EventHandler h) { event_handler_ = std::move(h); }
/// Application-level error listener. Receives the exception any node's
/// function throws, after that node's own handler (if any) declined it.
/// Return true to skip the failed invocation and keep the node running,
/// false to let it stop. Without a listener the exception is discarded
/// and only a Closed event survives, which reports that a node stopped
/// but not why.
void set_error_handler(NodeErrorHandler h) { error_handler_ = std::move(h); }
#ifdef KPN_WEB_DEBUG #ifdef KPN_WEB_DEBUG
void set_web_debug_port(uint16_t port) { web_debug_port_ = port; } void set_web_debug_port(uint16_t port) { web_debug_port_ = port; }
// Called by DebugHub::register_network() so the hub owns the debug server. // Called by DebugHub::register_network() so the hub owns the debug server.
@@ -231,6 +219,24 @@ public:
FanoutStorage& fanouts_storage() { return *fanouts_; } FanoutStorage& fanouts_storage() { return *fanouts_; }
// Make every node in this network push losslessly (block until the consumer
// drains) instead of dropping on a full channel. Includes the fanout nodes
// make_network() inserts automatically, which is the part user code cannot
// reach: they are unnamed, and they drop silently per-output, so a network
// whose own nodes are all lossless can still lose items at a fanout.
//
// Only safe when every consumer eventually drains. A branch that can stall
// indefinitely — a display node nobody is servicing, say — will block the
// whole pipeline through backpressure. Call before start().
void set_lossless(bool on = true) {
for (auto* n : user_nodes_topo_) if (n) n->set_lossless_output(on);
for (auto* n : fanout_nodes_ptr_) if (n) n->set_lossless_output(on);
}
// Block until every channel is empty. Useful before stop() so work already
// in flight completes rather than being discarded at teardown.
void drain() const { drain_all_channels(); }
private: private:
struct Snapshots { struct Snapshots {
std::vector<NodeSnapshot> nodes; std::vector<NodeSnapshot> nodes;
@@ -288,7 +294,6 @@ private:
std::vector<std::pair<std::string, IResourceProbe*>> resource_probes_; std::vector<std::pair<std::string, IResourceProbe*>> resource_probes_;
std::vector<std::pair<std::string, IPoolProbe*>> pool_probes_; std::vector<std::pair<std::string, IPoolProbe*>> pool_probes_;
EventHandler event_handler_; EventHandler event_handler_;
NodeErrorHandler error_handler_;
clock_t::time_point start_time_; clock_t::time_point start_time_;
#ifdef KPN_WEB_DEBUG #ifdef KPN_WEB_DEBUG
uint16_t web_debug_port_{9090}; uint16_t web_debug_port_{9090};
-1
View File
@@ -34,7 +34,6 @@ add_executable(kpn_tests
test_static_network.cpp test_static_network.cpp
test_shared_resource.cpp test_shared_resource.cpp
test_pool_node.cpp test_pool_node.cpp
test_backpressure_deadlock.cpp
test_scheduler.cpp test_scheduler.cpp
) )
-171
View File
@@ -1,171 +0,0 @@
// Regression: a blocking push must not park a pool worker.
//
// Node outputs use push_blocking so a full channel costs time rather than data
// (a dropped frame does not degrade a downstream result, it silently changes
// one). But push_blocking sleeps *inside* fire_once, which runs on a pool
// worker — and nodes are pinned to workers by index. Park enough workers in
// that retry loop and there is nobody left to run the consumer that would drain
// the channel, so the whole chain wedges.
//
// This is the failure channel.hpp:174 already warns about for sentinels
// ("a blocking push would park that thread and stop it draining its own input,
// cascading into a hold-and-wait deadlock under backpressure"). The warning
// applies to data pushes too.
//
// Observed in the field as an intermittent hang: frame_source, camera_pos,
// face_detector and face_aligner all asleep in push_blocking at once.
#include <catch2/catch_test_macros.hpp>
#include <kpn/kpn.hpp>
#include <atomic>
#include <chrono>
#include <thread>
namespace {
struct Produce {
static constexpr std::string_view label() { return "produce"; }
int n{0};
int operator()() { return n++; }
};
struct Relay {
static constexpr std::string_view label() { return "relay"; }
int operator()(int v) { return v; }
};
// Deliberately slower than the producer, so the channels between them fill.
struct SlowSink {
static constexpr std::string_view label() { return "slow_sink"; }
std::atomic<int>* seen;
void operator()(int) {
std::this_thread::sleep_for(std::chrono::milliseconds(2));
seen->fetch_add(1, std::memory_order_relaxed);
}
};
} // namespace
TEST_CASE("a saturated chain keeps draining", "[backpressure][deadlock]") {
std::atomic<int> seen{0};
Produce p_fn;
Relay r1_fn, r2_fn, r3_fn;
SlowSink s_fn{&seen};
// Small channels so they saturate immediately, and a chain longer than a
// modest pool — the shape that starves workers.
kpn::ObjectNode<Produce, kpn::in<>, kpn::out<"a">, "produce", 0> p (p_fn, 2);
kpn::ObjectNode<Relay, kpn::in<"a">, kpn::out<"b">, "relay1", 0> r1(r1_fn, 2);
kpn::ObjectNode<Relay, kpn::in<"b">, kpn::out<"c">, "relay2", 0> r2(r2_fn, 2);
kpn::ObjectNode<Relay, kpn::in<"c">, kpn::out<"d">, "relay3", 0> r3(r3_fn, 2);
kpn::ObjectNode<SlowSink, kpn::in<"d">, kpn::out<>, "slow_sink", 0> s (s_fn, 2);
auto net = kpn::make_network(
kpn::edge(p.output<"a">(), r1.input<"a">()),
kpn::edge(r1.output<"b">(), r2.input<"b">()),
kpn::edge(r2.output<"c">(), r3.input<"c">()),
kpn::edge(r3.output<"d">(), s.input<"d">())
);
net.start();
// The sink is the slowest stage at 2 ms/item, so 40 items is ~80 ms of real
// work. Anything approaching the timeout means the chain stopped draining
// rather than merely running slowly.
const auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(20);
while (seen.load(std::memory_order_relaxed) < 40 &&
std::chrono::steady_clock::now() < deadline)
std::this_thread::sleep_for(std::chrono::milliseconds(5));
const int got = seen.load(std::memory_order_relaxed);
net.stop();
INFO("items drained: " << got << " of 40");
CHECK(got >= 40);
}
// Regression: a saturated chain must never stall.
//
// push_outputs parks from two places: the retry at the top of fire_once, and
// the ordinary push after the node function returns. Both release the worker,
// so both face the same lost wakeup — a space callback firing while queued_ is
// still up is swallowed by try_submit's CAS. Only the retry path re-checked for
// space afterwards. The normal path fell through to on_input_ready(), which
// resubmits only if inputs are ready — and the firing that just parked had
// consumed its input, so they are not.
//
// The strand is permanent under saturation: the node holds its value, its
// consumer waits for exactly that value, and its producer fills the node's
// input channel and parks too. Nothing moves again.
//
// The test above cannot catch it — 40 items drain before any strand occurs.
// This one runs the chain saturated and watches for progress to *freeze*, which
// is the signature of the deadlock. It deliberately does not assert a total:
// capacity-1 channels are slow, and "slow" must never be reported as "wedged".
namespace {
struct FreeRun {
static constexpr std::string_view label() { return "free_run"; }
int n{0};
int operator()() { return n++; }
};
struct CountingSink {
static constexpr std::string_view label() { return "counting_sink"; }
std::atomic<int>* seen;
void operator()(int) { seen->fetch_add(1, std::memory_order_relaxed); }
};
} // namespace
TEST_CASE("a saturated chain never stalls", "[backpressure][deadlock]") {
std::atomic<int> seen{0};
FreeRun p_fn;
Relay r1_fn, r2_fn;
CountingSink s_fn{&seen};
// Capacity 1 everywhere: every push contends, so the park path is taken
// constantly and the race window is sampled millions of times.
kpn::ObjectNode<FreeRun, kpn::in<>, kpn::out<"a">, "free_run", 0> p (p_fn, 1);
kpn::ObjectNode<Relay, kpn::in<"a">, kpn::out<"b">, "relay1", 0> r1(r1_fn, 1);
kpn::ObjectNode<Relay, kpn::in<"b">, kpn::out<"c">, "relay2", 0> r2(r2_fn, 1);
kpn::ObjectNode<CountingSink, kpn::in<"c">, kpn::out<>, "sink", 0> s (s_fn, 1);
auto net = kpn::make_network(
kpn::edge(p.output<"a">(), r1.input<"a">()),
kpn::edge(r1.output<"b">(), r2.input<"b">()),
kpn::edge(r2.output<"c">(), s.input<"c">())
);
net.start();
// A live chain moves thousands of items a second, so 3 s with no movement
// at all is a wedge, not a slow patch. Sampling for 25 s gives the race
// ample opportunity: the pipeline hit it roughly twice in 30 runs.
const auto giveup = std::chrono::steady_clock::now() + std::chrono::seconds(25);
int last = 0;
auto last_move = std::chrono::steady_clock::now();
bool stalled = false;
int stall_at = 0;
while (std::chrono::steady_clock::now() < giveup) {
std::this_thread::sleep_for(std::chrono::milliseconds(100));
const int now_seen = seen.load(std::memory_order_relaxed);
if (now_seen != last) {
last = now_seen;
last_move = std::chrono::steady_clock::now();
} else if (std::chrono::steady_clock::now() - last_move >
std::chrono::seconds(3)) {
stalled = true;
stall_at = now_seen;
break;
}
}
net.stop();
INFO("chain stalled after " << stall_at << " items");
CHECK_FALSE(stalled);
// Guard against the test passing because nothing ever ran.
CHECK(last > 1000);
}
+7 -23
View File
@@ -243,13 +243,7 @@ TEST_CASE("interrupt node: trigger after stop is ignored", "[interrupt_node]") {
// ── Overflow callback ───────────────────────────────────────────────────────── // ── Overflow callback ─────────────────────────────────────────────────────────
// NOTE: output overflow is no longer reachable on the data path. A node whose TEST_CASE("pool node overflow callback fires on full output channel", "[pool_node][overflow]") {
// output channel is full now PARKS — it keeps the value in a hidden one-slot
// buffer, releases its scheduler worker, and is re-submitted when the consumer
// frees a slot. The overflow callback survives for other producers (a direct
// Channel::push by non-node code still throws), but a pool node cannot trigger
// it, so these cases assert the stronger property instead: nothing is dropped.
TEST_CASE("pool node parks instead of overflowing a full output", "[pool_node][overflow]") {
auto pool = std::make_shared<ThreadPool>(2); auto pool = std::make_shared<ThreadPool>(2);
pool->start(); pool->start();
@@ -270,13 +264,10 @@ TEST_CASE("pool node parks instead of overflowing a full output", "[pool_node][o
node.stop(); node.stop();
pool->stop(); pool->stop();
// Parked, not overflowed: the value is still owned by the node. REQUIRE(overflow_count.load() > 0);
REQUIRE(overflow_count.load() == 0);
// And it was never handed downstream, so nothing was lost or duplicated.
REQUIRE(full_ch.size() == 1);
} }
TEST_CASE("parking is per node, not shared", "[pool_node][overflow]") { TEST_CASE("pool node overflow callback is independent per instance", "[pool_node][overflow]") {
auto pool = std::make_shared<ThreadPool>(2); auto pool = std::make_shared<ThreadPool>(2);
pool->start(); pool->start();
@@ -305,10 +296,7 @@ TEST_CASE("parking is per node, not shared", "[pool_node][overflow]") {
nodeB.stop(); nodeB.stop();
pool->stop(); pool->stop();
// Neither overflows now: A parks on its full output, B runs normally. The REQUIRE(a_overflows.load() > 0);
// point of the case is unchanged — one node's backpressure must not leak
// into another's callbacks.
REQUIRE(a_overflows.load() == 0);
REQUIRE(b_overflows.load() == 0); REQUIRE(b_overflows.load() == 0);
} }
@@ -424,9 +412,7 @@ TEST_CASE("network_overflow_callback fires on overflow", "[pool_node][network]")
node.stop(); node.stop();
pool->stop(); pool->stop();
// Parking replaced overflow on the data path, so the network callback no REQUIRE(net_overflows.load() > 0);
// longer fires for a pool node's own output. See the note above.
REQUIRE(net_overflows.load() == 0);
} }
TEST_CASE("network_closed_callback fires on crash", "[pool_node][network]") { TEST_CASE("network_closed_callback fires on crash", "[pool_node][network]") {
@@ -473,8 +459,6 @@ TEST_CASE("per-node and network overflow callbacks both fire independently", "[p
node.stop(); node.stop();
pool->stop(); pool->stop();
// Both zero now: the node parks rather than overflowing. The case still REQUIRE(per_node.load() > 0);
// guards that the two callbacks are wired independently. REQUIRE(network.load() > 0);
REQUIRE(per_node.load() == 0);
REQUIRE(network.load() == 0);
} }