Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ac36159f37 | ||
|
|
4e81752838 | ||
|
|
75b34f31bb |
+25
-2
@@ -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};
|
||||||
|
|||||||
@@ -37,6 +37,10 @@ struct INode {
|
|||||||
// 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(); }
|
||||||
|
|||||||
@@ -133,6 +133,14 @@ public:
|
|||||||
void set_error_handler(NodeErrorHandler h) { error_handler_ = std::move(h); }
|
void set_error_handler(NodeErrorHandler h) { 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); }
|
||||||
@@ -400,6 +408,15 @@ private:
|
|||||||
ch->push_sentinel(std::move(val));
|
ch->push_sentinel(std::move(val));
|
||||||
return;
|
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 {
|
try {
|
||||||
ch->push(std::move(val));
|
ch->push(std::move(val));
|
||||||
} catch (const ChannelOverflowError&) {
|
} catch (const ChannelOverflowError&) {
|
||||||
@@ -422,6 +439,7 @@ 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_{};
|
||||||
@@ -498,6 +516,14 @@ public:
|
|||||||
void set_error_handler(NodeErrorHandler h) { error_handler_ = std::move(h); }
|
void set_error_handler(NodeErrorHandler h) { 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); }
|
||||||
@@ -706,6 +732,11 @@ private:
|
|||||||
ch->push_sentinel(std::move(val));
|
ch->push_sentinel(std::move(val));
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
// Lossless mode: block until the consumer drains instead of dropping.
|
||||||
|
if (lossless_) {
|
||||||
|
ch->push_blocking(std::move(val));
|
||||||
|
return;
|
||||||
|
}
|
||||||
try {
|
try {
|
||||||
ch->push(std::move(val));
|
ch->push(std::move(val));
|
||||||
} catch (const ChannelOverflowError&) {
|
} catch (const ChannelOverflowError&) {
|
||||||
@@ -717,6 +748,7 @@ private:
|
|||||||
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_{};
|
||||||
|
|||||||
@@ -219,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;
|
||||||
|
|||||||
Reference in New Issue
Block a user