#[cfg(feature = "playback")]
use std::sync::atomic::{AtomicBool, Ordering};
#[cfg(feature = "playback")]
use std::sync::Arc;
#[cfg(feature = "playback")]
use cpal::traits::{DeviceTrait, StreamTrait};
#[cfg(feature = "playback")]
use crossbeam_channel::{Receiver, Sender};
use crate::device::DeviceSelector;
use crate::error::DecibriError;
#[derive(Debug, Clone)]
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")]
pub struct SpeakerStream {
_stream: cpal::Stream,
sender: Sender<Vec<f32>>,
running: Arc<AtomicBool>,
drain_complete: Arc<AtomicBool>,
}
#[cfg(feature = "playback")]
impl SpeakerStream {
pub fn send(&self, samples: Vec<f32>) -> Result<(), DecibriError> {
if samples.is_empty() {
return Ok(());
}
self.sender
.send(samples)
.map_err(|_| DecibriError::SpeakerStreamClosed)
}
pub fn drain(&self) {
let _ = self.sender.send(Vec::new());
while !self.drain_complete.load(Ordering::Relaxed) {
std::thread::sleep(std::time::Duration::from_millis(10));
}
self.running.store(false, Ordering::Relaxed);
}
pub fn is_playing(&self) -> bool {
self.running.load(Ordering::Relaxed)
}
pub fn stop(&self) {
self.running.store(false, Ordering::Relaxed);
while self.sender.try_send(Vec::new()).is_ok() {}
}
pub fn sink(&self) -> SpeakerSink {
SpeakerSink {
sender: self.sender.clone(),
running: self.running.clone(),
drain_complete: self.drain_complete.clone(),
}
}
}
#[cfg(feature = "playback")]
#[derive(Clone)]
pub struct SpeakerSink {
sender: Sender<Vec<f32>>,
running: Arc<AtomicBool>,
drain_complete: Arc<AtomicBool>,
}
#[cfg(feature = "playback")]
impl SpeakerSink {
pub fn send(&self, samples: Vec<f32>) -> Result<(), DecibriError> {
if samples.is_empty() {
return Ok(());
}
self.sender
.send(samples)
.map_err(|_| DecibriError::SpeakerStreamClosed)
}
pub fn drain(&self) {
let _ = self.sender.send(Vec::new());
while !self.drain_complete.load(Ordering::Relaxed) {
std::thread::sleep(std::time::Duration::from_millis(10));
}
self.running.store(false, Ordering::Relaxed);
}
}
#[cfg(feature = "playback")]
pub struct Speaker {
config: SpeakerConfig,
device: cpal::Device,
}
#[cfg(feature = "playback")]
impl Speaker {
pub fn new(config: SpeakerConfig) -> Result<Self, DecibriError> {
config.validate()?;
let device = crate::device::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 drain_complete = Arc::new(AtomicBool::new(false));
let drain_complete_clone = drain_complete.clone();
let running_clone = running.clone();
let sample_rate = self.config.sample_rate;
let channels = self.config.channels;
let stream_config = cpal::StreamConfig {
channels,
sample_rate,
buffer_size: cpal::BufferSize::Default,
};
let mut accum: Vec<f32> = Vec::new();
let stream = self
.device
.build_output_stream(
&stream_config,
move |data: &mut [f32], _: &cpal::OutputCallbackInfo| {
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() {
for sample in &mut data[written..] {
*sample = 0.0;
}
drain_complete_clone.store(true, Ordering::Relaxed);
return;
}
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(_) => {
for sample in &mut data[written..] {
*sample = 0.0;
}
return;
}
}
}
},
move |err| {
eprintln!("decibri: audio output error: {err}");
running_clone.store(false, Ordering::Relaxed);
},
None,
)
.map_err(|e| DecibriError::StreamOpenFailed(e.to_string()))?;
stream
.play()
.map_err(|e| DecibriError::StreamStartFailed(e.to_string()))?;
Ok(SpeakerStream {
_stream: stream,
sender,
running,
drain_complete,
})
}
}