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.
74    scratch: Mutex<Vec<S>>,
75    /// The dedicated [`StreamWorker`] thread, if one drains this tap.
76    /// `push_frames` unparks it (a lock-free atomic load + unpark) so the
77    /// consumer wakes without polling. Unset when the tap is drained by
78    /// the shared pool instead.
79    waker: OnceLock<Thread>,
80}
81
82/// Builds a stereo tap at a default capacity (32768 frames, ~170 ms at
83/// 192 kHz), so a plugin can hold an `Arc<AudioTap<S>>` as a `#[skip]`
84/// field and still `#[derive(Default)]` its params. Use [`AudioTap::new`]
85/// when the size or channel count needs to differ from the default.
86impl<S: Copy + Send + 'static> Default for AudioTap<S> {
87    fn default() -> Self {
88        Self::new(DEFAULT_TAP_FRAMES, DEFAULT_TAP_CHANNELS)
89    }
90}
91
92impl<S: Copy + Send + 'static> AudioTap<S> {
93    /// Build a tap holding up to `frame_capacity` interleaved frames of
94    /// `channels` each. Size `frame_capacity` for the worst realistic
95    /// consumer-scheduling gap; drop-on-full is the safety net beyond it
96    /// (e.g. the analyzer uses ~32k frames, ~170 ms at 192 kHz).
97    ///
98    /// # Panics
99    ///
100    /// Panics if `channels` is zero.
101    #[must_use]
102    pub fn new(frame_capacity: usize, channels: usize) -> Self {
103        assert!(channels > 0, "AudioTap needs at least one channel");
104        Self {
105            ring: ArrayQueue::new(frame_capacity * channels),
106            channels,
107            scratch: Mutex::new(Vec::new()),
108            waker: OnceLock::new(),
109        }
110    }
111
112    /// Channel count each frame carries.
113    #[must_use]
114    pub fn channels(&self) -> usize {
115        self.channels
116    }
117
118    /// Push interleaved frames from the audio thread. Wait-free. Drops
119    /// whole frames (never a partial one, so a drop can't desync
120    /// channels) once the ring is full.
121    pub fn push_frames(&self, interleaved: &[S]) {
122        for frame in interleaved.chunks_exact(self.channels) {
123            // The producer is the sole pusher and the consumer only frees
124            // space, so free room seen here is a lower bound - if a whole
125            // frame fits now it still fits when we push it.
126            if self.ring.capacity() - self.ring.len() < self.channels {
127                break;
128            }
129            for &sample in frame {
130                let _ = self.ring.push(sample);
131            }
132        }
133        // Wake a dedicated worker, if one is attached. One unpark per
134        // block; the park token means a push landing between the worker's
135        // drain and its next park is never lost.
136        if let Some(worker) = self.waker.get() {
137            worker.unpark();
138        }
139    }
140
141    /// Discard everything buffered. Called off the audio thread (e.g.
142    /// from `reset` on a sample-rate change, so frames captured at the
143    /// old rate aren't analyzed against the new one). Serialized with
144    /// [`Self::drain_with`]; if a drain is mid-flight this is a no-op and
145    /// the drain finishes the stale frames - `reset` should follow with
146    /// the consumer's own state reset.
147    pub fn clear(&self) {
148        let Ok(_guard) = self.scratch.try_lock() else {
149            return;
150        };
151        while self.ring.pop().is_some() {}
152    }
153
154    /// Drain everything buffered and hand it to `f` as one interleaved
155    /// slice. Runs on the consumer (a pool worker); safe to call from a
156    /// `BackgroundTask::run`. If another worker is already draining
157    /// this tap it returns without double-draining - that worker sees the
158    /// same data. `f` is not called when nothing is buffered.
159    pub fn drain_with(&self, mut f: impl FnMut(&[S])) {
160        let Ok(mut scratch) = self.scratch.try_lock() else {
161            return;
162        };
163        scratch.clear();
164        while let Some(sample) = self.ring.pop() {
165            scratch.push(sample);
166        }
167        if !scratch.is_empty() {
168            f(&scratch);
169        }
170    }
171
172    /// Spawn a dedicated thread that drains this tap sequentially. The
173    /// thread parks until [`Self::push_frames`] unparks it, then hands
174    /// every buffered frame to `on_drain` as one interleaved slice, in
175    /// order. Consumer state lives inside the closure - a single owner,
176    /// so no lock. The returned [`StreamWorker`] joins the thread on drop.
177    ///
178    /// This is the dedicated-thread alternative to draining on the shared
179    /// task pool: pick it when the consumer runs continuously and would
180    /// otherwise contend with unrelated pool work. It spends one thread
181    /// per worker, so unlike the pool it does not bound thread growth.
182    ///
183    /// Attach **at most one** worker per tap: the per-sample ring has a
184    /// single consumer, and only one thread can be registered for the
185    /// wake-on-push. Drain a tap with either a `StreamWorker` or the shared
186    /// pool, not both.
187    ///
188    /// # Panics
189    ///
190    /// Panics if a worker is already attached to this tap, or if the OS
191    /// refuses to spawn the worker thread.
192    #[must_use]
193    pub fn spawn_worker(
194        self: Arc<Self>,
195        name: &str,
196        mut on_drain: impl FnMut(&[S]) + Send + 'static,
197    ) -> StreamWorker {
198        // Fail loud on a second attach rather than silently registering no
199        // waker (which would fall back to the park timeout and race the
200        // first worker on the drain lock). Checked before spawning so a
201        // rejected call leaks no thread.
202        assert!(
203            self.waker.get().is_none(),
204            "AudioTap already has a StreamWorker; attach at most one worker per tap",
205        );
206        let tap = Arc::clone(&self);
207        let shutdown = Arc::new(AtomicBool::new(false));
208        let shutdown_thread = Arc::clone(&shutdown);
209        let handle = thread::Builder::new()
210            .name(name.to_owned())
211            .spawn(move || {
212                while !shutdown_thread.load(Ordering::Acquire) {
213                    tap.drain_with(&mut on_drain);
214                    thread::park_timeout(WORKER_PARK_TIMEOUT);
215                }
216            })
217            .expect("spawn truce stream worker");
218        // Publish the thread so `push_frames` can unpark it. Set from the
219        // spawner (not the worker itself) so it is in place before the
220        // first push; a push before it lands just relies on the timeout.
221        let _ = self.waker.set(handle.thread().clone());
222        StreamWorker {
223            shutdown,
224            thread: handle.thread().clone(),
225            handle: Some(handle),
226        }
227    }
228}
229
230/// Handle to a dedicated [`AudioTap`] consumer thread spawned by
231/// [`AudioTap::spawn_worker`]. Dropping it asks the thread to stop and
232/// joins it, so the worker's lifetime is tied to whatever owns the handle
233/// (typically a `#[skip]` field alongside the tap).
234pub struct StreamWorker {
235    shutdown: Arc<AtomicBool>,
236    thread: Thread,
237    handle: Option<JoinHandle<()>>,
238}
239
240impl Drop for StreamWorker {
241    fn drop(&mut self) {
242        self.shutdown.store(true, Ordering::Release);
243        // Wake the thread out of its park so it sees the flag now rather
244        // than after the timeout.
245        self.thread.unpark();
246        if let Some(handle) = self.handle.take() {
247            let _ = handle.join();
248        }
249    }
250}
251
252#[cfg(test)]
253mod tests {
254    use super::*;
255
256    #[test]
257    fn round_trips_frames() {
258        let tap = AudioTap::<f32>::new(16, 2);
259        tap.push_frames(&[1.0, 2.0, 3.0, 4.0]); // two stereo frames
260        let mut got = Vec::new();
261        tap.drain_with(|chunk| got.extend_from_slice(chunk));
262        assert_eq!(got, vec![1.0, 2.0, 3.0, 4.0]);
263    }
264
265    #[test]
266    fn drop_on_full_stays_frame_aligned() {
267        // Capacity 2 frames (4 samples). Push 4 frames; the last 2 drop.
268        let tap = AudioTap::<i32>::new(2, 2);
269        tap.push_frames(&[1, 1, 2, 2, 3, 3, 4, 4]);
270        let mut got = Vec::new();
271        tap.drain_with(|chunk| got.extend_from_slice(chunk));
272        // Whatever survived is a whole number of frames (even length),
273        // and never a split frame.
274        assert_eq!(got.len() % 2, 0);
275        assert_eq!(got, vec![1, 1, 2, 2]);
276    }
277
278    #[test]
279    fn drain_of_empty_tap_does_not_call_f() {
280        let tap = AudioTap::<f32>::new(4, 1);
281        let mut called = false;
282        tap.drain_with(|_| called = true);
283        assert!(!called, "no callback when nothing is buffered");
284    }
285
286    #[test]
287    fn stream_worker_drains_in_order() {
288        use std::sync::mpsc;
289
290        let tap = Arc::new(AudioTap::<i32>::new(64, 2));
291        let (tx, rx) = mpsc::channel();
292        let worker = tap
293            .clone()
294            .spawn_worker("test-stream-worker", move |chunk| {
295                for &sample in chunk {
296                    let _ = tx.send(sample);
297                }
298            });
299        tap.push_frames(&[1, 2, 3, 4]);
300
301        // The worker drains FIFO, so the samples arrive in push order.
302        let mut got = Vec::new();
303        for _ in 0..4 {
304            got.push(
305                rx.recv_timeout(Duration::from_secs(5))
306                    .expect("worker drained the pushed frames"),
307            );
308        }
309        assert_eq!(got, vec![1, 2, 3, 4]);
310        drop(worker);
311    }
312
313    #[test]
314    fn default_is_a_usable_stereo_tap() {
315        let tap = AudioTap::<f32>::default();
316        assert_eq!(tap.channels(), DEFAULT_TAP_CHANNELS);
317        tap.push_frames(&[1.0, 2.0]);
318        let mut got = Vec::new();
319        tap.drain_with(|chunk| got.extend_from_slice(chunk));
320        assert_eq!(got, vec![1.0, 2.0]);
321    }
322
323    #[test]
324    #[should_panic(expected = "at most one worker per tap")]
325    fn second_worker_panics() {
326        let tap = Arc::new(AudioTap::<i32>::new(16, 2));
327        let _first = tap.clone().spawn_worker("first", |_| {});
328        // A second worker on the same tap is misuse and must fail loudly.
329        let _second = tap.clone().spawn_worker("second", |_| {});
330    }
331}