Skip to main content

oximedia_audio/
stream_buffer.rs

1//! Lock-free ring-queue stream buffer for audio frame pipelines.
2//!
3//! This module provides two buffer implementations:
4//!
5//! - [`StreamBuffer`]: A simple FIFO queue suitable for single-threaded use.
6//! - [`LockFreeRingBuffer`]: A single-producer single-consumer (SPSC) lock-free
7//!   ring buffer designed for real-time audio threading.  Uses `AtomicUsize`
8//!   sequence numbers to coordinate access without a mutex.
9//!
10//! # Lock-free design
11//!
12//! [`LockFreeRingBuffer`] can safely be shared between exactly **one writer
13//! thread** (audio callback or capture thread) and **one reader thread**
14//! (processing or playback thread) without any locking.  The implementation
15//! follows the classic SPSC ring-buffer pattern:
16//!
17//! 1. `head` is updated only by the **producer** (writer).
18//! 2. `tail` is updated only by the **consumer** (reader).
19//! 3. Both indices are `AtomicUsize` accessed with `Acquire`/`Release`
20//!    ordering to ensure the data written by the producer is visible to the
21//!    consumer.
22//!
23//! The effective capacity is `capacity - 1` samples to distinguish between
24//! full and empty states without an extra flag.
25//!
26//! # Example
27//!
28//! ```
29//! use oximedia_audio::stream_buffer::LockFreeRingBuffer;
30//! use std::sync::Arc;
31//!
32//! let buf = Arc::new(LockFreeRingBuffer::new(1024));
33//!
34//! // Producer side (audio callback thread)
35//! let samples = vec![0.0_f32; 256];
36//! buf.write_samples(&samples);
37//!
38//! // Consumer side (processing thread)
39//! let mut out = vec![0.0_f32; 256];
40//! let n = buf.read_samples(&mut out);
41//! assert_eq!(n, 256);
42//! ```
43#![allow(dead_code)]
44
45use std::collections::VecDeque;
46use std::sync::atomic::{AtomicU32, AtomicUsize, Ordering};
47
48/// Configuration for a [`StreamBuffer`].
49#[derive(Debug, Clone, Copy)]
50pub struct StreamBufferConfig {
51    /// Maximum number of frames the buffer may hold.
52    pub max_frames: usize,
53    /// Sample rate in Hz (used for duration calculations).
54    pub sample_rate: u32,
55    /// Number of channels per frame.
56    pub channels: u16,
57}
58
59impl StreamBufferConfig {
60    /// Create a new configuration.
61    #[must_use]
62    pub fn new(max_frames: usize, sample_rate: u32, channels: u16) -> Self {
63        Self {
64            max_frames,
65            sample_rate,
66            channels,
67        }
68    }
69
70    /// Maximum queue depth expressed in frames.
71    #[must_use]
72    pub fn max_frames(&self) -> usize {
73        self.max_frames
74    }
75}
76
77impl Default for StreamBufferConfig {
78    fn default() -> Self {
79        Self {
80            max_frames: 64,
81            sample_rate: 48_000,
82            channels: 2,
83        }
84    }
85}
86
87/// A single audio frame inside the stream buffer.
88#[derive(Debug, Clone)]
89pub struct StreamFrame {
90    /// Interleaved PCM samples (f32).
91    pub samples: Vec<f32>,
92    /// Presentation timestamp in samples since stream start.
93    pub pts_samples: u64,
94    /// Number of channels in this frame.
95    pub channels: u16,
96    /// Sample rate of the frame (Hz).
97    pub sample_rate: u32,
98}
99
100impl StreamFrame {
101    /// Create a new frame.
102    #[must_use]
103    pub fn new(samples: Vec<f32>, pts_samples: u64, channels: u16, sample_rate: u32) -> Self {
104        Self {
105            samples,
106            pts_samples,
107            channels,
108            sample_rate,
109        }
110    }
111
112    /// Number of multi-channel audio samples (frames) in this buffer.
113    ///
114    /// That is, the length of `samples` divided by the channel count.
115    #[must_use]
116    pub fn sample_count(&self) -> usize {
117        if self.channels == 0 {
118            return 0;
119        }
120        self.samples.len() / self.channels as usize
121    }
122
123    /// Duration of this frame in milliseconds.
124    #[allow(clippy::cast_precision_loss)]
125    #[must_use]
126    pub fn duration_ms(&self) -> f64 {
127        if self.sample_rate == 0 {
128            return 0.0;
129        }
130        self.sample_count() as f64 / self.sample_rate as f64 * 1_000.0
131    }
132}
133
134/// FIFO queue of [`StreamFrame`]s with a configurable capacity.
135#[derive(Debug)]
136pub struct StreamBuffer {
137    queue: VecDeque<StreamFrame>,
138    config: StreamBufferConfig,
139    /// Total number of frames ever pushed (monotonically increasing).
140    total_pushed: u64,
141    /// Total number of frames ever popped.
142    total_popped: u64,
143}
144
145impl StreamBuffer {
146    /// Create a new stream buffer with the given configuration.
147    #[must_use]
148    pub fn new(config: StreamBufferConfig) -> Self {
149        Self {
150            queue: VecDeque::with_capacity(config.max_frames),
151            config,
152            total_pushed: 0,
153            total_popped: 0,
154        }
155    }
156
157    /// Push a frame into the buffer.
158    ///
159    /// Returns `false` (and discards the frame) when the buffer is full.
160    pub fn push_frame(&mut self, frame: StreamFrame) -> bool {
161        if self.queue.len() >= self.config.max_frames {
162            return false;
163        }
164        self.queue.push_back(frame);
165        self.total_pushed += 1;
166        true
167    }
168
169    /// Pop the oldest frame from the buffer, or `None` if empty.
170    pub fn pop_frame(&mut self) -> Option<StreamFrame> {
171        let frame = self.queue.pop_front();
172        if frame.is_some() {
173            self.total_popped += 1;
174        }
175        frame
176    }
177
178    /// Approximate total buffered audio in milliseconds.
179    #[allow(clippy::cast_precision_loss)]
180    #[must_use]
181    pub fn duration_ms(&self) -> f64 {
182        self.queue.iter().map(|f| f.duration_ms()).sum()
183    }
184
185    /// Number of frames currently queued.
186    #[must_use]
187    pub fn len(&self) -> usize {
188        self.queue.len()
189    }
190
191    /// Returns `true` when no frames are queued.
192    #[must_use]
193    pub fn is_empty(&self) -> bool {
194        self.queue.is_empty()
195    }
196
197    /// Returns `true` when the queue has reached its configured maximum.
198    #[must_use]
199    pub fn is_full(&self) -> bool {
200        self.queue.len() >= self.config.max_frames
201    }
202
203    /// Total frames pushed since creation.
204    #[must_use]
205    pub fn total_pushed(&self) -> u64 {
206        self.total_pushed
207    }
208
209    /// Total frames popped since creation.
210    #[must_use]
211    pub fn total_popped(&self) -> u64 {
212        self.total_popped
213    }
214}
215
216// ─────────────────────────────────────────────────────────────────────────────
217// Lock-free SPSC ring buffer
218// ─────────────────────────────────────────────────────────────────────────────
219
220/// Single-producer single-consumer lock-free ring buffer for `f32` samples.
221///
222/// Designed for real-time audio pipelines where one thread writes audio data
223/// (the audio callback or capture thread) and another thread reads it (the
224/// processing or playback thread).
225///
226/// Samples are stored as their `u32` bit patterns in `AtomicU32` cells so
227/// that the entire structure is `Send + Sync` without any `unsafe` code.
228/// The SPSC ring-buffer protocol using `AtomicUsize` head/tail indices
229/// ensures correct ordering between producer and consumer.
230///
231/// ## Capacity
232///
233/// The actual number of samples that can be buffered is `capacity - 1`.
234/// Choose a power-of-two capacity (e.g., 2048, 4096) for best performance.
235///
236/// ## Thread safety
237///
238/// Only one writer and one reader are supported.  Using more than one writer
239/// or more than one reader concurrently results in incorrect data ordering.
240pub struct LockFreeRingBuffer {
241    /// Internal sample storage as atomic u32 (f32 bit patterns).
242    data: Vec<AtomicU32>,
243    /// Write index (producer-owned).
244    head: AtomicUsize,
245    /// Read index (consumer-owned).
246    tail: AtomicUsize,
247    /// Capacity (length of `data`).
248    cap: usize,
249}
250
251// `Vec<AtomicU32>` is already `Send + Sync`, so `LockFreeRingBuffer` is too.
252// No `unsafe impl` blocks are required.
253
254impl LockFreeRingBuffer {
255    /// Create a new ring buffer that can hold up to `capacity - 1` samples.
256    ///
257    /// `capacity` should be a power of two for best performance.  A minimum
258    /// capacity of 2 is enforced.
259    #[must_use]
260    pub fn new(capacity: usize) -> Self {
261        let cap = capacity.max(2);
262        let data = (0..cap).map(|_| AtomicU32::new(0)).collect();
263        Self {
264            data,
265            head: AtomicUsize::new(0),
266            tail: AtomicUsize::new(0),
267            cap,
268        }
269    }
270
271    /// Returns the total capacity of the buffer (number of samples that can
272    /// ever be stored).  The usable capacity is `capacity() - 1`.
273    #[must_use]
274    pub fn capacity(&self) -> usize {
275        self.cap
276    }
277
278    /// Number of samples currently available for reading.
279    #[must_use]
280    pub fn available(&self) -> usize {
281        let head = self.head.load(Ordering::Acquire);
282        let tail = self.tail.load(Ordering::Acquire);
283        if head >= tail {
284            head - tail
285        } else {
286            self.cap - tail + head
287        }
288    }
289
290    /// Number of free slots available for writing.
291    #[must_use]
292    pub fn free(&self) -> usize {
293        self.cap - 1 - self.available()
294    }
295
296    /// Returns `true` when the buffer contains no readable samples.
297    #[must_use]
298    pub fn is_empty(&self) -> bool {
299        self.head.load(Ordering::Acquire) == self.tail.load(Ordering::Acquire)
300    }
301
302    /// Returns `true` when the buffer is full (cannot accept more writes).
303    #[must_use]
304    pub fn is_full(&self) -> bool {
305        let head = self.head.load(Ordering::Acquire);
306        let tail = self.tail.load(Ordering::Acquire);
307        (head + 1) % self.cap == tail
308    }
309
310    /// Write a single sample.
311    ///
312    /// Returns `true` when successful, `false` when the buffer is full.
313    ///
314    /// **Must only be called from the producer thread.**
315    pub fn write(&self, sample: f32) -> bool {
316        let head = self.head.load(Ordering::Relaxed);
317        let next_head = (head + 1) % self.cap;
318        if next_head == self.tail.load(Ordering::Acquire) {
319            return false; // full
320        }
321        // Store the f32 bit pattern atomically.  Only the producer writes to
322        // data[head], and head < cap, so the index is always in bounds.
323        self.data[head].store(sample.to_bits(), Ordering::Relaxed);
324        self.head.store(next_head, Ordering::Release);
325        true
326    }
327
328    /// Read a single sample.
329    ///
330    /// Returns `Some(sample)` when data is available, `None` when empty.
331    ///
332    /// **Must only be called from the consumer thread.**
333    pub fn read(&self) -> Option<f32> {
334        let tail = self.tail.load(Ordering::Relaxed);
335        if tail == self.head.load(Ordering::Acquire) {
336            return None; // empty
337        }
338        let bits = self.data[tail].load(Ordering::Relaxed);
339        self.tail.store((tail + 1) % self.cap, Ordering::Release);
340        Some(f32::from_bits(bits))
341    }
342
343    /// Write a block of samples.
344    ///
345    /// Returns the number of samples actually written (may be less than
346    /// `samples.len()` when the buffer does not have enough free space).
347    ///
348    /// **Must only be called from the producer thread.**
349    pub fn write_samples(&self, samples: &[f32]) -> usize {
350        let mut written = 0;
351        for &s in samples {
352            if !self.write(s) {
353                break;
354            }
355            written += 1;
356        }
357        written
358    }
359
360    /// Read samples into `dst`.
361    ///
362    /// Returns the number of samples actually read (may be less than
363    /// `dst.len()` when the buffer does not have enough data).
364    ///
365    /// **Must only be called from the consumer thread.**
366    pub fn read_samples(&self, dst: &mut [f32]) -> usize {
367        let mut read = 0;
368        for slot in dst.iter_mut() {
369            match self.read() {
370                Some(s) => {
371                    *slot = s;
372                    read += 1;
373                }
374                None => break,
375            }
376        }
377        read
378    }
379
380    /// Clear all samples from the buffer.
381    ///
382    /// **Must only be called when no concurrent reads/writes are in progress.**
383    pub fn clear(&self) {
384        let head = self.head.load(Ordering::Relaxed);
385        self.tail.store(head, Ordering::Release);
386    }
387}
388
389impl std::fmt::Debug for LockFreeRingBuffer {
390    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
391        f.debug_struct("LockFreeRingBuffer")
392            .field("cap", &self.cap)
393            .field("available", &self.available())
394            .field("free", &self.free())
395            .finish()
396    }
397}
398
399// ─────────────────────────────────────────────────────────────────────────────
400// Unit tests
401// ─────────────────────────────────────────────────────────────────────────────
402
403#[cfg(test)]
404mod tests {
405    use super::*;
406    use std::sync::Arc;
407
408    fn make_frame(n_samples: usize, pts: u64) -> StreamFrame {
409        StreamFrame::new(
410            vec![0.0_f32; n_samples * 2], // stereo
411            pts,
412            2,
413            48_000,
414        )
415    }
416
417    #[test]
418    fn test_config_max_frames() {
419        let cfg = StreamBufferConfig::new(32, 48_000, 2);
420        assert_eq!(cfg.max_frames(), 32);
421    }
422
423    #[test]
424    fn test_config_default() {
425        let cfg = StreamBufferConfig::default();
426        assert_eq!(cfg.max_frames, 64);
427        assert_eq!(cfg.sample_rate, 48_000);
428    }
429
430    #[test]
431    fn test_frame_sample_count() {
432        let frame = make_frame(480, 0);
433        assert_eq!(frame.sample_count(), 480);
434    }
435
436    #[test]
437    fn test_frame_duration_ms() {
438        let frame = make_frame(480, 0); // 480/48000 = 10ms
439        let dur = frame.duration_ms();
440        assert!((dur - 10.0).abs() < 0.001);
441    }
442
443    #[test]
444    fn test_frame_duration_zero_rate() {
445        let frame = StreamFrame::new(vec![0.0; 4], 0, 2, 0);
446        assert_eq!(frame.duration_ms(), 0.0);
447    }
448
449    #[test]
450    fn test_buffer_push_and_pop() {
451        let cfg = StreamBufferConfig::default();
452        let mut buf = StreamBuffer::new(cfg);
453        let f = make_frame(480, 0);
454        assert!(buf.push_frame(f));
455        assert_eq!(buf.len(), 1);
456        let popped = buf.pop_frame();
457        assert!(popped.is_some());
458        assert!(buf.is_empty());
459    }
460
461    #[test]
462    fn test_buffer_fifo_order() {
463        let cfg = StreamBufferConfig::default();
464        let mut buf = StreamBuffer::new(cfg);
465        buf.push_frame(make_frame(480, 0));
466        buf.push_frame(make_frame(480, 480));
467        let first = buf.pop_frame().expect("should succeed");
468        assert_eq!(first.pts_samples, 0);
469        let second = buf.pop_frame().expect("should succeed");
470        assert_eq!(second.pts_samples, 480);
471    }
472
473    #[test]
474    fn test_buffer_full_rejects_push() {
475        let cfg = StreamBufferConfig::new(2, 48_000, 2);
476        let mut buf = StreamBuffer::new(cfg);
477        assert!(buf.push_frame(make_frame(480, 0)));
478        assert!(buf.push_frame(make_frame(480, 480)));
479        assert!(buf.is_full());
480        assert!(!buf.push_frame(make_frame(480, 960)));
481    }
482
483    #[test]
484    fn test_buffer_is_empty_initially() {
485        let buf = StreamBuffer::new(StreamBufferConfig::default());
486        assert!(buf.is_empty());
487    }
488
489    #[test]
490    fn test_buffer_pop_empty_returns_none() {
491        let mut buf = StreamBuffer::new(StreamBufferConfig::default());
492        assert!(buf.pop_frame().is_none());
493    }
494
495    #[test]
496    fn test_buffer_duration_ms() {
497        let cfg = StreamBufferConfig::default();
498        let mut buf = StreamBuffer::new(cfg);
499        buf.push_frame(make_frame(480, 0)); // 10ms
500        buf.push_frame(make_frame(480, 480)); // 10ms
501        let dur = buf.duration_ms();
502        assert!((dur - 20.0).abs() < 0.01);
503    }
504
505    #[test]
506    fn test_buffer_total_counters() {
507        let cfg = StreamBufferConfig::default();
508        let mut buf = StreamBuffer::new(cfg);
509        buf.push_frame(make_frame(480, 0));
510        buf.push_frame(make_frame(480, 480));
511        buf.pop_frame();
512        assert_eq!(buf.total_pushed(), 2);
513        assert_eq!(buf.total_popped(), 1);
514    }
515
516    #[test]
517    fn test_frame_zero_channels() {
518        let frame = StreamFrame::new(vec![0.0; 10], 0, 0, 48_000);
519        assert_eq!(frame.sample_count(), 0);
520    }
521
522    // ── LockFreeRingBuffer tests ───────────────────────────────────────────────
523
524    #[test]
525    fn test_ringbuf_initially_empty() {
526        let rb = LockFreeRingBuffer::new(16);
527        assert!(rb.is_empty());
528        assert_eq!(rb.available(), 0);
529    }
530
531    #[test]
532    fn test_ringbuf_write_and_read_single() {
533        let rb = LockFreeRingBuffer::new(16);
534        assert!(rb.write(0.5));
535        assert!(!rb.is_empty());
536        let s = rb.read().expect("should have data");
537        assert!((s - 0.5).abs() < 1e-7);
538        assert!(rb.is_empty());
539    }
540
541    #[test]
542    fn test_ringbuf_write_block_read_block() {
543        let rb = LockFreeRingBuffer::new(64);
544        let data: Vec<f32> = (0..32).map(|i| i as f32 * 0.1).collect();
545        let written = rb.write_samples(&data);
546        assert_eq!(written, 32);
547        let mut dst = vec![0.0_f32; 32];
548        let read = rb.read_samples(&mut dst);
549        assert_eq!(read, 32);
550        for (a, b) in data.iter().zip(dst.iter()) {
551            assert!((a - b).abs() < 1e-6);
552        }
553    }
554
555    #[test]
556    fn test_ringbuf_full_rejects_write() {
557        let rb = LockFreeRingBuffer::new(4); // usable cap = 3
558        assert!(rb.write(1.0));
559        assert!(rb.write(2.0));
560        assert!(rb.write(3.0));
561        assert!(rb.is_full());
562        assert!(!rb.write(4.0)); // should fail
563    }
564
565    #[test]
566    fn test_ringbuf_read_empty_returns_none() {
567        let rb = LockFreeRingBuffer::new(16);
568        assert!(rb.read().is_none());
569    }
570
571    #[test]
572    fn test_ringbuf_wrap_around() {
573        let rb = LockFreeRingBuffer::new(8); // usable cap = 7
574                                             // Fill then partially drain, then fill again — exercises wrap-around
575        for i in 0..7 {
576            rb.write(i as f32);
577        }
578        for _ in 0..4 {
579            rb.read();
580        }
581        for i in 0..4 {
582            assert!(rb.write(i as f32 + 10.0));
583        }
584        // Remaining: 4 from original + 4 new = 7; but usable cap = 7, so last write should fail
585        assert_eq!(rb.available(), 7);
586    }
587
588    #[test]
589    fn test_ringbuf_clear() {
590        let rb = LockFreeRingBuffer::new(16);
591        rb.write_samples(&[1.0, 2.0, 3.0]);
592        assert_eq!(rb.available(), 3);
593        rb.clear();
594        assert!(rb.is_empty());
595    }
596
597    #[test]
598    fn test_ringbuf_capacity() {
599        let rb = LockFreeRingBuffer::new(32);
600        assert_eq!(rb.capacity(), 32);
601    }
602
603    #[test]
604    fn test_ringbuf_free() {
605        let rb = LockFreeRingBuffer::new(16); // usable = 15
606        rb.write(0.5);
607        assert_eq!(rb.free(), 14);
608    }
609
610    #[test]
611    fn test_ringbuf_fifo_order() {
612        let rb = LockFreeRingBuffer::new(16);
613        rb.write(1.0);
614        rb.write(2.0);
615        rb.write(3.0);
616        assert_eq!(rb.read().expect("1"), 1.0);
617        assert_eq!(rb.read().expect("2"), 2.0);
618        assert_eq!(rb.read().expect("3"), 3.0);
619    }
620
621    #[test]
622    fn test_ringbuf_arc_shared() {
623        let rb = Arc::new(LockFreeRingBuffer::new(64));
624        let rb2 = Arc::clone(&rb);
625        // Simulate producer/consumer in the same thread for determinism
626        rb.write_samples(&[0.1, 0.2, 0.3]);
627        let mut out = vec![0.0_f32; 3];
628        rb2.read_samples(&mut out);
629        assert!((out[0] - 0.1).abs() < 1e-6);
630    }
631
632    #[test]
633    fn test_ringbuf_minimum_capacity_enforced() {
634        let rb = LockFreeRingBuffer::new(0); // should be clamped to 2
635        assert_eq!(rb.capacity(), 2);
636    }
637
638    #[test]
639    fn test_ringbuf_debug_format() {
640        let rb = LockFreeRingBuffer::new(16);
641        let s = format!("{rb:?}");
642        assert!(s.contains("LockFreeRingBuffer"));
643    }
644}