pub mod decode;
#[cfg(feature = "opus")]
pub mod opus;
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};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use cpal::Device;
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 output::{Output, Plan, Shared};
use resample::Resample;
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>),
TogglePause,
Next,
Prev,
Stop,
Seek(Duration),
SeekBy(i64),
SetVolume(f32),
SpeedBy(i32),
SpeedReset,
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,
}
impl Status {
pub fn current(&self) -> Option<&PathBuf> {
self.queue.get(self.index)
}
}
pub struct Player {
tx: Sender<Cmd>,
status: Arc<Mutex<Status>>,
shared: Arc<Shared>,
handle: Option<std::thread::JoinHandle<()>>,
}
impl Player {
pub fn new() -> Result<Self, output::OutputError> {
let device = output::default_device()?;
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(device, 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 _ = self.tx.send(cmd);
}
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()
}
}
impl Drop for Player {
fn drop(&mut self) {
let _ = self.tx.send(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);
struct Staged {
stream: AudioStream,
spec: Spec,
plan: Plan,
index: usize,
first: Vec<f32>,
}
struct Engine {
device: Device,
rx: Receiver<Cmd>,
status: Arc<Mutex<Status>>,
shared: Arc<Shared>,
device_errors: (Sender<String>, Receiver<String>),
out: Option<Output>,
stream: Option<AudioStream>,
resampler: Option<Resample>,
src: Spec,
queue: Arc<[PathBuf]>,
index: usize,
state: State,
semitones: i32,
written: u64,
marks: VecDeque<(u64, usize, Option<Duration>)>,
staged: Option<Staged>,
carry: Vec<f32>,
scratch: Vec<f32>,
}
impl Engine {
fn new(
device: Device,
rx: Receiver<Cmd>,
status: Arc<Mutex<Status>>,
shared: Arc<Shared>,
) -> Self {
Engine {
device,
rx,
status,
shared,
device_errors: std::sync::mpsc::channel(),
out: None,
stream: None,
resampler: None,
src: Spec {
rate: 0,
channels: 0,
},
queue: Arc::default(),
index: 0,
state: State::Stopped,
semitones: 0,
written: 0,
marks: VecDeque::new(),
staged: None,
carry: Vec::new(),
scratch: Vec::new(),
}
}
fn run(mut self) {
loop {
let wait = if self.state == State::Playing {
Duration::from_millis(5)
} else {
Duration::from_millis(100)
};
match self.rx.recv_timeout(wait) {
Ok(Cmd::Quit) => break,
Ok(cmd) => {
self.handle(cmd);
while let Ok(next) = self.rx.try_recv() {
if matches!(next, Cmd::Quit) {
return;
}
self.handle(next);
}
}
Err(RecvTimeoutError::Timeout) => {}
Err(RecvTimeoutError::Disconnected) => break,
}
while let Ok(e) = self.device_errors.1.try_recv() {
self.fail(format!("audio device: {e}"));
}
if self.state == State::Playing {
self.pump();
}
self.advance_marks();
self.publish();
}
}
fn handle(&mut self, cmd: Cmd) {
match cmd {
Cmd::Play(paths, index) => {
self.queue = paths.into();
self.index = index.min(self.queue.len().saturating_sub(1));
self.teardown();
if !self.queue.is_empty() {
self.start(self.index);
}
}
Cmd::Enqueue(paths) => {
let first_new = self.queue.len();
self.queue = self.queue.iter().cloned().chain(paths).collect();
if self.state == State::Stopped {
self.start(first_new);
} else if self.stream.is_none() && self.staged.is_none() {
self.stage_from(first_new);
}
}
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);
}
}
},
Cmd::Next => {
let next = self.index + 1;
if next < self.queue.len() {
self.teardown();
self.index = next;
self.start(next);
} else {
self.teardown();
self.state = State::Stopped;
}
}
Cmd::Prev => {
let restart = self.elapsed() > Duration::from_secs(3) || self.index == 0;
let target = if restart { self.index } else { self.index - 1 };
self.teardown();
self.index = target;
self.start(target);
}
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(v) => self.shared.set_volume(v),
Cmd::SpeedBy(delta) => self.set_semitones(self.semitones + delta),
Cmd::SpeedReset => self.set_semitones(0),
Cmd::Quit => {}
}
}
fn elapsed(&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(),
)
}
fn teardown(&mut self) {
self.stream = None;
self.resampler = None;
self.staged = 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, mut i: usize) {
let (stream, first, spec) = loop {
let Some(path) = self.queue.get(i).cloned() else {
self.state = State::Stopped;
return;
};
match self.open_track(&path) {
Some(opened) => break opened,
None => i += 1,
}
};
let plan = match output::negotiate(&self.device, spec) {
Ok(p) => p,
Err(e) => {
self.fail(e.to_string());
self.state = State::Stopped;
return;
}
};
if let Err(e) = self.rebuild_output(plan) {
self.fail(e.to_string());
self.state = State::Stopped;
return;
}
self.index = i;
self.src = spec;
self.resampler = self.make_resampler(spec, plan);
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();
}
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.device,
plan,
self.shared.clone(),
self.device_errors.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 make_resampler(&self, spec: Spec, plan: Plan) -> Option<Resample> {
if !plan.needs_resample(spec) && self.semitones == 0 {
return None;
}
Resample::new(spec.rate, plan.rate, plan.channels, self.speed())
}
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.stream.is_some() {
self.seek(at);
}
}
fn seek(&mut self, pos: Duration) {
let Some(stream) = self.stream.as_mut() else {
return;
};
if stream.seek(pos).is_err() {
return;
}
if let Some(plan) = self.out.as_ref().map(|o| o.plan) {
self.resampler = self.make_resampler(self.src, plan);
}
self.carry.clear();
self.staged = None;
let plan = self.out.as_ref().map(|o| o.plan);
if let Some(plan) = plan {
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;
}
}
}
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 offset = (pos.as_secs_f64() * rate).max(0.0) as u64;
self.shared.frames_out.store(0, Ordering::Relaxed);
self.shared.track_start.store(0, Ordering::Relaxed);
self.shared.position_offset.store(offset, Ordering::Relaxed);
}
fn convert_and_carry(&mut self, decoded: &[f32]) {
let Some(out) = &self.out else { return };
let dst_ch = out.plan.channels as usize;
let src_ch = self.src.channels as usize;
if let Some(r) = self.resampler.as_mut() {
self.scratch.clear();
r.push(decoded, &mut self.scratch);
let resampled = std::mem::take(&mut self.scratch);
output::remap_channels(&resampled, src_ch, dst_ch, &mut self.carry);
self.scratch = resampled;
} else {
output::remap_channels(decoded, src_ch, dst_ch, &mut self.carry);
}
}
fn pump(&mut self) {
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) => {
if let Some(r) = self.resampler.as_mut() {
let mut tail = Vec::new();
r.flush(&mut tail);
if !tail.is_empty() {
let (src_ch, dst_ch) = (
self.src.channels as usize,
self.out
.as_ref()
.map(|o| o.plan.channels as usize)
.unwrap_or(2),
);
output::remap_channels(&tail, src_ch, dst_ch, &mut self.carry);
}
}
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 next = self
.marks
.back()
.map(|(_, i, _)| i + 1)
.unwrap_or(self.index + 1);
self.stage_from(next);
}
fn stage_from(&mut self, mut next: usize) {
let (stream, first, spec) = loop {
let Some(path) = self.queue.get(next).cloned() else {
return;
};
match self.open_track(&path) {
Some(opened) => break opened,
None => next += 1,
}
};
let plan = match output::negotiate(&self.device, spec) {
Ok(p) => p,
Err(e) => {
self.fail(e.to_string());
return;
}
};
let duration = stream.duration();
if self.out.as_ref().map(|o| o.plan) == Some(plan) {
self.src = spec;
self.resampler = self.make_resampler(spec, plan);
self.marks.push_back((self.written, next, duration));
self.stream = Some(stream);
self.convert_and_carry(&first);
} else {
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.state = State::Stopped;
return;
}
self.src = staged.spec;
self.resampler = self.make_resampler(staged.spec, staged.plan);
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.queue = self.queue.clone();
s.index = self.index;
s.duration = self.marks.front().and_then(|(_, _, d)| *d);
s.source = if self.src.rate == 0 {
None
} else {
Some(self.src)
};
s.output_rate = self.out.as_ref().map(|o| o.plan.rate).unwrap_or(0);
s.resampling = self.resampler.is_some();
s.semitones = self.semitones;
}
}