use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use crossbeam_channel::{Receiver, SendTimeoutError, Sender, bounded};
use super::state::{PlayMode, PlaylistItem, QueueItemId, Repeat, SleepTimer};
#[derive(Debug)]
pub enum PlayerCommand {
Play(QueueItemId),
Cue {
id: QueueItemId,
position_ms: u64,
play: bool,
},
Pause,
PauseAndReport(Sender<u64>),
Barrier(Sender<()>),
Resume,
Stop,
Seek(u64), NextTrack,
PrevTrack,
AddToPlaylist(Vec<PlaylistItem>),
RemoveFromPlaylist(QueueItemId),
RemoveFromPlaylistBatch(Vec<QueueItemId>),
MoveInPlaylist {
id: QueueItemId,
target: QueueItemId,
after: bool,
},
MoveItemsInPlaylist {
ids: Vec<QueueItemId>,
target: QueueItemId,
after: bool,
},
ReorderPlaylist(Vec<QueueItemId>),
UpdatePaths(Vec<(QueueItemId, PathBuf)>),
InsertInPlaylist {
items: Vec<PlaylistItem>,
after: QueueItemId,
},
ClearPlaylist,
ReplacePlaylist {
items: Vec<PlaylistItem>,
start: usize,
position_ms: u64,
play: bool,
},
TrackReady(QueueItemId),
TrackStreamReady(QueueItemId),
StreamProbed {
id: QueueItemId,
info: Box<crate::audio::buffer::StreamInfo>,
mode: crate::audio::streaming::ProbeMode,
},
TrackFailed(QueueItemId),
CacheTracks(Vec<i64>),
DecodeFinished(u64),
TrackQueued,
Undo,
Redo,
BeginUndoBatch,
EndUndoBatch,
SetOutputDevice(String),
ClearOutputDevice,
ReloadDsp,
RestartOutput,
UseRenderer(Option<Box<crate::upnp::Connection>>),
ResumeRenderer(Box<crate::upnp::Connection>),
ResumeRendererMissed,
ReleaseRenderer(crossbeam_channel::Sender<()>),
SetRendererVolume(u8),
SetShuffle(bool),
SetRepeat(Repeat),
SetSleepTimer(Option<SleepTimer>),
RestorePlayMode(PlayMode),
Renderer {
session: u64,
event: crate::upnp::session::Event,
},
}
pub struct CommandChannel {
pub tx: Sender<PlayerCommand>,
pub rx: Receiver<PlayerCommand>,
}
impl Default for CommandChannel {
fn default() -> Self {
Self::new()
}
}
impl CommandChannel {
pub fn new() -> Self {
let (tx, rx) = bounded(16);
Self { tx, rx }
}
}
pub fn send_unless_stopped(tx: &Sender<PlayerCommand>, command: PlayerCommand, stop: &AtomicBool) {
let mut command = command;
while !stop.load(Ordering::Relaxed) {
match tx.send_timeout(command, Duration::from_millis(10)) {
Err(SendTimeoutError::Timeout(unsent)) => command = unsent,
Ok(()) | Err(SendTimeoutError::Disconnected(_)) => return,
}
}
}
impl PlayerCommand {
pub fn asks_to_play(&self) -> bool {
matches!(
self,
Self::Play(_)
| Self::Cue { play: true, .. }
| Self::Resume
| Self::NextTrack
| Self::PrevTrack
| Self::ReplacePlaylist { play: true, .. }
)
}
}
pub fn release_renderer(
player: &crossbeam_channel::Sender<PlayerCommand>,
timeout: std::time::Duration,
) {
let (tx, rx) = crossbeam_channel::bounded(1);
if player.send(PlayerCommand::ReleaseRenderer(tx)).is_ok() {
let _ = rx.recv_timeout(timeout);
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use std::time::{Duration, Instant};
use super::*;
use crate::audio::buffer::{self, PlaybackTimeline, Processing, SourceEntry};
#[test]
fn a_full_channel_never_holds_up_stopping_the_decoder() {
let channel = CommandChannel::new();
while channel.tx.try_send(PlayerCommand::Stop).is_ok() {}
let tx = channel.tx.clone();
let (producer, _consumer) = rtrb::RingBuffer::new(1024);
let mut decode = buffer::start_decode(
SourceEntry::from_file(QueueItemId::new(), "/nonexistent/koan.flac".into()),
producer,
0,
|| None,
PlaybackTimeline::new(),
None,
Processing::default(),
move |stop| send_unless_stopped(&tx, PlayerCommand::DecodeFinished(1), stop),
)
.unwrap();
std::thread::sleep(Duration::from_millis(50));
let (done_tx, done_rx) = crossbeam_channel::bounded(1);
std::thread::spawn(move || {
decode.stop();
done_tx.send(()).ok();
});
let started = Instant::now();
done_rx
.recv_timeout(Duration::from_secs(5))
.expect("stopping the decoder waited on the full channel");
assert!(started.elapsed() < Duration::from_secs(1));
assert_eq!(channel.rx.len(), 16);
}
#[test]
fn a_natural_end_waits_for_room() {
let channel = CommandChannel::new();
while channel.tx.try_send(PlayerCommand::Stop).is_ok() {}
let stop = Arc::new(AtomicBool::new(false));
let tx = channel.tx.clone();
let flag = stop.clone();
let sender = std::thread::spawn(move || {
send_unless_stopped(&tx, PlayerCommand::DecodeFinished(7), &flag)
});
std::thread::sleep(Duration::from_millis(50));
for _ in 0..16 {
channel.rx.recv().unwrap();
}
sender.join().unwrap();
assert!(matches!(
channel.rx.try_recv(),
Ok(PlayerCommand::DecodeFinished(7))
));
}
}