use std::fs::File;
use std::future::Future;
use std::io::{BufReader, Read, Seek};
use std::path::PathBuf;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::task::{Context, Poll};
use std::thread::sleep;
use std::time::Duration;
use rodio::{Decoder, DeviceSinkBuilder, Player};
use thiserror::Error;
use tokio::task::{self, JoinHandle};
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum PlaybackError {
#[error(transparent)]
Io(#[from] std::io::Error),
#[error("failed to open audio output: {0}")]
Output(String),
#[error("failed to decode audio data")]
Decode(#[source] rodio::decoder::DecoderError),
#[error("playback task failed")]
Task,
}
#[derive(Debug)]
pub struct PlaybackHandle {
join_handle: JoinHandle<Result<(), PlaybackError>>,
cancelled: Arc<AtomicBool>,
detached: bool,
}
impl PlaybackHandle {
pub fn detach(mut self) {
self.detached = true;
}
}
impl Future for PlaybackHandle {
type Output = Result<(), PlaybackError>;
fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
Pin::new(&mut self.join_handle)
.poll(cx)
.map(|result| result.unwrap_or(Err(PlaybackError::Task)))
}
}
impl Drop for PlaybackHandle {
fn drop(&mut self) {
if !self.detached {
self.cancelled.store(true, Ordering::Relaxed);
}
}
}
pub fn play_file(path: impl Into<PathBuf>, volume: f32) -> PlaybackHandle {
let path = path.into();
play_with(volume, move || Ok(BufReader::new(File::open(path)?)))
}
pub fn play_reader<R>(reader: R, volume: f32) -> PlaybackHandle
where
R: Read + Seek + Send + Sync + 'static,
{
play_with(volume, move || Ok(reader))
}
fn play_with<R>(
volume: f32,
open: impl FnOnce() -> Result<R, PlaybackError> + Send + 'static,
) -> PlaybackHandle
where
R: Read + Seek + Send + Sync + 'static,
{
let cancelled = Arc::new(AtomicBool::new(false));
let cancel_flag = cancelled.clone();
let join_handle = task::spawn_blocking(move || {
if cancel_flag.load(Ordering::Relaxed) {
return Ok(());
}
let device_sink = DeviceSinkBuilder::open_default_sink()
.map_err(|error| PlaybackError::Output(error.to_string()))?;
let player = Player::connect_new(device_sink.mixer());
player.set_volume(volume);
let source = Decoder::new(open()?).map_err(PlaybackError::Decode)?;
if cancel_flag.load(Ordering::Relaxed) {
return Ok(());
}
player.append(source);
while !player.empty() {
sleep(Duration::from_millis(50));
if cancel_flag.load(Ordering::Relaxed) {
return Ok(());
}
}
Ok(())
});
PlaybackHandle {
join_handle,
cancelled,
detached: false,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn missing_file_surfaces_as_error() {
let error = play_file("/nonexistent/audio.mp3", 1.0).await.unwrap_err();
assert!(matches!(
error,
PlaybackError::Io(_) | PlaybackError::Output(_)
));
}
#[tokio::test]
async fn dropping_and_detaching_do_not_panic() {
drop(play_file("/nonexistent/audio.mp3", 1.0));
play_file("/nonexistent/audio.mp3", 1.0).detach();
tokio::task::yield_now().await;
}
}