#[cfg(feature = "playback")]
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
#[cfg(feature = "playback")]
use std::sync::{Arc, Mutex};
#[cfg(feature = "playback")]
use crossbeam_channel::{Receiver, Sender};
#[cfg(feature = "playback")]
use crate::backend::{
AudioBackend, BackendDevice, BackendStream, CpalBackend, OutputDataCallback,
StreamErrorCallback, StreamParams,
};
use crate::device::DeviceSelector;
use crate::error::DecibriError;
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct SpeakerConfig {
pub sample_rate: u32,
pub channels: u16,
pub device: DeviceSelector,
}
impl Default for SpeakerConfig {
fn default() -> Self {
Self {
sample_rate: 16000,
channels: 1,
device: DeviceSelector::Default,
}
}
}
impl SpeakerConfig {
pub fn validate(&self) -> Result<(), DecibriError> {
if !(1000..=384000).contains(&self.sample_rate) {
return Err(DecibriError::SampleRateOutOfRange);
}
if !(1..=32).contains(&self.channels) {
return Err(DecibriError::ChannelsOutOfRange);
}
Ok(())
}
}
#[cfg(feature = "playback")]
fn render_output(
data: &mut [f32],
accum: &mut Vec<f32>,
receiver: &Receiver<Vec<f32>>,
running: &AtomicBool,
sentinels_played: &AtomicUsize,
) {
if !running.load(Ordering::Relaxed) {
accum.clear();
while let Ok(chunk) = receiver.try_recv() {
if chunk.is_empty() {
sentinels_played.fetch_add(1, Ordering::Relaxed);
}
}
data.fill(0.0);
return;
}
let mut written = 0;
if !accum.is_empty() {
let take = accum.len().min(data.len());
data[..take].copy_from_slice(&accum[..take]);
accum.drain(..take);
written += take;
}
while written < data.len() {
match receiver.try_recv() {
Ok(samples) => {
if samples.is_empty() {
sentinels_played.fetch_add(1, Ordering::Relaxed);
continue;
}
let need = data.len() - written;
if samples.len() <= need {
data[written..written + samples.len()].copy_from_slice(&samples);
written += samples.len();
} else {
data[written..].copy_from_slice(&samples[..need]);
accum.extend_from_slice(&samples[need..]);
written = data.len();
}
}
Err(_) => {
data[written..].fill(0.0);
return;
}
}
}
}
#[cfg(feature = "playback")]
fn send_impl(
sender: &Sender<Vec<f32>>,
running: &AtomicBool,
samples: Vec<f32>,
) -> Result<(), DecibriError> {
if samples.is_empty() {
return Ok(());
}
if !running.load(Ordering::Relaxed) {
return Err(DecibriError::SpeakerStreamClosed);
}
sender
.send(samples)
.map_err(|_| DecibriError::SpeakerStreamClosed)
}
#[cfg(feature = "playback")]
fn drain_impl(
sender: &Sender<Vec<f32>>,
running: &AtomicBool,
sentinel_seq: &AtomicUsize,
sentinels_played: &AtomicUsize,
) {
let target = sentinel_seq.fetch_add(1, Ordering::Relaxed) + 1;
if sender.send(Vec::new()).is_err() {
return;
}
while sentinels_played.load(Ordering::Relaxed) < target {
if !running.load(Ordering::Relaxed) {
return;
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
}
#[cfg(feature = "playback")]
pub struct SpeakerStream {
_stream: BackendStream,
sender: Sender<Vec<f32>>,
running: Arc<AtomicBool>,
sentinel_seq: Arc<AtomicUsize>,
sentinels_played: Arc<AtomicUsize>,
last_error: Arc<Mutex<Option<DecibriError>>>,
}
#[cfg(feature = "playback")]
impl SpeakerStream {
pub fn send(&self, samples: Vec<f32>) -> Result<(), DecibriError> {
send_impl(&self.sender, &self.running, samples)
}
pub fn drain(&self) {
drain_impl(
&self.sender,
&self.running,
&self.sentinel_seq,
&self.sentinels_played,
);
}
pub fn is_playing(&self) -> bool {
self.running.load(Ordering::Relaxed)
}
pub fn take_last_error(&self) -> Option<DecibriError> {
self.last_error.lock().ok().and_then(|mut slot| slot.take())
}
pub fn stop(&self) {
self.running.store(false, Ordering::Relaxed);
self._stream.stop();
}
pub fn sink(&self) -> SpeakerSink {
SpeakerSink {
sender: self.sender.clone(),
running: self.running.clone(),
sentinel_seq: self.sentinel_seq.clone(),
sentinels_played: self.sentinels_played.clone(),
}
}
}
#[cfg(feature = "playback")]
impl Drop for SpeakerStream {
fn drop(&mut self) {
self.running.store(false, Ordering::Relaxed);
}
}
#[cfg(feature = "playback")]
#[derive(Clone)]
pub struct SpeakerSink {
sender: Sender<Vec<f32>>,
running: Arc<AtomicBool>,
sentinel_seq: Arc<AtomicUsize>,
sentinels_played: Arc<AtomicUsize>,
}
#[cfg(feature = "playback")]
impl SpeakerSink {
pub fn send(&self, samples: Vec<f32>) -> Result<(), DecibriError> {
send_impl(&self.sender, &self.running, samples)
}
pub fn drain(&self) {
drain_impl(
&self.sender,
&self.running,
&self.sentinel_seq,
&self.sentinels_played,
);
}
}
#[cfg(feature = "playback")]
pub struct Speaker {
config: SpeakerConfig,
device: BackendDevice,
}
#[cfg(feature = "playback")]
impl Speaker {
pub fn new(config: SpeakerConfig) -> Result<Self, DecibriError> {
config.validate()?;
let device = CpalBackend.resolve_output_device(&config.device)?;
Ok(Self { config, device })
}
pub fn devices() -> Result<Vec<crate::device::SpeakerInfo>, DecibriError> {
crate::device::output_devices()
}
pub fn start(&self) -> Result<SpeakerStream, DecibriError> {
let (sender, receiver): (Sender<Vec<f32>>, Receiver<Vec<f32>>) =
crossbeam_channel::bounded(32);
let running = Arc::new(AtomicBool::new(true));
let sentinel_seq = Arc::new(AtomicUsize::new(0));
let sentinels_played = Arc::new(AtomicUsize::new(0));
let running_cb = running.clone();
let sentinels_played_cb = sentinels_played.clone();
let running_clone = running.clone();
let last_error = Arc::new(Mutex::new(None));
let err_last_error = last_error.clone();
let sample_rate = self.config.sample_rate;
let channels = self.config.channels;
let mut accum: Vec<f32> = Vec::new();
let on_data: OutputDataCallback = Box::new(move |data: &mut [f32]| {
render_output(
data,
&mut accum,
&receiver,
&running_cb,
&sentinels_played_cb,
);
});
let on_error: StreamErrorCallback = Box::new(move |err: DecibriError| {
eprintln!("{err}");
if let Ok(mut slot) = err_last_error.lock() {
*slot = Some(err);
}
running_clone.store(false, Ordering::Relaxed);
});
let params = StreamParams {
channels,
sample_rate,
frames_per_buffer: None,
};
let _stream = CpalBackend.open_output_stream(&self.device, ¶ms, on_data, on_error)?;
Ok(SpeakerStream {
_stream,
sender,
running,
sentinel_seq,
sentinels_played,
last_error,
})
}
}
#[cfg(all(test, feature = "playback"))]
mod tests {
use super::*;
use std::thread;
use std::time::{Duration, Instant};
fn test_stream() -> (SpeakerStream, Receiver<Vec<f32>>, Arc<AtomicUsize>) {
let (sender, receiver) = crossbeam_channel::bounded::<Vec<f32>>(32);
let sentinels_played = Arc::new(AtomicUsize::new(0));
let stream = SpeakerStream {
_stream: BackendStream::empty(),
sender,
running: Arc::new(AtomicBool::new(true)),
sentinel_seq: Arc::new(AtomicUsize::new(0)),
sentinels_played: sentinels_played.clone(),
last_error: Arc::new(Mutex::new(None)),
};
(stream, receiver, sentinels_played)
}
fn render_once(
stream: &SpeakerStream,
accum: &mut Vec<f32>,
receiver: &Receiver<Vec<f32>>,
len: usize,
) -> Vec<f32> {
let mut data = vec![0.0_f32; len];
render_output(
&mut data,
accum,
receiver,
&stream.running,
&stream.sentinels_played,
);
data
}
fn consume_until(
stream: &SpeakerStream,
accum: &mut Vec<f32>,
receiver: &Receiver<Vec<f32>>,
played: &AtomicUsize,
target: usize,
) {
let mut data = vec![0.0_f32; 8];
let deadline = Instant::now() + Duration::from_secs(2);
while played.load(Ordering::Relaxed) < target {
render_output(
&mut data,
accum,
receiver,
&stream.running,
&stream.sentinels_played,
);
assert!(
Instant::now() <= deadline,
"timed out waiting for {target} sentinel(s)"
);
thread::sleep(Duration::from_millis(1));
}
}
#[test]
fn drain_returns_when_stream_dropped_without_stop() {
let (stream, receiver, _played) = test_stream();
let sink = stream.sink();
let (done_tx, done_rx) = crossbeam_channel::bounded::<()>(1);
let drain_thread = thread::spawn(move || {
sink.drain();
let _ = done_tx.send(());
});
thread::sleep(Duration::from_millis(50));
drop(stream);
assert!(
done_rx.recv_timeout(Duration::from_secs(2)).is_ok(),
"drain() did not return after the stream was dropped without stop() \
(regression: would hang forever without impl Drop for SpeakerStream)"
);
drain_thread.join().unwrap();
drop(receiver);
}
#[test]
fn test_render_plays_queued_audio_then_silence() {
let (stream, receiver, _played) = test_stream();
stream.send(vec![0.5_f32, 0.5, 0.5]).unwrap();
let mut accum = Vec::new();
let data = render_once(&stream, &mut accum, &receiver, 5);
assert_eq!(&data[..3], &[0.5, 0.5, 0.5], "queued audio is played");
assert_eq!(&data[3..], &[0.0, 0.0], "the remainder is silence");
}
#[test]
fn test_render_stashes_overflow_into_accumulator() {
let (stream, receiver, _played) = test_stream();
stream.send(vec![1.0_f32; 6]).unwrap();
let mut accum = Vec::new();
let first = render_once(&stream, &mut accum, &receiver, 4);
assert_eq!(first, vec![1.0; 4]);
assert_eq!(
accum,
vec![1.0; 2],
"overflow is stashed for the next cycle"
);
let second = render_once(&stream, &mut accum, &receiver, 4);
assert_eq!(&second[..2], &[1.0, 1.0]);
assert_eq!(&second[2..], &[0.0, 0.0]);
assert!(accum.is_empty());
}
#[test]
fn test_stop_discards_queued_audio() {
let (stream, receiver, _played) = test_stream();
stream.send(vec![0.5_f32; 50]).unwrap();
stream.send(vec![0.25_f32; 50]).unwrap();
stream.stop();
let mut accum = Vec::new();
let data = render_once(&stream, &mut accum, &receiver, 64);
assert!(
data.iter().all(|&s| s == 0.0),
"a stopped stream must emit only silence"
);
assert!(
receiver.try_recv().is_err(),
"stop() must discard queued audio, leaving the channel empty"
);
assert!(accum.is_empty(), "the accumulator must be cleared on stop");
}
#[test]
fn test_stop_discards_when_channel_full() {
let (stream, receiver, _played) = test_stream();
for _ in 0..32 {
stream.send(vec![1.0_f32; 8]).unwrap();
}
stream.stop();
let mut accum = Vec::new();
let data = render_once(&stream, &mut accum, &receiver, 16);
assert!(
data.iter().all(|&s| s == 0.0),
"a full channel must still go silent after stop (regression: stop() \
used to no-op when the channel was full)"
);
assert!(
receiver.try_recv().is_err(),
"a full channel must be fully discarded after stop"
);
}
#[test]
fn test_is_playing_false_immediately_after_stop() {
let (stream, _receiver, _played) = test_stream();
assert!(stream.is_playing(), "a fresh stream reports playing");
stream.stop();
assert!(
!stream.is_playing(),
"is_playing() must be false the instant stop() returns"
);
}
#[test]
fn test_send_after_stop_returns_closed() {
let (stream, _receiver, _played) = test_stream();
stream.stop();
let err = stream.send(vec![1.0_f32; 4]).unwrap_err();
assert!(
matches!(err, DecibriError::SpeakerStreamClosed),
"a non-empty send() after stop() must fail with SpeakerStreamClosed"
);
assert!(
stream.send(Vec::new()).is_ok(),
"an empty send is a documented no-op even after stop"
);
}
#[test]
fn test_second_drain_waits_for_its_own_audio() {
let (stream, receiver, played) = test_stream();
let mut accum = Vec::new();
stream.send(vec![1.0_f32; 4]).unwrap();
let sink1 = stream.sink();
let h1 = thread::spawn(move || sink1.drain());
consume_until(&stream, &mut accum, &receiver, &played, 1);
h1.join().unwrap();
stream.send(vec![2.0_f32; 4]).unwrap();
let sink2 = stream.sink();
let done = Arc::new(AtomicBool::new(false));
let done_writer = done.clone();
let h2 = thread::spawn(move || {
sink2.drain();
done_writer.store(true, Ordering::Relaxed);
});
thread::sleep(Duration::from_millis(60));
assert!(
!done.load(Ordering::Relaxed),
"second drain must not return before its audio is consumed \
(stale-sentinel regression)"
);
consume_until(&stream, &mut accum, &receiver, &played, 2);
h2.join().unwrap();
assert!(
done.load(Ordering::Relaxed),
"second drain returns once its audio is consumed"
);
}
#[test]
fn test_send_after_drain_succeeds_and_stream_stays_open() {
let (stream, receiver, played) = test_stream();
let mut accum = Vec::new();
stream.send(vec![1.0_f32; 4]).unwrap();
let sink = stream.sink();
let h = thread::spawn(move || sink.drain());
consume_until(&stream, &mut accum, &receiver, &played, 1);
h.join().unwrap();
assert!(stream.is_playing(), "drain() must not end the stream");
assert!(
stream.send(vec![2.0_f32; 4]).is_ok(),
"send() after drain() must still succeed"
);
}
#[test]
fn test_sink_drain_is_repeatable() {
let (stream, receiver, played) = test_stream();
let mut accum = Vec::new();
for round in 1..=2usize {
let sink = stream.sink();
sink.send(vec![round as f32; 4]).unwrap();
let h = thread::spawn(move || sink.drain());
consume_until(&stream, &mut accum, &receiver, &played, round);
h.join().unwrap();
}
assert_eq!(
played.load(Ordering::Relaxed),
2,
"each drain through the sink waits for its own audio"
);
}
#[test]
fn test_speaker_stream_is_send_and_sync() {
fn assert_send_sync<T: Send + Sync>() {}
assert_send_sync::<SpeakerStream>();
}
#[test]
fn test_stop_empties_stream_and_is_idempotent() {
let (stream, _receiver, _played) = test_stream();
stream.stop();
assert!(!stream.is_playing(), "stop() clears the running flag");
assert!(
!stream._stream.is_active(),
"stop() leaves the stream slot empty (device released)"
);
stream.stop(); }
}