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, 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 self.layout.speakers().is_none() {
return Err(Error::Unsupported(format!(
"playback needs speaker positions to remix (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,
channel: Arc<Mutex<Channel>>,
control: Control,
overflowing: bool,
shared: Arc<Shared>,
engine: Arc<super::Handle>,
}
impl Sink {
pub fn write(&mut self, samples: &[u8]) -> Result<Write, Error> {
let channels = self.input.layout.channels();
let pcm = self.input.format.as_interleaved_f32(samples, channels)?;
let requested_sample_frames = pcm.len() / channels as usize;
let mut channel = self.channel.lock().unwrap();
let mixed;
let pcm = match &channel.remix {
Some(remix) => {
mixed = remix.process(&pcm);
&mixed
}
None => pcm.as_ref(),
};
let accepted_sample_frames = match channel.prod.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.channel.lock().unwrap().prod.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()
}
}
struct Channel {
prod: ResamplingProd<f32>,
remix: Option<Remix>,
}
pub(super) struct Registration {
pub(super) id: u64,
input: Input,
channel: Arc<Mutex<Channel>>,
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, bus: Layout) {
let (channel, cons) = channel(&self.input, rate, bus);
*self.channel.lock().unwrap() = channel;
self.pending = Some(cons);
}
}
pub(super) fn new(
id: u64,
rate: u32,
bus: Layout,
input: Input,
shared: Arc<Shared>,
engine: Arc<super::Handle>,
) -> Result<(Sink, Registration), Error> {
input.validate()?;
let (channel, cons) = self::channel(&input, rate, bus);
let channel = Arc::new(Mutex::new(channel));
let gain = Arc::new(Gain::new());
let sink = Sink {
id,
input,
channel: channel.clone(),
control: Control { gain: gain.clone() },
overflowing: false,
shared,
engine,
};
let registration = Registration {
id,
input: sink.input.clone(),
channel,
gain,
pending: Some(cons),
};
Ok((sink, registration))
}
fn channel(input: &Input, rate: u32, bus: Layout) -> (Channel, ResamplingCons<f32>) {
let remix =
(input.layout != bus).then(|| Remix::new(input.layout, bus).expect("sink and bus layouts name their speakers"));
let latency = input.latency.as_secs_f64();
let (prod, cons) = resampling_channel::<f32>(
bus.channels() as usize,
input.sample_rate,
rate,
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()
},
);
(Channel { prod, remix }, cons)
}
#[cfg(test)]
mod tests {
use super::*;
fn sink(input: Input, output_rate: u32) -> (Sink, ResamplingCons<f32>) {
sink_into(input, output_rate, Layout::Stereo)
}
fn sink_into(input: Input, output_rate: u32, bus: Layout) -> (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, bus, input, shared, engine).unwrap();
(sink, registration.pending.take().unwrap())
}
fn mix(input: Layout, bus: Layout, frame: &[f32]) -> Vec<f32> {
let input = Input {
layout: input,
..Default::default()
};
let (mut sink, mut cons) = sink_into(input, 48_000, bus);
let channels = bus.channels() as usize;
cons.read_interleaved(&mut vec![0.0; channels], false);
let pcm: Vec<u8> = frame.repeat(4800).iter().flat_map(|s| s.to_le_bytes()).collect();
assert_eq!(sink.write(&pcm).unwrap().dropped_sample_frames, 0);
let mut out = vec![0.0; 4800 * channels];
cons.read_interleaved(&mut out, false);
out[out.len() - channels..].to_vec()
}
fn close(got: &[f32], want: &[f32]) {
assert_eq!(got.len(), want.len(), "{got:?} vs {want:?}");
for (g, w) in got.iter().zip(want) {
assert!((g - w).abs() < 1e-4, "{got:?} vs {want:?}");
}
}
#[test]
fn a_surround_sink_downmixes_into_a_stereo_bus() {
let h = std::f32::consts::FRAC_1_SQRT_2;
let got = mix(Layout::FivePointOne, Layout::Stereo, &[0.1, 0.2, 0.3, 0.4, 0.05, 0.06]);
close(&got, &[0.1 + h * 0.3 + h * 0.05, 0.2 + h * 0.3 + h * 0.06]);
}
#[test]
fn a_stereo_sink_fills_the_front_of_a_surround_bus() {
let got = mix(Layout::Stereo, Layout::FivePointOne, &[0.25, 0.75]);
close(&got, &[0.25, 0.75, 0.0, 0.0, 0.0, 0.0]);
}
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; 2], 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(2), 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_named_layouts() {
for layout in [
Layout::Mono,
Layout::Stereo,
Layout::FivePointOne,
Layout::SevenPointOne,
] {
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();
}
}