1#![allow(dead_code)]
44
45use std::collections::VecDeque;
46use std::sync::atomic::{AtomicU32, AtomicUsize, Ordering};
47
48#[derive(Debug, Clone, Copy)]
50pub struct StreamBufferConfig {
51 pub max_frames: usize,
53 pub sample_rate: u32,
55 pub channels: u16,
57}
58
59impl StreamBufferConfig {
60 #[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 #[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#[derive(Debug, Clone)]
89pub struct StreamFrame {
90 pub samples: Vec<f32>,
92 pub pts_samples: u64,
94 pub channels: u16,
96 pub sample_rate: u32,
98}
99
100impl StreamFrame {
101 #[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 #[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 #[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#[derive(Debug)]
136pub struct StreamBuffer {
137 queue: VecDeque<StreamFrame>,
138 config: StreamBufferConfig,
139 total_pushed: u64,
141 total_popped: u64,
143}
144
145impl StreamBuffer {
146 #[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 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 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 #[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 #[must_use]
187 pub fn len(&self) -> usize {
188 self.queue.len()
189 }
190
191 #[must_use]
193 pub fn is_empty(&self) -> bool {
194 self.queue.is_empty()
195 }
196
197 #[must_use]
199 pub fn is_full(&self) -> bool {
200 self.queue.len() >= self.config.max_frames
201 }
202
203 #[must_use]
205 pub fn total_pushed(&self) -> u64 {
206 self.total_pushed
207 }
208
209 #[must_use]
211 pub fn total_popped(&self) -> u64 {
212 self.total_popped
213 }
214}
215
216pub struct LockFreeRingBuffer {
241 data: Vec<AtomicU32>,
243 head: AtomicUsize,
245 tail: AtomicUsize,
247 cap: usize,
249}
250
251impl LockFreeRingBuffer {
255 #[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 #[must_use]
274 pub fn capacity(&self) -> usize {
275 self.cap
276 }
277
278 #[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 #[must_use]
292 pub fn free(&self) -> usize {
293 self.cap - 1 - self.available()
294 }
295
296 #[must_use]
298 pub fn is_empty(&self) -> bool {
299 self.head.load(Ordering::Acquire) == self.tail.load(Ordering::Acquire)
300 }
301
302 #[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 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; }
321 self.data[head].store(sample.to_bits(), Ordering::Relaxed);
324 self.head.store(next_head, Ordering::Release);
325 true
326 }
327
328 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; }
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 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 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 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#[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], 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); 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)); buf.push_frame(make_frame(480, 480)); 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 #[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); 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)); }
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); 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 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); 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 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); 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}