Author SHA1 Message Date
dtourolleandClaude Opus 5 6595e6e925 fix: node outputs block instead of dropping on a full channel
🚦 CI / changes (push) Successful in 30s
🚦 CI / docker (push) Has been skipped
🚦 CI / test (push) Failing after 1h44m2s
🚦 CI / docs (push) Has been skipped
🚦 CI / tsan (push) Failing after 3h14m57s
Every node output used the throwing push(), so a consumer falling behind cost
values rather than time. push_blocking() already existed on Channel and
OutputPort — "wait for the consumer to drain instead of dropping; the producer
just runs slower" — but nothing called it.

A dropped frame does not degrade a downstream result, it silently changes one,
and the consumer has no way to tell it happened. For any pipeline whose output
is a claim about its input, that is corruption rather than degradation.

Safe because sentinels are already handled out-of-band, above this path: only
data blocks, so the EOF token that unwinds the network can always overtake a
stalled data path. That is exactly the hold-and-wait deadlock the push_sentinel
comment warns about, and the reason it is not reachable here.

Measured on a downstream consumer (face pipeline, 77s clip at 5 fps, expected
385 sampled frames):

  before  65 frames written, 320 dropped at one node, 29s
  after   385 frames written, 0 dropped, 17s

Faster, not slower — a dropped frame has already cost its decode, and the
overflow exception cost more. Two consecutive runs now produce byte-identical
output, which they did not before: what got dropped depended on timing, so the
same command could yield different results.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-31 11:31:53 +02:00
dtourolle 75b34f31bb Merge pull request 'feat: persistent-pipeline reuse — push_blocking, node introspection, stateful wrapper' (#2) from feature/persistent-pipeline-reuse into master
🚦 CI / changes (push) Successful in 4s
🚦 CI / docker (push) Has been skipped
🚦 CI / test (push) Successful in 4m39s
🚦 CI / tsan (push) Successful in 2m57s
🚦 CI / docs (push) Has been skipped
Reviewed-on: #2
2026-07-19 16:11:27 +00:00
dtourolle 4b6e498ba7 feat: persistent-pipeline reuse — push_blocking, node introspection, stateful wrapper
🚦 CI / changes (pull_request) Successful in 17s
🚦 CI / docker (pull_request) Has been skipped
🚦 CI / test (pull_request) Successful in 4m42s
🚦 CI / tsan (pull_request) Successful in 2m56s
🚦 CI / docs (pull_request) Has been skipped
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.
2026-07-19 16:55:36 +02:00
dtourolle 5ecf3cde4f Merge pull request 'spec-and-tsan' (#1) from spec-and-tsan into master
🚦 CI / changes (push) Successful in 4s
🚦 CI / docker (push) Has been skipped
🚦 CI / test (push) Successful in 4m36s
🚦 CI / tsan (push) Successful in 2m54s
🚦 CI / docs (push) Has been skipped
Reviewed-on: #1
2026-07-17 18:11:51 +00:00
dtourolleandClaude Opus 4.8 ec19137ed9 ci: fix TSan aborting at init on the nested-LXC runner
🚦 CI / changes (pull_request) Successful in 6s
🚦 CI / docker (pull_request) Has been skipped
🚦 CI / test (pull_request) Successful in 4m32s
🚦 CI / tsan (pull_request) Successful in 2m54s
🚦 CI / docs (pull_request) Has been skipped
The ThreadSanitizer job runs on Docker nested in an unprivileged LXC
container, whose kernel randomizes mmap addresses beyond the range TSan's
fixed shadow mapping expects. TSan aborted at init with "unexpected memory
mapping" before any test ran.

Disable ASLR per-process with `setarch -R`, which needs the personality(2)
syscall that Docker's default seccomp profile blocks; seccomp=unconfined on
the container permits it. Verified on the runner that both are required:
setarch -R alone gets EPERM, seccomp alone still aborts, both together run
clean. Scoped to the tsan job, which runs only our own test binaries.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-17 20:01:57 +02:00
6 changed files with 228 additions and 17 deletions
+12 -2
View File
@@ -18,6 +18,15 @@ jobs:
runs-on: linux/amd64
container:
image: gitea.tourolle.paris/dtourolle/kpnpp-builder:latest
# This runner is Docker nested in an unprivileged LXC container, whose
# kernel randomizes mmap addresses beyond the range TSan's fixed shadow
# mapping expects, so TSan aborts at init with "unexpected memory
# mapping". The fix is to disable ASLR per-process with `setarch -R`
# (below), which needs the personality(2) syscall that Docker's default
# seccomp profile blocks. seccomp=unconfined permits it. Verified on the
# runner: setarch -R alone gets EPERM, seccomp alone still aborts, both
# together run clean. Scoped to this job, which runs only our own tests.
options: --security-opt seccomp=unconfined
steps:
- name: Checkout repository
uses: actions/checkout@v4
@@ -54,13 +63,14 @@ jobs:
# the full picture for lock-order issues.
env:
TSAN_OPTIONS: "halt_on_error=1 second_deadlock_stack=1"
run: ./build/tests/kpn_tests_stress
# setarch -R disables ASLR for this process; see the container comment.
run: setarch -R ./build/tests/kpn_tests_stress
- name: Run unit tests under TSan
working-directory: tsan-${{ github.run_id }}
env:
TSAN_OPTIONS: "halt_on_error=1 second_deadlock_stack=1"
run: ./build/tests/kpn_tests
run: setarch -R ./build/tests/kpn_tests
- name: Cleanup
if: always()
+29
View File
@@ -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
+11 -12
View File
@@ -400,12 +400,15 @@ private:
ch->push_sentinel(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, not loss. A full downstream channel means the consumer
// is behind, and the correct response is for this producer to run
// slower — not to discard a value. A dropped frame does not degrade a
// result, it silently changes one, and the caller has no way to tell.
//
// Safe here because sentinels are handled above, out-of-band: this
// blocks only on data, so the EOF token that unwinds the network can
// always overtake a stalled data path.
ch->push_blocking(std::move(val));
}
template<std::size_t I>
@@ -706,12 +709,8 @@ private:
ch->push_sentinel(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: block rather than drop.
ch->push_blocking(std::move(val));
}
Obj& obj_;
+22 -3
View File
@@ -217,6 +217,20 @@ public:
return it->second;
}
// Raw node handle by name — lets a binding dynamic_cast to a concrete wrapper
// type and call its functor's runtime setters (persistent-pipeline reuse).
VNode* node_ptr(const std::string& name) { return &node_at(name); }
// Per-node timing snapshot for profiling where a replay spends its time.
std::map<std::string, double> node_stats(const std::string& name) {
auto& n = node_at(name);
NodeSnapshot s = n.node_snapshot(name, 0.0);
return {{"frames", double(s.frames_processed)},
{"exec_ms", s.ema_exec_ms}, {"max_ms", s.max_exec_ms},
{"blocked_ms", s.total_blocked_ms}, {"fps", s.throughput_fps},
{"cpu_ms", s.total_cpu_ms}, {"cpu_util_pct", s.cpu_util_pct}};
}
private:
VNode& node_at(const std::string& name) {
auto it = nodes_.find(name);
@@ -422,13 +436,17 @@ private:
for (std::size_t i = 0; i < out_channels_.size(); ++i) {
if (out_channels_[i])
out_channels_[i]->push(std::move(outputs[i]));
// Lossless: wait for space rather than drop. A dropped frame
// silently corrupts a replay's score; backpressure just slows
// the producer. (Was push() + "drop on overflow".)
out_channels_[i]->push_blocking(std::move(outputs[i]));
}
} catch (const ChannelClosedError&) {
break;
} catch (const ChannelOverflowError&) {
// drop and continue
// no longer reachable with push_blocking, kept for safety
break;
}
}
}
@@ -530,7 +548,8 @@ void register_py_network(nb::module_& m, const char* class_name = "Network") {
.def("read", &Net::read,
nb::arg("node"), nb::arg("out_idx") = std::size_t(0))
.def("write", &Net::write,
nb::arg("node"), nb::arg("in_idx"), nb::arg("value"));
nb::arg("node"), nb::arg("in_idx"), nb::arg("value"))
.def("node_stats", &Net::node_stats, nb::arg("node"));
}
} // namespace kpn::python
+149
View File
@@ -0,0 +1,149 @@
#pragma once
// ObjectVariantNodeWrapper — variant-node adapter for *stateful* functors.
//
// VariantNodeWrapper (variant_node.hpp) wraps Node<Func,...>, where Func is a
// default-constructible NTTP callable. That doesn't fit nodes whose functor must
// be constructed with runtime state (a Config, a loaded gallery, etc.) — those use
// ObjectNode<Obj>, which takes `Obj& obj` at construction.
//
// This wrapper owns an Obj instance and exposes the same IVariantNode surface so a
// stateful C++ node can live inside a PyNetwork. Build one via a factory that
// constructs the functor from Python-supplied config, e.g.:
//
// auto n = std::make_shared<ObjectVariantNodeWrapper<
// IdentityMatcherFunc, Variant, in<"tracked">, out<"matched">>>(
// fifo_cap, gallery, cfg); // Obj ctor args forwarded
// net.add("identity_matcher", n);
//
// The wrapper mirrors VariantNodeWrapper's channel plumbing exactly; only the
// underlying node type (PoolObjectNode, holding Obj&) differs.
#include "../channel.hpp"
#include "../node.hpp"
#include "../variant_node.hpp"
#include <memory>
#include <stdexcept>
#include <string>
#include <tuple>
#include <typeindex>
#include <utility>
#include <vector>
namespace kpn {
template<typename Obj, typename Variant,
typename InputTag = in<>,
typename OutputTag = out<>>
class ObjectVariantNodeWrapper;
template<typename Obj, typename Variant,
fixed_string... InNames, fixed_string... OutNames>
class ObjectVariantNodeWrapper<Obj, Variant, in<InNames...>, out<OutNames...>>
: public IVariantNode<Variant>
{
using NodeT = ObjectNode<Obj, in<InNames...>, out<OutNames...>>;
public:
using args_tuple = typename NodeT::args_tuple;
using return_tuple = typename NodeT::return_tuple;
static constexpr std::size_t n_in = NodeT::input_count;
static constexpr std::size_t n_out = NodeT::output_count;
// Owns the functor; forwards remaining args to Obj's constructor.
template<typename... ObjArgs>
explicit ObjectVariantNodeWrapper(std::size_t fifo_capacity, ObjArgs&&... obj_args)
: obj_(std::forward<ObjArgs>(obj_args)...)
, node_(obj_, fifo_capacity)
, in_channels_(n_in)
, out_channels_(n_out)
, out_type_indices_(n_out, std::type_index(typeid(void)))
{
init_inputs(std::make_index_sequence<n_in>{}, fifo_capacity);
init_out_types(std::make_index_sequence<n_out>{});
}
// Access the owned functor so callers can invoke its runtime setters (e.g. to
// change a threshold on a persistent pipeline without rebuilding the node).
Obj& functor() { return obj_; }
// ── INode ─────────────────────────────────────────────────────────────────
void start() override { node_.start(); }
void stop() override { node_.stop(); }
bool running() const override { return node_.running(); }
const NodeStats& stats() const override { return node_.stats(); }
void set_name(std::string name) override { node_.set_name(std::move(name)); }
NodeSnapshot node_snapshot(const std::string& name, double elapsed_s) const override {
return node_.node_snapshot(name, elapsed_s);
}
// ── IVariantNode ──────────────────────────────────────────────────────────
std::size_t input_count() const override { return n_in; }
std::size_t output_count() const override { return n_out; }
std::type_index input_type(std::size_t i) const override {
return in_channels_[i]->type_index();
}
std::type_index output_type(std::size_t i) const override {
return out_type_indices_[i];
}
std::shared_ptr<IVariantChannel<Variant>> input_channel(std::size_t i) override {
return in_channels_[i];
}
void set_output_channel(std::size_t i,
std::shared_ptr<IVariantChannel<Variant>> ch) override {
set_output_impl(i, std::move(ch), std::make_index_sequence<n_out>{});
}
private:
template<std::size_t... Is>
void init_inputs(std::index_sequence<Is...>, std::size_t cap) {
((init_one_input<Is>(cap)), ...);
}
template<std::size_t I>
void init_one_input(std::size_t cap) {
using T = std::tuple_element_t<I, args_tuple>;
auto shared_ch = std::make_shared<Channel<T>>(cap);
node_.template set_input_channel<I>(shared_ch);
in_channels_[I] = std::make_shared<VariantChannel<T, Variant>>(std::move(shared_ch));
}
template<std::size_t... Is>
void init_out_types(std::index_sequence<Is...>) {
((out_type_indices_[Is] =
std::type_index(typeid(std::tuple_element_t<Is, return_tuple>))), ...);
}
template<std::size_t... Is>
void set_output_impl(std::size_t port,
std::shared_ptr<IVariantChannel<Variant>> ch,
std::index_sequence<Is...>) {
bool matched = false;
((Is == port && (set_output_at<Is>(std::move(ch)), matched = true)), ...);
if (!matched)
throw std::out_of_range("set_output_channel: port index out of range");
}
template<std::size_t I>
void set_output_at(std::shared_ptr<IVariantChannel<Variant>> ch) {
using T = std::tuple_element_t<I, return_tuple>;
auto* typed = dynamic_cast<VariantChannel<T, Variant>*>(ch.get());
if (!typed)
throw std::runtime_error(
"set_output_channel: type mismatch at output port " + std::to_string(I));
node_.template set_output_channel<I>(typed->raw_ptr());
out_channels_[I] = std::move(ch);
}
Obj obj_; // owned; node_ holds Obj& — declaration order keeps obj_ alive first
NodeT node_;
std::vector<std::shared_ptr<IVariantChannel<Variant>>> in_channels_;
std::vector<std::shared_ptr<IVariantChannel<Variant>>> out_channels_;
std::vector<std::type_index> out_type_indices_;
};
} // namespace kpn
+5
View File
@@ -55,6 +55,8 @@ class IVariantChannel {
public:
virtual ~IVariantChannel() = default;
virtual void push(Variant v) = 0;
// Lossless push with backpressure (waits instead of dropping when full).
virtual void push_blocking(Variant v) = 0;
virtual Variant pop() = 0;
virtual std::type_index type_index() const = 0;
virtual std::string type_name() const = 0;
@@ -76,6 +78,9 @@ public:
void push(Variant v) override {
channel_->push(std::get<T>(std::move(v)));
}
void push_blocking(Variant v) override {
channel_->push_blocking(std::get<T>(std::move(v)));
}
Variant pop() override {
return Variant{ channel_->pop() };
}