use std::borrow::Cow;
use std::sync::mpsc::{SyncSender, TrySendError};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use fixed_resample::{PushStatus, ResamplingChannelConfig, ResamplingCons, ResamplingProd, resampling_channel};
use super::driver::Shared;
use super::mixer::{self, BUS_CHANNELS, Gain};
use crate::resample::remix;
use crate::{Error, Format, Layout};
const LATENCY: Duration = Duration::from_millis(50);
const HEADROOM: f64 = 3.0;
#[derive(Clone, Debug)]
#[non_exhaustive]
pub struct Input {
pub format: Format,
pub sample_rate: u32,
pub layout: Layout,
pub latency: Duration,
}
impl Default for Input {
fn default() -> Self {
Self {
format: Format::F32,
sample_rate: 48_000,
layout: Layout::Stereo,
latency: LATENCY,
}
}
}
impl Input {
pub const LATENCY_MAX: Duration = Duration::from_secs(10);
fn validate(&self) -> Result<(), Error> {
if self.sample_rate == 0 {
return Err(Error::Unsupported("sample rate must be > 0".into()));
}
if !matches!(self.layout, Layout::Mono | Layout::Stereo) {
return Err(Error::Unsupported(format!(
"playback accepts named mono or stereo input (got {:?})",
self.layout
)));
}
if self.latency.is_zero() || self.latency > Self::LATENCY_MAX {
return Err(Error::Unsupported(format!(
"playback latency must be non-zero and at most {:?} (got {:?})",
Self::LATENCY_MAX,
self.latency
)));
}
Ok(())
}
}
#[must_use = "inspect dropped_sample_frames; retrying dropped live audio adds latency"]
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
#[non_exhaustive]
pub struct Write {
pub accepted_sample_frames: usize,
pub dropped_sample_frames: usize,
}
impl Write {
fn from_accepted(requested_sample_frames: usize, accepted_sample_frames: usize) -> Self {
Self {
accepted_sample_frames,
dropped_sample_frames: requested_sample_frames - accepted_sample_frames,
}
}
}
pub struct Sink {
id: u64,
input: Input,
prod: Arc<Mutex<ResamplingProd<f32>>>,
control: Control,
overflowing: bool,
shared: Arc<Shared>,
engine: Arc<super::Handle>,
}
impl Sink {
pub fn write(&mut self, samples: &[u8]) -> Result<Write, Error> {
let pcm = self
.input
.format
.as_interleaved_f32(samples, self.input.layout.channels())?;
let pcm = match self.input.layout.channels() as usize {
BUS_CHANNELS => pcm,
_ => Cow::Owned(remix(&pcm, self.input.layout, Layout::Stereo)?),
};
let requested_sample_frames = pcm.len() / BUS_CHANNELS;
let accepted_sample_frames = match self.prod.lock().unwrap().push_interleaved(&pcm) {
PushStatus::Ok => {
self.overflowing = false;
requested_sample_frames
}
PushStatus::OutputNotReady => {
self.overflowing = false;
0
}
PushStatus::OverflowOccurred { num_frames_pushed } => {
if !self.overflowing {
tracing::warn!(num_frames_pushed, "audio playback overflow, dropping samples");
self.overflowing = true;
}
num_frames_pushed
}
PushStatus::UnderflowCorrected { num_zero_frames_pushed } => {
self.overflowing = false;
tracing::debug!(num_zero_frames_pushed, "audio playback underflow, padded with silence");
requested_sample_frames
}
};
Ok(Write::from_accepted(requested_sample_frames, accepted_sample_frames))
}
pub fn buffered(&self) -> Duration {
Duration::from_secs_f64(self.prod.lock().unwrap().occupied_seconds().max(0.0))
}
pub fn input(&self) -> &Input {
&self.input
}
pub fn control(&self) -> Control {
self.control.clone()
}
pub fn set_volume(&self, volume: f32) {
self.control.set_volume(volume);
}
pub fn volume(&self) -> f32 {
self.control.volume()
}
pub fn peak(&self) -> f32 {
self.control.peak()
}
}
impl Drop for Sink {
fn drop(&mut self) {
self.shared.remove(self.id);
self.engine.wake();
}
}
#[derive(Clone, Debug)]
pub struct Control {
gain: Arc<Gain>,
}
impl Control {
pub fn set_volume(&self, volume: f32) {
self.gain.set_volume(volume);
}
pub fn volume(&self) -> f32 {
self.gain.volume()
}
pub fn peak(&self) -> f32 {
self.gain.peak()
}
}
pub(super) struct Registration {
pub(super) id: u64,
rate: u32,
latency: Duration,
prod: Arc<Mutex<ResamplingProd<f32>>>,
gain: Arc<Gain>,
pending: Option<ResamplingCons<f32>>,
}
impl Registration {
pub(super) fn attached(&self) -> bool {
self.pending.is_none()
}
pub(super) fn attach(&mut self, mixer: &SyncSender<mixer::Command>) {
let Some(cons) = self.pending.take() else { return };
let command = mixer::Command::Add {
id: self.id,
cons,
gain: self.gain.clone(),
};
if let Err(err) = mixer.try_send(command) {
let (TrySendError::Full(rejected) | TrySendError::Disconnected(rejected)) = err;
if let mixer::Command::Add { cons, .. } = rejected {
self.pending = Some(cons);
}
}
}
pub(super) fn rebuild(&mut self, rate: u32) {
let (prod, cons) = channel(self.rate, rate, self.latency);
*self.prod.lock().unwrap() = prod;
self.pending = Some(cons);
}
}
pub(super) fn new(
id: u64,
rate: u32,
input: Input,
shared: Arc<Shared>,
engine: Arc<super::Handle>,
) -> Result<(Sink, Registration), Error> {
input.validate()?;
let (prod, cons) = channel(input.sample_rate, rate, input.latency);
let prod = Arc::new(Mutex::new(prod));
let gain = Arc::new(Gain::new());
let sink = Sink {
id,
input,
prod: prod.clone(),
control: Control { gain: gain.clone() },
overflowing: false,
shared,
engine,
};
let registration = Registration {
id,
rate: sink.input.sample_rate,
latency: sink.input.latency,
prod,
gain,
pending: Some(cons),
};
Ok((sink, registration))
}
fn channel(from: u32, to: u32, latency: Duration) -> (ResamplingProd<f32>, ResamplingCons<f32>) {
let latency = latency.as_secs_f64();
resampling_channel::<f32>(
BUS_CHANNELS,
from,
to,
true,
ResamplingChannelConfig {
latency_seconds: latency,
capacity_seconds: latency + HEADROOM,
underflow_autocorrect_percent_threshold: Some(25.0),
overflow_autocorrect_percent_threshold: Some(75.0),
..Default::default()
},
)
}
#[cfg(test)]
mod tests {
use super::*;
fn sink(input: Input, output_rate: u32) -> (Sink, ResamplingCons<f32>) {
let shared = Arc::new(Shared::default());
let engine = Arc::new(super::super::Handle {
commands: super::super::driver::Commands::default(),
});
let (sink, mut registration) = new(0, output_rate, input, shared, engine).unwrap();
(sink, registration.pending.take().unwrap())
}
fn s16(frames: usize, channels: usize) -> Vec<u8> {
vec![0; frames * channels * 2]
}
fn ready(cons: &mut ResamplingCons<f32>) {
cons.read_interleaved(&mut [0.0; BUS_CHANNELS], false);
}
#[test]
fn reports_output_not_ready_as_dropped() {
let input = Input {
format: Format::S16,
..Default::default()
};
let (mut sink, _cons) = sink(input, 48_000);
assert_eq!(
sink.write(&s16(10, 2)).unwrap(),
Write {
accepted_sample_frames: 0,
dropped_sample_frames: 10,
}
);
}
#[test]
fn reports_output_that_stops_as_dropped() {
let input = Input {
format: Format::S16,
..Default::default()
};
let (mut sink, mut cons) = sink(input, 48_000);
ready(&mut cons);
cons.set_output_stream_ready(false);
assert_eq!(
sink.write(&s16(10, 2)).unwrap(),
Write {
accepted_sample_frames: 0,
dropped_sample_frames: 10,
}
);
}
#[test]
fn reports_input_frame_units_through_conversion_and_resampling() {
let input = Input {
format: Format::S16,
sample_rate: 44_100,
layout: Layout::Mono,
..Default::default()
};
let (mut sink, mut cons) = sink(input, 48_000);
ready(&mut cons);
assert_eq!(
sink.write(&s16(441, 1)).unwrap(),
Write {
accepted_sample_frames: 441,
dropped_sample_frames: 0,
}
);
}
#[test]
fn reports_partial_acceptance_on_overflow() {
let input = Input {
format: Format::S16,
latency: Duration::from_millis(1),
..Default::default()
};
let (mut sink, mut cons) = sink(input, 48_000);
ready(&mut cons);
let requested = 4 * 48_000;
let write = sink.write(&s16(requested, 2)).unwrap();
assert!(write.accepted_sample_frames > 0);
assert!(write.dropped_sample_frames > 0);
assert_eq!(write.accepted_sample_frames + write.dropped_sample_frames, requested);
}
#[test]
fn keeps_invalid_input_distinct_from_dropping() {
let input = Input {
format: Format::S16,
..Default::default()
};
let (mut sink, _cons) = sink(input, 48_000);
assert!(matches!(sink.write(&[0]), Err(Error::Misaligned { .. })));
}
#[test]
fn rejects_layouts_it_cannot_mix() {
for layout in [Layout::Discrete(0), Layout::Discrete(6)] {
let input = Input {
layout,
..Default::default()
};
assert!(matches!(input.validate(), Err(Error::Unsupported(_))), "{layout:?}");
}
let input = Input {
sample_rate: 0,
..Default::default()
};
assert!(matches!(input.validate(), Err(Error::Unsupported(_))));
}
#[test]
fn accepts_mono_and_stereo() {
for layout in [Layout::Mono, Layout::Stereo] {
let input = Input {
layout,
..Default::default()
};
input.validate().unwrap();
}
}
#[test]
fn rejects_a_latency_it_cannot_buffer() {
for latency in [Duration::ZERO, Input::LATENCY_MAX + Duration::from_secs(1)] {
let input = Input {
latency,
..Default::default()
};
assert!(matches!(input.validate(), Err(Error::Unsupported(_))), "{latency:?}");
}
Input {
latency: Input::LATENCY_MAX,
..Default::default()
}
.validate()
.unwrap();
}
}