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. 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}