pub mod commands;
pub mod history;
mod renderer;
pub mod state;
pub mod undo;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::thread;
use thiserror::Error;
use crate::audio::{
analyzer::VizAnalyzer,
backend::{self, AudioBackend, AudioEngineHandle, BackendError, SampleRateWatch},
buffer, streaming,
viz::{VizBuffer, VizSnapshot},
};
use crate::remote::client::PlaybackReportState;
use buffer::PlaybackTimeline;
use commands::{CommandChannel, PlayerCommand};
use history::{InFlight, PlayEvent, PlayRecorder, PlaybackReport};
use state::{
ItemState, PlayMode, PlaybackSource, PlaybackState, QueueItemId, Repeat, SharedPlayerState,
Sleep, SleepTimer, TrackInfo,
};
use undo::{UndoEntry, UndoStack};
pub(crate) const RING_BUFFER_SIZE: usize = 192_000 * 2;
const SEEK_END_GUARD_MS: u64 = 500;
const BOUNDARY_SLACK: std::time::Duration = std::time::Duration::from_millis(5);
const FADE_CHECK: std::time::Duration = std::time::Duration::from_millis(50);
#[derive(Debug, Error)]
pub enum PlayerError {
#[error("backend error: {0}")]
Backend(#[from] BackendError),
#[error("decode error: {0}")]
Decode(#[from] buffer::DecodeError),
#[error("renderer: {0}")]
Renderer(String),
#[error("{0}")]
Unplayable(String),
}
#[derive(Clone)]
struct StreamSource {
path: PathBuf,
bytes_written: Arc<crate::remote::downloads::ByteFeed>,
total: u64,
mode: streaming::ProbeMode,
}
fn media_extension(path: &Path) -> Option<String> {
crate::remote::download::strip_part_suffix(path)
.extension()
.and_then(|e| e.to_str())
.map(str::to_ascii_lowercase)
}
fn hint_for(path: &Path) -> symphonia::core::formats::probe::Hint {
let mut hint = symphonia::core::formats::probe::Hint::new();
if let Some(ext) = media_extension(path) {
hint.with_extension(&ext);
}
hint
}
fn lengthless_mode_for(path: &Path) -> streaming::ProbeMode {
match media_extension(path).as_deref() {
Some("ogg" | "oga" | "opus" | "spx" | "m4a" | "m4b" | "mp4" | "mov") => {
streaming::ProbeMode::LengthlessWholeEnd
}
_ => streaming::ProbeMode::Lengthless,
}
}
pub struct Player {
shared_state: Arc<SharedPlayerState>,
commands: CommandChannel,
transport: Transport,
timeline: Arc<PlaybackTimeline>,
viz_buffer: Arc<VizBuffer>,
viz_snapshot: Arc<VizSnapshot>,
_viz_analyzer: VizAnalyzer,
undo_stack: UndoStack,
batch_buffer: Option<Vec<UndoEntry>>,
output_device_name: Option<String>,
backend: Box<dyn AudioBackend>,
stream_mode: streaming::ProbeMode,
history: Option<PlayRecorder>,
downloads: Option<crate::remote::queue::DownloadQueue>,
in_flight: Option<InFlight>,
lead_in_ends: Option<std::time::Instant>,
session: u64,
silence_waiters: Vec<crossbeam_channel::Sender<u64>>,
dsp: Option<DspCache>,
renderer: Option<renderer::RendererLink>,
mode: PlayMode,
sleep: Option<SleepSet>,
#[cfg(test)]
playback_starts: usize,
#[cfg(test)]
dsp_override: Option<Arc<crate::audio::dsp::Setup>>,
#[cfg(test)]
renderer_dsp_override: Option<Arc<crate::audio::dsp::Setup>>,
resume_renderer: bool,
held_run: Option<Run>,
}
#[derive(Clone, Copy)]
struct SleepSet {
sleep: Sleep,
at: Option<std::time::Instant>,
}
struct DspCache {
config: Arc<crate::config::Config>,
device: String,
setup: Option<Arc<crate::audio::dsp::Setup>>,
}
#[derive(Clone, Copy)]
enum Fade {
Cut,
Short,
Slow,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Run {
Playing,
Paused,
}
enum Transport {
Idle,
Waiting(Waiting),
Loaded(Session),
}
#[derive(Clone, Copy)]
struct Waiting {
id: QueueItemId,
position_ms: u64,
start: Run,
}
impl Waiting {
fn may_stream(&self, id: QueueItemId) -> bool {
self.id == id && self.position_ms == 0
}
}
struct Session {
track: TrackInfo,
run: Run,
lookahead: Arc<parking_lot::Mutex<Vec<state::Lookahead>>>,
output: Output,
}
enum Output {
Local(Local),
Renderer(Box<renderer::Play>),
}
struct Local {
engine: Box<dyn AudioEngineHandle>,
decode_handle: buffer::DecodeHandle,
stream: Option<LiveStream>,
_rate_watch: Option<Box<dyn SampleRateWatch>>,
dsp: Option<crate::audio::dsp::DspStatus>,
}
impl Session {
fn engine(&self) -> Option<&dyn AudioEngineHandle> {
match &self.output {
Output::Local(local) => Some(local.engine.as_ref()),
Output::Renderer(_) => None,
}
}
}
enum Source {
File(PathBuf),
Stream(StreamSource),
}
impl Source {
fn path(&self) -> &Path {
match self {
Source::File(path) => path,
Source::Stream(source) => &source.path,
}
}
}
struct LiveStream {
feed: Arc<crate::remote::downloads::ByteFeed>,
abandoned: Arc<std::sync::atomic::AtomicBool>,
}
impl LiveStream {
fn abandon(&self) {
self.abandoned
.store(true, std::sync::atomic::Ordering::Release);
self.feed.done();
}
}
impl Default for Player {
fn default() -> Self {
Self::new()
}
}
impl Player {
pub fn new() -> Self {
let viz_buffer = VizBuffer::new();
let viz_snapshot = VizSnapshot::new();
let timeline = PlaybackTimeline::new();
let cfg = crate::config::Config::cached();
let viz_analyzer = VizAnalyzer::spawn_with_snapshot(
Arc::clone(&viz_buffer),
&cfg.visualizer,
Arc::clone(&viz_snapshot),
timeline.samples_played_counter(),
);
let shared_state = SharedPlayerState::new();
shared_state.attach_timeline(timeline.clone());
let commands = CommandChannel::new();
let tx = commands.tx.clone();
timeline.on_queued(move || {
let _ = tx.try_send(PlayerCommand::TrackQueued);
});
Self {
shared_state,
commands,
transport: Transport::Idle,
lead_in_ends: None,
dsp: None,
session: 0,
silence_waiters: Vec::new(),
renderer: None,
mode: PlayMode::default(),
sleep: None,
timeline,
viz_buffer,
viz_snapshot,
_viz_analyzer: viz_analyzer,
undo_stack: UndoStack::new(),
batch_buffer: None,
output_device_name: cfg.playback.output_device.clone(),
backend: crate::audio::platform_backend(),
stream_mode: streaming::ProbeMode::Full,
history: None,
downloads: None,
in_flight: None,
#[cfg(test)]
playback_starts: 0,
#[cfg(test)]
dsp_override: None,
#[cfg(test)]
renderer_dsp_override: None,
resume_renderer: false,
held_run: None,
}
}
pub fn shared_state(&self) -> Arc<SharedPlayerState> {
self.shared_state.clone()
}
pub fn timeline(&self) -> Arc<PlaybackTimeline> {
self.timeline.clone()
}
pub fn viz_buffer(&self) -> Arc<VizBuffer> {
self.viz_buffer.clone()
}
pub fn viz_snapshot(&self) -> Arc<VizSnapshot> {
self.viz_snapshot.clone()
}
pub fn undo_stack(&self) -> &UndoStack {
&self.undo_stack
}
#[allow(clippy::type_complexity)]
fn create_engine_for(
&mut self,
info: &buffer::StreamInfo,
consumer: rtrb::Consumer<f32>,
) -> Result<(Box<dyn AudioEngineHandle>, Option<Box<dyn SampleRateWatch>>), PlayerError> {
let device = self.resolve_device()?;
let device_rate = self.backend.get_device_sample_rate(&device)?;
let dsp = match &self.dsp {
Some(cache) if cache.device == device.name => cache.setup.clone(),
_ => self.dsp_for(&device.name),
};
let source_rate =
dsp.as_ref()
.map_or(info.sample_rate, |d| d.output_rate(info.sample_rate)) as f64;
self.shared_state.clear_output_sample_rate();
let settled = if (device_rate - source_rate).abs() > 0.1 {
log::info!(
"switching device sample rate: {}Hz → {}Hz",
device_rate,
source_rate
);
match self.backend.set_device_sample_rate(&device, source_rate) {
Ok(rate) => rate,
Err(e) => {
log::warn!("failed to set device sample rate: {}", e);
device_rate
}
}
} else {
device_rate
};
if (settled - source_rate).abs() > 0.1 {
log::warn!(
"device stayed at {}Hz (wanted {}Hz) — output is resampled, not bit-perfect",
settled,
source_rate
);
}
self.shared_state
.set_output_sample_rate(settled.round() as u32);
let watch_state = self.shared_state.clone();
let watch_name = device.name.clone();
let rate_watch = self.backend.watch_device_sample_rate(
&device,
Box::new(move |rate| {
log::info!("device sample rate changed externally: {rate}Hz on '{watch_name}'");
watch_state.set_output_sample_rate(rate.round() as u32);
}),
);
let engine = self.backend.create_engine(
&device,
source_rate,
info.channels as u32,
consumer,
self.timeline.samples_played_counter(),
)?;
let now = std::time::Instant::now();
let lead_in = if device_rate > 0.0 && (settled - device_rate).abs() > 0.1 {
std::time::Duration::from_millis(
crate::config::Config::cached()
.playback
.rate_switch_lead_in_ms as u64,
)
} else {
self.lead_in_ends
.map(|end| end.saturating_duration_since(now))
.unwrap_or_default()
};
self.lead_in_ends = None;
if !lead_in.is_zero() {
engine.lead_in((settled * lead_in.as_secs_f64()) as u64);
self.lead_in_ends = Some(now + lead_in);
}
Ok((engine, rate_watch))
}
fn dsp_for(&mut self, device: &str) -> Option<Arc<crate::audio::dsp::Setup>> {
let config = crate::config::Config::cached();
#[cfg(test)]
if let Some(setup) = if self
.renderer
.as_ref()
.is_some_and(|l| l.device_name() == device)
{
self.renderer_dsp_override.clone()
} else {
self.dsp_override.clone()
} {
self.dsp = Some(DspCache {
config,
device: device.to_string(),
setup: Some(setup.clone()),
});
return Some(setup);
}
if let Some(cache) = &self.dsp
&& Arc::ptr_eq(&cache.config, &config)
&& cache.device == device
{
return cache.setup.clone();
}
let setup = config.dsp.profile_for(device).and_then(|profile| {
crate::audio::dsp::Setup::load(profile, &crate::config::config_dir())
.inspect_err(|e| {
log::error!(
"dsp: profile '{}' not loaded, playing without it: {e}",
profile.name
)
})
.ok()
.flatten()
.map(Arc::new)
});
self.dsp = Some(DspCache {
config,
device: device.to_string(),
setup: setup.clone(),
});
setup
}
fn reload_dsp(&mut self) {
let device = match &self.renderer {
Some(link) => Ok(link.device_name().to_string()),
None => self.resolve_device().map(|d| d.name),
};
let Ok(device) = device else {
self.dsp = None;
return;
};
let was = self
.dsp
.take()
.filter(|c| c.device == device)
.map(|c| c.setup);
let now = self.dsp_for(&device);
let changed = match was {
Some(was) => was != now,
None => now.is_some(),
};
if changed {
self.restart_on_current_track();
}
}
fn processing(&mut self) -> buffer::Processing {
let cfg = crate::config::Config::cached();
let dsp = match self.resolve_device() {
Ok(device) => self.dsp_for(&device.name),
Err(_) => None,
};
buffer::Processing {
rg_mode: cfg.playback.replaygain,
pre_amp_db: cfg.playback.pre_amp_db,
dsp,
}
}
fn resolve_device(&self) -> Result<backend::DeviceInfo, PlayerError> {
if let Some(ref name) = self.output_device_name {
match self.backend.list_devices() {
Ok(devices) => {
if let Some(dev) = devices.into_iter().find(|d| d.name == *name) {
return Ok(dev);
}
log::warn!(
"configured output device '{}' not found, falling back to default",
name,
);
}
Err(e) => {
log::warn!("failed to list devices while resolving '{}': {}", name, e);
}
}
}
Ok(self.backend.default_device()?)
}
pub fn set_output_device(&mut self, name: String) {
log::info!("switching output device to: {}", name);
self.output_device_name = Some(name.clone());
if let Err(e) = crate::config::Config::persist(|cfg| {
cfg.playback.output_device = Some(name);
cfg.playback.renderer = None;
cfg.playback.renderer_name = None;
}) {
log::error!("failed to save output device config: {}", e);
}
if self.renderer.is_some() {
self.use_renderer(None);
} else {
self.restart_on_current_track();
}
}
pub fn clear_output_device(&mut self) {
log::info!("reverting to system default output device");
self.output_device_name = None;
if let Err(e) = crate::config::Config::persist(|cfg| {
cfg.playback.output_device = None;
cfg.playback.renderer = None;
cfg.playback.renderer_name = None;
}) {
log::error!("failed to save output device config: {}", e);
}
if self.renderer.is_some() {
self.use_renderer(None);
} else {
self.restart_on_current_track();
}
}
fn restart_on_current_track(&mut self) {
let position_ms = self.shared_state.position_ms();
if let Err(e) = self.restart_current(position_ms) {
log::error!("failed to restart playback on device switch: {}", e);
}
}
pub fn output_device_name(&self) -> Option<&str> {
self.output_device_name.as_deref()
}
pub fn command_sender(&self) -> crossbeam_channel::Sender<PlayerCommand> {
self.commands.tx.clone()
}
fn session(&self) -> Option<&Session> {
match &self.transport {
Transport::Loaded(session) => Some(session),
_ => None,
}
}
fn waiting(&self) -> Option<Waiting> {
match self.transport {
Transport::Waiting(waiting) => Some(waiting),
_ => None,
}
}
fn may_stream(&self, waiting: &Waiting, id: QueueItemId) -> bool {
waiting.may_stream(id) && self.renderer.is_none()
}
fn forget_waiting(&mut self) {
if matches!(self.transport, Transport::Waiting(_)) {
self.transport = Transport::Idle;
}
}
fn publish(&self) {
let state = &self.shared_state;
let (playback, waiting, track) = match &self.transport {
Transport::Idle => (PlaybackState::Stopped, false, None),
Transport::Waiting(waiting) => (
match waiting.start {
Run::Playing => PlaybackState::Stopped,
Run::Paused => PlaybackState::Paused,
},
true,
None,
),
Transport::Loaded(session) => (
match session.run {
Run::Playing => PlaybackState::Playing,
Run::Paused => PlaybackState::Paused,
},
false,
Some(&session.track),
),
};
if state.track_info().as_ref() != track {
state.set_track_info(track.cloned());
}
state.set_transport(playback, waiting);
let dsp = match &self.transport {
Transport::Loaded(Session {
output: Output::Local(local),
..
}) => local.dsp.clone(),
Transport::Loaded(Session {
output: Output::Renderer(play),
..
}) => play.dsp(),
_ => None,
};
if state.dsp() != dsp {
state.set_dsp(dsp);
}
state.set_play_mode(self.mode);
state.set_sleep(self.sleep.map(|s| s.sleep));
}
fn intent(&self) -> Option<Run> {
match &self.transport {
Transport::Idle => None,
Transport::Waiting(waiting) => Some(waiting.start),
Transport::Loaded(session) => Some(session.run),
}
}
fn carry_on(&mut self, to: Option<QueueItemId>, intent: Option<Run>) {
match (to, intent) {
(Some(id), Some(start)) => self.cue(id, 0, start),
(Some(id), None) => {
self.stop_playback_and_clear_state();
self.shared_state.set_cursor(Some(id));
}
(None, _) => self.stop_playback_and_clear_state(),
}
}
pub fn play(&mut self, id: QueueItemId) {
if self.in_flight.as_ref().is_some_and(|f| f.item == id) {
self.finish_play();
}
self.forget_waiting();
self.shared_state.set_cursor(Some(id));
match self.shared_state.item_playback_source(id) {
Some(PlaybackSource::Ready(path)) => {
if let Err(e) = self.start_playback(id, &path, 0, Run::Playing) {
log::error!("play failed: {}", e);
}
}
Some(PlaybackSource::Streaming {
path,
bytes_written,
total,
}) => {
self.park(id, 0, Run::Playing);
if self.renderer.is_none() {
self.probe_stream_for_playback(id, &path, bytes_written, total);
}
}
None => {
self.park(id, 0, Run::Playing);
log::info!("play: item {:?} not ready, waiting for TrackReady", id);
}
}
}
fn cue(&mut self, id: QueueItemId, position_ms: u64, start: Run) {
self.forget_waiting();
self.shared_state.set_cursor(Some(id));
let Some(PlaybackSource::Ready(path)) = self.shared_state.item_playback_source(id) else {
if position_ms == 0 && start == Run::Playing {
self.play(id);
return;
}
self.park(id, position_ms, start);
log::info!("cue: {id:?} not on disk yet, opening at {position_ms}ms once it is");
return;
};
if let Err(e) = self.start_playback(id, &path, position_ms, start) {
log::error!("cue failed: {}", e);
return;
}
self.report(match start {
Run::Playing => PlaybackReportState::Playing,
Run::Paused => PlaybackReportState::Paused,
});
}
fn park(&mut self, id: QueueItemId, position_ms: u64, start: Run) {
self.report(PlaybackReportState::Stopped);
self.finish_play();
self.stop_engine();
self.timeline.reset();
self.shared_state.set_position_ms(position_ms);
self.transport = Transport::Waiting(Waiting {
id,
position_ms,
start,
});
}
fn start_playback(
&mut self,
id: QueueItemId,
path: &Path,
seek_ms: u64,
start: Run,
) -> Result<(), PlayerError> {
self.open_session(id, Source::File(path.to_path_buf()), None, seek_ms, start)
}
fn open_session(
&mut self,
id: QueueItemId,
source: Source,
info: Option<buffer::StreamInfo>,
seek_ms: u64,
start: Run,
) -> Result<(), PlayerError> {
#[cfg(test)]
{
self.playback_starts += 1;
}
self.forget_waiting();
let mut result = self.try_open_session(id, source, info, seek_ms, start);
let mut skipped = id;
while matches!(result, Err(PlayerError::Unplayable(_)))
&& self.shared_state.is_cursor(skipped)
{
let Some(next) = self.shared_state.advance_cursor_loadable() else {
log::info!("upnp: nothing further the renderer can play");
self.stop_playback_and_clear_state();
return Ok(());
};
match self.shared_state.item_playback_source(next) {
Some(PlaybackSource::Ready(path)) => {
skipped = next;
result = self.try_open_session(next, Source::File(path), None, 0, start);
}
_ => {
self.cue(next, 0, start);
return Ok(());
}
}
}
if result.is_err() {
self.stop_playback_and_clear_state();
}
self.wake_analyzer();
result
}
fn try_open_session(
&mut self,
id: QueueItemId,
source: Source,
info: Option<buffer::StreamInfo>,
seek_ms: u64,
start: Run,
) -> Result<(), PlayerError> {
if self.renderer.is_some() {
return self.try_open_on_renderer(id, source, info, seek_ms, start);
}
self.stop_engine();
let info = match info {
Some(info) => info,
None => buffer::probe_file(source.path())?,
};
let path = source.path().to_path_buf();
let streaming = matches!(source, Source::Stream(_));
let (first, stream) = match source {
Source::File(path) => (buffer::SourceEntry::from_file(id, path), None),
Source::Stream(source) => {
self.stream_mode = source.mode;
let live = LiveStream {
feed: source.bytes_written.clone(),
abandoned: Default::default(),
};
let status = {
let downloading = self.stream_status_fn(id);
let abandoned = live.abandoned.clone();
Arc::new(move || {
if abandoned.load(std::sync::atomic::Ordering::Acquire) {
streaming::StreamStatus::Failed
} else {
downloading()
}
}) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
};
let entry = buffer::SourceEntry {
id,
path: source.path.clone(),
hint: hint_for(&source.path),
make_mss: Box::new(move || {
let partial = streaming::PartialFileSource::open(
&source.path,
source.bytes_written.clone(),
source.total,
status.clone(),
source.mode,
)?;
Ok(symphonia::core::io::MediaSourceStream::new(
Box::new(partial),
Default::default(),
))
}),
};
(entry, Some(live))
}
};
let track = TrackInfo {
id,
path: path.clone(),
codec: info.codec.clone(),
sample_rate: info.sample_rate,
bit_depth: info.bit_depth,
bitrate_kbps: info.bitrate_kbps,
channels: info.channels,
duration_ms: info.duration_ms,
};
self.shared_state.set_position_ms(seek_ms);
self.on_track_changed(id, seek_ms);
log::info!(
"{}: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
if streaming { "streaming" } else { "playing" },
path.display(),
id,
info.codec,
info.sample_rate,
info.channels,
info.duration_ms,
if seek_ms > 0 {
format!(" @{}ms", seek_ms)
} else {
String::new()
}
);
let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
self.timeline.reset();
let lookahead = Arc::new(parking_lot::Mutex::new(Vec::new()));
let next_track = self.decode_cursor(id, lookahead.clone());
let processing = self.processing();
self.session += 1;
let session = self.session;
let finish_tx = self.commands.tx.clone();
let decode_handle = buffer::start_decode(
first,
producer,
seek_ms,
move || {
let (next_id, next_path) = next_track()?;
Some(buffer::SourceEntry::from_file(next_id, next_path))
},
self.timeline.clone(),
Some(self.viz_buffer.clone()),
processing,
move || {
finish_tx.send(PlayerCommand::DecodeFinished(session)).ok();
},
)?;
let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
let dsp = self
.dsp
.as_ref()
.and_then(|c| c.setup.as_ref())
.map(|s| s.status(info.sample_rate));
if start == Run::Playing {
engine.start()?;
}
self.transport = Transport::Loaded(Session {
track,
run: start,
lookahead,
output: Output::Local(Local {
engine,
decode_handle,
stream,
_rate_watch: rate_watch,
dsp,
}),
});
Ok(())
}
fn decode_cursor(
&self,
id: QueueItemId,
steps: Arc<parking_lot::Mutex<Vec<state::Lookahead>>>,
) -> impl Fn() -> Option<(QueueItemId, PathBuf)> + Send + 'static {
let state = self.shared_state.clone();
let timeline = self.timeline.clone();
move || {
let mut steps = steps.lock();
let current = match steps.last() {
None => id,
Some(step) => step.chosen.as_ref()?.0,
};
let mut step = state.lookahead_after(current)?;
step.boundary = timeline.boundary_count();
let next = step.chosen.clone();
steps.push(step);
next
}
}
fn probe_stream_for_playback(
&self,
id: QueueItemId,
path: &Path,
bytes_written: Arc<crate::remote::downloads::ByteFeed>,
total: u64,
) {
let path = path.to_path_buf();
let tx = self.commands.tx.clone();
let hint = hint_for(&path);
let status = {
let downloading = self.stream_status_fn(id);
let state = self.shared_state.clone();
Arc::new(move || {
if state.is_cursor(id) {
downloading()
} else {
streaming::StreamStatus::Failed
}
}) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
};
let spawned = thread::Builder::new()
.name("koan-stream-probe".into())
.spawn(move || {
let attempt = |mode, wait: bool| {
let open = if wait {
streaming::PartialFileSource::open(
&path,
bytes_written.clone(),
total,
status.clone(),
mode,
)
} else {
streaming::PartialFileSource::open_for_probe(
&path,
bytes_written.clone(),
total,
status.clone(),
mode,
)
};
open.map_err(buffer::DecodeError::Io).and_then(|source| {
let mss = symphonia::core::io::MediaSourceStream::new(
Box::new(source),
Default::default(),
);
buffer::probe_source(mss, &hint)
})
};
let info = match attempt(streaming::ProbeMode::Full, false) {
Ok(info) => Some((info, streaming::ProbeMode::Full)),
Err(e) => {
log::info!(
"stream probe: {} needs more than has arrived ({}), opening without a length",
path.display(),
e
);
let lengthless = lengthless_mode_for(&path);
attempt(lengthless, true)
.ok()
.map(|info| (info, lengthless))
}
};
match info {
Some((info, mode)) => {
tx.send(PlayerCommand::StreamProbed {
id,
info: Box::new(info),
mode,
})
.ok();
}
None => log::info!(
"stream probe: {} cannot start early, waiting for the download",
path.display()
),
}
});
if let Err(e) = spawned {
log::warn!("stream probe: could not spawn for {:?}: {}", id, e);
}
}
fn stream_probed(
&mut self,
id: QueueItemId,
info: buffer::StreamInfo,
mode: streaming::ProbeMode,
) {
let Some(waiting) = self.waiting().filter(|w| self.may_stream(w, id)) else {
return;
};
if !self.shared_state.is_cursor(id) {
return;
}
match self.shared_state.item_playback_source(id) {
Some(PlaybackSource::Ready(path)) => {
if let Err(e) = self.start_playback(id, &path, 0, waiting.start) {
log::error!("stream probe: playback failed: {}", e);
}
}
Some(PlaybackSource::Streaming {
path,
bytes_written,
total,
}) => {
let source = StreamSource {
path,
bytes_written,
total,
mode,
};
if let Err(e) =
self.open_session(id, Source::Stream(source), Some(info), 0, waiting.start)
{
log::error!("stream probe: streaming playback failed: {}", e);
}
}
None => {}
}
}
fn stream_status_fn(
&self,
id: QueueItemId,
) -> Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync> {
let state = self.shared_state.clone();
Arc::new(move || match state.item_state(id) {
Some(ItemState::Ready) => streaming::StreamStatus::Complete,
Some(ItemState::Failed(_)) => streaming::StreamStatus::Failed,
_ => streaming::StreamStatus::Downloading,
})
}
pub fn seek(&mut self, position_ms: u64) {
let Some(id) = self.session().map(|s| s.track.id) else {
return;
};
let seekable = self.shared_state.seekable_ms();
if seekable == 0 {
log::debug!("seek declined: {:?} is not seekable yet", id);
return;
}
let ceiling = seekable.min(
self.shared_state
.duration_ms()
.saturating_sub(SEEK_END_GUARD_MS),
);
let clamped = position_ms.min(ceiling);
if self.seek_on_renderer(clamped) {
return;
}
if let Err(e) = self.restart_current(clamped) {
log::error!("seek failed: {}", e);
}
}
fn restart_current(&mut self, position_ms: u64) -> Result<(), PlayerError> {
let Some(session) = self.session() else {
return Ok(());
};
let (info, start) = (session.track.clone(), session.run);
match self.shared_state.item_playback_source(info.id) {
Some(PlaybackSource::Streaming { .. }) if self.renderer.is_some() => {
self.park(info.id, position_ms, start);
return Ok(());
}
Some(PlaybackSource::Streaming {
path,
bytes_written,
total,
}) => {
let known = buffer::StreamInfo {
codec: info.codec.clone(),
sample_rate: info.sample_rate,
channels: info.channels,
bit_depth: info.bit_depth,
bitrate_kbps: info.bitrate_kbps,
duration_ms: info.duration_ms,
};
let source = StreamSource {
path,
bytes_written,
total,
mode: self.stream_mode,
};
self.open_session(
info.id,
Source::Stream(source),
Some(known),
position_ms,
start,
)?;
}
Some(PlaybackSource::Ready(path)) => {
self.start_playback(info.id, &path, position_ms, start)?;
}
None => return Ok(()),
}
self.report(match start {
Run::Playing => PlaybackReportState::Playing,
Run::Paused => PlaybackReportState::Paused,
});
Ok(())
}
pub fn next_track(&mut self) {
match self.shared_state.advance_cursor_loadable() {
Some(id) => self.play(id),
None => {
log::info!("no more tracks in playlist");
self.stop_playback_and_clear_state();
}
}
}
pub fn prev_track(&mut self) {
match self.shared_state.retreat_cursor() {
Some((id, _)) => self.play(id),
None => {
if let Err(e) = self.restart_current(0) {
log::error!("restart failed: {}", e);
}
}
}
}
pub fn pause(&mut self) {
let fade = crate::config::Config::cached().playback.fade_on_pause;
self.pause_with(if fade { Fade::Short } else { Fade::Cut });
}
fn pause_with(&mut self, fade: Fade) {
match &mut self.transport {
Transport::Idle => return,
Transport::Waiting(waiting) => {
waiting.start = Run::Paused;
return;
}
Transport::Loaded(session) => {
if let Output::Local(local) = &session.output {
match fade {
Fade::Short => local.engine.fade_out(),
Fade::Slow => local.engine.fade_out_slowly(),
Fade::Cut => {
if let Err(e) = local.engine.stop() {
log::error!("pause failed: {}", e);
return;
}
}
}
}
session.run = Run::Paused;
}
}
self.pause_renderer();
self.report(PlaybackReportState::Paused);
}
pub fn resume(&mut self) {
self.answer_silence();
let session = match &mut self.transport {
Transport::Idle => {
if let Some(id) = self.shared_state.cursor() {
self.play(id);
}
return;
}
Transport::Waiting(waiting) => {
let waiting = *waiting;
self.cue(waiting.id, waiting.position_ms, Run::Playing);
return;
}
Transport::Loaded(session) => session,
};
let Output::Local(local) = &session.output else {
self.resume_renderer();
return;
};
let engine = &local.engine;
let resumed = if engine.is_running() || engine.is_silent() {
self.lead_in_ends = None;
engine.fade_in()
} else {
engine.start()
};
if let Err(e) = resumed {
log::error!("resume failed: {}", e);
return;
}
session.run = Run::Playing;
self.wake_analyzer();
self.report(PlaybackReportState::Playing);
}
fn wake_analyzer(&self) {
self.viz_snapshot.wake();
}
pub fn stop(&mut self) {
self.shared_state.clear_playlist();
self.stop_playback_and_clear_state();
}
fn stop_engine(&mut self) {
let playback = match std::mem::replace(&mut self.transport, Transport::Idle) {
Transport::Loaded(playback) => playback,
other => {
self.transport = other;
return;
}
};
self.session += 1;
self.bank_listening();
let local = match playback.output {
Output::Local(local) => local,
Output::Renderer(play) => {
self.halt_renderer(*play);
self.answer_silence();
return;
}
};
let Local {
engine,
mut decode_handle,
stream,
..
} = local;
let _ = engine.stop();
decode_handle.signal_stop();
if let Some(stream) = stream {
stream.abandon();
}
decode_handle.stop();
drop(engine);
self.answer_silence();
}
fn answer_silence(&mut self) {
let position_ms = self.shared_state.position_ms();
for reply in self.silence_waiters.drain(..) {
let _ = reply.send(position_ms);
}
}
fn stop_playback_and_clear_state(&mut self) {
self.forget_waiting();
self.report(PlaybackReportState::Stopped);
self.finish_play();
self.stop_engine();
self.timeline.reset();
self.shared_state.set_position_ms(0);
}
pub fn remove_from_playlist(&mut self, id: QueueItemId) {
let was_cursor = self.shared_state.is_cursor(id);
let resume_after = was_cursor
.then(|| self.shared_state.item_before(id))
.flatten();
self.shared_state.remove_item(id);
if was_cursor {
self.shared_state.set_cursor(resume_after);
let next = self.shared_state.advance_cursor_loadable();
self.carry_on(next, self.intent());
}
}
pub fn track_ready(&mut self, id: QueueItemId) {
if !self.shared_state.is_cursor(id) {
return;
}
if let Some(waiting) = self.waiting().filter(|w| w.id == id) {
log::info!("track_ready: opening {:?}", id);
self.cue(id, waiting.position_ms, waiting.start);
return;
}
if self.session().is_some_and(|s| s.track.id == id) {
log::info!(
"track_ready: download complete while streaming {:?}, refreshing metadata",
id
);
self.refresh_track_metadata(id);
}
}
pub fn track_stream_ready(&mut self, id: QueueItemId) {
let Some(waiting) = self.waiting().filter(|w| self.may_stream(w, id)) else {
return;
};
if !self.shared_state.is_cursor(id) {
return;
}
match self.shared_state.item_playback_source(id) {
Some(PlaybackSource::Streaming {
path,
bytes_written,
total,
}) => {
log::info!("track_stream_ready: probing partial file for {:?}", id);
self.probe_stream_for_playback(id, &path, bytes_written, total);
}
Some(PlaybackSource::Ready(path)) => {
log::info!(
"track_stream_ready: track already ready, starting normal playback for {:?}",
id
);
if let Err(e) = self.start_playback(id, &path, 0, waiting.start) {
log::error!("track_stream_ready playback failed: {}", e);
}
}
None => {} }
}
fn refresh_track_metadata(&mut self, id: QueueItemId) {
use crate::index::metadata;
let path = match self.shared_state.item_path_if_ready(id) {
Some(p) => p,
None => return,
};
match metadata::read_metadata(&path) {
Ok(meta) => {
self.shared_state.update_item_metadata(
id,
meta.title,
meta.artist,
meta.album_artist.unwrap_or_default(),
meta.album,
meta.duration_ms.map(|d| d as u64),
);
if let Transport::Loaded(session) = &mut self.transport
&& session.track.id == id
{
let current = &mut session.track;
let probed = buffer::probe_file(&path).ok();
let duration_ms = probed
.as_ref()
.map(|s| s.duration_ms)
.filter(|d| *d > current.duration_ms)
.unwrap_or(current.duration_ms);
if duration_ms != current.duration_ms {
log::info!(
"track_ready: duration corrected {}ms → {}ms",
current.duration_ms,
duration_ms
);
}
current.duration_ms = duration_ms;
current.path = path.clone();
}
self.shared_state.signal_metadata_refresh();
log::info!("track_ready: metadata refreshed for {:?}", id);
}
Err(e) => {
log::warn!("track_ready: metadata refresh failed for {:?}: {}", id, e);
}
}
}
fn on_track_changed(&mut self, id: QueueItemId, position_ms: u64) {
if let Some(f) = self.in_flight.as_mut().filter(|f| f.item == id) {
f.jump(position_ms);
f.boundary = 0;
return;
}
self.begin_play(id, position_ms, 0);
}
fn begin_play(&mut self, id: QueueItemId, position_ms: u64, boundary: usize) {
self.finish_play();
let track_id = self.shared_state.item_db_id(id);
let mut flight = InFlight::new(id, track_id, position_ms);
flight.boundary = boundary;
self.in_flight = Some(flight);
if let (Some(track_id), Some(recorder)) = (track_id, self.history.as_ref()) {
recorder.record(PlayEvent::Started {
track_id,
position_ms,
});
}
}
fn bank_listening(&mut self) {
let renderer_at = self
.shared_state
.renderer_clock()
.map(|_| self.shared_state.position_ms());
if let Some(f) = self.in_flight.as_mut()
&& let Some(at) = self.timeline.position_in(f.boundary).or(renderer_at)
{
f.advance(at);
}
}
fn report(&self, state: PlaybackReportState) {
let Some(track_id) = self.in_flight.as_ref().and_then(InFlight::track_id) else {
return;
};
if let Some(recorder) = self.history.as_ref() {
recorder.record(PlayEvent::Playback(PlaybackReport {
track_id,
state,
position_ms: self.shared_state.position_ms(),
}));
}
}
fn finish_play(&mut self) -> Option<PlayEvent> {
self.bank_listening();
let flight = self.in_flight.take()?;
let event = PlayEvent::Finished {
track_id: flight.track_id()?,
listened_ms: flight.listened_ms(),
};
if let Some(recorder) = self.history.as_ref() {
recorder.record(event);
}
Some(event)
}
pub fn update_playback_state(&mut self) {
if let Some(SleepSet { at: Some(at), .. }) = self.sleep
&& std::time::Instant::now() >= at
{
self.sleep = None;
log::info!("sleep timer: time's up");
if self.intent() == Some(Run::Playing) {
self.pause_with(Fade::Slow);
}
}
if self.session().is_some()
&& self
.lead_in_ends
.is_some_and(|end| std::time::Instant::now() >= end)
{
self.lead_in_ends = None;
self.shared_state.changed();
}
if let Some(session) = self.session()
&& session.run == Run::Paused
&& let Some(engine) = session.engine()
&& engine.is_running()
&& engine.is_silent()
{
if let Err(e) = engine.stop() {
log::error!("stopping after fade failed: {}", e);
}
self.answer_silence();
}
self.follow_playhead();
self.renderer_tick();
self.publish();
}
fn follow_playhead(&mut self) {
let before = self.in_flight.as_ref().map(|f| (f.item, f.boundary));
self.move_with_playhead();
let after = self.in_flight.as_ref().map(|f| (f.item, f.boundary));
if let (Some((ended, _)), Some((next, _))) = (before, after)
&& before != after
&& self.sleeps_between(Some(ended), Some(next))
{
self.pause_with(Fade::Cut);
}
}
fn move_with_playhead(&mut self) {
if self.session().is_none() {
return;
}
let Some(playhead) = self.timeline.playhead() else {
return;
};
let id = playhead.id;
if self
.in_flight
.as_ref()
.is_none_or(|f| f.item != id || f.boundary != playhead.boundary)
{
self.begin_play(id, playhead.position_ms, playhead.boundary);
}
let Transport::Loaded(session) = &mut self.transport else {
return;
};
if session.track.id == id {
return;
}
let Some((id, path, info, _)) = self.timeline.current_playback() else {
return;
};
log::info!("timeline: now playing {:?}", id);
session.track = TrackInfo {
id,
path,
codec: info.codec,
sample_rate: info.sample_rate,
bit_depth: info.bit_depth,
bitrate_kbps: info.bitrate_kbps,
channels: info.channels,
duration_ms: info.duration_ms,
};
self.shared_state.set_cursor(Some(id));
}
fn sleeps_between(&mut self, ended: Option<QueueItemId>, next: Option<QueueItemId>) -> bool {
let record = |id: Option<QueueItemId>| {
id.and_then(|id| self.shared_state.get_item(id))
.map(|i| (i.album, i.album_artist))
};
let due = match self.sleep.map(|s| s.sleep) {
Some(Sleep::EndOfTrack) => true,
Some(Sleep::EndOfRecord) => next.is_none() || record(ended) != record(next),
_ => false,
};
if due {
log::info!(
"sleep timer: stopping at the end of the {}",
match self.sleep {
Some(SleepSet {
sleep: Sleep::EndOfTrack,
..
}) => "track",
_ => "record",
}
);
self.sleep = None;
}
due
}
fn set_sleep_timer(&mut self, timer: Option<SleepTimer>) {
self.sleep = timer.map(|timer| match timer {
SleepTimer::After { minutes } => {
let after = std::time::Duration::from_secs(u64::from(minutes) * 60);
let unix_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
+ after;
SleepSet {
sleep: Sleep::At {
unix_ms: unix_ms.as_millis() as u64,
},
at: Some(std::time::Instant::now() + after),
}
}
SleepTimer::EndOfTrack => SleepSet {
sleep: Sleep::EndOfTrack,
at: None,
},
SleepTimer::EndOfRecord => SleepSet {
sleep: Sleep::EndOfRecord,
at: None,
},
});
log::info!("sleep timer: {:?}", self.sleep.map(|s| s.sleep));
}
pub fn track_failed(&mut self, id: QueueItemId) {
let Some(waiting) = self.waiting().filter(|w| w.id == id) else {
return;
};
if !self.shared_state.is_cursor(id) {
return;
}
log::info!("track {:?} cannot load, moving on", id);
let next = self.shared_state.advance_cursor_loadable();
self.carry_on(next, Some(waiting.start));
}
fn on_decode_finished(&mut self, session: u64) {
if session != self.session || self.renderer_stream_finished() {
return;
}
log::info!("decode finished, checking for next track");
self.track_ended();
}
fn track_ended(&mut self) {
self.bank_listening();
let ended = self.in_flight.as_ref().map(|f| f.item);
let heard = self.in_flight.as_ref().is_some_and(|f| f.listened_ms() > 0);
if self.mode.repeat != Repeat::Off && !heard {
log::info!("track ended with nothing heard; not repeating it");
self.stop_playback_and_clear_state();
return;
}
let again = (self.mode.repeat == Repeat::One)
.then(|| self.shared_state.cursor())
.flatten()
.filter(|id| {
self.shared_state
.item_state(*id)
.is_some_and(|s| !matches!(s, ItemState::Failed(_)))
});
let next = again.or_else(|| self.shared_state.advance_cursor_loadable());
if next.is_some() && next == ended {
self.finish_play();
}
let intent = match self.intent() {
Some(_) if self.sleeps_between(ended, next) => Some(Run::Paused),
intent => intent,
};
self.carry_on(next, intent);
}
fn revoke_stale_lookahead(&mut self) {
if self.session().is_none() {
return;
}
self.update_playback_state();
let Some(playhead) = self.timeline.playhead() else {
return;
};
let position_ms = playhead.position_ms;
let Some(session) = self.session().filter(|s| s.track.id == playhead.id) else {
return;
};
let stale = {
let steps = session.lookahead.lock();
steps
.iter()
.filter(|step| step.boundary > playhead.boundary)
.any(|step| !self.shared_state.still_follows(step))
};
if !stale {
return;
}
log::info!("queue changed under the lookahead, restarting at {position_ms}ms");
if let Err(e) = self.restart_current(position_ms) {
log::error!("restart after a queue edit failed: {}", e);
}
}
fn snapshot_for_undo(
&self,
ids: &[QueueItemId],
) -> Vec<(Box<state::PlaylistItem>, Option<QueueItemId>)> {
self.shared_state
.items_before(ids)
.into_iter()
.filter_map(|(id, after)| Some((Box::new(self.shared_state.get_item(id)?), after)))
.collect()
}
fn push_undo(&mut self, entry: UndoEntry) {
if let Some(ref mut batch) = self.batch_buffer {
batch.push(entry);
} else {
self.undo_stack.push(entry);
}
}
pub fn process_command(&mut self, cmd: PlayerCommand) {
if cmd.asks_to_play() && !matches!(cmd, PlayerCommand::Cue { .. })
|| matches!(
cmd,
PlayerCommand::UseRenderer(_)
| PlayerCommand::SetOutputDevice(_)
| PlayerCommand::ClearOutputDevice
| PlayerCommand::Pause
| PlayerCommand::PauseAndReport(_)
| PlayerCommand::Stop
| PlayerCommand::ClearPlaylist
| PlayerCommand::ReplacePlaylist { .. }
)
{
self.resume_renderer = false;
self.held_run = None;
}
let cmd = match cmd {
PlayerCommand::Cue {
id,
position_ms,
play: true,
} if self.resume_renderer => {
self.held_run = Some(Run::Playing);
PlayerCommand::Cue {
id,
position_ms,
play: false,
}
}
cmd => cmd,
};
let edits_queue = matches!(
cmd,
PlayerCommand::AddToPlaylist(_)
| PlayerCommand::InsertInPlaylist { .. }
| PlayerCommand::RemoveFromPlaylist(_)
| PlayerCommand::RemoveFromPlaylistBatch(_)
| PlayerCommand::MoveInPlaylist { .. }
| PlayerCommand::MoveItemsInPlaylist { .. }
| PlayerCommand::ReorderPlaylist(_)
| PlayerCommand::Undo
| PlayerCommand::Redo
| PlayerCommand::SetShuffle(_)
| PlayerCommand::SetRepeat(_)
| PlayerCommand::RestorePlayMode(_)
);
self.apply_command(cmd);
if edits_queue {
self.revoke_stale_lookahead();
}
self.publish();
}
fn apply_command(&mut self, cmd: PlayerCommand) {
match cmd {
PlayerCommand::Play(id) => self.play(id),
PlayerCommand::Cue {
id,
position_ms,
play,
} => self.cue(
id,
position_ms,
if play { Run::Playing } else { Run::Paused },
),
PlayerCommand::Pause => self.pause(),
PlayerCommand::PauseAndReport(reply) => {
self.pause();
self.silence_waiters.push(reply);
if self
.session()
.is_none_or(|p| p.engine().is_none_or(|e| !e.is_running()))
{
self.answer_silence();
}
}
PlayerCommand::Resume => self.resume(),
PlayerCommand::Stop => self.stop(),
PlayerCommand::Seek(pos) => self.seek(pos),
PlayerCommand::NextTrack => self.next_track(),
PlayerCommand::PrevTrack => self.prev_track(),
PlayerCommand::AddToPlaylist(items) => {
let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
let whole = self.shared_state.is_empty();
self.shared_state.add_items(items);
if whole
&& self.mode.shuffle
&& let Some(&first) = ids.first()
{
self.shared_state.shuffle_from(first);
}
self.push_undo(UndoEntry::Added { ids });
}
PlayerCommand::UpdatePaths(updates) => {
self.shared_state.update_paths(&updates);
if let Transport::Loaded(session) = &mut self.transport
&& let Some((_, new_path)) =
updates.iter().find(|(id, _)| *id == session.track.id)
{
session.track.path = new_path.clone();
}
}
PlayerCommand::InsertInPlaylist { items, after } => {
let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
self.shared_state.insert_items_after(items, after);
self.push_undo(UndoEntry::Inserted { ids });
}
PlayerCommand::ClearPlaylist => self.clear_playlist(),
PlayerCommand::ReplacePlaylist {
items,
start,
position_ms,
play,
} => {
if items.is_empty() {
self.clear_playlist();
return;
}
let start_id = items.get(start).unwrap_or(&items[0]).id;
self.stop_playback_and_clear_state();
let (old_items, cursor) = self.shared_state.replace_playlist(items);
if self.mode.shuffle {
self.shared_state.shuffle_from(start_id);
}
self.push_undo(UndoEntry::Replaced {
items: old_items,
cursor,
});
self.cue(
start_id,
position_ms,
if play { Run::Playing } else { Run::Paused },
);
}
PlayerCommand::RemoveFromPlaylist(id) => {
let item = self.shared_state.get_item(id);
let after = self.shared_state.item_before(id);
self.remove_from_playlist(id);
if let Some(item) = item {
self.push_undo(UndoEntry::Removed {
items: vec![(Box::new(item), after)],
});
}
}
PlayerCommand::RemoveFromPlaylistBatch(ids) => {
let items_with_pos = self.snapshot_for_undo(&ids);
let intent = self.intent();
let resume_after = match self.shared_state.cursor() {
Some(cursor) if ids.contains(&cursor) => {
Some(self.shared_state.surviving_item_before(cursor, &ids))
}
_ => None,
};
self.shared_state.remove_items(&ids);
if let Some(resume_after) = resume_after {
self.shared_state.set_cursor(resume_after);
let next = self.shared_state.advance_cursor_loadable();
self.carry_on(next, intent);
}
if !items_with_pos.is_empty() {
self.push_undo(UndoEntry::Removed {
items: items_with_pos,
});
}
}
PlayerCommand::MoveInPlaylist { id, target, after } => {
let was_after = self.shared_state.item_before(id);
self.shared_state.move_item(id, target, after);
self.push_undo(UndoEntry::Moved { id, was_after });
}
PlayerCommand::MoveItemsInPlaylist { ids, target, after } => {
let entries = self.shared_state.items_before(&ids);
self.shared_state.move_items(&ids, target, after);
self.push_undo(UndoEntry::MovedBatch { entries });
}
PlayerCommand::ReorderPlaylist(order) => {
let entries = self.shared_state.items_before(&order);
self.shared_state.reorder_to(&order);
self.push_undo(UndoEntry::MovedBatch { entries });
}
PlayerCommand::TrackReady(id) => self.track_ready(id),
PlayerCommand::DecodeFinished(session) => self.on_decode_finished(session),
PlayerCommand::TrackQueued => {}
PlayerCommand::TrackStreamReady(id) => self.track_stream_ready(id),
PlayerCommand::StreamProbed { id, info, mode } => self.stream_probed(id, *info, mode),
PlayerCommand::TrackFailed(id) => self.track_failed(id),
PlayerCommand::CacheTracks(ids) => {
if let Some(downloads) = &self.downloads {
downloads.cache(ids);
}
}
PlayerCommand::Undo => self.execute_undo(),
PlayerCommand::Redo => self.execute_redo(),
PlayerCommand::BeginUndoBatch => {
self.batch_buffer = Some(Vec::new());
}
PlayerCommand::EndUndoBatch => {
if let Some(entries) = self.batch_buffer.take() {
if entries.len() == 1 {
self.undo_stack.push(entries.into_iter().next().unwrap());
} else if !entries.is_empty() {
self.undo_stack.push(UndoEntry::Batch(entries));
}
}
}
PlayerCommand::SetOutputDevice(name) => self.set_output_device(name),
PlayerCommand::RestartOutput => {
log::info!("restarting audio output");
self.restart_on_current_track();
}
PlayerCommand::ClearOutputDevice => self.clear_output_device(),
PlayerCommand::ReloadDsp => self.reload_dsp(),
PlayerCommand::UseRenderer(connection) => {
self.remember_renderer(connection.as_deref());
self.use_renderer(connection);
}
PlayerCommand::ResumeRenderer(connection) => {
if std::mem::take(&mut self.resume_renderer) && self.renderer.is_none() {
log::info!(
"upnp: back to {}, as last time",
connection.session.renderer().name
);
self.use_renderer(Some(connection));
if self.held_run.take() == Some(Run::Playing) {
self.resume();
}
} else {
log::info!(
"upnp: not going back to {}: playback or the output moved first",
connection.session.renderer().name
);
}
}
PlayerCommand::ResumeRendererMissed => {
if std::mem::take(&mut self.resume_renderer)
&& self.held_run.take() == Some(Run::Playing)
{
log::info!("upnp: playing here what was held for the renderer");
self.resume();
}
}
PlayerCommand::ReleaseRenderer(reply) => {
self.release_renderer();
let _ = reply.send(());
}
PlayerCommand::SetRendererVolume(volume) => self.set_renderer_volume(volume),
PlayerCommand::Renderer { session, event } => self.on_renderer_event(session, event),
PlayerCommand::SetShuffle(on) => self.set_shuffle(on),
PlayerCommand::SetRepeat(repeat) => self.mode.repeat = repeat,
PlayerCommand::RestorePlayMode(mode) => self.mode = mode,
PlayerCommand::SetSleepTimer(timer) => self.set_sleep_timer(timer),
}
}
fn set_shuffle(&mut self, on: bool) {
if self.mode.shuffle == on {
return;
}
let order = self.shared_state.shuffle_order();
if on {
self.shared_state.shuffle_after_cursor();
} else {
self.shared_state.unshuffle();
}
self.push_undo(UndoEntry::Shuffled {
shuffle: self.mode.shuffle,
order,
});
self.mode.shuffle = on;
}
fn clear_playlist(&mut self) {
self.stop_playback_and_clear_state();
let (items, cursor) = self.shared_state.snapshot_playlist();
self.shared_state.clear_playlist();
self.push_undo(UndoEntry::Replaced { items, cursor });
}
fn apply_entry(&mut self, entry: UndoEntry) -> UndoEntry {
match entry {
UndoEntry::Added { ids } | UndoEntry::Inserted { ids } => {
let items_with_pos = self.snapshot_for_undo(&ids);
self.shared_state.remove_items(&ids);
UndoEntry::Removed {
items: items_with_pos,
}
}
UndoEntry::Removed { items } => {
let mut ids = Vec::with_capacity(items.len());
for (item, after) in items {
ids.push(item.id);
self.shared_state.insert_item_at(*item, after);
}
UndoEntry::Added { ids }
}
UndoEntry::Moved { id, was_after } => {
let current_after = self.shared_state.item_before(id);
self.shared_state.move_item_to(id, was_after);
UndoEntry::Moved {
id,
was_after: current_after,
}
}
UndoEntry::MovedBatch { entries } => {
let ids: Vec<QueueItemId> = entries.iter().map(|(id, _)| *id).collect();
let current_positions = self.shared_state.items_before(&ids);
self.shared_state.move_items_to(&entries);
UndoEntry::MovedBatch {
entries: current_positions,
}
}
UndoEntry::Replaced { items, cursor } => {
let (current_items, current_cursor) = self.shared_state.snapshot_playlist();
self.shared_state.restore_playlist(items, cursor);
UndoEntry::Replaced {
items: current_items,
cursor: current_cursor,
}
}
UndoEntry::Shuffled { shuffle, order } => {
let inverse = UndoEntry::Shuffled {
shuffle: self.mode.shuffle,
order: self.shared_state.shuffle_order(),
};
self.shared_state.restore_shuffle_order(&order);
self.mode.shuffle = shuffle;
inverse
}
UndoEntry::Batch(entries) => {
let mut inverses: Vec<_> = entries
.into_iter()
.rev()
.map(|e| self.apply_entry(e))
.collect();
inverses.reverse();
UndoEntry::Batch(inverses)
}
}
}
fn reconcile_playback(&mut self) {
let orphaned = match &self.transport {
Transport::Idle => false,
Transport::Waiting(waiting) => self.shared_state.get_item(waiting.id).is_none(),
Transport::Loaded(session) => self.shared_state.get_item(session.track.id).is_none(),
};
if !orphaned {
return;
}
let cursor = self.shared_state.cursor();
self.carry_on(cursor, self.intent());
}
fn execute_undo(&mut self) {
let Some(entry) = self.undo_stack.pop_undo() else {
return;
};
let inverse = self.apply_entry(entry);
self.undo_stack.push_redo(inverse);
self.reconcile_playback();
}
fn execute_redo(&mut self) {
let Some(entry) = self.undo_stack.pop_redo() else {
return;
};
let inverse = self.apply_entry(entry);
self.undo_stack.push_undo_keep_redo(inverse);
self.reconcile_playback();
}
pub fn run(&mut self) {
use crossbeam_channel::RecvTimeoutError;
let rx = self.commands.rx.clone();
loop {
let received = match self.next_wake() {
Some(at) => rx.recv_deadline(at),
None => rx.recv().map_err(|_| RecvTimeoutError::Disconnected),
};
match received {
Ok(cmd) => self.process_command(cmd),
Err(RecvTimeoutError::Timeout) => {}
Err(RecvTimeoutError::Disconnected) => break,
}
self.update_playback_state();
}
self.stop();
}
fn next_wake(&self) -> Option<std::time::Instant> {
let sleep = self.sleep.and_then(|s| s.at);
match (self.next_event(), sleep) {
(Some(a), Some(b)) => Some(a.min(b)),
(a, b) => a.or(b),
}
}
fn next_event(&self) -> Option<std::time::Instant> {
let session = self.session()?;
let now = std::time::Instant::now();
let Output::Local(local) = &session.output else {
let next_track = (session.run == Run::Playing && self.streaming_to_renderer())
.then(|| self.timeline.until_next_track())
.flatten()
.map(|left| now + left + BOUNDARY_SLACK);
return match (next_track, self.renderer_deadline()) {
(Some(a), Some(b)) => Some(a.min(b)),
(a, b) => a.or(b),
};
};
match session.run {
Run::Playing => {
let next_track = self
.timeline
.until_next_track()
.map(|left| now + left + BOUNDARY_SLACK);
match (next_track, self.lead_in_ends) {
(Some(a), Some(b)) => Some(a.min(b)),
(a, b) => a.or(b),
}
}
Run::Paused if local.engine.is_running() => Some(now + FADE_CHECK),
Run::Paused => None,
}
}
fn remember_renderer(&self, connection: Option<&crate::upnp::Connection>) {
let renderer = connection.map(|c| c.session.renderer());
if let Err(e) = crate::config::Config::persist(|cfg| {
cfg.playback.renderer = renderer.map(|r| r.udn.clone());
cfg.playback.renderer_name = renderer.map(|r| r.name.clone());
}) {
log::error!("failed to save the output: {e}");
}
}
pub fn spawn() -> (
Arc<SharedPlayerState>,
Arc<PlaybackTimeline>,
Arc<VizSnapshot>,
crossbeam_channel::Sender<PlayerCommand>,
) {
Self::spawn_with(false)
}
pub fn spawn_for_listening() -> (
Arc<SharedPlayerState>,
Arc<PlaybackTimeline>,
Arc<VizSnapshot>,
crossbeam_channel::Sender<PlayerCommand>,
) {
Self::spawn_with(true)
}
fn renderer_to_resume(listening: bool) -> Option<String> {
listening
.then(|| crate::config::Config::cached().playback.renderer.clone())
.flatten()
}
fn spawn_with(
listening: bool,
) -> (
Arc<SharedPlayerState>,
Arc<PlaybackTimeline>,
Arc<VizSnapshot>,
crossbeam_channel::Sender<PlayerCommand>,
) {
let mut player = Self::new();
player.history = PlayRecorder::spawn();
let state = player.shared_state();
let timeline = player.timeline();
let viz_snapshot = player.viz_snapshot();
let tx = player.command_sender();
player.downloads = Some(crate::remote::queue::DownloadQueue::spawn(
tx.clone(),
state.clone(),
));
if let Some(udn) = Self::renderer_to_resume(listening) {
player.resume_renderer = true;
crate::upnp::resume(udn, &tx);
}
thread::Builder::new()
.name("koan-player".into())
.spawn(move || player.run())
.expect("failed to spawn player thread");
(state, timeline, viz_snapshot, tx)
}
}
#[cfg(test)]
mod tests {
#[test]
fn a_headless_player_does_not_go_back_to_a_renderer() {
crate::config::isolate_config_for_tests();
crate::config::Config::persist(|cfg| {
cfg.playback.renderer = Some("uuid:headless-test".into());
})
.unwrap();
assert_eq!(Player::renderer_to_resume(false), None);
}
#[test]
fn a_download_in_progress_is_known_by_its_own_extension() {
use std::path::Path;
assert_eq!(
lengthless_mode_for(Path::new("/c/t.m4a.part")),
streaming::ProbeMode::LengthlessWholeEnd
);
assert_eq!(
lengthless_mode_for(Path::new("/c/t.OPUS.part")),
streaming::ProbeMode::LengthlessWholeEnd
);
assert_eq!(
lengthless_mode_for(Path::new("/c/t.flac.part")),
streaming::ProbeMode::Lengthless
);
assert_eq!(
media_extension(Path::new("/c/t.m4a.part")).as_deref(),
Some("m4a")
);
assert_eq!(
media_extension(Path::new("/c/t.mp3")).as_deref(),
Some("mp3")
);
}
use super::*;
use state::PlaylistItem;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
fn make_item(title: &str) -> PlaylistItem {
PlaylistItem {
playlist_entry_id: None,
id: QueueItemId::new(),
db_id: None,
path: PathBuf::from(format!("/music/{title}.flac")),
title: title.to_string(),
artist: String::new(),
album_artist: String::new(),
album: String::new(),
year: None,
codec: None,
track_number: None,
disc: None,
duration_ms: None,
state: ItemState::Ready,
pre_shuffle: None,
}
}
pub(super) fn playlist_ids(player: &Player) -> Vec<QueueItemId> {
let (items, _) = player.shared_state.snapshot_playlist();
items.iter().map(|i| i.id).collect()
}
fn playlist_titles(player: &Player) -> Vec<String> {
let (items, _) = player.shared_state.snapshot_playlist();
items.iter().map(|i| i.title.clone()).collect()
}
fn pending_item(title: &str) -> PlaylistItem {
PlaylistItem {
playlist_entry_id: None,
state: ItemState::Pending,
..make_item(title)
}
}
fn test_session(id: QueueItemId, engine: Box<dyn AudioEngineHandle>) -> Session {
Session {
track: TrackInfo {
id,
path: PathBuf::from("/music/t.flac"),
codec: String::new(),
sample_rate: 44_100,
bit_depth: None,
bitrate_kbps: None,
channels: 2,
duration_ms: 1_000,
},
run: Run::Playing,
lookahead: Default::default(),
output: Output::Local(Local {
engine,
decode_handle: buffer::DecodeHandle::new_for_test(Default::default()),
stream: None,
_rate_watch: None,
dsp: None,
}),
}
}
fn pretend_playing(player: &mut Player, id: QueueItemId) {
assert!(
player.shared_state.get_item(id).is_some(),
"item is in the queue"
);
player.stop_engine();
player.transport = Transport::Loaded(test_session(
id,
Box::new(NullEngine {
starts: Default::default(),
running: Default::default(),
lead_in: Default::default(),
}),
));
player.shared_state.set_cursor(Some(id));
player.publish();
}
fn playing_id(player: &Player) -> Option<QueueItemId> {
player.shared_state.track_info().map(|t| t.id)
}
fn seed(player: &mut Player, n: usize) -> Vec<QueueItemId> {
let items: Vec<_> = (0..n).map(|i| make_item(&format!("t{i}"))).collect();
let ids = items.iter().map(|i| i.id).collect();
player.process_command(PlayerCommand::AddToPlaylist(items));
ids
}
fn listen(player: &mut Player, from_ms: u64, to_ms: u64) {
let mut at = from_ms;
if let Some(f) = player.in_flight.as_mut() {
f.advance(at); }
while at < to_ms {
at = (at + 50).min(to_ms);
if let Some(f) = player.in_flight.as_mut() {
f.advance(at);
}
}
}
fn start(player: &mut Player, track_id: i64) -> QueueItemId {
let id = QueueItemId::new();
player.on_track_changed(id, 0);
player
.in_flight
.as_mut()
.unwrap()
.track_id_for_test(track_id);
id
}
#[test]
fn a_gapless_transition_closes_the_outgoing_track_and_opens_the_next() {
let mut player = Player::new();
start(&mut player, 11);
listen(&mut player, 0, 200_000);
let b = QueueItemId::new();
player.on_track_changed(b, 0);
let f = player
.in_flight
.as_ref()
.expect("the next track is counting");
assert_eq!(f.item, b);
assert_eq!(f.listened_ms(), 0, "and starts from nothing");
}
#[test]
fn a_track_skipped_seconds_in_is_still_history() {
let mut player = Player::new();
start(&mut player, 7);
listen(&mut player, 0, 2_000);
let event = player
.finish_play()
.expect("putting something on is a thing you did, however briefly");
assert!(matches!(
event,
history::PlayEvent::Finished {
track_id: 7,
listened_ms: 2_000
}
));
}
#[test]
fn a_track_is_closed_out_once() {
let mut player = Player::new();
start(&mut player, 7);
listen(&mut player, 0, 200_000);
assert!(player.finish_play().is_some());
assert!(player.finish_play().is_none());
}
#[test]
fn seeking_around_a_track_does_not_enter_it_twice() {
let mut player = Player::new();
let id = start(&mut player, 7);
listen(&mut player, 0, 120_000);
player.on_track_changed(id, 30_000);
assert_eq!(
player.in_flight.as_ref().unwrap().listened_ms(),
120_000,
"the seek kept the count rather than restarting it"
);
listen(&mut player, 30_000, 40_000);
let Some(history::PlayEvent::Finished { listened_ms, .. }) = player.finish_play() else {
panic!("still one play");
};
assert_eq!(listened_ms, 130_000);
assert!(player.finish_play().is_none());
}
#[test]
fn a_track_that_is_not_in_the_library_is_not_recorded() {
let mut player = Player::new();
let id = QueueItemId::new();
player.on_track_changed(id, 0);
listen(&mut player, 0, 200_000);
assert!(player.finish_play().is_none());
}
#[test]
fn stopping_closes_out_what_was_heard() {
let mut player = Player::new();
start(&mut player, 7);
listen(&mut player, 0, 150_000);
player.stop_playback_and_clear_state();
assert!(player.in_flight.is_none(), "the stop consumed it");
}
#[test]
fn resume_with_nothing_loaded_plays_the_cursor() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: Default::default(),
});
let item = PlaylistItem {
db_id: Some(5),
path,
..make_item("t")
};
let id = item.id;
player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
player.shared_state.set_cursor(Some(id));
assert!(player.session().is_none());
player.process_command(PlayerCommand::Resume);
assert!(player.session().is_some(), "the cursor's track started");
assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_player_with_nothing_coming_has_nothing_to_wake_for() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: Default::default(),
});
let item = PlaylistItem {
path,
..make_item("t")
};
let id = item.id;
player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
assert_eq!(player.next_wake(), None, "stopped");
player.process_command(PlayerCommand::Play(id));
assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
assert_eq!(
player.next_wake(),
None,
"playing, with no track queued after it"
);
if let Transport::Loaded(session) = &mut player.transport {
session.engine().unwrap().stop().unwrap();
session.run = Run::Paused;
}
assert_eq!(player.next_wake(), None, "paused");
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_track_cued_paused_never_starts_the_output() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: starts.clone(),
});
let item = PlaylistItem {
path,
..make_item("t")
};
let id = item.id;
player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
player.process_command(PlayerCommand::Cue {
id,
position_ms: 2_000,
play: false,
});
assert!(player.session().is_some(), "loaded");
assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
loop {
match player
.commands
.rx
.recv_timeout(std::time::Duration::from_secs(5))
{
Ok(PlayerCommand::TrackQueued) => break,
Ok(_) => {}
Err(e) => panic!("the decoder never queued the track: {e}"),
}
}
let at = player.shared_state.position_ms();
assert!((1_750..=2_000).contains(&at), "cued at {at}ms");
assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
player.process_command(PlayerCommand::Seek(4_000));
assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
assert_eq!(starts.load(Ordering::Relaxed), 0);
player.process_command(PlayerCommand::Resume);
assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
assert_eq!(starts.load(Ordering::Relaxed), 1);
player.process_command(PlayerCommand::Stop);
}
fn downloading_wav(dir: &Path) -> (Player, QueueItemId, Arc<std::sync::atomic::AtomicUsize>) {
let path = dir.join("t.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: starts.clone(),
});
let item = PlaylistItem {
path,
state: ItemState::Pending,
..make_item("t")
};
let id = item.id;
player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
(player, id, starts)
}
fn land(player: &mut Player, id: QueueItemId) {
player.shared_state.update_item_state(id, ItemState::Ready);
player.process_command(PlayerCommand::TrackReady(id));
}
#[test]
fn track_ready_never_resurrects_an_item_that_failed() {
let mut player = Player::new();
let item = PlaylistItem {
state: ItemState::Failed("gone".into()),
..make_item("t")
};
let id = item.id;
player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
player.process_command(PlayerCommand::TrackReady(id));
assert!(matches!(
player.shared_state.get_item(id).map(|i| i.state),
Some(ItemState::Failed(_))
));
}
fn await_queued(player: &Player) {
loop {
match player
.commands
.rx
.recv_timeout(std::time::Duration::from_secs(5))
{
Ok(PlayerCommand::TrackQueued) => return,
Ok(_) => {}
Err(e) => panic!("the decoder never queued the track: {e}"),
}
}
}
#[test]
fn a_cue_on_a_downloading_track_opens_at_its_position_once_it_lands() {
let dir = tempfile::tempdir().unwrap();
let (mut player, id, starts) = downloading_wav(dir.path());
player.process_command(PlayerCommand::Cue {
id,
position_ms: 6_000,
play: true,
});
assert!(player.session().is_none(), "nothing opened early");
assert_eq!(player.shared_state.playback_state(), PlaybackState::Stopped);
assert_eq!(starts.load(Ordering::Relaxed), 0);
land(&mut player, id);
assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
assert_eq!(player.playback_starts, 1, "opened once, at the position");
assert_eq!(starts.load(Ordering::Relaxed), 1);
await_queued(&player);
let at = player.shared_state.position_ms();
assert!((5_750..=6_000).contains(&at), "opened at {at}ms");
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_paused_cue_on_a_downloading_track_lands_paused() {
let dir = tempfile::tempdir().unwrap();
let (mut player, id, starts) = downloading_wav(dir.path());
player.process_command(PlayerCommand::Cue {
id,
position_ms: 3_000,
play: false,
});
land(&mut player, id);
assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
await_queued(&player);
let at = player.shared_state.position_ms();
assert!((2_750..=3_000).contains(&at), "cued at {at}ms");
player.process_command(PlayerCommand::Stop);
}
#[test]
fn playing_something_else_forgets_a_waiting_cue() {
let dir = tempfile::tempdir().unwrap();
let (mut player, id, _) = downloading_wav(dir.path());
let other = seed(&mut player, 1)[0];
player.process_command(PlayerCommand::Cue {
id,
position_ms: 6_000,
play: true,
});
player.process_command(PlayerCommand::Play(other));
player.process_command(PlayerCommand::Play(id));
assert!(player.waiting().is_some_and(|w| w.position_ms == 0));
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_paused_track_whose_download_lands_stays_paused_where_it_was() {
let dir = tempfile::tempdir().unwrap();
let (mut player, id, starts) = downloading_wav(dir.path());
player.shared_state.update_item_state(id, ItemState::Ready);
player.process_command(PlayerCommand::Cue {
id,
position_ms: 3_000,
play: false,
});
await_queued(&player);
player.process_command(PlayerCommand::TrackReady(id));
assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
assert_eq!(player.playback_starts, 1, "not reopened");
assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
player.process_command(PlayerCommand::Stop);
}
#[test]
fn pausing_a_track_on_its_way_opens_it_paused() {
let dir = tempfile::tempdir().unwrap();
let (mut player, id, starts) = downloading_wav(dir.path());
player.process_command(PlayerCommand::Play(id));
assert!(player.shared_state.is_waiting());
assert!(player.shared_state.wants_to_play(), "a toggle pauses it");
assert!(!player.shared_state.is_idle(), "adding tracks leaves it be");
player.process_command(PlayerCommand::Pause);
assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
assert!(!player.shared_state.wants_to_play());
player.shared_state.update_item_state(id, ItemState::Ready);
player.process_command(PlayerCommand::TrackReady(id));
assert!(player.session().is_some(), "loaded");
assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_cue_paused_and_resumed_on_its_way_keeps_its_position() {
let dir = tempfile::tempdir().unwrap();
let (mut player, id, starts) = downloading_wav(dir.path());
player.process_command(PlayerCommand::Cue {
id,
position_ms: 6_000,
play: true,
});
player.process_command(PlayerCommand::Pause);
player.process_command(PlayerCommand::Resume);
assert!(player.session().is_none(), "still on its way");
assert_eq!(player.shared_state.playback_state(), PlaybackState::Stopped);
player.shared_state.update_item_state(id, ItemState::Ready);
player.process_command(PlayerCommand::TrackReady(id));
assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
assert_eq!(starts.load(Ordering::Relaxed), 1);
await_queued(&player);
let at = player.shared_state.position_ms();
assert!((5_750..=6_000).contains(&at), "opened at {at}ms");
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_restored_cursor_does_not_start_when_its_download_lands() {
let dir = tempfile::tempdir().unwrap();
let (mut player, id, starts) = downloading_wav(dir.path());
player.shared_state.set_cursor(Some(id));
player.shared_state.update_item_state(id, ItemState::Ready);
player.process_command(PlayerCommand::TrackReady(id));
assert!(player.session().is_none());
assert_eq!(starts.load(Ordering::Relaxed), 0);
}
#[test]
fn replacing_the_queue_keeps_a_transfer_both_queues_want() {
let pending = |title: &str| PlaylistItem {
db_id: Some(7),
state: ItemState::Pending,
..make_item(title)
};
let mut player = Player::new();
let old = pending("old");
let old_id = old.id;
player.process_command(PlayerCommand::AddToPlaylist(vec![old]));
let store = player.shared_state.downloads().clone();
store.claim(7, Some(old_id));
let before = player.shared_state.pending_version();
player.process_command(PlayerCommand::ReplacePlaylist {
items: vec![pending("again")],
start: 0,
position_ms: 0,
play: false,
});
assert_eq!(
player.shared_state.pending_version(),
before + 1,
"one change, so no reader can see the playlist between two"
);
store.resync(&player.shared_state.pending_downloads());
assert!(!store.abandoned(7));
}
#[test]
fn a_hand_off_opens_paused_at_its_position_in_one_command() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: starts.clone(),
});
seed(&mut player, 2);
let item = PlaylistItem {
path,
..make_item("t")
};
let id = item.id;
player.process_command(PlayerCommand::ReplacePlaylist {
items: vec![make_item("before"), item],
start: 1,
position_ms: 4_000,
play: false,
});
assert_eq!(
playlist_titles(&player),
vec!["before", "t"],
"replaced, not added to"
);
assert_eq!(player.shared_state.cursor(), Some(id));
assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
assert_eq!(player.playback_starts, 1, "opened once, at the position");
assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
await_queued(&player);
let at = player.shared_state.position_ms();
assert!((3_750..=4_000).contains(&at), "opened at {at}ms");
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["t0", "t1"], "one undo step");
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_decode_end_from_a_replaced_session_is_ignored() {
let dir = tempfile::tempdir().unwrap();
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: Default::default(),
});
let items: Vec<_> = ["a", "b", "c"]
.iter()
.map(|name| {
let path = dir.path().join(format!("{name}.wav"));
crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
PlaylistItem {
path,
..make_item(name)
}
})
.collect();
let ids: Vec<_> = items.iter().map(|i| i.id).collect();
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::Play(ids[0]));
let stale = player.session;
player.process_command(PlayerCommand::Play(ids[1]));
player.process_command(PlayerCommand::DecodeFinished(stale));
assert_eq!(player.shared_state.cursor(), Some(ids[1]));
assert_eq!(player.playback_starts, 2);
player.process_command(PlayerCommand::Stop);
}
fn queued_wavs(dir: &Path, names: &[&str]) -> (Player, Vec<QueueItemId>) {
let (player, ids) = wavs_playing(dir, names, 1.0);
for _ in &ids {
await_queued(&player);
}
(player, ids)
}
fn wavs_playing(dir: &Path, names: &[&str], seconds: f32) -> (Player, Vec<QueueItemId>) {
wavs_in(dir, names, seconds, Repeat::Off)
}
fn wavs_in(
dir: &Path,
names: &[&str],
seconds: f32,
repeat: Repeat,
) -> (Player, Vec<QueueItemId>) {
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: Default::default(),
});
let items: Vec<_> = names
.iter()
.zip(1..)
.map(|(name, track)| {
let path = dir.join(format!("{name}.wav"));
crate::test_utils::generate_wav(&path, 8_000, 1, seconds, 16);
PlaylistItem {
path,
db_id: Some(track),
..make_item(name)
}
})
.collect();
let ids: Vec<_> = items.iter().map(|i| i.id).collect();
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::SetRepeat(repeat));
player.process_command(PlayerCommand::Play(ids[0]));
(player, ids)
}
fn queued_at_least(player: &Player, n: usize) -> Vec<QueueItemId> {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
loop {
let queued = player.timeline.queued_after_playhead();
if queued.len() >= n {
return queued;
}
assert!(
std::time::Instant::now() < deadline,
"the decoder queued {queued:?}, never {n}"
);
std::thread::yield_now();
}
}
#[test]
fn repeating_the_queue_runs_the_last_track_into_the_first() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = wavs_in(dir.path(), &["a", "b"], 1.0, Repeat::Queue);
let queued = queued_at_least(&player, 3);
assert_eq!(
queued[..3],
[ids[1], ids[0], ids[1]],
"gapless, round again"
);
let steps = player.session().unwrap().lookahead.lock().clone();
assert!(!steps[0].wrapped);
let wrap = steps.iter().find(|s| s.after == ids[1]).unwrap();
assert_eq!(wrap.next, Some(ids[0]));
assert!(wrap.wrapped);
assert_eq!(player.playback_starts, 1);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn turning_repeat_off_takes_back_a_wrap_already_queued() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = wavs_in(dir.path(), &["a"], 1.0, Repeat::Queue);
assert_eq!(queued_at_least(&player, 1)[0], ids[0]);
player.process_command(PlayerCommand::SetRepeat(Repeat::Off));
assert_eq!(player.playback_starts, 2, "restarted at the playhead");
wait_for_lookahead(&player);
assert!(player.timeline.queued_after_playhead().is_empty());
player.process_command(PlayerCommand::Stop);
}
#[test]
fn repeating_one_track_queues_it_again_and_each_pass_is_a_play() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = wavs_in(dir.path(), &["a", "b"], 1.0, Repeat::One);
let (recorder, events) = history::PlayRecorder::capture();
player.history = Some(recorder);
player.in_flight = None;
player.on_track_changed(ids[0], 0);
assert_eq!(queued_at_least(&player, 2)[..2], [ids[0], ids[0]]);
player
.timeline
.samples_played
.store(8_400, Ordering::Relaxed);
player.update_playback_state();
let flight = player.in_flight.as_ref().unwrap();
assert_eq!((flight.item, flight.boundary), (ids[0], 1));
assert_eq!(player.shared_state.cursor(), Some(ids[0]));
let events: Vec<_> = events.try_iter().collect();
let started = events
.iter()
.filter(|e| matches!(e, PlayEvent::Started { track_id: 1, .. }))
.count();
assert_eq!(started, 2, "two plays: {events:?}");
assert!(
events.contains(&PlayEvent::Finished {
track_id: 1,
listened_ms: 1_000
}),
"the first pass banked whole: {events:?}"
);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn next_moves_on_from_a_track_repeating() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = wavs_in(dir.path(), &["a", "b"], 1.0, Repeat::One);
player.process_command(PlayerCommand::NextTrack);
assert_eq!(playing_id(&player), Some(ids[1]));
player.process_command(PlayerCommand::NextTrack);
assert_eq!(
playing_id(&player),
Some(ids[0]),
"round from the last, as repeating does"
);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_repeat_of_nothing_heard_stops_rather_than_going_round() {
let dir = tempfile::tempdir().unwrap();
let (mut player, _) = wavs_in(dir.path(), &["a"], 1.0, Repeat::One);
player.process_command(PlayerCommand::DecodeFinished(player.session));
assert!(matches!(player.transport, Transport::Idle));
}
fn session_run(player: &Player) -> Option<Run> {
player.session().map(|s| s.run)
}
#[test]
fn a_sleep_timer_fades_out_at_its_time_and_keeps_the_queue() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = wavs_playing(dir.path(), &["a", "b"], 1.0);
player.process_command(PlayerCommand::SetSleepTimer(Some(SleepTimer::After {
minutes: 30,
})));
let Some(Sleep::At { unix_ms }) = player.shared_state.sleep() else {
panic!("published as a time");
};
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
assert!(unix_ms.abs_diff(now + 30 * 60_000) < 5_000);
let at = player.sleep.and_then(|s| s.at).unwrap();
assert!(player.next_wake().is_some_and(|w| w <= at), "woken for it");
player.sleep.as_mut().unwrap().at = Some(std::time::Instant::now());
player.update_playback_state();
assert_eq!(session_run(&player), Some(Run::Paused));
assert!(
player.session().unwrap().engine().unwrap().is_running(),
"fading, not cut off"
);
assert_eq!(player.shared_state.sleep(), None, "spent");
assert_eq!(playlist_ids(&player), ids, "the queue as it was");
assert_eq!(player.shared_state.cursor(), Some(ids[0]));
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_sleep_timer_going_off_while_paused_only_ends() {
let dir = tempfile::tempdir().unwrap();
let (mut player, _) = wavs_playing(dir.path(), &["a"], 1.0);
player.process_command(PlayerCommand::Pause);
player.process_command(PlayerCommand::SetSleepTimer(Some(SleepTimer::After {
minutes: 1,
})));
player.sleep.as_mut().unwrap().at = Some(std::time::Instant::now());
player.update_playback_state();
assert_eq!(player.shared_state.sleep(), None);
player.process_command(PlayerCommand::Resume);
assert_eq!(
session_run(&player),
Some(Run::Playing),
"nothing left to go off"
);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn cancelling_a_sleep_timer_leaves_playback_alone() {
let dir = tempfile::tempdir().unwrap();
let (mut player, _) = wavs_playing(dir.path(), &["a"], 1.0);
player.process_command(PlayerCommand::SetSleepTimer(Some(SleepTimer::After {
minutes: 1,
})));
player.process_command(PlayerCommand::SetSleepTimer(None));
assert_eq!(player.shared_state.sleep(), None);
assert!(player.sleep.is_none());
player.update_playback_state();
assert_eq!(session_run(&player), Some(Run::Playing));
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_sleep_timer_for_the_end_of_the_track_stops_at_the_boundary() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = queued_wavs(dir.path(), &["a", "b"]);
player.process_command(PlayerCommand::SetSleepTimer(Some(SleepTimer::EndOfTrack)));
assert_eq!(player.shared_state.sleep(), Some(Sleep::EndOfTrack));
player
.timeline
.samples_played
.store(4_000, Ordering::Relaxed);
player.update_playback_state();
assert_eq!(session_run(&player), Some(Run::Playing));
assert_eq!(player.shared_state.sleep(), Some(Sleep::EndOfTrack));
player
.timeline
.samples_played
.store(8_400, Ordering::Relaxed);
player.update_playback_state();
assert_eq!(player.shared_state.cursor(), Some(ids[1]));
assert_eq!(session_run(&player), Some(Run::Paused));
assert!(
!player.session().unwrap().engine().unwrap().is_running(),
"stopped at once, nothing of the next track faded through"
);
assert_eq!(player.shared_state.sleep(), None);
assert_eq!(playlist_ids(&player), ids);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_sleep_timer_for_the_end_of_the_track_stops_a_track_repeating() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = wavs_in(dir.path(), &["a", "b"], 1.0, Repeat::One);
assert_eq!(queued_at_least(&player, 1)[0], ids[0]);
player.process_command(PlayerCommand::SetSleepTimer(Some(SleepTimer::EndOfTrack)));
player
.timeline
.samples_played
.store(8_400, Ordering::Relaxed);
player.update_playback_state();
assert_eq!(player.shared_state.cursor(), Some(ids[0]));
assert_eq!(session_run(&player), Some(Run::Paused));
assert_eq!(player.shared_state.sleep(), None);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_sleep_timer_for_the_end_of_the_track_opens_the_next_paused() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = wavs_playing(dir.path(), &["a", "b"], 1.0);
player.process_command(PlayerCommand::SetSleepTimer(Some(SleepTimer::EndOfTrack)));
player.process_command(PlayerCommand::DecodeFinished(player.session));
assert_eq!(player.shared_state.cursor(), Some(ids[1]));
assert_eq!(player.intent(), Some(Run::Paused));
assert_eq!(player.shared_state.sleep(), None);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_sleep_timer_for_the_end_of_the_record_plays_the_record_out() {
let dir = tempfile::tempdir().unwrap();
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: Default::default(),
});
let items: Vec<_> = [("a", "Low"), ("b", "Low"), ("c", "Heroes")]
.iter()
.map(|(name, album)| {
let path = dir.path().join(format!("{name}.wav"));
crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
PlaylistItem {
path,
album: album.to_string(),
album_artist: "David Bowie".into(),
..make_item(name)
}
})
.collect();
let ids: Vec<_> = items.iter().map(|i| i.id).collect();
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::Play(ids[0]));
player.process_command(PlayerCommand::SetSleepTimer(Some(SleepTimer::EndOfRecord)));
player.process_command(PlayerCommand::DecodeFinished(player.session));
assert_eq!(player.shared_state.cursor(), Some(ids[1]));
assert_eq!(player.intent(), Some(Run::Playing), "the same record");
assert_eq!(player.shared_state.sleep(), Some(Sleep::EndOfRecord));
player.process_command(PlayerCommand::DecodeFinished(player.session));
assert_eq!(player.shared_state.cursor(), Some(ids[2]));
assert_eq!(player.intent(), Some(Run::Paused), "the next record waits");
assert_eq!(player.shared_state.sleep(), None);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn shuffle_on_and_off_again_puts_the_queue_back() {
let mut player = Player::new();
let ids = seed(&mut player, 20);
pretend_playing(&mut player, ids[3]);
player.process_command(PlayerCommand::SetShuffle(true));
let shuffled = playlist_ids(&player);
assert!(player.shared_state.play_mode().shuffle);
assert_eq!(
shuffled[..4],
ids[..4],
"nothing up to the playing track moves"
);
assert_ne!(shuffled, ids);
let mut sorted = shuffled.clone();
sorted.sort_by_key(|id| ids.iter().position(|i| i == id));
assert_eq!(sorted, ids, "the same items");
let extra = make_item("extra");
let extra_id = extra.id;
player.process_command(PlayerCommand::InsertInPlaylist {
items: vec![extra],
after: ids[3],
});
player.process_command(PlayerCommand::RemoveFromPlaylist(ids[10]));
player.process_command(PlayerCommand::SetShuffle(false));
let mut expected = ids.clone();
expected.remove(10);
expected.insert(4, extra_id);
assert_eq!(playlist_ids(&player), expected, "added since stays put");
assert!(!player.shared_state.play_mode().shuffle);
assert!(
player
.shared_state
.shuffle_order()
.iter()
.all(|(_, pre)| pre.is_none())
);
}
#[test]
fn a_queue_replaced_while_shuffled_plays_shuffled_from_its_start() {
let mut player = Player::new();
seed(&mut player, 3);
player.process_command(PlayerCommand::SetShuffle(true));
let items: Vec<_> = (0..20).map(|i| make_item(&format!("n{i}"))).collect();
let given: Vec<_> = items.iter().map(|i| i.id).collect();
player.process_command(PlayerCommand::ReplacePlaylist {
items,
start: 5,
position_ms: 0,
play: false,
});
let shuffled = playlist_ids(&player);
assert_eq!(shuffled[0], given[5], "the start first");
assert_eq!(player.shared_state.cursor(), Some(given[5]));
assert_ne!(shuffled, given);
player.process_command(PlayerCommand::SetShuffle(false));
assert_eq!(
playlist_ids(&player),
given,
"off gives the queue as it came"
);
}
#[test]
fn a_queue_added_to_an_empty_one_while_shuffled_plays_shuffled() {
let mut player = Player::new();
player.process_command(PlayerCommand::SetShuffle(true));
let given = seed(&mut player, 20);
assert_eq!(playlist_ids(&player)[0], given[0]);
assert_ne!(playlist_ids(&player), given);
player.process_command(PlayerCommand::SetShuffle(false));
assert_eq!(playlist_ids(&player), given);
}
#[test]
fn shuffle_is_one_undo_step() {
let mut player = Player::new();
let ids = seed(&mut player, 10);
pretend_playing(&mut player, ids[0]);
player.process_command(PlayerCommand::SetShuffle(true));
let shuffled = playlist_ids(&player);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_ids(&player), ids);
assert!(!player.shared_state.play_mode().shuffle);
player.process_command(PlayerCommand::Redo);
assert_eq!(playlist_ids(&player), shuffled);
assert!(player.shared_state.play_mode().shuffle);
player.process_command(PlayerCommand::SetShuffle(false));
assert_eq!(
playlist_ids(&player),
ids,
"the redone shuffle still unwinds"
);
}
#[test]
fn removing_the_last_track_playing_carries_on_from_the_top_when_repeating() {
let mut player = Player::new();
let ids = seed(&mut player, 3);
player.process_command(PlayerCommand::SetRepeat(Repeat::Queue));
player.shared_state.set_cursor(Some(ids[2]));
player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
assert_eq!(player.shared_state.cursor(), Some(ids[0]));
}
#[test]
fn removing_a_track_the_decoder_queued_takes_it_back() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = queued_wavs(dir.path(), &["a", "b", "c"]);
assert_eq!(player.timeline.queued_after_playhead(), ids[1..]);
player.process_command(PlayerCommand::RemoveFromPlaylist(ids[1]));
assert_eq!(player.playback_starts, 2, "restarted at the playhead");
assert_eq!(player.shared_state.cursor(), Some(ids[0]));
await_queued(&player);
await_queued(&player);
assert_eq!(player.timeline.queued_after_playhead(), vec![ids[2]]);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_track_inserted_to_play_next_is_not_skipped() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = queued_wavs(dir.path(), &["a", "b"]);
let path = dir.path().join("next.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
let next = PlaylistItem {
path,
..make_item("next")
};
let next_id = next.id;
player.process_command(PlayerCommand::InsertInPlaylist {
items: vec![next],
after: ids[0],
});
assert_eq!(player.playback_starts, 2);
for _ in 0..3 {
await_queued(&player);
}
assert_eq!(
player.timeline.queued_after_playhead(),
vec![next_id, ids[1]]
);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_track_the_decoder_could_not_open_does_not_make_every_edit_restart() {
let dir = tempfile::tempdir().unwrap();
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: Default::default(),
});
let wav = |name: &str| {
let path = dir.path().join(format!("{name}.wav"));
crate::test_utils::generate_wav(&path, 8_000, 1, 30.0, 16);
PlaylistItem {
path,
..make_item(name)
}
};
let items = vec![wav("a"), make_item("missing"), wav("c")];
let ids: Vec<_> = items.iter().map(|i| i.id).collect();
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::Play(ids[0]));
await_queued(&player);
await_queued(&player);
player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("later")]));
assert_eq!(
player.playback_starts, 1,
"nothing the decoder decided changed"
);
assert_eq!(player.timeline.queued_after_playhead(), vec![ids[2]]);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn a_track_added_after_the_decoder_reached_the_end_follows_gaplessly() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = queued_wavs(dir.path(), &["a"]);
wait_for_lookahead(&player);
let path = dir.path().join("b.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
let b = PlaylistItem {
path,
..make_item("b")
};
let b_id = b.id;
player.process_command(PlayerCommand::AddToPlaylist(vec![b]));
assert_eq!(player.playback_starts, 2, "restarted to queue it");
assert_eq!(player.shared_state.cursor(), Some(ids[0]));
await_queued(&player);
await_queued(&player);
assert_eq!(player.timeline.queued_after_playhead(), vec![b_id]);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn playing_next_a_track_still_downloading_takes_back_the_lookahead() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = queued_wavs(dir.path(), &["a", "b"]);
let remote = pending_item("remote");
let remote_id = remote.id;
player.process_command(PlayerCommand::InsertInPlaylist {
items: vec![remote],
after: ids[0],
});
assert_eq!(player.playback_starts, 2, "b is taken back out of the ring");
await_queued(&player);
wait_for_lookahead(&player);
let session = player.session().unwrap();
let step = session.lookahead.lock().last().cloned().unwrap();
assert_eq!(step.next, Some(remote_id), "the decoder waits for it");
assert!(player.timeline.queued_after_playhead().is_empty());
player.process_command(PlayerCommand::Stop);
}
#[test]
fn gapless_playback_waits_for_a_track_still_downloading() {
let dir = tempfile::tempdir().unwrap();
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: Default::default(),
});
let wav = |name: &str| {
let path = dir.path().join(format!("{name}.wav"));
crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
PlaylistItem {
path,
..make_item(name)
}
};
let items = vec![wav("a"), pending_item("arriving"), wav("c")];
let ids: Vec<_> = items.iter().map(|i| i.id).collect();
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::Play(ids[0]));
await_queued(&player);
wait_for_lookahead(&player);
assert!(
player.timeline.queued_after_playhead().is_empty(),
"c is not queued over the track before it"
);
player.process_command(PlayerCommand::DecodeFinished(player.session));
assert_eq!(player.shared_state.cursor(), Some(ids[1]));
assert!(player.waiting().is_some());
player.process_command(PlayerCommand::Stop);
}
fn wait_for_lookahead(player: &Player) {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while player
.session()
.is_some_and(|s| s.lookahead.lock().is_empty())
{
assert!(
std::time::Instant::now() < deadline,
"the decoder never looked ahead"
);
std::thread::yield_now();
}
}
#[test]
fn an_edit_after_what_the_decoder_queued_leaves_playback_alone() {
let dir = tempfile::tempdir().unwrap();
let (mut player, ids) = wavs_playing(dir.path(), &["a", "b"], 30.0);
await_queued(&player);
await_queued(&player);
let path = dir.path().join("later.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
player.process_command(PlayerCommand::AddToPlaylist(vec![PlaylistItem {
path,
..make_item("later")
}]));
assert_eq!(player.playback_starts, 1);
assert_eq!(player.timeline.queued_after_playhead(), vec![ids[1]]);
player.process_command(PlayerCommand::Stop);
}
#[test]
fn undoing_the_add_of_a_track_on_its_way_forgets_it() {
let dir = tempfile::tempdir().unwrap();
let (mut player, _, starts) = downloading_wav(dir.path());
let path = dir.path().join("added.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
let added = PlaylistItem {
path,
state: ItemState::Pending,
..make_item("added")
};
let added_id = added.id;
player.process_command(PlayerCommand::AddToPlaylist(vec![added]));
player.process_command(PlayerCommand::Play(added_id));
player.process_command(PlayerCommand::Undo);
assert!(player.waiting().is_none());
assert_eq!(player.shared_state.playback_state(), PlaybackState::Stopped);
player.process_command(PlayerCommand::Redo);
player
.shared_state
.update_item_state(added_id, ItemState::Ready);
player.process_command(PlayerCommand::TrackReady(added_id));
assert!(player.session().is_none());
assert_eq!(starts.load(Ordering::Relaxed), 0);
}
pub(super) struct Rng(pub(super) u64);
impl Rng {
pub(super) fn next(&mut self) -> u64 {
self.0 ^= self.0 << 13;
self.0 ^= self.0 >> 7;
self.0 ^= self.0 << 17;
self.0
}
pub(super) fn below(&mut self, n: usize) -> usize {
(self.next() % n.max(1) as u64) as usize
}
pub(super) fn coin(&mut self) -> bool {
self.next() & 1 == 1
}
}
pub(super) fn asks_to_play(cmd: &PlayerCommand) -> bool {
cmd.asks_to_play()
}
pub(super) fn check_invariants(
player: &Player,
wanted_before: bool,
asked: bool,
) -> Result<(), String> {
let state = &player.shared_state;
let ids = playlist_ids(player);
let listed = |id: QueueItemId| ids.contains(&id);
let mut seen = std::collections::HashSet::new();
if !ids.iter().all(|id| seen.insert(*id)) {
return Err("an item is in the playlist twice".into());
}
let cursor = state.cursor();
if cursor.is_some_and(|c| !listed(c)) {
return Err("the cursor is on an item not in the playlist".into());
}
let info = state.track_info();
match (&player.transport, &info) {
(Transport::Loaded(_), None) => return Err("loaded, but no track published".into()),
(Transport::Loaded(_), Some(info)) => {
if !listed(info.id) {
return Err("playing an item not in the playlist".into());
}
if cursor != Some(info.id) {
return Err("the cursor is not on what is playing".into());
}
if !matches!(
state.playback_state(),
PlaybackState::Playing | PlaybackState::Paused
) {
return Err("loaded, but published as stopped".into());
}
}
(_, Some(_)) => return Err("a track is published with nothing loaded".into()),
(Transport::Waiting(w), None) => {
if cursor != Some(w.id) || !listed(w.id) {
return Err("waiting on an item that is not the cursor's".into());
}
let expected = match w.start {
Run::Playing => PlaybackState::Stopped,
Run::Paused => PlaybackState::Paused,
};
if state.playback_state() != expected {
return Err("a wait published in the wrong state".into());
}
}
(Transport::Idle, None) => {
if state.playback_state() != PlaybackState::Stopped {
return Err("idle, but not published as stopped".into());
}
}
}
if state.is_waiting() != player.waiting().is_some() {
return Err("the published wait disagrees with the player".into());
}
if player
.timeline
.queued_after_playhead()
.into_iter()
.any(|id| !listed(id))
{
return Err("the decoder has queued an item no longer in the playlist".into());
}
if !asked && !wanted_before && state.wants_to_play() {
return Err("started without being asked".into());
}
if state.play_mode() != player.mode {
return Err("the published mode disagrees with the player".into());
}
if !player.mode.shuffle && state.shuffle_order().iter().any(|(_, pre)| pre.is_some()) {
return Err("shuffle is off, but an item remembers a place to go back to".into());
}
if player.mode.repeat == Repeat::Off
&& let Some(session) = player.session()
&& let Some(playhead) = player.timeline.playhead()
&& session.track.id == playhead.id
&& session
.lookahead
.lock()
.iter()
.filter(|step| step.boundary > playhead.boundary)
.any(|step| step.wrapped || step.next == Some(step.after))
{
return Err("repeat is off, but the decoder has queued a wrap".into());
}
Ok(())
}
#[test]
fn random_use_keeps_the_player_honest() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 0.5, 16);
for seed in 1..=40u64 {
let mut rng = Rng(seed.wrapping_mul(0x9E37_79B9_7F4A_7C15) | 1);
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: Default::default(),
});
let fresh = |rng: &mut Rng| {
let state = match rng.below(3) {
0 => ItemState::Pending,
_ => ItemState::Ready,
};
PlaylistItem {
path: path.clone(),
state,
..make_item("t")
}
};
let start: Vec<_> = (0..5).map(|_| fresh(&mut rng)).collect();
player.process_command(PlayerCommand::AddToPlaylist(start));
let mut history = Vec::new();
for step in 0..150 {
let ids = playlist_ids(&player);
let pick = |rng: &mut Rng| ids.get(rng.below(ids.len())).copied();
let cmd = match rng.below(23) {
0 => pick(&mut rng).map(PlayerCommand::Play),
1 => pick(&mut rng).map(|id| PlayerCommand::Cue {
id,
position_ms: if rng.coin() { 0 } else { 200 },
play: rng.coin(),
}),
2 => Some(PlayerCommand::Pause),
3 => Some(PlayerCommand::Resume),
4 => Some(PlayerCommand::NextTrack),
5 => Some(PlayerCommand::PrevTrack),
6 => Some(PlayerCommand::Seek(100)),
7 => pick(&mut rng).map(PlayerCommand::RemoveFromPlaylist),
8 => Some(PlayerCommand::RemoveFromPlaylistBatch(
(0..2).filter_map(|_| pick(&mut rng)).collect(),
)),
9 => pick(&mut rng).zip(pick(&mut rng)).map(|(id, target)| {
PlayerCommand::MoveInPlaylist {
id,
target,
after: rng.coin(),
}
}),
10 => pick(&mut rng).map(|after| PlayerCommand::InsertInPlaylist {
items: vec![fresh(&mut rng)],
after,
}),
11 => Some(PlayerCommand::AddToPlaylist(vec![fresh(&mut rng)])),
12 => Some(PlayerCommand::Undo),
13 => Some(PlayerCommand::Redo),
14 | 15 => {
let (items, _) = player.shared_state.snapshot_playlist();
let pending: Vec<_> = items
.iter()
.filter(|i| matches!(i.state, ItemState::Pending))
.map(|i| i.id)
.collect();
pending.get(rng.below(pending.len())).map(|&id| {
if rng.below(4) == 0 {
player
.shared_state
.update_item_state(id, ItemState::Failed("gone".into()));
PlayerCommand::TrackFailed(id)
} else {
player.shared_state.update_item_state(id, ItemState::Ready);
PlayerCommand::TrackReady(id)
}
})
}
16 => Some(PlayerCommand::DecodeFinished(player.session)),
17 => Some(PlayerCommand::DecodeFinished(
player.session.wrapping_sub(1),
)),
18 => Some(PlayerCommand::ClearPlaylist),
20 => Some(PlayerCommand::SetShuffle(rng.coin())),
21 => Some(PlayerCommand::SetRepeat(
[Repeat::Off, Repeat::Queue, Repeat::One][rng.below(3)],
)),
_ => Some(PlayerCommand::ReplacePlaylist {
items: (0..3).map(|_| fresh(&mut rng)).collect(),
start: rng.below(4),
position_ms: if rng.coin() { 0 } else { 200 },
play: rng.coin(),
}),
};
let Some(cmd) = cmd else { continue };
let label = format!("{cmd:?}");
let asked = asks_to_play(&cmd);
let replaced_shuffled = player.mode.shuffle
&& matches!(&cmd, PlayerCommand::ReplacePlaylist { items, .. } if items.len() > 1);
let wanted_before = player.shared_state.wants_to_play();
player.process_command(cmd);
while let Ok(sent) = player.commands.rx.try_recv() {
player.process_command(sent);
}
player.update_playback_state();
history.push(label);
let broken = check_invariants(&player, wanted_before, asked)
.err()
.or_else(|| {
(replaced_shuffled
&& player
.shared_state
.shuffle_order()
.iter()
.all(|(_, pre)| pre.is_none()))
.then(|| "a queue replaced while shuffled plays in order".to_string())
});
if let Some(broken) = broken {
let tail = history[history.len().saturating_sub(8)..].join("\n ");
panic!("seed {seed}, step {step}: {broken}\nlast commands:\n {tail}");
}
}
player.process_command(PlayerCommand::Stop);
}
}
#[test]
fn a_pause_reports_where_the_fade_went_silent() {
use std::sync::atomic::AtomicBool;
struct FadingEngine {
running: Arc<AtomicBool>,
silent: Arc<AtomicBool>,
}
impl AudioEngineHandle for FadingEngine {
fn start(&self) -> Result<(), BackendError> {
self.running.store(true, Ordering::Relaxed);
Ok(())
}
fn stop(&self) -> Result<(), BackendError> {
self.running.store(false, Ordering::Relaxed);
Ok(())
}
fn is_running(&self) -> bool {
self.running.load(Ordering::Relaxed)
}
fn fade_out(&self) {}
fn fade_in(&self) -> Result<(), BackendError> {
Ok(())
}
fn is_silent(&self) -> bool {
self.silent.load(Ordering::Relaxed)
}
}
let running = Arc::new(AtomicBool::new(true));
let silent = Arc::new(AtomicBool::new(false));
let mut player = Player::new();
player.transport = Transport::Loaded(test_session(
QueueItemId::new(),
Box::new(FadingEngine {
running: running.clone(),
silent: silent.clone(),
}),
));
player
.shared_state
.set_playback_state(PlaybackState::Playing);
player.shared_state.set_position_ms(5_000);
let (reply, answer) = crossbeam_channel::bounded(1);
player.process_command(PlayerCommand::PauseAndReport(reply));
if crate::config::Config::cached().playback.fade_on_pause {
assert!(answer.try_recv().is_err(), "not while the fade is audible");
player.shared_state.set_position_ms(5_150);
player.update_playback_state();
assert!(answer.try_recv().is_err());
silent.store(true, Ordering::Relaxed);
player.update_playback_state();
assert!(!running.load(Ordering::Relaxed));
assert_eq!(answer.try_recv().unwrap(), 5_150);
} else {
assert_eq!(answer.try_recv().unwrap(), 5_000);
}
}
#[test]
fn the_server_hears_each_turn_playback_takes() {
use PlaybackReportState::{Paused, Playing, Stopped};
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.wav");
crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: 8_000.0,
asked: Default::default(),
starts: Default::default(),
});
let (recorder, events) = PlayRecorder::capture();
player.history = Some(recorder);
let item = PlaylistItem {
db_id: Some(5),
path,
..make_item("t")
};
let id = item.id;
player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
player.process_command(PlayerCommand::Play(id));
player.process_command(PlayerCommand::Pause);
player.process_command(PlayerCommand::Seek(4_000));
await_queued(&player);
player.process_command(PlayerCommand::Resume);
player.process_command(PlayerCommand::Stop);
let events: Vec<_> = events.try_iter().collect();
let near_seek = |position_ms: u64| (3_750..=4_000).contains(&position_ms);
let reports: Vec<_> = events
.iter()
.filter_map(|e| match e {
PlayEvent::Playback(r) => Some((r.track_id, r.state, r.position_ms)),
_ => None,
})
.collect();
assert_eq!(
events.first(),
Some(&PlayEvent::Started {
track_id: 5,
position_ms: 0
})
);
assert_eq!(reports.len(), 4, "{events:?}");
assert_eq!(reports[0], (5, Paused, 0));
for (report, state) in reports[1..].iter().zip([Paused, Playing, Stopped]) {
assert_eq!((report.0, report.1), (5, state), "{events:?}");
assert!(near_seek(report.2), "{events:?}");
}
assert_eq!(
events.last(),
Some(&PlayEvent::Finished {
track_id: 5,
listened_ms: 0
})
);
}
#[test]
fn removing_the_playing_track_resumes_at_its_successor() {
let mut player = Player::new();
let ids = seed(&mut player, 5);
pretend_playing(&mut player, ids[2]);
player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
assert_eq!(
player.shared_state.cursor(),
Some(ids[3]),
"playback must continue at the next track, not restart the queue"
);
assert_eq!(player.playback_starts, 1);
}
#[test]
fn removing_the_paused_track_moves_on_paused() {
let mut player = Player::new();
let ids = seed(&mut player, 3);
pretend_playing(&mut player, ids[1]);
if let Transport::Loaded(session) = &mut player.transport {
session.run = Run::Paused;
}
player.process_command(PlayerCommand::RemoveFromPlaylist(ids[1]));
assert_eq!(player.shared_state.cursor(), Some(ids[2]));
assert!(!player.shared_state.wants_to_play(), "still paused");
}
#[test]
fn removing_the_cursor_with_nothing_loaded_starts_nothing() {
let mut player = Player::new();
let ids = seed(&mut player, 3);
player.shared_state.set_cursor(Some(ids[1]));
player.process_command(PlayerCommand::RemoveFromPlaylist(ids[1]));
assert_eq!(player.shared_state.cursor(), Some(ids[2]));
assert_eq!(player.playback_starts, 0);
}
#[test]
fn removing_the_first_playing_track_resumes_at_the_new_first() {
let mut player = Player::new();
let ids = seed(&mut player, 3);
pretend_playing(&mut player, ids[0]);
player.process_command(PlayerCommand::RemoveFromPlaylist(ids[0]));
assert_eq!(player.shared_state.cursor(), Some(ids[1]));
}
#[test]
fn next_track_parks_on_a_track_that_has_not_downloaded_yet() {
let mut player = Player::new();
let playing = make_item("playing");
let waiting = pending_item("waiting");
let later = make_item("later");
let (playing_id, waiting_id) = (playing.id, waiting.id);
player.process_command(PlayerCommand::AddToPlaylist(vec![playing, waiting, later]));
pretend_playing(&mut player, playing_id);
player.process_command(PlayerCommand::DecodeFinished(player.session));
assert_eq!(
player.shared_state.cursor(),
Some(waiting_id),
"the cursor parks on the track being fetched"
);
assert_eq!(
player.playback_starts, 0,
"nothing to play until its bytes land"
);
player
.shared_state
.update_item_state(waiting_id, ItemState::Ready);
player.process_command(PlayerCommand::TrackReady(waiting_id));
assert_eq!(player.playback_starts, 1);
assert_eq!(player.shared_state.cursor(), Some(waiting_id));
}
#[test]
fn a_download_that_cannot_land_moves_the_cursor_on() {
let mut player = Player::new();
let waiting = pending_item("waiting");
let later = make_item("later");
let (waiting_id, later_id) = (waiting.id, later.id);
player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, later]));
player.process_command(PlayerCommand::Play(waiting_id));
assert_eq!(player.playback_starts, 0, "nothing to play yet");
player
.shared_state
.update_item_state(waiting_id, ItemState::Failed("remote unavailable".into()));
player.process_command(PlayerCommand::TrackFailed(waiting_id));
assert_eq!(
player.shared_state.cursor(),
Some(later_id),
"the queue moves past a track that can never load"
);
assert_eq!(player.playback_starts, 1);
}
#[test]
fn a_queue_that_can_never_load_stops_rather_than_waiting() {
let mut player = Player::new();
let first = pending_item("first");
let second = pending_item("second");
let (first_id, second_id) = (first.id, second.id);
player.process_command(PlayerCommand::AddToPlaylist(vec![first, second]));
player.process_command(PlayerCommand::Play(first_id));
for id in [first_id, second_id] {
player
.shared_state
.update_item_state(id, ItemState::Failed("remote unavailable".into()));
player.process_command(PlayerCommand::TrackFailed(id));
}
assert_eq!(player.playback_starts, 0);
assert_eq!(
player.shared_state.playback_state(),
PlaybackState::Stopped,
"a stop the UI can see, not an indefinite wait for TrackReady"
);
}
#[test]
fn a_failure_elsewhere_in_the_queue_leaves_the_cursor_alone() {
let mut player = Player::new();
let waiting = pending_item("waiting");
let other = pending_item("other");
let (waiting_id, other_id) = (waiting.id, other.id);
player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, other]));
player.process_command(PlayerCommand::Play(waiting_id));
player
.shared_state
.update_item_state(other_id, ItemState::Failed("remote unavailable".into()));
player.process_command(PlayerCommand::TrackFailed(other_id));
assert_eq!(
player.shared_state.cursor(),
Some(waiting_id),
"a track still downloading keeps the cursor"
);
}
#[test]
fn batch_delete_containing_the_cursor_restarts_the_engine_once() {
let mut player = Player::new();
let ids = seed(&mut player, 5);
pretend_playing(&mut player, ids[2]);
player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![
ids[1], ids[2], ids[3],
]));
assert_eq!(playlist_titles(&player), vec!["t0", "t4"]);
assert_eq!(player.shared_state.cursor(), Some(ids[4]));
assert_eq!(
player.playback_starts, 1,
"one resume for the whole selection, not one per deleted track"
);
}
#[test]
fn batch_delete_below_the_cursor_leaves_playback_alone() {
let mut player = Player::new();
let ids = seed(&mut player, 4);
player.shared_state.set_cursor(Some(ids[0]));
player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![ids[2], ids[3]]));
assert_eq!(player.shared_state.cursor(), Some(ids[0]));
assert_eq!(player.playback_starts, 0);
}
#[test]
fn undo_of_a_batch_delete_restores_the_original_order() {
let mut player = Player::new();
let items = vec![
make_item("A"),
make_item("B"),
make_item("C"),
make_item("D"),
];
let (b_id, c_id) = (items[1].id, items[2].id);
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![c_id, b_id]));
assert_eq!(playlist_titles(&player), vec!["A", "D"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
}
#[test]
fn undo_add_removes_items() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B")];
let ids: Vec<_> = items.iter().map(|i| i.id).collect();
player.process_command(PlayerCommand::AddToPlaylist(items));
assert_eq!(playlist_ids(&player), ids);
assert!(player.undo_stack().can_undo());
player.process_command(PlayerCommand::Undo);
assert!(playlist_ids(&player).is_empty());
assert!(player.undo_stack().can_redo());
}
#[test]
fn redo_add_restores_items() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B")];
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::Undo);
assert!(playlist_ids(&player).is_empty());
player.process_command(PlayerCommand::Redo);
assert_eq!(playlist_titles(&player), vec!["A", "B"]);
}
#[test]
fn undo_remove_restores_item_at_position() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B"), make_item("C")];
let b_id = items[1].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
assert_eq!(playlist_titles(&player), vec!["A", "C"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
}
#[test]
fn undo_remove_first_item() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B")];
let a_id = items[0].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::RemoveFromPlaylist(a_id));
assert_eq!(playlist_titles(&player), vec!["B"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B"]);
}
#[test]
fn undo_batch_remove_restores_all() {
let mut player = Player::new();
let items = vec![
make_item("A"),
make_item("B"),
make_item("C"),
make_item("D"),
];
let b_id = items[1].id;
let c_id = items[2].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
let version_before = player.shared_state.playlist_version();
player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![b_id, c_id]));
assert_eq!(playlist_titles(&player), vec!["A", "D"]);
assert_eq!(
player.shared_state.playlist_version(),
version_before + 1,
"batch removal must bump the playlist version exactly once"
);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
}
#[test]
fn redo_batch_remove() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B"), make_item("C")];
let a_id = items[0].id;
let b_id = items[1].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![a_id, b_id]));
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
player.process_command(PlayerCommand::Redo);
assert_eq!(playlist_titles(&player), vec!["C"]);
}
#[test]
fn redo_remove() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B"), make_item("C")];
let b_id = items[1].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
player.process_command(PlayerCommand::Redo);
assert_eq!(playlist_titles(&player), vec!["A", "C"]);
}
#[test]
fn undo_insert_removes_inserted_items() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("C")];
let a_id = items[0].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
let inserted = vec![make_item("B")];
player.process_command(PlayerCommand::InsertInPlaylist {
items: inserted,
after: a_id,
});
assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "C"]);
}
#[test]
fn undo_move_restores_position() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B"), make_item("C")];
let a_id = items[0].id;
let c_id = items[2].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::MoveInPlaylist {
id: a_id,
target: c_id,
after: true,
});
assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
}
#[test]
fn redo_move() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B"), make_item("C")];
let a_id = items[0].id;
let c_id = items[2].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::MoveInPlaylist {
id: a_id,
target: c_id,
after: true,
});
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
player.process_command(PlayerCommand::Redo);
assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
}
#[test]
fn undo_batch_move() {
let mut player = Player::new();
let items = vec![
make_item("A"),
make_item("B"),
make_item("C"),
make_item("D"),
];
let a_id = items[0].id;
let b_id = items[1].id;
let d_id = items[3].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::MoveItemsInPlaylist {
ids: vec![a_id, b_id],
target: d_id,
after: true,
});
assert_eq!(playlist_titles(&player), vec!["C", "D", "A", "B"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
}
#[test]
fn undo_clear_restores_playlist() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B"), make_item("C")];
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::ClearPlaylist);
assert!(playlist_ids(&player).is_empty());
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
}
#[test]
fn undoing_a_replace_does_not_leave_the_engine_on_an_orphaned_track() {
let mut player = Player::new();
let original = seed(&mut player, 3);
player.shared_state.set_cursor(Some(original[0]));
pretend_playing(&mut player, original[0]);
let replacement = vec![make_item("something else")];
let orphan = replacement[0].id;
player.process_command(PlayerCommand::ReplacePlaylist {
items: replacement,
start: 0,
position_ms: 0,
play: true,
});
pretend_playing(&mut player, orphan);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_ids(&player), original, "the queue comes back");
assert!(
player.shared_state.get_item(orphan).is_none(),
"and the replacement is gone from it"
);
assert!(
playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
"so nothing may still be playing out of it"
);
}
#[test]
fn undoing_an_add_does_not_leave_the_engine_on_a_removed_track() {
let mut player = Player::new();
seed(&mut player, 2);
let added = seed(&mut player, 1);
pretend_playing(&mut player, added[0]);
player.process_command(PlayerCommand::Undo);
assert!(player.shared_state.get_item(added[0]).is_none());
assert!(
playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
"the engine cannot be left on the item the undo removed"
);
}
#[test]
fn undoing_a_move_leaves_playback_alone() {
let mut player = Player::new();
let ids = seed(&mut player, 3);
player.shared_state.set_cursor(Some(ids[0]));
pretend_playing(&mut player, ids[0]);
let starts = player.playback_starts;
player.process_command(PlayerCommand::MoveInPlaylist {
id: ids[2],
target: ids[0],
after: false,
});
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_ids(&player), ids);
assert_eq!(playing_id(&player), Some(ids[0]), "still on the same track");
assert_eq!(player.playback_starts, starts, "and not restarted");
}
#[test]
fn redo_clear() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B")];
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::ClearPlaylist);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B"]);
player.process_command(PlayerCommand::Redo);
assert!(playlist_ids(&player).is_empty());
}
#[test]
fn multiple_undos_in_sequence() {
let mut player = Player::new();
player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("A")]));
player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("C")]));
assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A"]);
player.process_command(PlayerCommand::Undo);
assert!(playlist_ids(&player).is_empty());
}
#[test]
fn undo_redo_undo_cycle() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B")];
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::Undo);
assert!(playlist_ids(&player).is_empty());
player.process_command(PlayerCommand::Redo);
assert_eq!(playlist_titles(&player), vec!["A", "B"]);
player.process_command(PlayerCommand::Undo);
assert!(playlist_ids(&player).is_empty());
}
#[test]
fn new_action_clears_redo_stack() {
let mut player = Player::new();
let items = vec![make_item("A")];
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::Undo);
assert!(player.undo_stack().can_redo());
player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
assert!(!player.undo_stack().can_redo());
}
#[test]
fn undo_on_empty_stack_is_noop() {
let mut player = Player::new();
player.process_command(PlayerCommand::Undo);
assert!(playlist_ids(&player).is_empty());
}
#[test]
fn redo_on_empty_stack_is_noop() {
let mut player = Player::new();
player.process_command(PlayerCommand::Redo);
assert!(playlist_ids(&player).is_empty());
}
#[test]
fn playback_commands_not_undoable() {
let mut player = Player::new();
player.process_command(PlayerCommand::Pause);
player.process_command(PlayerCommand::Resume);
player.process_command(PlayerCommand::NextTrack);
player.process_command(PlayerCommand::PrevTrack);
assert!(!player.undo_stack().can_undo());
}
#[test]
fn update_paths_not_undoable() {
let mut player = Player::new();
let items = vec![make_item("A")];
let id = items[0].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
let undo_count = player.undo_stack().undo_len();
player.process_command(PlayerCommand::UpdatePaths(vec![(
id,
PathBuf::from("/new/path.flac"),
)]));
assert_eq!(player.undo_stack().undo_len(), undo_count);
}
#[test]
fn add_remove_undo_undo_produces_original() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B"), make_item("C")];
let b_id = items[1].id;
let original_titles = vec!["A", "B", "C"];
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
assert_eq!(playlist_titles(&player), vec!["A", "C"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), original_titles);
player.process_command(PlayerCommand::Undo);
assert!(playlist_ids(&player).is_empty());
}
#[test]
fn interleaved_adds_and_moves_undo() {
let mut player = Player::new();
let items = vec![make_item("A"), make_item("B"), make_item("C")];
let a_id = items[0].id;
let c_id = items[2].id;
player.process_command(PlayerCommand::AddToPlaylist(items));
player.process_command(PlayerCommand::MoveInPlaylist {
id: a_id,
target: c_id,
after: true,
});
assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("D")]));
assert_eq!(playlist_titles(&player), vec!["B", "C", "A", "D"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
player.process_command(PlayerCommand::Undo);
assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
}
#[test]
fn stop_engine_drops_engine_synchronously() {
use std::sync::atomic::{AtomicBool, Ordering};
struct MockEngine {
dropped: Arc<AtomicBool>,
}
impl AudioEngineHandle for MockEngine {
fn start(&self) -> Result<(), BackendError> {
Ok(())
}
fn stop(&self) -> Result<(), BackendError> {
Ok(())
}
fn is_running(&self) -> bool {
false
}
fn fade_out(&self) {}
fn fade_in(&self) -> Result<(), BackendError> {
Ok(())
}
fn is_silent(&self) -> bool {
false
}
}
impl Drop for MockEngine {
fn drop(&mut self) {
self.dropped.store(true, Ordering::SeqCst);
}
}
let dropped = Arc::new(AtomicBool::new(false));
let mut player = Player::new();
player.transport = Transport::Loaded(test_session(
QueueItemId::new(),
Box::new(MockEngine {
dropped: dropped.clone(),
}),
));
player.stop_engine();
assert!(
dropped.load(Ordering::SeqCst),
"AudioEngine must be dropped synchronously in stop_engine (GitHub #89)"
);
}
#[test]
fn abandoning_a_stream_wakes_a_reader_parked_at_the_write_head() {
let live = LiveStream {
feed: crate::remote::downloads::ByteFeed::new(),
abandoned: Default::default(),
};
let feed = live.feed.clone();
let started = std::time::Instant::now();
let reader = thread::spawn(move || {
feed.wait_past(
0,
std::time::Instant::now() + std::time::Duration::from_secs(30),
)
});
thread::sleep(std::time::Duration::from_millis(50));
live.abandon();
reader.join().unwrap();
assert!(live.abandoned.load(std::sync::atomic::Ordering::Acquire));
assert!(started.elapsed() < std::time::Duration::from_secs(5));
}
pub(super) struct StuckBackend {
pub(super) rate: f64,
pub(super) asked: Arc<std::sync::Mutex<Option<(f64, u32)>>>,
pub(super) starts: Arc<std::sync::atomic::AtomicUsize>,
}
struct NullEngine {
starts: Arc<std::sync::atomic::AtomicUsize>,
running: std::sync::atomic::AtomicBool,
lead_in: Arc<AtomicU64>,
}
impl AudioEngineHandle for NullEngine {
fn start(&self) -> Result<(), BackendError> {
self.starts.fetch_add(1, Ordering::Relaxed);
self.running.store(true, Ordering::Relaxed);
Ok(())
}
fn stop(&self) -> Result<(), BackendError> {
self.running.store(false, Ordering::Relaxed);
Ok(())
}
fn is_running(&self) -> bool {
self.running.load(Ordering::Relaxed)
}
fn fade_out(&self) {}
fn fade_in(&self) -> Result<(), BackendError> {
Ok(())
}
fn is_silent(&self) -> bool {
false
}
fn lead_in(&self, frames: u64) {
self.lead_in.store(frames, Ordering::Relaxed);
}
}
impl AudioBackend for StuckBackend {
fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
Ok(vec![self.default_device()?])
}
fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
Ok(backend::DeviceInfo {
name: "Stuck DAC".into(),
sample_rates: vec![self.rate],
platform_id: 0,
kind: Default::default(),
})
}
fn supported_sample_rates(
&self,
_device: &backend::DeviceInfo,
) -> Result<Vec<f64>, BackendError> {
Ok(vec![self.rate])
}
fn get_device_sample_rate(
&self,
_device: &backend::DeviceInfo,
) -> Result<f64, BackendError> {
Ok(self.rate)
}
fn set_device_sample_rate(
&self,
_device: &backend::DeviceInfo,
rate: f64,
) -> Result<f64, BackendError> {
Err(BackendError::UnsupportedSampleRate(rate))
}
fn create_engine(
&self,
_device: &backend::DeviceInfo,
sample_rate: f64,
channels: u32,
_consumer: rtrb::Consumer<f32>,
_samples_played: Arc<AtomicU64>,
) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
*self.asked.lock().unwrap() = Some((sample_rate, channels));
Ok(Box::new(NullEngine {
starts: self.starts.clone(),
running: Default::default(),
lead_in: Default::default(),
}))
}
}
struct CaptureBackend {
consumers: Arc<std::sync::Mutex<Vec<rtrb::Consumer<f32>>>>,
}
impl AudioBackend for CaptureBackend {
fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
Ok(vec![self.default_device()?])
}
fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
Ok(backend::DeviceInfo {
name: "Capture DAC".into(),
sample_rates: vec![44100.0],
platform_id: 0,
kind: Default::default(),
})
}
fn supported_sample_rates(
&self,
_device: &backend::DeviceInfo,
) -> Result<Vec<f64>, BackendError> {
Ok(vec![44100.0])
}
fn get_device_sample_rate(
&self,
_device: &backend::DeviceInfo,
) -> Result<f64, BackendError> {
Ok(44100.0)
}
fn set_device_sample_rate(
&self,
_device: &backend::DeviceInfo,
rate: f64,
) -> Result<f64, BackendError> {
Ok(rate)
}
fn create_engine(
&self,
_device: &backend::DeviceInfo,
_sample_rate: f64,
_channels: u32,
consumer: rtrb::Consumer<f32>,
_samples_played: Arc<AtomicU64>,
) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
self.consumers.lock().unwrap().push(consumer);
Ok(Box::new(NullEngine {
starts: Default::default(),
running: Default::default(),
lead_in: Default::default(),
}))
}
}
fn peak_reaching_the_device(
consumers: &std::sync::Mutex<Vec<rtrb::Consumer<f32>>>,
want: usize,
) -> f32 {
let mut consumer = consumers.lock().unwrap().pop().expect("an engine was made");
let mut peak = 0.0f32;
let mut got = 0;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
while got < want && std::time::Instant::now() < deadline {
let n = consumer.slots();
if n == 0 {
std::thread::sleep(std::time::Duration::from_millis(2));
continue;
}
let chunk = consumer.read_chunk(n).unwrap();
let (a, b) = chunk.as_slices();
peak = a.iter().chain(b).fold(peak, |p, s| p.max(s.abs()));
got += n;
chunk.commit_all();
}
assert!(got >= want, "only {got} of {want} samples arrived");
peak
}
#[test]
fn a_profile_processes_files_and_streams_alike() {
let dir = tempfile::tempdir().unwrap();
let tone = dir.path().join("tone.wav");
crate::test_utils::generate_wav_tone(&tone, 44100, 440.0, 0.5);
let want = 22050;
let consumers = Arc::new(std::sync::Mutex::new(Vec::new()));
let mut player = Player::new();
player.backend = Box::new(CaptureBackend {
consumers: consumers.clone(),
});
let mut item = make_item("tone");
item.path = tone.clone();
let id = item.id;
player.shared_state.add_items(vec![item]);
player.shared_state.set_cursor(Some(id));
let stream = || {
let feed = crate::remote::downloads::ByteFeed::new();
feed.set(std::fs::metadata(&tone).unwrap().len());
Source::Stream(StreamSource {
path: tone.clone(),
bytes_written: feed,
total: 0,
mode: streaming::ProbeMode::Full,
})
};
let peaks = |player: &mut Player| {
let mut out = Vec::new();
for source in [Source::File(tone.clone()), stream()] {
player
.try_open_session(id, source, None, 0, Run::Playing)
.unwrap();
player.publish();
out.push((
peak_reaching_the_device(&consumers, want),
player.shared_state.dsp().map(|d| d.profile),
));
player.stop_engine();
}
out
};
let untouched = peaks(&mut player);
player.dsp_override = Some(Arc::new(
crate::audio::dsp::Setup::new(vec![], vec![]).with_preamp(-6.0206),
));
let processed = peaks(&mut player);
for (kind, ((before, none), (after, half))) in ["file", "stream"]
.iter()
.zip(untouched.iter().zip(&processed))
{
assert_eq!(
none, &None,
"{kind}: nothing is published without a profile"
);
assert_eq!(half.as_deref(), Some("test"), "{kind}: the badge names it");
assert!(*before > 0.1, "{kind}: the tone reached the device");
assert!(
(after / before - 0.5).abs() < 0.01,
"{kind}: {before} → {after}, not halved"
);
}
}
fn engine_format_for(source_rate: u32, channels: u16, device_rate: f64) -> (f64, u32) {
let asked = Arc::new(std::sync::Mutex::new(None));
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: device_rate,
asked: asked.clone(),
starts: Default::default(),
});
let info = buffer::StreamInfo {
codec: "MP3".into(),
sample_rate: source_rate,
channels,
bit_depth: Some(16),
bitrate_kbps: None,
duration_ms: 1000,
};
let (_producer, consumer) = rtrb::RingBuffer::new(16);
player
.create_engine_for(&info, consumer)
.expect("engine creation should succeed");
let asked = *asked.lock().unwrap();
asked.expect("engine was never created")
}
#[test]
fn engine_uses_source_rate_when_device_refuses_switch() {
assert_eq!(engine_format_for(22050, 2, 48000.0), (22050.0, 2));
assert_eq!(engine_format_for(32000, 2, 44100.0), (32000.0, 2));
}
#[test]
fn engine_uses_source_channel_count() {
assert_eq!(engine_format_for(44100, 1, 44100.0), (44100.0, 1));
}
fn settled_rate_for(source_rate: u32, device_rate: f64) -> Option<u32> {
let mut player = Player::new();
player.backend = Box::new(StuckBackend {
rate: device_rate,
asked: Arc::new(std::sync::Mutex::new(None)),
starts: Default::default(),
});
let state = player.shared_state.clone();
let info = buffer::StreamInfo {
codec: "MP3".into(),
sample_rate: source_rate,
channels: 2,
bit_depth: Some(16),
bitrate_kbps: None,
duration_ms: 1000,
};
let (_producer, consumer) = rtrb::RingBuffer::new(16);
player
.create_engine_for(&info, consumer)
.expect("engine creation should succeed");
state.output_sample_rate()
}
#[test]
fn settled_device_rate_reaches_the_shared_state() {
assert_eq!(settled_rate_for(22050, 48000.0), Some(48000));
assert_eq!(settled_rate_for(44100, 44100.0), Some(44100));
}
struct SlowBackend {
observed: Arc<std::sync::Mutex<Vec<Option<u32>>>>,
state: Arc<SharedPlayerState>,
lead_in: Arc<AtomicU64>,
}
impl AudioBackend for SlowBackend {
fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
Ok(vec![self.default_device()?])
}
fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
Ok(backend::DeviceInfo {
name: "Slow DAC".into(),
sample_rates: vec![44100.0, 48000.0],
platform_id: 0,
kind: Default::default(),
})
}
fn supported_sample_rates(
&self,
_device: &backend::DeviceInfo,
) -> Result<Vec<f64>, BackendError> {
Ok(vec![44100.0, 48000.0])
}
fn get_device_sample_rate(
&self,
_device: &backend::DeviceInfo,
) -> Result<f64, BackendError> {
Ok(48000.0)
}
fn set_device_sample_rate(
&self,
_device: &backend::DeviceInfo,
rate: f64,
) -> Result<f64, BackendError> {
self.observed
.lock()
.unwrap()
.push(self.state.output_sample_rate());
Ok(rate)
}
fn create_engine(
&self,
_device: &backend::DeviceInfo,
_sample_rate: f64,
_channels: u32,
_consumer: rtrb::Consumer<f32>,
_samples_played: Arc<AtomicU64>,
) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
Ok(Box::new(NullEngine {
starts: Default::default(),
running: Default::default(),
lead_in: self.lead_in.clone(),
}))
}
}
#[test]
fn the_previous_rate_is_not_published_while_the_device_reclocks() {
let mut player = Player::new();
let state = player.shared_state.clone();
state.set_output_sample_rate(48000);
let observed = Arc::new(std::sync::Mutex::new(Vec::new()));
player.backend = Box::new(SlowBackend {
observed: observed.clone(),
state: state.clone(),
lead_in: Default::default(),
});
let info = buffer::StreamInfo {
codec: "FLAC".into(),
sample_rate: 44100,
channels: 2,
bit_depth: Some(16),
bitrate_kbps: None,
duration_ms: 1000,
};
let (_producer, consumer) = rtrb::RingBuffer::new(16);
player
.create_engine_for(&info, consumer)
.expect("engine creation should succeed");
assert_eq!(
*observed.lock().unwrap(),
vec![None],
"mid-switch the output rate must read as unknown, not as the last track's"
);
assert_eq!(state.output_sample_rate(), Some(44100));
}
fn lead_in_for(source_rate: u32) -> u64 {
let mut player = Player::new();
let lead_in = Arc::new(AtomicU64::new(0));
player.backend = Box::new(SlowBackend {
observed: Default::default(),
state: player.shared_state.clone(),
lead_in: lead_in.clone(),
});
let info = buffer::StreamInfo {
codec: "FLAC".into(),
sample_rate: source_rate,
channels: 2,
bit_depth: Some(16),
bitrate_kbps: None,
duration_ms: 1000,
};
let (_producer, consumer) = rtrb::RingBuffer::new(16);
player
.create_engine_for(&info, consumer)
.expect("engine creation should succeed");
lead_in.load(Ordering::Relaxed)
}
#[test]
fn an_engine_made_inside_the_silence_keeps_the_rest_of_it() {
let mut player = Player::new();
let lead_in = Arc::new(AtomicU64::new(0));
player.backend = Box::new(SlowBackend {
observed: Default::default(),
state: player.shared_state.clone(),
lead_in: lead_in.clone(),
});
let info = |sample_rate| buffer::StreamInfo {
codec: "FLAC".into(),
sample_rate,
channels: 2,
bit_depth: Some(16),
bitrate_kbps: None,
duration_ms: 1000,
};
let (_p, consumer) = rtrb::RingBuffer::new(16);
player.create_engine_for(&info(44100), consumer).unwrap();
let first = lead_in.swap(0, Ordering::Relaxed);
assert!(first > 0);
let (_p, consumer) = rtrb::RingBuffer::new(16);
player.create_engine_for(&info(48000), consumer).unwrap();
let carried = lead_in.load(Ordering::Relaxed);
assert!(
carried > 0 && carried < first * 48000 / 44100,
"carried {carried} of {first}"
);
}
#[test]
fn silence_leads_in_only_after_the_device_changed_rate() {
assert!(
lead_in_for(44100) > 0,
"the device is relocking, so the start of the track would be lost"
);
assert_eq!(lead_in_for(48000), 0, "no switch, nothing to wait for");
}
struct WatchedBackend {
inner: StuckBackend,
#[allow(clippy::type_complexity)]
captured: Arc<std::sync::Mutex<Option<Box<dyn Fn(f64) + Send + Sync>>>>,
}
struct NullWatch;
impl backend::SampleRateWatch for NullWatch {}
impl AudioBackend for WatchedBackend {
fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
self.inner.list_devices()
}
fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
self.inner.default_device()
}
fn supported_sample_rates(
&self,
device: &backend::DeviceInfo,
) -> Result<Vec<f64>, BackendError> {
self.inner.supported_sample_rates(device)
}
fn get_device_sample_rate(
&self,
device: &backend::DeviceInfo,
) -> Result<f64, BackendError> {
self.inner.get_device_sample_rate(device)
}
fn set_device_sample_rate(
&self,
device: &backend::DeviceInfo,
rate: f64,
) -> Result<f64, BackendError> {
self.inner.set_device_sample_rate(device, rate)
}
fn watch_device_sample_rate(
&self,
_device: &backend::DeviceInfo,
on_change: Box<dyn Fn(f64) + Send + Sync>,
) -> Option<Box<dyn backend::SampleRateWatch>> {
*self.captured.lock().unwrap() = Some(on_change);
Some(Box::new(NullWatch))
}
fn create_engine(
&self,
device: &backend::DeviceInfo,
sample_rate: f64,
channels: u32,
consumer: rtrb::Consumer<f32>,
samples_played: Arc<AtomicU64>,
) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
self.inner
.create_engine(device, sample_rate, channels, consumer, samples_played)
}
}
#[test]
fn external_rate_change_reaches_the_shared_state() {
let captured = Arc::new(std::sync::Mutex::new(None));
let mut player = Player::new();
player.backend = Box::new(WatchedBackend {
inner: StuckBackend {
rate: 44100.0,
asked: Arc::new(std::sync::Mutex::new(None)),
starts: Default::default(),
},
captured: captured.clone(),
});
let state = player.shared_state.clone();
let info = buffer::StreamInfo {
codec: "FLAC".into(),
sample_rate: 44100,
channels: 2,
bit_depth: Some(16),
bitrate_kbps: None,
duration_ms: 1000,
};
let (_producer, consumer) = rtrb::RingBuffer::new(16);
player
.create_engine_for(&info, consumer)
.expect("engine creation should succeed");
assert_eq!(state.output_sample_rate(), Some(44100));
let on_change = captured.lock().unwrap().take().expect("watch registered");
on_change(48000.0);
assert_eq!(state.output_sample_rate(), Some(48000));
}
}