use std::any::Any;
use std::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering};
use std::sync::mpsc::Sender;
use std::sync::Arc;
use cpal::traits::{DeviceTrait, HostTrait, StreamTrait};
use cpal::{
Device, ErrorKind, SampleFormat, StreamConfig, SupportedStreamConfig,
SupportedStreamConfigRange,
};
use super::decode::Spec;
use super::meter::Meter;
const BUFFER_SECONDS: u32 = 2;
#[derive(Debug)]
pub enum OutputError {
NoDevice,
NoConfig,
Build(String),
}
impl std::fmt::Display for OutputError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
OutputError::NoDevice => write!(f, "no audio output device"),
OutputError::NoConfig => write!(f, "device offers no usable output format"),
OutputError::Build(s) => write!(f, "could not open audio stream: {s}"),
}
}
}
impl std::error::Error for OutputError {}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum DeviceEvent {
Lost(String),
Rerouted,
Error(String),
}
impl From<cpal::Error> for DeviceEvent {
fn from(e: cpal::Error) -> Self {
match e.kind() {
ErrorKind::DeviceChanged => DeviceEvent::Rerouted,
ErrorKind::DeviceNotAvailable
| ErrorKind::HostUnavailable
| ErrorKind::StreamInvalidated => DeviceEvent::Lost(e.to_string()),
_ => DeviceEvent::Error(e.to_string()),
}
}
}
pub trait Backend: Send + 'static {
fn negotiate(&self, src: Spec) -> Result<Plan, OutputError>;
fn start(
&self,
plan: Plan,
consumer: rtrb::Consumer<f32>,
shared: Arc<Shared>,
events: Sender<DeviceEvent>,
) -> Result<Box<dyn Any>, OutputError>;
}
pub struct Cpal(pub Device);
impl Backend for Cpal {
fn negotiate(&self, src: Spec) -> Result<Plan, OutputError> {
negotiate(&self.0, src)
}
fn start(
&self,
plan: Plan,
consumer: rtrb::Consumer<f32>,
shared: Arc<Shared>,
events: Sender<DeviceEvent>,
) -> Result<Box<dyn Any>, OutputError> {
let config = StreamConfig {
channels: plan.channels,
sample_rate: plan.rate,
buffer_size: cpal::BufferSize::Default,
};
let device = &self.0;
let stream = match plan.format {
SampleFormat::F32 => build::<f32>(device, &config, consumer, shared, events, |v| v),
SampleFormat::F64 => {
build::<f64>(device, &config, consumer, shared, events, |v| v as f64)
}
SampleFormat::I32 => build::<i32>(device, &config, consumer, shared, events, |v| {
(v.clamp(-1.0, 1.0) as f64 * i32::MAX as f64) as i32
}),
SampleFormat::I24 => {
build::<cpal::I24>(device, &config, consumer, shared, events, |v| {
cpal::I24::new_unchecked((v.clamp(-1.0, 1.0) as f64 * 8_388_607.0) as i32)
})
}
SampleFormat::I16 => build::<i16>(device, &config, consumer, shared, events, |v| {
(v.clamp(-1.0, 1.0) * i16::MAX as f32) as i16
}),
SampleFormat::U16 => build::<u16>(device, &config, consumer, shared, events, |v| {
((v.clamp(-1.0, 1.0) * 0.5 + 0.5) * u16::MAX as f32) as u16
}),
other => {
return Err(OutputError::Build(format!(
"unsupported sample format {other:?}"
)))
}
}?;
stream
.play()
.map_err(|e| OutputError::Build(e.to_string()))?;
Ok(Box::new(stream))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Plan {
pub rate: u32,
pub channels: u16,
pub format: SampleFormat,
}
impl Plan {
pub fn needs_resample(&self, src: Spec) -> bool {
src.rate != 0 && src.rate != self.rate
}
}
pub fn negotiate(device: &Device, src: Spec) -> Result<Plan, OutputError> {
let ranges: Vec<SupportedStreamConfigRange> = device
.supported_output_configs()
.map_err(|_| OutputError::NoConfig)?
.collect();
choose(&ranges, device.default_output_config().ok(), src).ok_or(OutputError::NoConfig)
}
pub fn choose(
ranges: &[SupportedStreamConfigRange],
default: Option<SupportedStreamConfig>,
src: Spec,
) -> Option<Plan> {
let want_ch = if src.channels == 0 { 2 } else { src.channels };
let supports = |r: &SupportedStreamConfigRange, rate: u32| {
r.min_sample_rate() <= rate && rate <= r.max_sample_rate()
};
let usable: Vec<_> = ranges
.iter()
.filter_map(|r| format_rank(r.sample_format()).map(|rank| (rank, r)))
.collect();
let plan = |r: &SupportedStreamConfigRange, rate| Plan {
rate,
channels: r.channels(),
format: r.sample_format(),
};
if src.rate != 0 {
let exact = usable
.iter()
.filter(|(_, r)| r.channels() == want_ch && supports(r, src.rate))
.min_by_key(|(rank, _)| *rank);
if let Some((_, r)) = exact {
return Some(plan(r, src.rate));
}
let by_rate = usable
.iter()
.filter(|(_, r)| supports(r, src.rate))
.min_by_key(|(rank, r)| (*rank, r.channels().abs_diff(want_ch)));
if let Some((_, r)) = by_rate {
return Some(plan(r, src.rate));
}
}
let default = default?;
if format_rank(default.sample_format()).is_some() {
return Some(Plan {
rate: default.sample_rate(),
channels: default.channels(),
format: default.sample_format(),
});
}
usable
.iter()
.filter(|(_, r)| supports(r, default.sample_rate()))
.min_by_key(|(rank, r)| (*rank, r.channels().abs_diff(default.channels())))
.map(|(_, r)| plan(r, default.sample_rate()))
}
fn format_rank(f: SampleFormat) -> Option<u8> {
match f {
SampleFormat::F32 => Some(0),
SampleFormat::F64 => Some(1),
SampleFormat::I32 => Some(2),
SampleFormat::I24 => Some(3),
SampleFormat::I16 => Some(4),
SampleFormat::U16 => Some(5),
_ => None,
}
}
pub fn default_device() -> Result<Device, OutputError> {
cpal::default_host()
.default_output_device()
.ok_or(OutputError::NoDevice)
}
pub struct Shared {
pub frames_out: AtomicU64,
volume: AtomicU32,
pub paused: AtomicBool,
pub position_rate: AtomicU32,
pub track_start: AtomicU64,
pub position_offset: AtomicU64,
speed_bits: AtomicU32,
pub flush_requested: AtomicU64,
pub flush_done: AtomicU64,
momentary_bits: AtomicU32,
peak_bits: AtomicU32,
}
impl Shared {
pub fn new() -> Self {
Shared {
frames_out: AtomicU64::new(0),
volume: AtomicU32::new(1.0f32.to_bits()),
paused: AtomicBool::new(false),
position_rate: AtomicU32::new(0),
track_start: AtomicU64::new(0),
position_offset: AtomicU64::new(0),
speed_bits: AtomicU32::new(1.0f32.to_bits()),
flush_requested: AtomicU64::new(0),
flush_done: AtomicU64::new(0),
momentary_bits: AtomicU32::new(f32::NEG_INFINITY.to_bits()),
peak_bits: AtomicU32::new(0),
}
}
pub fn loudness(&self) -> Option<f32> {
let lufs = f32::from_bits(self.momentary_bits.load(Ordering::Relaxed));
(lufs >= super::meter::SILENCE_LUFS).then_some(lufs)
}
pub fn take_peak(&self) -> f32 {
f32::from_bits(self.peak_bits.swap(0, Ordering::Relaxed))
}
pub fn volume(&self) -> f32 {
f32::from_bits(self.volume.load(Ordering::Relaxed))
}
pub fn set_volume(&self, v: f32) {
self.volume
.store(v.clamp(0.0, 1.0).to_bits(), Ordering::Relaxed);
}
pub fn speed(&self) -> f64 {
f32::from_bits(self.speed_bits.load(Ordering::Relaxed)) as f64
}
pub fn set_speed(&self, v: f64) {
self.speed_bits
.store((v as f32).to_bits(), Ordering::Relaxed);
}
}
impl Default for Shared {
fn default() -> Self {
Self::new()
}
}
pub struct Output {
_stream: Box<dyn Any>,
pub producer: rtrb::Producer<f32>,
pub plan: Plan,
pub capacity: usize,
pub shared: Arc<Shared>,
}
impl Output {
pub fn open(
backend: &dyn Backend,
plan: Plan,
shared: Arc<Shared>,
events: Sender<DeviceEvent>,
) -> Result<Self, OutputError> {
let capacity = (plan.rate * BUFFER_SECONDS) as usize * plan.channels as usize;
let (producer, consumer) = rtrb::RingBuffer::<f32>::new(capacity);
let stream = backend.start(plan, consumer, shared.clone(), events)?;
Ok(Output {
_stream: stream,
producer,
plan,
capacity,
shared,
})
}
pub fn play(&self) {
self.shared.paused.store(false, Ordering::Relaxed);
}
pub fn pause(&self) {
self.shared.paused.store(true, Ordering::Relaxed);
}
pub fn buffered_frames(&self) -> usize {
(self.capacity - self.producer.slots()) / self.plan.channels as usize
}
pub fn is_drained(&self) -> bool {
self.buffered_frames() == 0
}
}
fn build<T>(
device: &Device,
config: &StreamConfig,
mut consumer: rtrb::Consumer<f32>,
shared: Arc<Shared>,
events: Sender<DeviceEvent>,
conv: fn(f32) -> T,
) -> Result<cpal::Stream, OutputError>
where
T: cpal::SizedSample + Send + 'static,
{
let channels = config.channels as u64;
let mut meter = Meter::new(config.sample_rate, config.channels);
device
.build_output_stream::<T, _, _>(
*config,
move |out: &mut [T], _| render(out, &mut consumer, &shared, &mut meter, channels, conv),
move |e| {
let _ = events.send(e.into());
},
None,
)
.map_err(|e| OutputError::Build(e.to_string()))
}
pub fn render<T>(
out: &mut [T],
consumer: &mut rtrb::Consumer<f32>,
shared: &Shared,
meter: &mut Meter,
channels: u64,
conv: fn(f32) -> T,
) {
let requested = shared.flush_requested.load(Ordering::Relaxed);
if requested != shared.flush_done.load(Ordering::Relaxed) {
if let Ok(chunk) = consumer.read_chunk(consumer.slots()) {
chunk.commit_all();
}
shared.flush_done.store(requested, Ordering::Relaxed);
}
if shared.paused.load(Ordering::Relaxed) {
for slot in out.iter_mut() {
*slot = conv(0.0);
}
return;
}
let gain = shared.volume();
let mut filled = 0usize;
for slot in out.iter_mut() {
let s = match consumer.pop() {
Ok(s) => {
filled += 1;
s
}
Err(_) => 0.0,
};
if let Some(lufs) = meter.sample(s) {
shared
.momentary_bits
.store(lufs.to_bits(), Ordering::Relaxed);
}
*slot = conv(s * gain);
}
let peak = meter.take_peak();
let _ = shared
.peak_bits
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |bits| {
(peak > f32::from_bits(bits)).then_some(peak.to_bits())
});
shared
.frames_out
.fetch_add(filled as u64 / channels, Ordering::Relaxed);
}
pub fn remap_channels(input: &[f32], src_ch: usize, dst_ch: usize, out: &mut Vec<f32>) {
if src_ch == dst_ch {
out.extend_from_slice(input);
return;
}
if src_ch == 0 || dst_ch == 0 {
return;
}
for frame in input.chunks_exact(src_ch) {
if src_ch == 1 {
out.extend(std::iter::repeat_n(frame[0], dst_ch));
} else {
for c in 0..dst_ch {
out.push(frame.get(c).copied().unwrap_or(0.0));
}
}
}
}