From e2a12fb4c3fb5a42df4abe44214143b169a8c351 Mon Sep 17 00:00:00 2001 From: Timothy Place Date: Wed, 30 Sep 2026 17:46:52 -0500 Subject: [PATCH 1/4] Plan: a design for worker mode (2.5), for review before code Python on a thread of its own, a fixed latency of whole vectors behind the audio thread, which never takes the GIL: a lock-free ring of slots in the core, the existing processor::process() run on the worker, underruns output as silence with the latency kept, @mode and @latency applied at the next chain compile, and a read-only latency attribute (Max documents no call to report latency to the host). Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_019rVV9SkjkWmiBXFMT4whkz --- docs/PRODUCTION-PLAN.md | 40 +++++++++++++++++++++++++++++++++++++++- 1 file changed, 39 insertions(+), 1 deletion(-) diff --git a/docs/PRODUCTION-PLAN.md b/docs/PRODUCTION-PLAN.md index cd2745d..35f6097 100644 --- a/docs/PRODUCTION-PLAN.md +++ b/docs/PRODUCTION-PLAN.md @@ -131,7 +131,45 @@ while audio ran segfaulted in 5 of 5 runs. that `thispatcher` connects to the new inlet and outlet carry signal (this check fails against (b)); shrunk to one, the removed outlets' cords go with them. - [ ] **2.5 Worker mode (D1)** — `@mode worker`: Python on a worker thread, lock-free FIFO, - latency reported to Max, underrun → silence. + latency reported to Max, underrun → silence. *Design (proposed, 2026-09-30):* + - **What it buys.** In direct mode the audio thread takes the GIL, so anything else holding it — + a reload compiling a large file (2.6's limit), a message handler, a GC pass — delays the + buffer. In worker mode the audio thread never takes the GIL or calls CPython: it copies + vectors in and out of a queue, and Python runs on a thread of its own, a fixed latency behind. + - **In the core** (D6), beside `processor`: a `worker` that owns the thread and a ring of slots, + each one vector of every host input and output channel. The audio thread writes vector *k*'s + inputs into a slot and reads vector *k − L*'s outputs; the worker takes slots in order and runs + the existing `processor::process()` on them — per sample or per vector, with 2.4's channel + matching — so the class contract does not change. Slot states are atomics (single producer, + single consumer, no locks); the audio thread wakes the worker with a C++20 atomic notify, which + never blocks. The worker's thread is created and joined on the main thread. + - **Latency** *L* is whole vectors, set by `@latency` (in vectors, default 2), reported by a + read-only `latency` attribute in samples (*L* × vector size) that a patch can read to align + other paths. Max has no documented call for an MSP object to report latency to the host, so + none is used; the output is primed with *L* vectors of silence. + - **Underruns.** A slot not done when its outputs are due is output as silence and counted; the + worker still processes it when it gets there (so the class's state stays continuous) and + discards its outputs, catching up through the backlog, so the latency stays *L*. If it falls + a whole ring behind (the ring holds *L* + a margin), the oldest inputs are dropped, which the + class sees as a gap. Both are reported once per load from the main thread (`flush_reports()`), + never printed on the audio thread. + - **`prepare()` and the mode.** The ring is sized when the chain compiles (`dspsetup`: vector + size, and the object's inlet and outlet counts), with the worker stopped and restarted around + it; `prepare()` runs as now, on the main thread. `@mode` and `@latency` take effect at the next + compile — setting them marks the chain broken (as 2.4 does), so that is at once while audio + runs. + - **Reloads, attributes and messages** are unchanged: they take the GIL on the main or scheduler + thread, and the worker picks up a new binding at its next vector, as the audio thread does + now. An attribute change is heard *L* vectors later. + - **Tests.** Core battery (Linux, under TSan too): worker output equals direct output delayed by + *L* vectors, per sample and per vector, several channels; a worker stalled by a class that + sleeps underruns to silence, is reported once, and comes back at the same latency; reload + under a running worker; resizing in `prepare()`; destruction joins the thread. Runtime test in + Max: `@mode worker` against `delay~` of *L* vectors, a reload under audio, and the `latency` + attribute. + - **Open:** the worker calls Python once per host vector; batching several vectors per call + would cut the per-call overhead further (worker mode's other win) at more latency — a later + option, `@block`, if measurements show it pays. - [x] **2.6 Shorter reload stalls** — compile outside the swap and hold the GIL only for the swap. *Measured first* (`core/bench/reload_bench.cpp`: an audio thread with Core Audio's real-time scheduling computes 512-sample buffers of 64-sample vectors at 96 kHz while the main thread saves From 4fa81fe19e52d31f4c3e70387ace5171bbc5392a Mon Sep 17 00:00:00 2001 From: Timothy Place Date: Wed, 30 Sep 2026 20:02:03 -0500 Subject: [PATCH 2/4] Plan: the worker-mode design is decided (2.5) Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_019rVV9SkjkWmiBXFMT4whkz --- docs/PRODUCTION-PLAN.md | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/docs/PRODUCTION-PLAN.md b/docs/PRODUCTION-PLAN.md index 35f6097..6f4acc6 100644 --- a/docs/PRODUCTION-PLAN.md +++ b/docs/PRODUCTION-PLAN.md @@ -131,7 +131,7 @@ while audio ran segfaulted in 5 of 5 runs. that `thispatcher` connects to the new inlet and outlet carry signal (this check fails against (b)); shrunk to one, the removed outlets' cords go with them. - [ ] **2.5 Worker mode (D1)** — `@mode worker`: Python on a worker thread, lock-free FIFO, - latency reported to Max, underrun → silence. *Design (proposed, 2026-09-30):* + latency reported to Max, underrun → silence. *Design (decided 2026-09-30):* - **What it buys.** In direct mode the audio thread takes the GIL, so anything else holding it — a reload compiling a large file (2.6's limit), a message handler, a GC pass — delays the buffer. In worker mode the audio thread never takes the GIL or calls CPython: it copies @@ -167,9 +167,10 @@ while audio ran segfaulted in 5 of 5 runs. under a running worker; resizing in `prepare()`; destruction joins the thread. Runtime test in Max: `@mode worker` against `delay~` of *L* vectors, a reload under audio, and the `latency` attribute. - - **Open:** the worker calls Python once per host vector; batching several vectors per call - would cut the per-call overhead further (worker mode's other win) at more latency — a later - option, `@block`, if measurements show it pays. + - **Decided:** latency in whole vectors, default 2; an underrun is silence with the latency kept; + one Python call per host vector — batching several per call would cut the per-call overhead + further (worker mode's other win) at more latency, a later option (`@block`) if measurements + show it pays. - [x] **2.6 Shorter reload stalls** — compile outside the swap and hold the GIL only for the swap. *Measured first* (`core/bench/reload_bench.cpp`: an audio thread with Core Audio's real-time scheduling computes 512-sample buffers of 64-sample vectors at 96 kHz while the main thread saves From 0541649e00471b615335c2e4b02aeecd6d48dbe2 Mon Sep 17 00:00:00 2001 From: Timothy Place Date: Wed, 30 Sep 2026 20:03:41 -0500 Subject: [PATCH 3/4] Plan: worker mode's two latency attributes get two names (2.5) @latency sets it in vectors; the read-only @latencysamples reports it in samples. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_019rVV9SkjkWmiBXFMT4whkz --- docs/PRODUCTION-PLAN.md | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/docs/PRODUCTION-PLAN.md b/docs/PRODUCTION-PLAN.md index 6f4acc6..973ca7d 100644 --- a/docs/PRODUCTION-PLAN.md +++ b/docs/PRODUCTION-PLAN.md @@ -143,8 +143,8 @@ while audio ran segfaulted in 5 of 5 runs. matching — so the class contract does not change. Slot states are atomics (single producer, single consumer, no locks); the audio thread wakes the worker with a C++20 atomic notify, which never blocks. The worker's thread is created and joined on the main thread. - - **Latency** *L* is whole vectors, set by `@latency` (in vectors, default 2), reported by a - read-only `latency` attribute in samples (*L* × vector size) that a patch can read to align + - **Latency** *L* is whole vectors, set by `@latency` (in vectors, default 2), and reported in + samples (*L* × vector size) by the read-only `@latencysamples`, which a patch can read to align other paths. Max has no documented call for an MSP object to report latency to the host, so none is used; the output is primed with *L* vectors of silence. - **Underruns.** A slot not done when its outputs are due is output as silence and counted; the @@ -165,8 +165,8 @@ while audio ran segfaulted in 5 of 5 runs. *L* vectors, per sample and per vector, several channels; a worker stalled by a class that sleeps underruns to silence, is reported once, and comes back at the same latency; reload under a running worker; resizing in `prepare()`; destruction joins the thread. Runtime test in - Max: `@mode worker` against `delay~` of *L* vectors, a reload under audio, and the `latency` - attribute. + Max: `@mode worker` against `delay~` of *L* vectors, a reload under audio, and + `@latencysamples`. - **Decided:** latency in whole vectors, default 2; an underrun is silence with the latency kept; one Python call per host vector — batching several per call would cut the per-call overhead further (worker mode's other win) at more latency, a later option (`@block`) if measurements From 65a06942d3f528fbd166d026bdb4454ccf88facd Mon Sep 17 00:00:00 2001 From: Timothy Place Date: Wed, 30 Sep 2026 20:16:57 -0500 Subject: [PATCH 4/4] =?UTF-8?q?2.5=20(a):=20worker=20mode=20in=20the=20cor?= =?UTF-8?q?e=20=E2=80=94=20process()=20on=20a=20thread=20of=20its=20own?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit worker.h runs a processor's process() on a worker thread a fixed number of vectors behind the audio thread, which only copies vectors through a lock-free ring and never takes the GIL. A vector whose outputs are not ready in time is output as silence; the worker still processes it and catches up, so the latency stays fixed and the class sees every vector in order. A worker a whole ring behind drops new inputs until it catches up. Both are reported once per load from the main thread (the processor gains load_count() for that). The audio thread borrows the ring by taking an atomic pointer, so start() and stop() can run while a host's old signal chain still does. Core battery: the output is the direct output L vectors later, per sample and per vector; channel matching; underruns, drops and their reports; restart and stop; reloads under an audio thread. Clean under ThreadSanitizer. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_019rVV9SkjkWmiBXFMT4whkz --- CLAUDE.md | 4 +- core/CMakeLists.txt | 3 + core/include/tap/python/processor.h | 18 +- core/include/tap/python/worker.h | 333 ++++++++++++++++++++++++++++ core/tests/CMakeLists.txt | 1 + core/tests/python/stalls.py | 13 ++ core/tests/test_worker.cpp | 323 +++++++++++++++++++++++++++ docs/PRODUCTION-PLAN.md | 18 +- 8 files changed, 704 insertions(+), 9 deletions(-) create mode 100644 core/include/tap/python/worker.h create mode 100644 core/tests/python/stalls.py create mode 100644 core/tests/test_worker.cpp diff --git a/CLAUDE.md b/CLAUDE.md index f686ca0..b930dc1 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -21,7 +21,9 @@ its phases, and the audit findings behind them. Tick its items (with the PR) as signatures) that `processor.h` builds its descriptions from), `value.h` (the Max-atom value model and its coercion rules, which reproduce `atom_getlong`/`atom_getfloat`/`atom_getsym`), `processor.h` (load/reload, class introspection, attribute and message dispatch, `prepare()`, and `process()` — per sample, or - per vector when the input is hinted `np.ndarray` — plus `flush_reports()`). New CPython-facing + per vector when the input is hinted `np.ndarray` — plus `flush_reports()`), `worker.h` (worker + mode, plan 2.5: `process()` on a thread of its own, a fixed number of vectors behind an audio + thread that only copies through a lock-free ring and never takes the GIL). New CPython-facing behavior goes here, never in the wrapper. CMake target `tap::python` (`core/CMakeLists.txt`). - **`core/tests/`** — the core's Catch2 battery against CPython 3.13, with Python fixtures in `core/tests/python/` and the shipped examples copied alongside. Runs on Linux, including under diff --git a/core/CMakeLists.txt b/core/CMakeLists.txt index bfcf337..12501b6 100644 --- a/core/CMakeLists.txt +++ b/core/CMakeLists.txt @@ -16,6 +16,9 @@ add_library(tap::python ALIAS tap_python_core) # the family convention: tap add_library(tap::python_core ALIAS tap_python_core) # the original spelling, kept target_include_directories(tap_python_core INTERFACE "${CMAKE_CURRENT_SOURCE_DIR}/include") target_compile_features(tap_python_core INTERFACE cxx_std_20) +# worker.h runs process() on a std::thread (plan 2.5) +find_package(Threads REQUIRED) +target_link_libraries(tap_python_core INTERFACE Threads::Threads) if (CMAKE_SOURCE_DIR STREQUAL CMAKE_CURRENT_SOURCE_DIR) enable_testing() diff --git a/core/include/tap/python/processor.h b/core/include/tap/python/processor.h index 917846e..29fdfa2 100644 --- a/core/include/tap/python/processor.h +++ b/core/include/tap/python/processor.h @@ -127,6 +127,10 @@ namespace tap::python { /// thread, after load(). const std::vector& input_names() const noexcept { return m_input_names; } + /// How many times load() has succeeded: what a host reporting once per load compares against + /// (worker mode does, plan 2.5). Any thread. + std::uint64_t load_count() const noexcept { return m_loads.load(std::memory_order_acquire); } + /// The most inputs, and the most outputs, a process() may declare. static constexpr std::size_t k_max_channels = 64; @@ -237,6 +241,7 @@ namespace tap::python { next.process_function = nullptr; next.prepare_function = nullptr; reset_warnings(); // each kind of audio-thread problem is reported once per load + m_loads.fetch_add(1, std::memory_order_release); release(previous); // may run finalizers, which may let other threads in: now safe if (executed) { @@ -528,12 +533,13 @@ namespace tap::python { std::vector m_attributes; std::vector m_messages; // The rest of the binding and the block buffer: read and written only under the GIL. - PyObject* m_prepare_function{}; // strong, or null - bool m_block_mode{}; - std::size_t m_inputs{1}; // the bound process()'s inputs and outputs - std::size_t m_outputs{1}; - std::vector m_input_names; // main thread only - bool m_announcing{}; // this load() ran the file: see announce() + PyObject* m_prepare_function{}; // strong, or null + bool m_block_mode{}; + std::size_t m_inputs{1}; // the bound process()'s inputs and outputs + std::size_t m_outputs{1}; + std::vector m_input_names; // main thread only + bool m_announcing{}; // this load() ran the file: see announce() + std::atomic m_loads{}; // successful load()s // The np.ndarray each input arrives in, reused every vector: strong, its buffer held so its // memory cannot move. struct block_buffer { diff --git a/core/include/tap/python/worker.h b/core/include/tap/python/worker.h new file mode 100644 index 0000000..64cbdfa --- /dev/null +++ b/core/include/tap/python/worker.h @@ -0,0 +1,333 @@ +/// @file worker.h +/// @brief Worker mode (plan 2.5): a processor's process() on a thread of its own, a fixed number of +/// vectors behind the audio thread, which never takes the GIL. +// SPDX-License-Identifier: MIT +// Copyright 2022-2026 Timothy Place. +// +// Host-independent. In direct mode the audio thread calls processor::process() and so takes the +// GIL: whatever else holds it — a reload compiling a large file, a message handler, a garbage +// collection — delays the audio. A worker moves that call to a thread of its own. The audio thread +// calls worker::process() once per vector, which copies the vector's inputs into a ring of slots +// and copies out the outputs of the vector `latency` vectors before it; the worker thread runs +// processor::process() on each slot in order — per sample or per vector, with the processor's +// channel matching — so the class contract does not change. The first `latency` vectors out are +// silence. +// +// A vector whose outputs are not ready when they are due (an underrun) is output as silence. The +// worker still processes it when it gets there, so the class sees every vector in order and the +// latency stays fixed; only that vector's output is lost. If the worker falls a whole ring behind +// (the latency plus a quarter of a second), the audio thread drops the inputs of new vectors until +// it catches up: the class sees a gap. Both are reported once per load, from the main thread +// (flush_reports()); the audio thread never prints, allocates, locks or calls Python. +// +// Threads: start(), stop(), flush_reports() and the destructor run on the host's main thread, +// never while it holds the GIL (stopping joins the worker, which may be waiting for it); process() +// on the audio thread, possibly at the same time as any of them (a host can run the old signal +// chain while it prepares the next). The audio thread borrows the ring for each vector by taking +// an atomic pointer, so start() and stop() wait at most one vector for it, and a vector that +// arrives while they hold it is silence. + +#pragma once + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "tap/python/processor.h" + +namespace tap::python { + + class worker { + public: + /// The latency, in vectors, unless the host says otherwise. + static constexpr std::size_t k_default_latency = 2; + + /// @param target the processor to run; it must outlive the worker + /// @param log receives the worker's reports (main thread) + /// @param report_ready called on the audio thread when something needs reporting; the host + /// must then call flush_reports() from its main thread. Must be + /// real-time safe, as the processor's is (it may be the same callback). + explicit worker(processor& target, log_function log = {}, std::function report_ready = {}) + : m_target{target} + , m_log{std::move(log)} + , m_report_ready{std::move(report_ready)} {} + + ~worker() { stop(); } + + worker(const worker&) = delete; + worker& operator=(const worker&) = delete; + + /// Start the worker, or restart it with new settings: a ring for `inputs` and `outputs` host + /// channels of `vector_size` frames, `latency` vectors deep (at least 1), with room for the + /// worker to fall a quarter of a second behind at `sample_rate` (at least 16 vectors). + /// Whatever the previous worker had not processed is discarded. Main thread. + void start(const std::size_t inputs, const std::size_t outputs, const std::size_t vector_size, + const std::size_t latency, const double sample_rate) { + stop(); + const auto vectors = std::max(1, latency); + const auto frames = std::max(1, vector_size); + const auto backlog = + std::isfinite(sample_rate) && sample_rate > 0.0 + ? static_cast(std::ceil(k_backlog_seconds * sample_rate / static_cast(frames))) + : std::size_t{0}; + auto next = + std::make_unique(inputs, outputs, frames, vectors, vectors + std::max(k_min_backlog, backlog)); + ring* r = next.get(); + r->thread = std::thread{[this, r] { run(*r); }}; + m_owned = std::move(next); + m_latency = vectors * frames; + m_ring.store(r, std::memory_order_release); + } + + /// Stop the worker and join its thread; process() outputs silence until the next start(). + /// Main thread. + void stop() { + if (!m_owned) { + return; + } + // take the ring back from the audio thread, which holds it for one vector at most + ring* r = m_ring.exchange(nullptr, std::memory_order_acq_rel); + while (!r) { + std::this_thread::yield(); + r = m_ring.exchange(nullptr, std::memory_order_acq_rel); + } + r->stopping.store(true, std::memory_order_release); + r->signal.fetch_add(1, std::memory_order_release); + r->signal.notify_one(); + r->thread.join(); + m_owned.reset(); + m_latency = 0; + } + + /// True between start() and stop(). Main thread. + bool running() const noexcept { return m_owned != nullptr; } + + /// The latency in samples: the vectors it was started with times the vector size; 0 when + /// stopped. Main thread. + std::size_t latency_samples() const noexcept { return m_latency; } + + /// How many vectors the worker has processed (or passed over, dropped) since start(). Main + /// thread; for tests and diagnostics. + std::uint64_t processed() const noexcept { return m_owned ? m_owned->done.load(std::memory_order_acquire) : 0; } + + /// The audio thread's side, once per vector: queue this vector's inputs and output those of + /// the vector `latency` before (silence if the worker has not finished it). Inputs the host + /// does not give read as silence, outputs it does not take are dropped, and the class's + /// own matching applies on the worker. Never blocks, allocates, or takes the GIL. + void process(const double* const* inputs, const std::size_t input_count, double* const* outputs, + const std::size_t output_count, const std::size_t frame_count) { + ring* r = m_ring.exchange(nullptr, std::memory_order_acquire); + if (!r) { + silence(outputs, output_count, 0, frame_count); + return; + } + const auto frames = std::min(frame_count, r->vector_size); + + // this vector's inputs, unless the worker is a whole ring behind + const auto k = r->written.load(std::memory_order_relaxed); // written only here + if (k - r->done.load(std::memory_order_acquire) < r->slot_count) { + const auto index = k % r->slot_count; + for (std::size_t c = 0; c < r->inputs; ++c) { + double* to = r->input(index, c); + if (c < input_count && inputs[c]) { + std::copy_n(inputs[c], frames, to); + } + else { + std::fill_n(to, frames, 0.0); + } + } + r->slots[index].frames.store(frames, std::memory_order_relaxed); + r->slots[index].in_seq.store(k, std::memory_order_release); + } + else { + record(m_dropped, m_dropped_load); + } + r->written.store(k + 1, std::memory_order_release); + r->signal.fetch_add(1, std::memory_order_release); + r->signal.notify_one(); + + // the outputs of the vector `latency` before (the first `latency` vectors have none) + std::size_t given = 0; // frames of output given; the rest is silence + if (k >= r->latency) { + const auto m = k - r->latency; + const auto index = m % r->slot_count; + auto& slot = r->slots[index]; + if (r->done.load(std::memory_order_acquire) > m && slot.out_seq.load(std::memory_order_acquire) == m) { + given = std::min(slot.frames.load(std::memory_order_relaxed), frame_count); + for (std::size_t c = 0; c < output_count; ++c) { + if (c < r->outputs) { + std::copy_n(r->output(index, c), given, outputs[c]); + } + else { + std::fill_n(outputs[c], given, 0.0); + } + } + } + else if (slot.in_seq.load(std::memory_order_relaxed) == m) { + record(m_late, m_late_load); // its inputs were queued: the worker is late + } + } + silence(outputs, output_count, given, frame_count); + + m_ring.store(r, std::memory_order_release); + } + + /// Print what the audio thread recorded since the last flush: vectors output as silence + /// because the worker was late, and inputs dropped because it fell a whole ring behind. + /// Main thread. + void flush_reports() { + if (const auto late = m_late.exchange(0, std::memory_order_acq_rel); late != 0) { + log(log_level::error, "process() was late for " + std::to_string(late) + + " vector(s) in worker mode — output as silence, the latency kept " + "(reported once per load)"); + } + if (const auto dropped = m_dropped.exchange(0, std::memory_order_acq_rel); dropped != 0) { + log(log_level::error, "process() fell too far behind in worker mode — the input of " + + std::to_string(dropped) + + " vector(s) was dropped (reported once per load)"); + } + } + + private: + static constexpr std::uint64_t k_none = std::numeric_limits::max(); + static constexpr double k_backlog_seconds = 0.25; + static constexpr std::size_t k_min_backlog = 16; + + struct slot { + std::atomic in_seq{k_none}; // the vector whose inputs it holds + std::atomic out_seq{k_none}; // the vector whose outputs it holds + std::atomic frames{}; + }; + + /// One run of the worker: the slots, and the thread that processes them. + struct ring { + ring(const std::size_t input_channels, const std::size_t output_channels, const std::size_t frames, + const std::size_t vectors, const std::size_t count) + : inputs{input_channels} + , outputs{output_channels} + , vector_size{frames} + , latency{vectors} + , slot_count{count} + , input_samples(count * input_channels * frames) + , output_samples(count * output_channels * frames) + , slots(count) {} + + double* input(const std::size_t index, const std::size_t channel) { + return input_samples.data() + ((index * inputs) + channel) * vector_size; + } + double* output(const std::size_t index, const std::size_t channel) { + return output_samples.data() + ((index * outputs) + channel) * vector_size; + } + + std::size_t inputs; + std::size_t outputs; + std::size_t vector_size; + std::size_t latency; + std::size_t slot_count; + std::vector input_samples; + std::vector output_samples; + std::vector slots; + std::atomic written{}; // vectors the audio thread has queued + std::atomic done{}; // vectors the worker has finished (or passed over) + std::atomic signal{}; // bumped to wake the worker + std::atomic stopping{}; + std::thread thread; + }; + + processor& m_target; + log_function m_log; + std::function m_report_ready; + + std::atomic m_ring{}; // the running ring, while the audio thread is not using it + std::unique_ptr m_owned; // the running ring (main thread) + std::size_t m_latency{}; // in samples (main thread) + + std::atomic m_late{}; // vectors output as silence since the last flush + std::atomic m_dropped{}; // vectors whose inputs were dropped + std::atomic m_late_load{k_none}; // the load for which late vectors were last reported + std::atomic m_dropped_load{k_none}; // the same, for dropped ones + + /// The worker thread: process each queued vector in order, sleeping while there is none. + void run(ring& r) { + std::vector in(r.inputs); + std::vector out(r.outputs); + std::uint64_t next = 0; + for (;;) { + const auto signalled = r.signal.load(std::memory_order_acquire); + if (r.stopping.load(std::memory_order_acquire)) { + return; + } + if (next == r.written.load(std::memory_order_acquire)) { + r.signal.wait(signalled, std::memory_order_acquire); + continue; + } + const auto index = next % r.slot_count; + auto& slot = r.slots[index]; + if (slot.in_seq.load(std::memory_order_acquire) == next) { // not dropped + for (std::size_t c = 0; c < r.inputs; ++c) { + in[c] = r.input(index, c); + } + for (std::size_t c = 0; c < r.outputs; ++c) { + out[c] = r.output(index, c); + } + const auto frames = slot.frames.load(std::memory_order_relaxed); + try { + m_target.process(in.data(), r.inputs, out.data(), r.outputs, frames); + } + catch (...) { // an exception must not end the thread (and with it, the host) + silence(out.data(), r.outputs, 0, frames); + } + slot.out_seq.store(next, std::memory_order_release); + } + ++next; + r.done.store(next, std::memory_order_release); + } + } + + /// Count one late or dropped vector (audio thread), and ask for a report if it is the first + /// since the class was last loaded; later ones are counted into that report until it is + /// flushed, and then not again until the next load. + void record(std::atomic& count, std::atomic& reported_load) { + const auto load = m_target.load_count(); + if (reported_load.load(std::memory_order_relaxed) == load) { + // reported already: count it into the report if that is still pending (never into one + // flushed meanwhile, which would carry it over to the next load's) + auto pending = count.load(std::memory_order_relaxed); + while (pending != 0 && !count.compare_exchange_weak(pending, pending + 1, std::memory_order_relaxed)) { + } + return; + } + reported_load.store(load, std::memory_order_relaxed); + count.fetch_add(1, std::memory_order_release); + if (m_report_ready) { + m_report_ready(); + } + } + + void log(const log_level level, const std::string& text) const { + if (m_log) { + m_log(level, text); + } + } + + /// Zero `count` channels of `out` from `first_frame` to `frame_count`. + static void silence(double* const* out, const std::size_t count, const std::size_t first_frame, + const std::size_t frame_count) { + for (std::size_t c = 0; c < count; ++c) { + if (first_frame < frame_count) { + std::fill(out[c] + first_frame, out[c] + frame_count, 0.0); + } + } + } + }; + +} // namespace tap::python diff --git a/core/tests/CMakeLists.txt b/core/tests/CMakeLists.txt index 30d74d4..380a962 100644 --- a/core/tests/CMakeLists.txt +++ b/core/tests/CMakeLists.txt @@ -35,6 +35,7 @@ add_executable(tap_python_core_tests test_loading.cpp test_types.cpp test_channels.cpp + test_worker.cpp ) target_link_libraries(tap_python_core_tests PRIVATE tap::python Python3::Python Catch2::Catch2WithMain) diff --git a/core/tests/python/stalls.py b/core/tests/python/stalls.py new file mode 100644 index 0000000..2c3d3b9 --- /dev/null +++ b/core/tests/python/stalls.py @@ -0,0 +1,13 @@ +# Test fixture: a process() that can be told to stall once, as a worker falling behind (plan 2.5). +# Identity otherwise, so each output sample says which input it came from. +import time + + +class stalls: + stall: float = 0.0 # seconds to sleep at the next sample, once + + def process(self, x: float) -> float: + if self.stall: + seconds, self.stall = self.stall, 0.0 + time.sleep(seconds) + return x diff --git a/core/tests/test_worker.cpp b/core/tests/test_worker.cpp new file mode 100644 index 0000000..b34e825 --- /dev/null +++ b/core/tests/test_worker.cpp @@ -0,0 +1,323 @@ +/// @file test_worker.cpp +/// @brief Worker mode (plan 2.5): process() on a thread of its own, a fixed latency behind the +/// audio thread. +// SPDX-License-Identifier: MIT +// Copyright 2022-2026 Timothy Place. + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include + +#include "support.h" +#include "tap/python/worker.h" + +using namespace tap::python; +using namespace tap::python::test; + +namespace { + + constexpr std::size_t k_frames = 64; + + using channel_data = std::vector>; + + /// A processor run by a worker, with one counter for both's report_ready. + struct harness { + log_capture log; + std::atomic notified{0}; + processor p; + worker w; // declared after p: stopped (and joined) before p goes + + explicit harness(const char* name) + : p{name, log.sink(), {}, [this] { ++notified; }} + , w{p, log.sink(), [this] { ++notified; }} {} + + /// One vector through the worker, as the audio thread calls it: each input channel holds + /// one value throughout; returns `output_count` channels. + channel_data push(const std::vector& values, const std::size_t output_count, + const std::size_t frame_count = k_frames) { + channel_data inputs; + channel_data outputs(output_count, std::vector(frame_count, 12345.0)); + std::vector in; + std::vector out; + for (const auto v : values) { + inputs.emplace_back(frame_count, v); + } + for (const auto& channel : inputs) { + in.push_back(channel.data()); + } + for (auto& channel : outputs) { + out.push_back(channel.data()); + } + w.process(in.data(), in.size(), out.data(), out.size(), frame_count); + return outputs; + } + + /// Wait until the worker has processed `count` vectors (a test may wait; the audio thread + /// never does). + void wait_for(const std::uint64_t count) { + const auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds{20}; + while (w.processed() < count && std::chrono::steady_clock::now() < deadline) { + std::this_thread::sleep_for(std::chrono::microseconds{100}); + } + REQUIRE(w.processed() >= count); + } + + std::size_t lines_containing(const std::string& fragment) const { + const auto lines = log.lines(); + return static_cast(std::count_if(lines.begin(), lines.end(), [&](const auto& line) { + return line.text.find(fragment) != std::string::npos; + })); + } + }; + + bool numpy_importable() { + gil_lock lock; + PyObject* numpy = PyImport_ImportModule("numpy"); + if (!numpy) { + PyErr_Clear(); + } + Py_XDECREF(numpy); + if (std::getenv("TAP_PYTHON_TEST_REQUIRE_EXAMPLES")) { + REQUIRE(numpy != nullptr); + } + return numpy != nullptr; + } + +} // namespace + +SCENARIO("In worker mode the output is the direct output, the latency later (plan 2.5)") { + ensure_runtime(); + const auto [name, block] = GENERATE(std::pair{"two_by_two", false}, std::pair{"block_two_by_two", true}); + if (block && !numpy_importable()) { + SKIP("numpy is not importable by the embedded interpreter"); + } + const std::size_t latency = GENERATE(1, 2, 3); + harness h{name}; + REQUIRE(h.p.load()); + h.p.prepare(48000.0, k_frames); + h.w.start(2, 2, k_frames, latency, 48000.0); + CHECK(h.w.running()); + CHECK(h.w.latency_samples() == latency * k_frames); + + for (std::size_t k = 0; k < 12; ++k) { + const auto a = static_cast(k + 1); + const auto outputs = h.push({a, a / 4.0}, 2); + h.wait_for(k + 1); + if (k < latency) { + CHECK(all_equal(outputs[0], 0.0)); // the first `latency` vectors are silence + CHECK(all_equal(outputs[1], 0.0)); + } + else { + const auto was = static_cast(k - latency + 1); // vector k - latency's input + CHECK(all_equal(outputs[0], was * 1.25)); + CHECK(all_equal(outputs[1], was * 0.75)); + } + } + CHECK(h.notified.load() == 0); +} + +SCENARIO("The host's channels are matched to the worker's ring (plan 2.5)") { + ensure_runtime(); + harness h{"two_by_two"}; + REQUIRE(h.p.load()); + h.w.start(2, 2, k_frames, 1, 48000.0); + + WHEN("the host gives one input and takes three outputs") { + h.push({2.0}, 3); + h.wait_for(1); + const auto outputs = h.push({2.0}, 3); + THEN("the missing input is silence and the extra output silent") { + CHECK(all_equal(outputs[0], 2.0)); // 2 + 0 + CHECK(all_equal(outputs[1], 2.0)); // 2 - 0 + CHECK(all_equal(outputs[2], 0.0)); + } + } + WHEN("the host gives a shorter vector") { + h.push({1.0, 1.0}, 2, 16); + h.wait_for(1); + const auto outputs = h.push({1.0, 1.0}, 2, 16); + THEN("that many frames come back") { + CHECK(all_equal(outputs[0], 2.0)); + } + } +} + +SCENARIO("A worker that falls behind outputs silence, keeps its latency, and says so once per load " + "(plan 2.5)") { + ensure_runtime(); + harness h{"stalls"}; + REQUIRE(h.p.load()); + h.w.start(1, 1, k_frames, 2, 48000.0); + std::uint64_t k = 0; + const auto paced = [&] { // one vector, waiting for the worker: never late + const auto outputs = h.push({static_cast(k + 1)}, 1); + h.wait_for(++k); + return outputs; + }; + const auto stall = [&] { // the worker stalls 0.2 s; the next vectors arrive meanwhile + REQUIRE(h.p.set_attribute("stall", value{0.2})); + channel_data outputs; + for (int i = 0; i < 5; ++i) { + outputs.push_back(h.push({static_cast(++k)}, 1)[0]); + } + h.wait_for(k); + return outputs; + }; + for (int i = 0; i < 4; ++i) { + paced(); + } + + const auto during = stall(); + THEN("a vector due while it stalled is output as silence") { + CHECK(all_equal(during[2], 0.0)); // due two vectors after the stalled one + } + THEN("once it has caught up, the output is the input two vectors before") { + const auto outputs = paced(); // vector k is in; vector k - 2 is out + CHECK(all_equal(outputs[0], static_cast(k - 2))); + } + THEN("it is reported, from the main thread, once") { + CHECK(h.notified.load() == 1); + h.w.flush_reports(); + CHECK(h.log.contains("late for", log_level::error)); + CHECK(h.lines_containing("late for") == 1); + + AND_WHEN("it stalls again") { + stall(); + h.w.flush_reports(); + THEN("nothing more is reported") { + CHECK(h.notified.load() == 1); + CHECK(h.lines_containing("late for") == 1); + } + } + AND_WHEN("the class is reloaded and it stalls again") { + REQUIRE(h.p.load()); + stall(); + h.w.flush_reports(); + THEN("it is reported again") { + CHECK(h.notified.load() == 2); + CHECK(h.lines_containing("late for") == 2); + } + } + } +} + +SCENARIO("A worker a whole ring behind drops new inputs until it catches up (plan 2.5)") { + ensure_runtime(); + harness h{"stalls"}; + REQUIRE(h.p.load()); + h.w.start(1, 1, k_frames, 2, 0.0); // no sample rate: the least backlog, 16 vectors (18 slots) + h.push({1.0}, 1); + h.push({2.0}, 1); + h.wait_for(2); + + REQUIRE(h.p.set_attribute("stall", value{0.3})); + for (std::uint64_t k = 2; k < 32; ++k) { // the worker is stuck on vector 2 + h.push({static_cast(k + 1)}, 1); + } + h.wait_for(32); + h.w.flush_reports(); + THEN("the vectors beyond the ring were dropped, and it says so") { + CHECK(h.log.contains("input of 12 vector(s) was dropped", log_level::error)); + } + THEN("a dropped vector's output is silence, without counting it late") { + h.w.flush_reports(); // what was late during the stall has been reported + const auto outputs = h.push({33.0}, 1); // vector 30's output is due: dropped + h.wait_for(33); + CHECK(all_equal(outputs[0], 0.0)); + h.w.flush_reports(); + CHECK(h.lines_containing("late for") == 1); + } + THEN("afterwards it runs at the same latency") { + h.push({33.0}, 1); + h.wait_for(33); + h.push({34.0}, 1); + h.wait_for(34); + const auto outputs = h.push({35.0}, 1); // vector 34 in, vector 32 (input 33.0) out + CHECK(all_equal(outputs[0], 33.0)); + } +} + +SCENARIO("A worker can be restarted with new settings, and stopped (plan 2.5)") { + ensure_runtime(); + harness h{"stalls"}; + REQUIRE(h.p.load()); + h.w.start(1, 1, k_frames, 2, 48000.0); + h.push({1.0}, 1); + h.push({2.0}, 1); + + WHEN("it is restarted with another vector size and latency") { + h.w.start(1, 1, 128, 3, 48000.0); + THEN("the latency is the new one, and what was queued is gone") { + CHECK(h.w.latency_samples() == 3 * 128); + for (std::uint64_t k = 0; k < 5; ++k) { + const auto outputs = h.push({static_cast(k + 10)}, 1, 128); + h.wait_for(k + 1); + CHECK(all_equal(outputs[0], k < 3 ? 0.0 : static_cast(k - 3 + 10))); + } + } + } + WHEN("it is stopped") { + h.w.stop(); + THEN("it outputs silence, and nothing is running") { + CHECK_FALSE(h.w.running()); + CHECK(h.w.latency_samples() == 0); + CHECK(all_equal(h.push({5.0}, 1)[0], 0.0)); + } + } + WHEN("it is destroyed while the class stalls") { + REQUIRE(h.p.set_attribute("stall", value{0.2})); + h.push({3.0}, 1); + std::this_thread::sleep_for(std::chrono::milliseconds{50}); // into the stall + THEN("stopping waits for the vector in progress") { + h.w.stop(); // joins: returns once the stall is over + CHECK_FALSE(h.w.running()); + } + } +} + +SCENARIO("A reload while the worker runs takes effect without a pause on the audio thread (plan 2.5)") { + ensure_runtime(); + write_script("worker_reload", + "class worker_reload:\n def process(self, x: float) -> float:\n return x\n"); + harness h{"worker_reload"}; + REQUIRE(h.p.load()); + h.w.start(1, 1, k_frames, 2, 48000.0); + + std::atomic running{true}; + std::atomic pushed{0}; + std::thread audio{[&] { // the audio thread, at about real time (64 frames at 48 kHz) + while (running.load()) { + h.push({1.0}, 1); + ++pushed; + std::this_thread::sleep_for(std::chrono::microseconds{1333}); + } + }}; + for (int i = 0; i < 10; ++i) { + write_script("worker_reload", + "class worker_reload:\n def process(self, x: float) -> float:\n return x * " + + std::to_string(i % 2 == 0 ? 2 : 1) + "\n"); + REQUIRE(h.p.load()); + std::this_thread::sleep_for(std::chrono::milliseconds{20}); + } + running = false; + audio.join(); + + const auto queued = pushed.load(); + h.wait_for(queued); // all it was given while reloading + h.push({1.0}, 1); + h.wait_for(queued + 1); + h.push({1.0}, 1); + h.wait_for(queued + 2); + const auto outputs = h.push({1.0}, 1); // vector `queued` out: the first after the reloads + THEN("the last reload is what plays") { + CHECK(all_equal(outputs[0], 1.0)); // the tenth load: x * 1 + } +} diff --git a/docs/PRODUCTION-PLAN.md b/docs/PRODUCTION-PLAN.md index 973ca7d..daddafc 100644 --- a/docs/PRODUCTION-PLAN.md +++ b/docs/PRODUCTION-PLAN.md @@ -150,8 +150,8 @@ while audio ran segfaulted in 5 of 5 runs. - **Underruns.** A slot not done when its outputs are due is output as silence and counted; the worker still processes it when it gets there (so the class's state stays continuous) and discards its outputs, catching up through the backlog, so the latency stays *L*. If it falls - a whole ring behind (the ring holds *L* + a margin), the oldest inputs are dropped, which the - class sees as a gap. Both are reported once per load from the main thread (`flush_reports()`), + a whole ring behind (the ring holds *L* + a margin), new inputs are dropped until it catches + up, which the class sees as a gap. Both are reported once per load from the main thread (`flush_reports()`), never printed on the audio thread. - **`prepare()` and the mode.** The ring is sized when the chain compiles (`dspsetup`: vector size, and the object's inlet and outlet counts), with the worker stopped and restarted around @@ -171,6 +171,20 @@ while audio ran segfaulted in 5 of 5 runs. one Python call per host vector — batching several per call would cut the per-call overhead further (worker mode's other win) at more latency, a later option (`@block`) if measurements show it pays. + - *(a) the core — done:* `core/include/tap/python/worker.h`. The ring is single producer, single + consumer, with a sequence number per slot: the audio thread queues a vector only into a slot + the worker has finished with, so a full ring drops the new vector rather than racing the + worker for an old one; the margin is a quarter of a second, at least 16 vectors. The audio + thread borrows the ring for each vector by taking an atomic pointer, so `start()` and `stop()` + — which a host may call while its old signal chain still runs — wait at most one vector, and + a vector arriving meanwhile is silence. Reports are once per load, by the processor's new + `load_count()`. Core battery: the output is the direct output *L* vectors later (per sample + and per vector, *L* = 1–3, two channels); the host's channels matched to the ring; a stall + goes silent, is reported once, comes back at the same latency, and is reported again after a + reload; a stall past the ring drops exactly the vectors beyond it; restart, stop, and stop + during a stall; ten reloads under an audio thread. Clean under TSan (locally, macOS) — the + Linux CI job runs it too. *(b) the Max object:* `@mode`, `@latency`, `@latencysamples`, the + ring sized in `dspsetup`; runtime test in Max. - [x] **2.6 Shorter reload stalls** — compile outside the swap and hold the GIL only for the swap. *Measured first* (`core/bench/reload_bench.cpp`: an audio thread with Core Audio's real-time scheduling computes 512-sample buffers of 64-sample vectors at 96 kHz while the main thread saves