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
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
dtourolleandClaude Opus 4.8 a0c4bf580e fix: deliver EOF sentinel only when the ring is freshly empty
🚦 CI / changes (pull_request) Successful in 4s
🚦 CI / docker (pull_request) Has been skipped
🚦 CI / test (pull_request) Successful in 7m30s
🚦 CI / tsan (pull_request) Failing after 2m21s
🚦 CI / docs (pull_request) Has been skipped
Channel<T>::pop() surfaced the out-of-band sentinel from its empty branch
using the tail_ snapshot taken at the top of the loop. Under contention the
producer can push more values *and* the sentinel in the window between that
snapshot and take_sentinel(), so pop() could return the sentinel while real
values still sat in the ring — the sentinel jumping ahead of values pushed
before it. No value was lost (a consumer that keeps draining still receives
them, and approx_size() keeps counting them so a PoolNode reschedules), but a
consumer treating the sentinel as a hard "last message" barrier would act on
EOF early.

Re-confirm emptiness against a fresh tail_ load before taking the sentinel.
Costs one acquire-load on the empty-ring path only; never runs in steady
state. The spin and post-spin takes already reload tail_ on the line above
them; try_pop_now() already reads tail_ fresh in the same branch — both were
correct and are unchanged.

The two sentinel stress cases now assert the strict "sentinel is last, after
every value" ordering (previously relaxed to avoid the flake this fixes).
Verified TSan-clean (2606 assertions, no data races) over repeated runs.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-14 23:14:45 +02:00
dtourolleandClaude Opus 4.8 3ac2242df1 docs: rewrite SPEC.md to match the implemented library
The spec had drifted far from the code. Key corrections:

- Execution model is reactive (PoolNode submits fire_once() to a
  ThreadPool when inputs are ready), not one blocking thread per node
- Channel<T> is a lock-free SPSC ring buffer (atomic wait/notify +
  spin-before-sleep), not a mutex+CV queue
- Remove latch<> ports (never implemented)
- NodeErrorHandler returns bool (skip vs stop); per-node
- Document new subsystems: scheduler, InterruptNode, Router/FilterNode,
  MainThreadNode, SharedResource, DebugHub, diagnostics/stats layer
- Update StaticNetwork, Python auto_bind layer, examples 01-16

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-14 23:14:45 +02:00
dtourolleandClaude Opus 4.8 399ee4cf9b test: add ThreadSanitizer verification for lock-free Channel<T>
Add a contended SPSC stress suite (tests/test_channel_stress.cpp) that
actually exercises the ring's memory-ordering pairing and spin/futex/
lost-wakeup logic, plus the CMake and CI plumbing to run it under TSan:

