Skip to main content

truce_core/
audio_tap.rs

1//! Audio tap: stream interleaved frames off the audio thread to a
2//! background consumer.
3//!
4//! The audio thread pushes frames wait-free from `process`; a background
5//! consumer drains them. This is the *stream + analyze* half of the
6//! worker pattern (a spectrum analyzer, a loudness meter, an
7//! oscilloscope), where [`crate::tasks`]'s discrete task pool is the
8//! *construction offload* half.
9//!
10//! Two ways to drain:
11//! - a `BackgroundTask` handler (in `truce_plugin`) on the shared pool,
12//!   woken by a coalescing task each block - bounded threads, best when
13//!   analysis is bursty or coalescable.
14//! - a dedicated [`StreamWorker`] via [`AudioTap::spawn_worker`] - one
15//!   owned thread that parks on the tap and drains sequentially. Consumer
16//!   state lives thread-local (no lock), and it never stalls on unrelated
17//!   pool work, at the cost of one thread per worker.
18//!
19//! The tap is a plain `#[skip]` field the plugin owns (usually
20//! `Arc<AudioTap<S>>`), so both `process` (via `&params`) and the
21//! handler (via `&params`) reach it - the same shared-`Arc` mechanism as
22//! any audio -> worker channel. It composes with the pool rather than
23//! introducing new threading:
24//!
25//! ```ignore
26//! #[derive(Params)]
27//! struct AnalyzerParams {
28//!     #[skip]
29//!     tap: Arc<AudioTap<f32>>,        // Default: a stereo tap at the default capacity
30//!     // ... published spectrum atoms (also #[skip]) ...
31//! }
32//!
33//! // process (audio thread):
34//! params.tap.push_frames(&interleaved);
35//! if let Some(t) = ctx.tasks::<Analyze>() { t.spawn_coalescing(Analyze); }
36//!
37//! // run (pool thread):
38//! params.tap.drain_with(|chunk| { /* run the FFT, publish */ });
39//! ```
40
41use std::sync::atomic::{AtomicBool, Ordering};
42use std::sync::{Arc, Mutex, OnceLock};
43use std::thread::{self, JoinHandle, Thread};
44use std::time::Duration;
45
46use crossbeam_queue::ArrayQueue;
47
48/// Safety-net wake interval for a dedicated [`StreamWorker`]. Normal
49/// wakeups come from `push_frames` unparking the thread; this only bounds
50/// how long a missed unpark or a pending shutdown can go unnoticed.
51const WORKER_PARK_TIMEOUT: Duration = Duration::from_millis(100);
52
53/// Frame capacity [`AudioTap::default`] builds: 32768 stereo frames, a
54/// generous consumer-scheduling margin (~170 ms at 192 kHz, ~740 ms at
55/// 44.1 kHz). Call [`AudioTap::new`] for a different size or channel count.
56const DEFAULT_TAP_FRAMES: usize = 32 * 1024;
57/// Channel count [`AudioTap::default`] builds. Stereo is the common case;
58/// call [`AudioTap::new`] for mono or higher channel counts.
59const DEFAULT_TAP_CHANNELS: usize = 2;
60
61/// A lock-free audio tap. The audio thread is the sole producer
62/// ([`Self::push_frames`], wait-free); a background consumer drains it
63/// ([`Self::drain_with`]). Whole-frame drop-on-full keeps a drop from
64/// desyncing channels.
65pub struct AudioTap<S> {
66    ring: ArrayQueue<S>,
67    channels: usize,
68    /// Reused pop buffer whose lock also serializes drains: with a shared
69    /// pool a second worker could pick up a coalesced drain while the
70    /// first is still running, and a per-sample ring must have one
71    /// consumer at a time. `try_lock` makes the second worker bow out
72    /// (the first is already draining the same data). Locked only on the
73    /// pool thread, never on the audio thread. Also retains any trailing
74    /// partial frame between drains (see [`Self::drain_with`]): the
75    /// producer pushes a frame's channels one at a time, so a concurrent
76    /// drain can catch it mid-frame, and the leftover reassembles next call.
77    scratch: Mutex<Vec<S>>,
78    /// The dedicated [`StreamWorker`] thread, if one drains this tap.
79    /// `push_frames` unparks it (a lock-free atomic load + unpark) so the
80    /// consumer wakes without polling. Unset when the tap is drained by
81    /// the shared pool instead.
82    waker: OnceLock<Thread>,
83}
84
85/// Builds a stereo tap at a default capacity (32768 frames, ~170 ms at
86/// 192 kHz), so a plugin can hold an `Arc<AudioTap<S>>` as a `#[skip]`
87/// field and still `#[derive(Default)]` its params. Use [`AudioTap::new`]
88/// when the size or channel count needs to differ from the default.
89impl<S: Copy + Send + 'static> Default for AudioTap<S> {
90    fn default() -> Self {
91        Self::new(DEFAULT_TAP_FRAMES, DEFAULT_TAP_CHANNELS)
92    }
93}
94
95impl<S: Copy + Send + 'static> AudioTap<S> {
96    /// Build a tap holding up to `frame_capacity` interleaved frames of
97    /// `channels` each. Size `frame_capacity` for the worst realistic
98    /// consumer-scheduling gap; drop-on-full is the safety net beyond it
99    /// (e.g. the analyzer uses ~32k frames, ~170 ms at 192 kHz).
100    ///
101    /// # Panics
102    ///
103    /// Panics if `channels` is zero.
104    #[must_use]
105    pub fn new(frame_capacity: usize, channels: usize) -> Self {
106        assert!(channels > 0, "AudioTap needs at least one channel");
107        Self {
108            ring: ArrayQueue::new(frame_capacity * channels),
109            channels,
110            scratch: Mutex::new(Vec::new()),
111            waker: OnceLock::new(),
112        }
113    }
114
115    /// Channel count each frame carries.
116    #[must_use]
117    pub fn channels(&self) -> usize {
118        self.channels
119    }
120
121    /// Push interleaved frames from the audio thread. Wait-free. Drops
122    /// whole frames (never a partial one, so a drop can't desync
123    /// channels) once the ring is full.
124    pub fn push_frames(&self, interleaved: &[S]) {
125        for frame in interleaved.chunks_exact(self.channels) {
126            // The producer is the sole pusher and the consumer only frees
127            // space, so free room seen here is a lower bound - if a whole
128            // frame fits now it still fits when we push it.
129            if self.ring.capacity() - self.ring.len() < self.channels {
130                break;
131            }
132            for &sample in frame {
133                let _ = self.ring.push(sample);
134            }
135        }
136        // Wake a dedicated worker, if one is attached. One unpark per
137        // block; the park token means a push landing between the worker's
138        // drain and its next park is never lost.
139        if let Some(worker) = self.waker.get() {
140            worker.unpark();
141        }
142    }
143
144    /// Discard everything buffered. Called off the audio thread (e.g.
145    /// from `reset` on a sample-rate change, so frames captured at the
146    /// old rate aren't analyzed against the new one). Serialized with
147    /// [`Self::drain_with`]; if a drain is mid-flight this is a no-op and
148    /// the drain finishes the stale frames - `reset` should follow with
149    /// the consumer's own state reset.
150    pub fn clear(&self) {
151        let Ok(mut guard) = self.scratch.try_lock() else {
152            return;
153        };
154        while self.ring.pop().is_some() {}
155        // Drop any partial frame carried from an interrupted drain, so it
156        // can't prepend stale samples onto post-clear frames.
157        guard.clear();
158    }
159
160    /// Drain the buffered whole frames and hand them to `f` as one
161    /// interleaved slice - always a whole number of frames. Runs on the
162    /// consumer (a pool worker); safe to call from a `BackgroundTask::run`.
163    /// If another worker is already draining this tap it returns without
164    /// double-draining - that worker sees the same data. `f` is not called
165    /// when no whole frame is buffered.
166    ///
167    /// The producer pushes a frame's channels one at a time, so a drain
168    /// running concurrently can catch it mid-frame. Any trailing partial
169    /// frame is held in `scratch` and reassembled on the next drain rather
170    /// than handed to `f` split - a consumer deinterleaving with
171    /// `chunks_exact(channels)` never sees a channel-swapped chunk.
172    pub fn drain_with(&self, mut f: impl FnMut(&[S])) {
173        let Ok(mut scratch) = self.scratch.try_lock() else {
174            return;
175        };
176        // Do NOT clear: `scratch` may hold a partial frame carried from a
177        // previous drain. Append after it so the frame reassembles.
178        while let Some(sample) = self.ring.pop() {
179            scratch.push(sample);
180        }
181        // Hand off only whole frames; keep any trailing partial for the
182        // next call. `drain` shifts the (sub-channel-count) remainder to
183        // the front and preserves the buffer's capacity for reuse.
184        let whole = scratch.len() - scratch.len() % self.channels;
185        if whole > 0 {
186            f(&scratch[..whole]);
187            scratch.drain(..whole);
188        }
189    }
190
191    /// Spawn a dedicated thread that drains this tap sequentially. The
192    /// thread parks until [`Self::push_frames`] unparks it, then hands
193    /// every buffered frame to `on_drain` as one interleaved slice, in
194    /// order. Consumer state lives inside the closure - a single owner,
195    /// so no lock. The returned [`StreamWorker`] joins the thread on drop.
196    ///
197    /// This is the dedicated-thread alternative to draining on the shared
198    /// task pool: pick it when the consumer runs continuously and would
199    /// otherwise contend with unrelated pool work. It spends one thread
200    /// per worker, so unlike the pool it does not bound thread growth.
201    ///
202    /// Attach **at most one** worker per tap: the per-sample ring has a
203    /// single consumer, and only one thread can be registered for the
204    /// wake-on-push. Drain a tap with either a `StreamWorker` or the shared
205    /// pool, not both.
206    ///
207    /// # Panics
208    ///
209    /// Panics if a worker is already attached to this tap, or if the OS
210    /// refuses to spawn the worker thread.
211    #[must_use]
212    pub fn spawn_worker(
213        self: Arc<Self>,
214        name: &str,
215        mut on_drain: impl FnMut(&[S]) + Send + 'static,
216    ) -> StreamWorker {
217        // Fail loud on a second attach rather than silently registering no
218        // waker (which would fall back to the park timeout and race the
219        // first worker on the drain lock). Checked before spawning so a
220        // rejected call leaks no thread.
221        assert!(
222            self.waker.get().is_none(),
223            "AudioTap already has a StreamWorker; attach at most one worker per tap",
224        );
225        let tap = Arc::clone(&self);
226        let shutdown = Arc::new(AtomicBool::new(false));
227        let shutdown_thread = Arc::clone(&shutdown);
228        let handle = thread::Builder::new()
229            .name(name.to_owned())
230            .spawn(move || {
231                while !shutdown_thread.load(Ordering::Acquire) {
232                    tap.drain_with(&mut on_drain);
233                    thread::park_timeout(WORKER_PARK_TIMEOUT);
234                }
235            })
236            .expect("spawn truce stream worker");
237        // Publish the thread so `push_frames` can unpark it. Set from the
238        // spawner (not the worker itself) so it is in place before the
239        // first push; a push before it lands just relies on the timeout.
240        let _ = self.waker.set(handle.thread().clone());
241        StreamWorker {
242            shutdown,
243            thread: handle.thread().clone(),
244            handle: Some(handle),
245        }
246    }
247}
248
249/// Handle to a dedicated [`AudioTap`] consumer thread spawned by
250/// [`AudioTap::spawn_worker`]. Dropping it asks the thread to stop and
251/// joins it, so the worker's lifetime is tied to whatever owns the handle
252/// (typically a `#[skip]` field alongside the tap).
253pub struct StreamWorker {
254    shutdown: Arc<AtomicBool>,
255    thread: Thread,
256    handle: Option<JoinHandle<()>>,
257}
258
259impl Drop for StreamWorker {
260    fn drop(&mut self) {
261        self.shutdown.store(true, Ordering::Release);
262        // Wake the thread out of its park so it sees the flag now rather
263        // than after the timeout.
264        self.thread.unpark();
265        if let Some(handle) = self.handle.take() {
266            let _ = handle.join();
267        }
268    }
269}
270
271#[cfg(test)]
272mod tests {
273    use super::*;
274
275    #[test]
276    fn round_trips_frames() {
277        let tap = AudioTap::<f32>::new(16, 2);
278        tap.push_frames(&[1.0, 2.0, 3.0, 4.0]); // two stereo frames
279        let mut got = Vec::new();
280        tap.drain_with(|chunk| got.extend_from_slice(chunk));
281        assert_eq!(got, vec![1.0, 2.0, 3.0, 4.0]);
282    }
283
284    #[test]
285    fn drop_on_full_stays_frame_aligned() {
286        // Capacity 2 frames (4 samples). Push 4 frames; the last 2 drop.
287        let tap = AudioTap::<i32>::new(2, 2);
288        tap.push_frames(&[1, 1, 2, 2, 3, 3, 4, 4]);
289        let mut got = Vec::new();
290        tap.drain_with(|chunk| got.extend_from_slice(chunk));
291        // Whatever survived is a whole number of frames (even length),
292        // and never a split frame.
293        assert_eq!(got.len() % 2, 0);
294        assert_eq!(got, vec![1, 1, 2, 2]);
295    }
296
297    #[test]
298    fn drain_of_empty_tap_does_not_call_f() {
299        let tap = AudioTap::<f32>::new(4, 1);
300        let mut called = false;
301        tap.drain_with(|_| called = true);
302        assert!(!called, "no callback when nothing is buffered");
303    }
304
305    #[test]
306    fn stream_worker_drains_in_order() {
307        use std::sync::mpsc;
308
309        let tap = Arc::new(AudioTap::<i32>::new(64, 2));
310        let (tx, rx) = mpsc::channel();
311        let worker = tap
312            .clone()
313            .spawn_worker("test-stream-worker", move |chunk| {
314                for &sample in chunk {
315                    let _ = tx.send(sample);
316                }
317            });
318        tap.push_frames(&[1, 2, 3, 4]);
319
320        // The worker drains FIFO, so the samples arrive in push order.
321        let mut got = Vec::new();
322        for _ in 0..4 {
323            got.push(
324                rx.recv_timeout(Duration::from_secs(5))
325                    .expect("worker drained the pushed frames"),
326            );
327        }
328        assert_eq!(got, vec![1, 2, 3, 4]);
329        drop(worker);
330    }
331
332    #[test]
333    fn default_is_a_usable_stereo_tap() {
334        let tap = AudioTap::<f32>::default();
335        assert_eq!(tap.channels(), DEFAULT_TAP_CHANNELS);
336        tap.push_frames(&[1.0, 2.0]);
337        let mut got = Vec::new();
338        tap.drain_with(|chunk| got.extend_from_slice(chunk));
339        assert_eq!(got, vec![1.0, 2.0]);
340    }
341
342    /// The producer pushes a frame's channels one at a time, so a drain
343    /// running concurrently routinely lands between them. Every chunk
344    /// handed to `f` must still be a whole number of frames - otherwise a
345    /// consumer deinterleaving with `chunks_exact(2)` gets a channel-swapped
346    /// chunk. Sizing the ring for the whole run rules out drops, so the
347    /// reassembled stream must also be the exact FIFO sequence.
348    ///
349    /// Several independent rounds: each producer/consumer start is its own
350    /// race, so a regressed drain (handing off partial frames) is caught
351    /// with high probability, while the correct one always passes.
352    ///
353    /// Skipped under Miri: 400k pushes across two threads per round are
354    /// hour-scale in the interpreter, and the mid-frame race it soaks for
355    /// depends on real thread timing Miri's scheduler doesn't reproduce.
356    #[allow(clippy::cast_possible_truncation, clippy::cast_possible_wrap)]
357    #[cfg_attr(
358        miri,
359        ignore = "concurrency soak - too slow under Miri, no timing repro"
360    )]
361    #[test]
362    fn concurrent_drain_hands_out_whole_frames() {
363        use std::time::Instant;
364
365        let frames = 100_000usize;
366        for _round in 0..4 {
367            let tap = Arc::new(AudioTap::<i32>::new(frames, 2));
368
369            let producer = {
370                let tap = Arc::clone(&tap);
371                thread::spawn(move || {
372                    // Frame i is [2i, 2i+1]: L even, R odd, values consecutive
373                    // across the stream, so any split or swap is visible.
374                    for i in 0..frames as i32 {
375                        tap.push_frames(&[2 * i, 2 * i + 1]);
376                    }
377                })
378            };
379
380            let mut got: Vec<i32> = Vec::with_capacity(frames * 2);
381            let start = Instant::now();
382            while got.len() < frames * 2 {
383                tap.drain_with(|chunk| {
384                    assert_eq!(
385                        chunk.len() % 2,
386                        0,
387                        "drain handed a partial frame (len {})",
388                        chunk.len()
389                    );
390                    got.extend_from_slice(chunk);
391                });
392                assert!(start.elapsed() < Duration::from_secs(30), "drain stalled");
393            }
394            producer.join().unwrap();
395
396            // No drops (ring sized for the run) means strict FIFO: sample j is j.
397            assert_eq!(got.len(), frames * 2);
398            for (j, &v) in got.iter().enumerate() {
399                assert_eq!(v, j as i32, "sample {j} out of place - frame misaligned");
400            }
401        }
402    }
403
404    #[test]
405    #[should_panic(expected = "at most one worker per tap")]
406    fn second_worker_panics() {
407        let tap = Arc::new(AudioTap::<i32>::new(16, 2));
408        let _first = tap.clone().spawn_worker("first", |_| {});
409        // A second worker on the same tap is misuse and must fail loudly.
410        let _second = tap.clone().spawn_worker("second", |_| {});
411    }
412}