pub mod convert;
pub mod decode;
pub mod meter;
#[cfg(feature = "opus")]
pub mod opus;
pub mod order;
pub mod output;
pub mod resample;
use std::collections::VecDeque;
use std::path::PathBuf;
use std::sync::atomic::Ordering;
use std::sync::mpsc::{Receiver, RecvTimeoutError, Sender, TryRecvError};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use decode::AudioStream;
pub use decode::Spec;
fn short(path: &std::path::Path) -> String {
path.file_name()
.map(|n| n.to_string_lossy().into_owned())
.unwrap_or_else(|| path.display().to_string())
}
use convert::Converter;
pub use order::Mode;
use order::Order;
use output::{Backend, DeviceEvent, Output, Plan, Shared};
pub const MAX_SEMITONES: i32 = 12;
pub fn speed_for(semitones: i32) -> f64 {
2f64.powf(semitones.clamp(-MAX_SEMITONES, MAX_SEMITONES) as f64 / 12.0)
}
pub fn track_position(
frames_out: u64,
track_start: u64,
offset: u64,
rate: u32,
speed: f64,
) -> Duration {
if rate == 0 || speed <= 0.0 {
return Duration::ZERO;
}
let played = frames_out.saturating_sub(track_start) as f64 * speed;
Duration::from_secs_f64((played + offset as f64) / rate as f64)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum State {
#[default]
Stopped,
Playing,
Paused,
}
#[derive(Debug, Clone)]
pub enum Cmd {
Play(Vec<PathBuf>, usize),
Enqueue(Vec<PathBuf>),
Jump(usize),
TogglePause,
Next,
Prev,
Stop,
Seek(Duration),
SeekBy(i64),
SetVolume(f32),
SpeedBy(i32),
SpeedReset,
SetSpeed(i32),
SetMode(Mode),
Quit,
}
#[derive(Debug, Clone, Default)]
pub struct Status {
pub state: State,
pub queue: Arc<[PathBuf]>,
pub index: usize,
pub duration: Option<Duration>,
pub source: Option<Spec>,
pub output_rate: u32,
pub resampling: bool,
pub error: Option<String>,
pub error_seq: u64,
pub semitones: i32,
pub mode: Mode,
}
impl Status {
pub fn current(&self) -> Option<&PathBuf> {
self.queue.get(self.index)
}
}
enum Msg {
Play(Arc<[PathBuf]>, usize),
Enqueue(Arc<[PathBuf]>),
Cmd(Cmd),
}
pub struct Player {
tx: Sender<Msg>,
status: Arc<Mutex<Status>>,
shared: Arc<Shared>,
handle: Option<std::thread::JoinHandle<()>>,
}
impl Player {
pub fn new() -> Result<Self, output::OutputError> {
Self::with_backend(output::Cpal(output::default_device()?))
}
pub fn with_backend(backend: impl Backend) -> Result<Self, output::OutputError> {
let (tx, rx) = std::sync::mpsc::channel();
let status = Arc::new(Mutex::new(Status::default()));
let shared = Arc::new(Shared::new());
let handle = {
let status = status.clone();
let shared = shared.clone();
std::thread::Builder::new()
.name("playr-audio".into())
.spawn(move || Engine::new(Box::new(backend), rx, status, shared).run())
.map_err(|e| output::OutputError::Build(e.to_string()))?
};
Ok(Player {
tx,
status,
shared,
handle: Some(handle),
})
}
pub fn send(&self, cmd: Cmd) {
let Ok(mut status) = self.status.lock() else {
return;
};
let msg = match cmd {
Cmd::Play(paths, index) => {
status.queue = paths.into();
Msg::Play(status.queue.clone(), index)
}
Cmd::Enqueue(paths) => {
status.queue = status.queue.iter().cloned().chain(paths).collect();
Msg::Enqueue(status.queue.clone())
}
Cmd::SetVolume(v) => return self.shared.set_volume(v),
Cmd::SetMode(mode) => {
status.mode = mode;
Msg::Cmd(Cmd::SetMode(mode))
}
cmd => Msg::Cmd(cmd),
};
let _ = self.tx.send(msg);
}
pub fn mode(&self) -> Mode {
self.status.lock().map(|s| s.mode).unwrap_or_default()
}
pub fn queue(&self) -> Arc<[PathBuf]> {
self.status
.lock()
.map(|s| s.queue.clone())
.unwrap_or_default()
}
pub fn status(&self) -> Status {
self.status.lock().map(|s| s.clone()).unwrap_or_default()
}
pub fn position(&self) -> Duration {
track_position(
self.shared.frames_out.load(Ordering::Relaxed),
self.shared.track_start.load(Ordering::Relaxed),
self.shared.position_offset.load(Ordering::Relaxed),
self.shared.position_rate.load(Ordering::Relaxed),
self.shared.speed(),
)
}
pub fn volume(&self) -> f32 {
self.shared.volume()
}
pub fn loudness(&self) -> Option<f32> {
self.shared.loudness()
}
pub fn take_peak(&self) -> f32 {
self.shared.take_peak()
}
}
impl Drop for Player {
fn drop(&mut self) {
let _ = self.tx.send(Msg::Cmd(Cmd::Quit));
if let Some(h) = self.handle.take() {
let _ = h.join();
}
}
}
const REFILL_BELOW: f64 = 0.5;
const PUMP_BUDGET: Duration = Duration::from_millis(8);
const FLUSH_TIMEOUT: Duration = Duration::from_secs(1);
#[derive(Clone, Copy)]
struct Flush {
generation: u64,
deadline: std::time::Instant,
at: Duration,
offset: u64,
}
struct Staged {
stream: AudioStream,
spec: Spec,
plan: Plan,
index: usize,
first: Vec<f32>,
}
struct Engine {
backend: Box<dyn Backend>,
rx: Receiver<Msg>,
status: Arc<Mutex<Status>>,
shared: Arc<Shared>,
device_events: (Sender<DeviceEvent>, Receiver<DeviceEvent>),
deferred: VecDeque<Msg>,
out: Option<Output>,
stream: Option<AudioStream>,
conv: Option<Converter>,
queue: Arc<[PathBuf]>,
order: Order,
index: usize,
state: State,
semitones: i32,
written: u64,
marks: VecDeque<(u64, usize, Option<Duration>)>,
staged: Option<Staged>,
flush: Option<Flush>,
carry: Vec<f32>,
}
impl Engine {
fn new(
backend: Box<dyn Backend>,
rx: Receiver<Msg>,
status: Arc<Mutex<Status>>,
shared: Arc<Shared>,
) -> Self {
Engine {
backend,
rx,
status,
shared,
device_events: std::sync::mpsc::channel(),
deferred: VecDeque::new(),
out: None,
stream: None,
conv: None,
queue: Arc::default(),
order: Order::new(0, Mode::Normal, 0, order::random_seed()),
index: 0,
state: State::Stopped,
semitones: 0,
written: 0,
marks: VecDeque::new(),
staged: None,
flush: None,
carry: Vec::new(),
}
}
fn run(mut self) {
loop {
let wait = if self.state == State::Playing {
Duration::from_millis(5)
} else {
Duration::from_millis(100)
};
let first = match self.deferred.pop_front() {
Some(cmd) => Ok(cmd),
None => self.rx.recv_timeout(wait),
};
match first {
Ok(Msg::Cmd(Cmd::Quit)) => break,
Ok(msg) => {
self.handle(msg);
while let Some(next) = self
.deferred
.pop_front()
.or_else(|| self.rx.try_recv().ok())
{
if matches!(next, Msg::Cmd(Cmd::Quit)) {
return;
}
self.handle(next);
}
}
Err(RecvTimeoutError::Timeout) => {}
Err(RecvTimeoutError::Disconnected) => break,
}
while let Ok(event) = self.device_events.1.try_recv() {
self.device_event(event);
}
self.poll_flush();
if self.state == State::Playing {
self.pump();
}
self.advance_marks();
self.publish();
}
}
fn handle(&mut self, msg: Msg) {
let cmd = match msg {
Msg::Play(queue, index) => {
self.queue = queue;
return self.jump(index);
}
Msg::Enqueue(queue) => {
let first_new = self.queue.len();
self.queue = queue;
return self.enqueued(first_new);
}
Msg::Cmd(cmd) => cmd,
};
match cmd {
Cmd::Play(..) | Cmd::Enqueue(..) => {}
Cmd::Jump(index) => self.jump(index),
Cmd::TogglePause => match self.state {
State::Playing => {
if let Some(o) = &self.out {
o.pause();
}
self.state = State::Paused;
}
State::Paused => {
if let Some(o) = &self.out {
o.play();
}
self.state = State::Playing;
}
State::Stopped => {
if !self.queue.is_empty() {
self.start(self.index, false);
}
}
},
Cmd::Next => match self.order.successor(self.index, true) {
Some(next) => {
self.teardown();
self.index = next;
self.start(next, false);
}
None => {
self.teardown();
self.state = State::Stopped;
}
},
Cmd::Prev => {
let current = self.index;
let before = if self.elapsed() > Duration::from_secs(3) {
None
} else {
self.order.predecessor(current)
};
self.teardown();
if !before.is_some_and(|b| self.start(b, true)) {
self.start(current, false);
}
}
Cmd::Stop => {
self.teardown();
self.state = State::Stopped;
}
Cmd::Seek(pos) => self.seek(pos),
Cmd::SeekBy(delta) => {
let cur = self.elapsed().as_secs_f64();
let target = (cur + delta as f64).max(0.0);
self.seek(Duration::from_secs_f64(target));
}
Cmd::SetVolume(_) => {}
Cmd::SpeedBy(delta) => self.set_semitones(self.semitones + delta),
Cmd::SpeedReset => self.set_semitones(0),
Cmd::SetSpeed(semitones) => self.set_semitones(semitones),
Cmd::SetMode(mode) => self.set_mode(mode),
Cmd::Quit => {}
}
}
fn device_event(&mut self, event: DeviceEvent) {
match event {
DeviceEvent::Rerouted => {}
DeviceEvent::Lost(e) => {
self.fail(format!("audio device lost: {e}"));
self.teardown();
self.state = State::Stopped;
}
DeviceEvent::Error(e) => self.fail(format!("audio device: {e}")),
}
}
fn jump(&mut self, index: usize) {
self.index = index.min(self.queue.len().saturating_sub(1));
self.order = Order::new(
self.queue.len(),
self.order.mode(),
self.index,
order::random_seed(),
);
self.teardown();
if self.queue.is_empty() {
self.state = State::Stopped;
} else {
self.start(self.index, false);
}
}
fn enqueued(&mut self, first_new: usize) {
self.order.extend(self.queue.len());
if self.state == State::Stopped {
self.start(first_new, false);
} else if self.stream.is_none() && self.staged.is_none() {
self.stage_from(first_new);
}
}
fn elapsed(&self) -> Duration {
if let Some(flush) = self.flush {
return flush.at;
}
track_position(
self.shared.frames_out.load(Ordering::Relaxed),
self.shared.track_start.load(Ordering::Relaxed),
self.shared.position_offset.load(Ordering::Relaxed),
self.shared.position_rate.load(Ordering::Relaxed),
self.shared.speed(),
)
}
fn teardown(&mut self) {
self.stream = None;
self.conv = None;
self.staged = None;
self.flush = None;
self.carry.clear();
self.marks.clear();
self.written = 0;
self.out = None;
self.shared.frames_out.store(0, Ordering::Relaxed);
self.shared.track_start.store(0, Ordering::Relaxed);
self.shared.position_offset.store(0, Ordering::Relaxed);
}
fn start(&mut self, i: usize, back: bool) -> bool {
let Some((i, stream, first, spec)) = self.open_from(i, back) else {
self.state = State::Stopped;
return false;
};
let plan = match self.backend.negotiate(spec) {
Ok(p) => p,
Err(e) => {
self.fail(e.to_string());
self.state = State::Stopped;
return true;
}
};
if let Err(e) = self.rebuild_output(plan) {
self.fail(e.to_string());
self.state = State::Stopped;
return true;
}
self.index = i;
self.conv = Some(Converter::new(spec, plan, self.speed()));
let duration = stream.duration();
self.stream = Some(stream);
self.marks.clear();
self.marks.push_back((0, i, duration));
self.shared.track_start.store(0, Ordering::Relaxed);
self.shared.position_offset.store(0, Ordering::Relaxed);
self.convert_and_carry(&first);
self.state = State::Playing;
if let Some(o) = &self.out {
o.play();
}
self.pump();
true
}
fn open_from(
&mut self,
mut i: usize,
back: bool,
) -> Option<(usize, AudioStream, Vec<f32>, Spec)> {
for _ in 0..self.queue.len() {
let path = self.queue.get(i)?.clone();
if let Some((stream, first, spec)) = self.open_track(&path) {
return Some((i, stream, first, spec));
}
if self.interrupted() {
return None;
}
i = if back {
self.order.predecessor(i)?
} else {
self.order.successor(i, true)?
};
}
None
}
fn interrupted(&mut self) -> bool {
let mut stop = false;
loop {
match self.rx.try_recv() {
Ok(msg) => {
stop |= matches!(
msg,
Msg::Play(..)
| Msg::Cmd(
Cmd::Quit | Cmd::Stop | Cmd::Jump(_) | Cmd::Next | Cmd::Prev
)
);
self.deferred.push_back(msg);
}
Err(TryRecvError::Empty) => return stop,
Err(TryRecvError::Disconnected) => return true,
}
}
}
fn open_track(&mut self, path: &std::path::Path) -> Option<(AudioStream, Vec<f32>, Spec)> {
let mut stream = match AudioStream::open(path) {
Ok(s) => s,
Err(e) => {
self.fail(format!("{}: {e}", short(path)));
return None;
}
};
let first = match stream.next_chunk() {
Ok(Some(c)) => c.to_vec(),
Ok(None) => Vec::new(),
Err(e) => {
self.fail(format!("{}: {e}", short(path)));
return None;
}
};
let spec = stream.spec();
if spec.rate == 0 || spec.channels == 0 {
self.fail(format!("{}: unknown stream format", short(path)));
return None;
}
Some((stream, first, spec))
}
fn open_output(&self, plan: Plan) -> Result<Output, output::OutputError> {
Output::open(
self.backend.as_ref(),
plan,
self.shared.clone(),
self.device_events.0.clone(),
)
}
fn fail(&mut self, msg: String) {
if let Ok(mut s) = self.status.lock() {
s.error = Some(msg);
s.error_seq += 1;
}
}
fn rebuild_output(&mut self, plan: Plan) -> Result<(), output::OutputError> {
if self.out.as_ref().map(|o| o.plan) == Some(plan) {
return Ok(());
}
self.out = None;
let out = self.open_output(plan)?;
self.written = 0;
self.shared.frames_out.store(0, Ordering::Relaxed);
self.shared
.position_rate
.store(plan.rate, Ordering::Relaxed);
self.out = Some(out);
Ok(())
}
fn speed(&self) -> f64 {
speed_for(self.semitones)
}
fn set_mode(&mut self, mode: Mode) {
if mode == self.order.mode() {
return;
}
let current = self.marks.front().map_or(self.index, |(_, i, _)| *i);
self.order.set_mode(mode, current);
let chosen = self.marks.len() > 1 || self.staged.is_some() || self.stream.is_none();
if chosen && !self.marks.is_empty() {
self.seek(self.elapsed());
}
}
fn set_semitones(&mut self, semitones: i32) {
let want = semitones.clamp(-MAX_SEMITONES, MAX_SEMITONES);
if want == self.semitones {
return;
}
let at = self.elapsed();
self.semitones = want;
self.shared.set_speed(self.speed());
if !self.marks.is_empty() {
self.seek(at);
}
}
fn seek(&mut self, pos: Duration) {
let Some(&(_, audible, _)) = self.marks.front() else {
return;
};
let src = if self.marks.len() > 1 || self.stream.is_none() {
let Some(path) = self.queue.get(audible).cloned() else {
return;
};
let Some((mut stream, _, spec)) = self.open_track(&path) else {
return;
};
if stream.seek(pos).is_err() {
return;
}
self.stream = Some(stream);
self.index = audible;
spec
} else {
let (Some(stream), Some(conv)) = (self.stream.as_mut(), self.conv.as_ref()) else {
return;
};
if stream.seek(pos).is_err() {
return;
}
conv.src()
};
if let Some(plan) = self.out.as_ref().map(|o| o.plan) {
self.conv = Some(Converter::new(src, plan, self.speed()));
}
self.carry.clear();
self.staged = None;
let dur = self.stream.as_ref().and_then(|s| s.duration());
let idx = self.index;
self.written = 0;
self.marks.clear();
self.marks.push_back((0, idx, dur));
let rate = self.shared.position_rate.load(Ordering::Relaxed) as f64;
let generation = self.shared.flush_requested.fetch_add(1, Ordering::Relaxed) + 1;
self.flush = Some(Flush {
generation,
deadline: std::time::Instant::now() + FLUSH_TIMEOUT,
at: pos,
offset: (pos.as_secs_f64() * rate).max(0.0) as u64,
});
self.poll_flush();
}
fn poll_flush(&mut self) {
let Some(flush) = self.flush else { return };
if self.shared.flush_done.load(Ordering::Relaxed) != flush.generation {
if std::time::Instant::now() < flush.deadline {
return;
}
let Some(plan) = self.out.as_ref().map(|o| o.plan) else {
self.flush = None;
return;
};
self.out = None;
match self.open_output(plan) {
Ok(o) => {
if self.state == State::Playing {
o.play();
}
self.out = Some(o);
}
Err(e) => {
self.fail(e.to_string());
self.teardown();
self.state = State::Stopped;
return;
}
}
}
self.flush = None;
self.shared.frames_out.store(0, Ordering::Relaxed);
self.shared.track_start.store(0, Ordering::Relaxed);
self.shared
.position_offset
.store(flush.offset, Ordering::Relaxed);
}
fn convert_and_carry(&mut self, decoded: &[f32]) {
if let Some(conv) = self.conv.as_mut() {
conv.push(decoded, &mut self.carry);
}
}
fn finish_track(&mut self) {
if let Some(conv) = self.conv.as_mut() {
conv.finish(&mut self.carry);
}
}
fn pump(&mut self) {
if self.flush.is_some() {
return;
}
let deadline = std::time::Instant::now() + PUMP_BUDGET;
if self.staged.is_some() {
let drained = self.out.as_ref().map(|o| o.is_drained()).unwrap_or(true);
if drained {
self.promote_staged();
} else {
self.push_carry();
return;
}
}
loop {
self.push_carry();
if std::time::Instant::now() >= deadline {
return;
}
let Some(out) = &self.out else { return };
let free = out.producer.slots();
if (free as f64) < out.capacity as f64 * REFILL_BELOW {
return;
}
if !self.carry.is_empty() {
return;
}
let Some(stream) = self.stream.as_mut() else {
if self.out.as_ref().map(|o| o.is_drained()).unwrap_or(true) {
self.state = State::Stopped;
self.teardown();
}
return;
};
match stream.next_chunk() {
Ok(Some(chunk)) => {
let chunk = chunk.to_vec();
self.convert_and_carry(&chunk);
}
Ok(None) => {
self.stream = None;
self.stage_next();
if self.staged.is_some() {
return;
}
}
Err(e) => {
let playing = self.marks.back().map(|(_, i, _)| *i).unwrap_or(self.index);
if let Some(path) = self.queue.get(playing).cloned() {
self.fail(format!("{}: {e}", short(&path)));
}
self.stream = None;
self.stage_next();
if self.staged.is_some() {
return;
}
}
}
}
}
fn stage_next(&mut self) {
let last = self.marks.back().map_or(self.index, |(_, i, _)| *i);
match self.order.successor(last, false) {
Some(next) => self.stage_from(next),
None => self.finish_track(),
}
}
fn stage_from(&mut self, next: usize) {
let Some((next, stream, first, spec)) = self.open_from(next, false) else {
self.finish_track();
return;
};
let plan = match self.backend.negotiate(spec) {
Ok(p) => p,
Err(e) => {
self.fail(e.to_string());
self.finish_track();
return;
}
};
let duration = stream.duration();
if self.out.as_ref().map(|o| o.plan) == Some(plan) {
let speed = self.speed();
if !self
.conv
.as_ref()
.is_some_and(|c| c.continues(spec, plan, speed))
{
self.finish_track();
self.conv = Some(Converter::new(spec, plan, speed));
}
self.marks.push_back((self.written, next, duration));
self.stream = Some(stream);
self.convert_and_carry(&first);
} else {
self.finish_track();
self.staged = Some(Staged {
stream,
spec,
plan,
index: next,
first,
});
}
}
fn promote_staged(&mut self) {
let Some(staged) = self.staged.take() else {
return;
};
if let Err(e) = self.rebuild_output(staged.plan) {
self.fail(e.to_string());
self.teardown();
self.index = staged.index;
self.state = State::Stopped;
return;
}
self.conv = Some(Converter::new(staged.spec, staged.plan, self.speed()));
let duration = staged.stream.duration();
self.index = staged.index;
self.stream = Some(staged.stream);
self.written = 0;
self.marks.clear();
self.marks.push_back((0, staged.index, duration));
self.shared.track_start.store(0, Ordering::Relaxed);
self.shared.position_offset.store(0, Ordering::Relaxed);
self.convert_and_carry(&staged.first);
if let Some(o) = &self.out {
o.play();
}
}
fn push_carry(&mut self) {
let Some(out) = self.out.as_mut() else { return };
if self.carry.is_empty() {
return;
}
let (written, _) = out.producer.push_partial_slice(&self.carry);
let n = written.len();
if n > 0 {
self.carry.drain(..n);
self.written += n as u64 / out.plan.channels as u64;
}
}
fn advance_marks(&mut self) {
let played = self.shared.frames_out.load(Ordering::Relaxed);
while self.marks.len() > 1 {
let next_start = self.marks[1].0;
if played < next_start {
break;
}
self.marks.pop_front();
let (start, idx, _) = self.marks[0];
self.index = idx;
self.shared.track_start.store(start, Ordering::Relaxed);
self.shared.position_offset.store(0, Ordering::Relaxed);
}
}
fn publish(&self) {
let Ok(mut s) = self.status.lock() else {
return;
};
s.state = self.state;
s.index = self.index;
s.duration = self.marks.front().and_then(|(_, _, d)| *d);
s.source = self.conv.as_ref().map(Converter::src);
s.output_rate = self.out.as_ref().map(|o| o.plan.rate).unwrap_or(0);
s.resampling = self.conv.as_ref().is_some_and(Converter::resampling);
s.semitones = self.semitones;
}
}