- KPN_SANITIZER cache var + kpn_sanitizer_flags() helper (no-op when unset)
- kpn_tests_stress executable, labelled "stress" for CTest
- reusable tsan.yaml workflow (gcc:14 builder image, already ships libtsan)
- ci.yaml gains a tsan job on the same code/dockerfile triggers as test

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-14 23:14:45 +02:00
14 changed files with 1284 additions and 1162 deletions
+7
View File
@@ -73,6 +73,13 @@ jobs:
if: ${{ !failure() && !cancelled() && (needs.changes.outputs.code == 'true' || needs.changes.outputs.dockerfile == 'true') }}
uses: ./.gitea/workflows/test.yaml
# ThreadSanitizer run for the lock-free Channel<T>. Same trigger conditions as
# test (code or image changed); runs in parallel with test.
tsan:
needs: [changes, docker]
if: ${{ !failure() && !cancelled() && (needs.changes.outputs.code == 'true' || needs.changes.outputs.dockerfile == 'true') }}
uses: ./.gitea/workflows/tsan.yaml
docs:
needs: [changes, docker]
if: ${{ !failure() && !cancelled() && github.ref == 'refs/heads/master' && (needs.changes.outputs.docs == 'true' || needs.changes.outputs.dockerfile == 'true') }}
+77
View File
@@ -0,0 +1,77 @@
name: '🧵 ThreadSanitizer'
# Reusable workflow: builds the channel stress suite with ThreadSanitizer and
# runs it. This is the dynamic half of verifying the lock-free SPSC Channel<T>
# (the static half is the CDSChecker model-check harness in verify/).
#
# Triggering and path filtering are owned by ci.yaml (the orchestrator), which
# calls this only when code changed. workflow_dispatch is kept for manual runs.
#
# Runs in the prebuilt builder image (gcc:14), which already ships libtsan — no
# package installs at job time.
on:
workflow_call:
workflow_dispatch:
jobs:
tsan:
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
with:
path: tsan-${{ github.run_id }}
- name: Cache FetchContent dependencies
uses: actions/cache@v3
with:
path: ~/.cmake/fetchcontent
key: cmake-fetchcontent-${{ hashFiles('**/CMakeLists.txt') }}
restore-keys: cmake-fetchcontent-
- name: Configure (TSan)
working-directory: tsan-${{ github.run_id }}
run: |
cmake -S . -B build \
-G Ninja \
-DCMAKE_BUILD_TYPE=Debug \
-DKPN_SANITIZER=thread \
-DKPN_BUILD_TESTS=ON \
-DKPN_BUILD_EXAMPLES=OFF \
-DKPN_BUILD_PYTHON=OFF \
-DFETCHCONTENT_BASE_DIR=$HOME/.cmake/fetchcontent
- name: Build (TSan)
working-directory: tsan-${{ github.run_id }}
run: cmake --build build --parallel --target kpn_tests kpn_tests_stress
- name: Run stress suite under TSan
working-directory: tsan-${{ github.run_id }}
# halt_on_error=1 makes the first detected race fail the job; the report
# (with both stacks) is printed to the log. second_deadlock_stack gives
# the full picture for lock-order issues.
env:
TSAN_OPTIONS: "halt_on_error=1 second_deadlock_stack=1"
# 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: setarch -R ./build/tests/kpn_tests
- name: Cleanup
if: always()
run: rm -rf tsan-${{ github.run_id }}
+24
View File
@@ -10,6 +10,30 @@ option(KPN_BUILD_PYTHON "Build Python bindings (requires nanobind)" ON)
option(KPN_BUILD_EXAMPLES "Build examples" ON)
option(KPN_WEB_DEBUG "Enable web debug UI (cpp-httplib)" OFF)
# Sanitizer build. Empty = off. Accepts "thread", "address", "undefined",
# or a combination like "address,undefined". Applied to all kpn targets via
# the kpn_sanitizer_flags() helper below.
#
# The lock-free SPSC Channel<T> (include/kpn/channel.hpp) has hand-reasoned
# acquire/release ordering; -DKPN_SANITIZER=thread + the channel stress test
# (tests/test_channel_stress.cpp) is the dynamic half of verifying it. The
# static half is the CDSChecker model-check harness (see verify/).
set(KPN_SANITIZER "" CACHE STRING
"Build with sanitizer: thread | address | undefined | <combo> (empty = off)")
# Translate KPN_SANITIZER into compile/link flags. No-op when empty.
function(kpn_sanitizer_flags out_var)
if(KPN_SANITIZER)
set(${out_var}
-fsanitize=${KPN_SANITIZER}
-fno-omit-frame-pointer
-g
PARENT_SCOPE)
else()
set(${out_var} "" PARENT_SCOPE)
endif()
endfunction()
# ── Core library (header-only) ────────────────────────────────────────────────
add_library(kpn INTERFACE)
target_include_directories(kpn INTERFACE
+561 -1145
View File
File diff suppressed because it is too large Load Diff
+39 -1
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
@@ -180,7 +209,16 @@ public:
if (h == t) {
// Ring drained — deliver any pending out-of-band sentinel (EOF)
// now, so it always arrives after the data pushed before it.
{ T s; if (take_sentinel(s)) return s; }
//
// Re-confirm emptiness against a fresh tail_ first: the snapshot
// at the top of the loop may be stale (the producer can push more
// values *and* the sentinel in the window since), and the sentinel
// must never jump ahead of ring values pushed before it. The spin
// and post-spin takes below already reload tail_ on the line above
// them; this is the one take that used the loop-top snapshot.
if (h == tail_.load(std::memory_order_acquire)) {
T s; if (take_sentinel(s)) return s;
}
if (!accepting_.load(std::memory_order_acquire))
throw ChannelClosedError{};
+24 -1
View File
@@ -115,6 +115,16 @@ public:
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:
void run_loop() {
while (!stop_flag_.load(std::memory_order_relaxed)) {
@@ -126,8 +136,19 @@ private:
for (std::size_t i = 0; i < N; ++i) {
if (out_channels_[i]) {
// Lossless: block until this consumer drains. Note the
// 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 for this output independently
catch (const ChannelOverflowError&) {} // drop independently
}
}
}
@@ -142,6 +163,8 @@ private:
std::string name_;
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::array<Channel<T>*, N> out_channels_{};
std::atomic<bool> stop_flag_{false};
+4
View File
@@ -37,6 +37,10 @@ struct INode {
// halt(): alias for stop() — immediate, discards in-flight work.
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
// back to stop(). Network and StaticNetwork override with topo-ordered drain.
virtual void shutdown() { stop(); }
+32
View File
@@ -133,6 +133,14 @@ public:
void set_error_handler(NodeErrorHandler h) { error_handler_ = std::move(h); }
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_network_overflow_callback(NodeEventCallback cb) override { event_callbacks_[1] = 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));
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&) {
@@ -422,6 +439,7 @@ private:
std::shared_ptr<IScheduler> scheduler_;
std::string name_;
bool lossless_{false};
std::size_t fifo_capacity_;
input_channels_t input_channels_;
output_channels_t output_channels_{};
@@ -498,6 +516,14 @@ public:
void set_error_handler(NodeErrorHandler h) { error_handler_ = std::move(h); }
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_network_overflow_callback(NodeEventCallback cb) override { event_callbacks_[1] = 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));
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&) {
@@ -717,6 +748,7 @@ private:
Obj& obj_;
std::shared_ptr<IScheduler> scheduler_;
std::string name_;
bool lossless_{false};
std::size_t fifo_capacity_;
input_channels_t input_channels_;
output_channels_t output_channels_{};
+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
+18
View File
@@ -219,6 +219,24 @@ public:
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:
struct Snapshots {
std::vector<NodeSnapshot> nodes;
+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() };
}
+30 -1
View File
@@ -43,6 +43,35 @@ target_link_libraries(kpn_tests PRIVATE
GTest::gtest
)
# Channel stress suite (separate executable)
# Contended SPSC tests for the lock-free Channel<T>. Kept out of kpn_tests
# because each case runs many reps / tens of thousands of items and is slow.
# Most valuable under -DKPN_SANITIZER=thread, but correct (and run) without it.
add_executable(kpn_tests_stress test_channel_stress.cpp)
target_link_libraries(kpn_tests_stress PRIVATE kpn Catch2::Catch2WithMain)
# Sanitizer flags
# kpn_sanitizer_flags() is defined in the top-level CMakeLists and is a no-op
# unless -DKPN_SANITIZER=... is set. Sanitizer must be on both compile and link.
kpn_sanitizer_flags(_kpn_san)
if(_kpn_san)
foreach(_t kpn_tests kpn_tests_stress)
target_compile_options(${_t} PRIVATE ${_kpn_san})
target_link_options(${_t} PRIVATE ${_kpn_san})
endforeach()
endif()
include(CTest)
include(Catch)
catch_discover_tests(kpn_tests)
# DISCOVERY_MODE PRE_TEST defers test enumeration to `ctest` run time. The
# default (POST_BUILD) runs each test binary during the build to list its
# cases which fails a sanitizer build: a TSan/ASan binary needs a fixed
# address-space layout and aborts on startup ("unexpected memory mapping")
# under the container's ASLR, breaking the build before any test runs. The
# tsan.yaml job invokes the binaries directly (not via ctest), so deferring
# discovery costs nothing there and keeps `ctest` working for normal builds.
catch_discover_tests(kpn_tests DISCOVERY_MODE PRE_TEST)
# Register the stress suite under its own label so CI can run / time it
# separately from the fast unit tests.
catch_discover_tests(kpn_tests_stress DISCOVERY_MODE PRE_TEST PROPERTIES LABELS "stress")
+281
View File
@@ -0,0 +1,281 @@
// Contended stress tests for the lock-free SPSC Channel<T>.
//
// The other channel tests (test_channel.cpp) are single-threaded or use a
// single 20 ms sleep to order two threads — they never actually contend on the
// ring, so they exercise neither the memory-ordering pairing nor the
// spin/futex/lost-wakeup logic in pop().
//
// These tests are written to be run under ThreadSanitizer:
//
// cmake -B build -DKPN_SANITIZER=thread -DKPN_BUILD_EXAMPLES=OFF -DKPN_BUILD_PYTHON=OFF
// cmake --build build --target kpn_tests_tsan
// ./build/tests/kpn_tests_tsan
//
// They are also valid (and meaningful) without a sanitizer: the value/sequence
// assertions catch lost or duplicated items regardless of build flags. TSan
// adds detection of the underlying data race even on runs where the race did
// not corrupt observable state.
//
// Channel<T> is SPSC: exactly one producer thread and one consumer thread per
// channel. Every scenario below honours that contract.
#include <catch2/catch_test_macros.hpp>
#include <atomic>
#include <chrono>
#include <kpn/channel.hpp>
#include <thread>
#include <vector>
using namespace kpn;
using namespace std::chrono_literals;
namespace {
// Repeat each scenario enough times that rare interleavings (spin window just
// missing / just catching the next push, disable landing inside the futex
// wait) actually occur across a run. Kept modest so a TSan run stays minutes,
// not hours.
constexpr int kReps = 200;
} // namespace
TEST_CASE("SPSC: every pushed item is popped exactly once, in order", "[channel][stress]") {
// Small capacity forces frequent full/empty transitions, so both the
// producer's overflow-retry and the consumer's spin->futex path are hit
// many times. The producer retries on overflow rather than dropping, so
// the consumer must observe a strictly contiguous 0..N-1 sequence.
constexpr int N = 50'000;
Channel<int> ch(/*capacity=*/4, /*spin_count=*/16);
std::thread producer([&] {
for (int i = 0; i < N; ++i) {
for (;;) {
try { ch.push(i); break; }
catch (const ChannelOverflowError&) { std::this_thread::yield(); }
}
}
});
int expected = 0;
bool in_order = true;
for (int i = 0; i < N; ++i) {
int v = ch.pop();
if (v != expected) in_order = false;
++expected;
}
producer.join();
REQUIRE(in_order);
REQUIRE(expected == N);
REQUIRE(ch.size() == 0);
}
TEST_CASE("SPSC: tight empty<->non-empty transitions exercise spin/futex boundary",
"[channel][stress]") {
// spin_count=0 forces every empty pop() straight into atomic::wait, so this
// hammers the lost-wakeup guard (snapshot wake_, re-check tail_, then wait).
// The producer pushes one item then waits to go empty again, maximising the
// number of empty->non-empty edges relative to item count.
constexpr int N = 20'000;
Channel<int> ch(/*capacity=*/2, /*spin_count=*/0);
std::thread producer([&] {
for (int i = 0; i < N; ++i) {
for (;;) {
try { ch.push(i); break; }
catch (const ChannelOverflowError&) { std::this_thread::yield(); }
}
}
});
long sum = 0;
for (int i = 0; i < N; ++i) sum += ch.pop();
producer.join();
// Sum of 0..N-1 — detects any lost or duplicated item.
REQUIRE(sum == static_cast<long>(N) * (N - 1) / 2);
}
TEST_CASE("SPSC: disable() while consumer is blocked in pop() unblocks cleanly",
"[channel][stress]") {
// The data race of record: consumer blocked in pop() (spinning or parked in
// the futex) while the owner thread calls disable(). pop() must observe the
// close and throw ChannelClosedError — it must not hang and must not read
// past the ring. Repeated so disable() lands at many points in pop()'s loop.
for (int rep = 0; rep < kReps; ++rep) {
Channel<int> ch(/*capacity=*/4, /*spin_count=*/8);
std::atomic<bool> threw{false};
std::atomic<bool> finished{false};
std::thread consumer([&] {
try {
ch.pop(); // empty channel: will block
} catch (const ChannelClosedError&) {
threw.store(true, std::memory_order_relaxed);
}
finished.store(true, std::memory_order_relaxed);
});
// Give the consumer a chance to reach the wait, then close.
std::this_thread::sleep_for(50us);
ch.disable();
consumer.join();
REQUIRE(finished.load());
REQUIRE(threw.load());
}
}
TEST_CASE("SPSC: producer racing a disable() never throws and never hangs",
"[channel][stress]") {
// Mirror of the above from the producer side: push() racing disable() must
// either enqueue or silently drop, never throw ChannelClosedError and never
// wedge. Overflow is still a legal outcome (full accepting channel) and is
// tolerated here.
for (int rep = 0; rep < kReps; ++rep) {
Channel<int> ch(/*capacity=*/8, /*spin_count=*/8);
std::atomic<bool> bad{false};
std::thread producer([&] {
for (int i = 0; i < 1000; ++i) {
try { ch.push(i); }
catch (const ChannelOverflowError&) { /* legal: full */ }
catch (...) { bad.store(true, std::memory_order_relaxed); break; }
}
});
std::this_thread::sleep_for(20us);
ch.disable(); // owner closes mid-stream
producer.join();
REQUIRE_FALSE(bad.load());
}
}
TEST_CASE("SPSC: push_callback fires on each empty->non-empty transition",
"[channel][stress]") {
// The empty->non-empty callback ([channel.hpp] was_empty branch) is read by
// the consumer-side notification path. Run it under contention to make sure
// the was_empty detection isn't torn by a concurrent pop().
Channel<int> ch(/*capacity=*/4, /*spin_count=*/4);
std::atomic<int> callbacks{0};
ch.set_push_callback([&] { callbacks.fetch_add(1, std::memory_order_relaxed); });
constexpr int N = 10'000;
std::thread producer([&] {
for (int i = 0; i < N; ++i) {
for (;;) {
try { ch.push(i); break; }
catch (const ChannelOverflowError&) { std::this_thread::yield(); }
}
}
});
for (int i = 0; i < N; ++i) (void)ch.pop();
producer.join();
// At least one transition, at most one per item; mainly we assert the run
// completed without TSan flagging a race on push_callback_/was_empty.
REQUIRE(callbacks.load() >= 1);
REQUIRE(callbacks.load() <= N);
}
// Ordering contract of the out-of-band sentinel under contention.
//
// push_sentinel() publishes has_eof_ (release) after the producer's N ring
// pushes; a consumer that observes has_eof_ (acquire) therefore also observes
// every value pushed before it. Both pop() and try_pop_now() only surface the
// sentinel once the ring is *freshly* observed empty, so the sentinel is the
// strictly last item received — it never jumps ahead of a ring value pushed
// before it. These tests treat the sentinel as a hard "last message" barrier
// (the consumer stops draining the moment it sees it) and assert that all N
// values arrived, in a contiguous 0..N-1 sequence, before it.
//
// Regression guard: an earlier version of pop() checked emptiness against a
// stale tail_ snapshot from the top of its loop, so under load the sentinel
// could surface with a few real values still queued — breaking in_order /
// values==N here. Under TSan these also cover the has_eof_/eof_value_
// acquire/release handshake and the spin/futex wakeup on push_sentinel().
TEST_CASE("SPSC: sentinel is strictly last, after every value (blocking pop)",
"[channel][stress]") {
constexpr int N = 20'000;
constexpr int SENTINEL = -1;
for (int rep = 0; rep < kReps; ++rep) {
// Small ring + tiny spin window so the ring is frequently empty exactly
// when the sentinel is published — the interleaving under test.
Channel<int> ch(/*capacity=*/4, /*spin_count=*/8);
std::thread producer([&] {
for (int i = 0; i < N; ++i) {
for (;;) {
try { ch.push(i); break; }
catch (const ChannelOverflowError&) { std::this_thread::yield(); }
}
}
ch.push_sentinel(SENTINEL); // must-deliver, never overflows/blocks
});
int expected = 0;
bool in_order = true;
bool saw_sentinel = false;
// Treat the sentinel as EOF: stop draining the instant it appears.
for (;;) {
int v = ch.pop();
if (v == SENTINEL) { saw_sentinel = true; break; }
if (v != expected) in_order = false;
++expected;
}
producer.join();
REQUIRE(saw_sentinel);
REQUIRE(in_order);
REQUIRE(expected == N); // all N values received before the sentinel
REQUIRE(ch.size() == 0);
REQUIRE(ch.approx_size() == 0);
}
}
TEST_CASE("SPSC: sentinel is strictly last, after every value (try_pop_now)",
"[channel][stress]") {
// The pool-node consume path is try_pop_now(), not pop(): it must surface
// the out-of-band sentinel only once the ring is freshly observed empty.
// The consumer spins with no sleeps, racing the producer at full tilt
// across the empty-ring boundary where take_sentinel() is reached.
constexpr int N = 20'000;
constexpr int SENTINEL = -1;
for (int rep = 0; rep < kReps; ++rep) {
Channel<int> ch(/*capacity=*/4, /*spin_count=*/0);
std::thread producer([&] {
for (int i = 0; i < N; ++i) {
for (;;) {
try { ch.push(i); break; }
catch (const ChannelOverflowError&) { std::this_thread::yield(); }
}
}
ch.push_sentinel(SENTINEL);
});
int expected = 0;
bool in_order = true;
bool saw_sentinel = false;
int v;
for (;;) {
if (!ch.try_pop_now(v)) { std::this_thread::yield(); continue; }
if (v == SENTINEL) { saw_sentinel = true; break; }
if (v != expected) in_order = false;
++expected;
}
producer.join();
REQUIRE(saw_sentinel);
REQUIRE(in_order);
REQUIRE(expected == N);
// Sentinel held no ring slot; once taken the channel is fully empty.
REQUIRE(ch.size() == 0);
REQUIRE(ch.approx_size() == 0);
}
}