use anyhow::Result;
use std::sync::mpsc::{self, Receiver, SyncSender, TrySendError};
use std::time::Duration;
use super::Playout;
const QUEUE_FRAMES: usize = 50;
pub struct AudioOutput {
frames: SyncSender<Vec<f32>>,
stats: std::sync::Arc<Stats>,
}
#[derive(Debug, Default)]
pub struct Stats {
pub played: std::sync::atomic::AtomicU64,
pub dropped: std::sync::atomic::AtomicU64,
}
impl std::fmt::Debug for AudioOutput {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AudioOutput")
.field("played", &self.stats.played)
.field("dropped", &self.stats.dropped)
.finish()
}
}
impl AudioOutput {
pub fn start<R: Send + 'static>(capacity: usize) -> Result<(Self, PlayedReceiver<R>)> {
let (frames_tx, frames_rx) = mpsc::sync_channel::<Vec<f32>>(QUEUE_FRAMES);
let (played_tx, played_rx) = mpsc::channel::<PlayedFrame<R>>();
let stats = std::sync::Arc::new(Stats::default());
let thread_stats = stats.clone();
let thread = std::thread::Builder::new()
.name("pcc-audio-out".into())
.spawn(move || {
let mut playout = Playout::<R>::new(capacity);
match playout.start_default_output() {
Ok(()) => {}
Err(e) => {
tracing::debug!("no audio output device: {e}");
}
}
while let Ok(pcm) = frames_rx.recv() {
playout.queue_audio(&pcm);
if let Some(now) = playout.audio_now() {
let due = playout.take_due(now);
if !due.is_empty() {
thread_stats
.played
.fetch_add(due.len() as u64, std::sync::atomic::Ordering::Relaxed);
if played_tx
.send(PlayedFrame {
revisions: due,
at: now,
})
.is_err()
{
break;
}
}
}
}
})
.map_err(|e| anyhow::anyhow!("could not start the audio output thread: {e}"))?;
drop(thread);
Ok((
Self {
frames: frames_tx,
stats,
},
played_rx,
))
}
pub fn push(&self, pcm: &[f32]) {
match self.frames.try_send(pcm.to_vec()) {
Ok(()) => {}
Err(TrySendError::Full(_)) => {
self.stats
.dropped
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
Err(TrySendError::Disconnected(_)) => {}
}
}
pub fn played(&self) -> u64 {
self.stats.played.load(std::sync::atomic::Ordering::Relaxed)
}
pub fn dropped(&self) -> u64 {
self.stats
.dropped
.load(std::sync::atomic::Ordering::Relaxed)
}
}
#[derive(Debug)]
pub struct PlayedFrame<R> {
pub revisions: Vec<R>,
pub at: std::time::Instant,
}
impl<R> PlayedFrame<R> {
pub fn next_poll() -> Duration {
Duration::from_millis(super::FRAME_MS)
}
}
pub type PlayedReceiver<R> = Receiver<PlayedFrame<R>>;
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn pushing_never_blocks_even_with_nothing_draining() {
let (tx, rx) = mpsc::sync_channel::<Vec<f32>>(1);
let out = AudioOutput {
frames: tx,
stats: std::sync::Arc::new(Stats::default()),
};
for _ in 0..1_000 {
out.push(&[0.0; 10]);
}
assert!(out.dropped() > 0, "an undrained queue must count its drops");
drop(rx);
}
}