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 `¶ms`) and the
21//! handler (via `¶ms`) 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}