feat: persistent-pipeline reuse — push_blocking, node introspection, stateful wrapper
Adds three pieces needed to build one KPN network and reuse it across many replays/configs instead of tearing down and rebuilding per run: - Channel<T>::push_blocking (+ IVariantChannel/VariantChannel forwarding): lossless backpressure push that waits for space instead of dropping when the ring is full. PyNode's run_loop now uses it so a downstream consumer lagging behind never silently drops a frame. - PyNetwork::node_ptr / node_stats: raw node handle by name (for a binding to dynamic_cast to a concrete wrapper and call functor-specific runtime setters) and a per-node timing snapshot for profiling. - ObjectVariantNodeWrapper: variant-node adapter for functors that need runtime-constructed state (a Config, a loaded gallery), mirroring VariantNodeWrapper's channel plumbing but backed by ObjectNode<Obj>. Built and used downstream in scene-actor-extraction's sae_kpn Python replay bindings for repeated threshold-sweep evaluation of the same pipeline.
This commit is contained in:
@@ -136,6 +136,35 @@ public:
|
||||
push_callback_();
|
||||
}
|
||||
|
||||
// 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.
|
||||
// Use when every value must be delivered (e.g. replaying a dump for scoring, where
|
||||
// a dropped frame silently corrupts the result). SPSC: only the sole producer may
|
||||
// call it. Returns false if the channel was disabled while waiting.
|
||||
bool push_blocking(T value) {
|
||||
for (;;) {
|
||||
if (!accepting_.load(std::memory_order_acquire)) {
|
||||
stats_.record_drop();
|
||||
return false;
|
||||
}
|
||||
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_) { // space available → normal push
|
||||
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;
|
||||
}
|
||||
// full: yield briefly and retry (consumer will drain)
|
||||
std::this_thread::sleep_for(std::chrono::microseconds(50));
|
||||
}
|
||||
}
|
||||
|
||||
// Lossless, non-blocking delivery for a must-deliver control token (EOF).
|
||||
//
|
||||
// A sentinel is stored out-of-band — in a dedicated slot that does NOT
|
||||
|
||||
Reference in New Issue
Block a user