crabidy/crabidy-server/src/playback.rs

921 lines
35 KiB
Rust

use crate::queue_store::{self, QueueSnapshot, QueueStore, SaveQueueError};
use crate::{PendingResolve, QueueManager, ResolveKind};
use crate::{PlaybackCommand, PlaybackMessage, ProviderCommand, ProviderMessage};
use audio_player::Player;
use crabidy_core::proto::crabidy::QueueModifiers;
use crabidy_core::proto::crabidy::{
get_update_stream_response::Update as StreamUpdate, InitResponse, PlayState,
Queue as ProtoQueue, QueueTrack, Track, TrackPosition,
};
use crabidy_core::ProviderError;
use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use tracing::{debug, debug_span, error, info, instrument, trace, warn, Instrument};
pub struct Playback {
update_tx: tokio::sync::broadcast::Sender<StreamUpdate>,
provider_tx: flume::Sender<ProviderMessage>,
pub playback_tx: flume::Sender<PlaybackMessage>,
playback_rx: flume::Receiver<PlaybackMessage>,
queue: Mutex<QueueManager>,
state: Mutex<PlayState>,
/// In-flight resolve operations by op id. Non-empty means the broadcast
/// `Queue` snapshots carry `resolving = true`. Only the playback loop
/// touches this map (same single-writer discipline as `queue`).
pending: Mutex<HashMap<u64, PendingResolve>>,
next_op_id: AtomicU64,
/// `None` when queue persistence is disabled (no usable queues
/// directory) — the queue then lives in memory only.
store: Option<Arc<QueueStore>>,
/// Feeds the persister task; latest snapshot wins, so the loop never
/// waits on disk (architecture/queue-persistence.md D4).
persist_tx: tokio::sync::watch::Sender<Option<QueueSnapshot>>,
pub player: Player,
}
impl Playback {
pub fn new(
update_tx: tokio::sync::broadcast::Sender<StreamUpdate>,
provider_tx: flume::Sender<ProviderMessage>,
store: Option<Arc<QueueStore>>,
) -> Self {
let (playback_tx, playback_rx) = flume::bounded(64);
let queue = Mutex::new(QueueManager::new());
let state = Mutex::new(PlayState::Stopped);
let (persist_tx, _) = tokio::sync::watch::channel(None);
let player = Player::default();
Self {
update_tx,
provider_tx,
playback_tx,
playback_rx,
queue,
state,
pending: Mutex::new(HashMap::new()),
next_op_id: AtomicU64::new(0),
store,
persist_tx,
player,
}
}
/// Reloads the persisted current queue: tracks, position, and
/// shuffle/repeat. Never starts playback — a restarted server stays
/// silent. Call before [`Self::run`] so nothing observes the empty
/// queue first.
pub async fn restore_current(&self) {
let Some(store) = &self.store else {
return;
};
let Some(snapshot) = store.load_current().await else {
debug!("no persisted queue, starting fresh");
return;
};
let Ok(mut queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
// No autoplay: the track the replace would start is ignored.
let _ = queue.replace_with_tracks(&snapshot.tracks);
queue.repeat = snapshot.repeat;
// An out-of-range position (edited folder) is refused by
// `set_current_position` and playback starts at the first track.
let _ = queue.set_current_position(snapshot.current_position);
if snapshot.shuffle {
// The play order is not persisted; restoring shuffle reshuffles
// around the restored current track.
queue.shuffle_on();
}
info!(
tracks = snapshot.tracks.len(),
position = snapshot.current_position,
"restored the persisted queue"
);
}
pub fn run(self) {
if let Some(store) = &self.store {
queue_store::spawn_persister(Arc::clone(store), self.persist_tx.subscribe());
}
tokio::spawn(async move {
while let Ok(PlaybackMessage { span, command }) = self.playback_rx.recv_async().await {
// Attribute all handler events to a span that is a child of
// the span that was current when the command was sent.
let handler_span =
debug_span!(parent: &span, "playback_command", command = command.name());
self.handle_command(command).instrument(handler_span).await;
}
warn!("playback message channel closed, loop exiting");
});
}
async fn handle_command(&self, command: PlaybackCommand) {
trace!("handling playback command");
match command {
PlaybackCommand::Init { result_tx } => {
let response = {
let Ok(queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
let queue_track = QueueTrack {
queue_position: queue.current_position() as u32,
track: queue.current_track(),
};
let position = TrackPosition {
duration: 0,
position: 0,
};
let play_state = {
let Ok(play_state) = self.state.lock() else {
error!("play state lock poisoned");
return;
};
*play_state
};
InitResponse {
// Snapshot with `resolving`: a client connecting
// mid-resolve must show the indicator right away.
queue: Some(self.queue_snapshot(&queue)),
queue_track: Some(queue_track),
play_state: play_state as i32,
volume: 0.0,
mute: false,
position: Some(position),
mods: Some(QueueModifiers {
repeat: queue.repeat,
shuffle: queue.shuffle,
}),
}
};
trace!(?response, "sending init response");
if let Err(err) = result_tx.send(response) {
error!("failed to send init response: {err}");
}
}
PlaybackCommand::Replace { paths } => {
// A replace obsoletes whatever earlier ops are still
// resolving; their late chunks must not land in the new
// queue.
self.cancel_pending_resolves();
self.start_resolve(ResolveKind::Replace, paths);
}
PlaybackCommand::Queue { paths } => {
let position = {
let Ok(queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
queue.current_position() as u32
};
self.start_resolve(ResolveKind::InsertAfter(position), paths);
}
PlaybackCommand::Append { paths } => {
self.start_resolve(ResolveKind::Append, paths);
}
PlaybackCommand::ApplyResolvedChunk { op_id, tracks } => {
self.apply_resolved_chunk(op_id, tracks).await;
}
PlaybackCommand::ResolveFinished { op_id } => {
self.finish_resolve(op_id);
}
PlaybackCommand::Remove { positions } => {
debug!(?positions, "removing tracks");
let (track, was_last) = {
let Ok(mut queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
let was_last = queue.is_last_track();
let track = queue.remove_tracks(&positions);
self.broadcast_queue(&queue);
(track, was_last)
};
let state = {
let Ok(state) = self.state.lock() else {
error!("play state lock poisoned");
return;
};
*state
};
if state == PlayState::Playing && track.is_some() {
// The playing track was removed: play its successor, or
// stop when it was the last one.
if was_last {
self.stop_player().await;
} else {
self.play(track).await;
}
}
}
PlaybackCommand::Insert { position, paths } => {
self.start_resolve(ResolveKind::InsertAfter(position), paths);
}
PlaybackCommand::Clear { exclude_current } => {
debug!(exclude_current, "clearing queue");
// Chunks still resolving would repopulate the queue the
// user just emptied.
self.cancel_pending_resolves();
let should_stop = {
let Ok(mut queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
let should_stop = queue.clear(exclude_current);
self.broadcast_queue(&queue);
should_stop
};
if should_stop {
self.stop_player().await;
}
}
PlaybackCommand::SetCurrent { position } => {
debug!(position, "jumping to queue position");
let track = {
let Ok(mut queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
queue.set_current_position(position);
queue.current_track()
};
self.play(track).await;
}
PlaybackCommand::SaveQueue { name, result_tx } => {
debug!(name, "saving the queue");
// Snapshot on the loop (single-writer discipline), write on
// a spawned task — the loop never waits on disk.
let snapshot = {
let Ok(queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
let proto: ProtoQueue = queue.clone().into();
QueueSnapshot {
tracks: proto.tracks,
current_position: proto.current_position,
repeat: queue.repeat,
shuffle: queue.shuffle,
}
};
let store = self.store.clone();
tokio::spawn(
async move {
let result = match &store {
Some(store) => store.save(&name, &snapshot).await,
None => Err(SaveQueueError::Disabled),
};
if let Err(err) = &result {
warn!(name, "cannot save queue: {err}");
}
let _ = result_tx.send_async(result).await;
}
.in_current_span(),
);
}
PlaybackCommand::ToggleShuffle => {
let (shuffle, repeat) = {
let Ok(mut queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
if queue.shuffle {
queue.shuffle_off()
} else {
queue.shuffle_on()
}
self.send_persist_snapshot(&queue);
(queue.shuffle, queue.repeat)
};
debug!(shuffle, "toggled shuffle");
self.broadcast(StreamUpdate::Mods(QueueModifiers { shuffle, repeat }));
}
PlaybackCommand::ToggleRepeat => {
let (shuffle, repeat) = {
let Ok(mut queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
queue.repeat = !queue.repeat;
self.send_persist_snapshot(&queue);
(queue.shuffle, queue.repeat)
};
debug!(repeat, "toggled repeat");
self.broadcast(StreamUpdate::Mods(QueueModifiers { shuffle, repeat }));
}
PlaybackCommand::TogglePlay => {
let state = {
let Ok(state) = self.state.lock() else {
error!("play state lock poisoned");
return;
};
*state
};
debug!(?state, "toggling play");
if state == PlayState::Playing {
if let Err(err) = self.player.pause().await {
warn!("pause failed: {err:?}");
}
} else if let Err(err) = self.player.unpause().await {
warn!("unpause failed: {err:?}");
}
}
PlaybackCommand::Stop => {
debug!("stopping playback");
self.stop_player().await;
}
PlaybackCommand::ChangeVolume { delta } => {
match self.player.volume().await {
Ok(volume) => {
debug!(volume, delta, "changing volume");
if let Err(err) = self.player.set_volume(volume + delta).await {
warn!("set_volume failed: {err:?}");
}
}
Err(err) => warn!("could not read volume: {err:?}"),
};
}
PlaybackCommand::ToggleMute => match self.player.toggle_mute().await {
Ok(muted) => {
debug!(muted, "toggled mute");
self.broadcast(StreamUpdate::Mute(muted));
}
Err(err) => warn!("toggle_mute failed: {err:?}"),
},
PlaybackCommand::Next => {
let track = {
let Ok(mut queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
queue.next_track()
};
debug!(
track = track.as_ref().map(|t| t.path.as_str()),
"advancing to next track"
);
self.play_or_stop(track).await;
}
PlaybackCommand::Prev => {
let track = {
let Ok(mut queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
queue.prev_track()
};
debug!(
track = track.as_ref().map(|t| t.path.as_str()),
"going back to previous track"
);
self.play_or_stop(track).await;
}
PlaybackCommand::StateChanged { state } => {
let play_state = {
let Ok(mut state_lock) = self.state.lock() else {
error!("play state lock poisoned");
return;
};
*state_lock = state;
state
};
debug!(?play_state, "player state changed");
self.broadcast(StreamUpdate::PlayState(play_state as i32));
}
PlaybackCommand::RestartTrack => {
debug!("restarting current track");
if let Err(err) = self.player.restart().await {
warn!("restart failed: {err:?}");
}
}
PlaybackCommand::VolumeChanged { volume } => {
trace!(volume, "volume changed");
self.broadcast(StreamUpdate::Volume(volume));
}
PlaybackCommand::MuteChanged { muted } => {
trace!(muted, "mute changed");
self.broadcast(StreamUpdate::Mute(muted));
}
PlaybackCommand::PositionChanged { duration, position } => {
trace!(duration, position, "position changed");
self.broadcast(StreamUpdate::Position(TrackPosition { duration, position }));
}
}
}
/// Sends an update to all connected clients. Having no subscribers is
/// normal and not an error.
fn broadcast(&self, update: StreamUpdate) {
if let Err(err) = self.update_tx.send(update) {
trace!("no update stream subscribers: {err}");
}
}
/// Registers a pending resolve operation and spawns its forwarder task.
///
/// The forwarder resolves `paths` one after the other (preserving the
/// request's path order): for each path it sends
/// `ProviderCommand::ResolveTracks` with a fresh bounded chunk channel
/// and forwards every chunk to the playback loop as
/// `PlaybackCommand::ApplyResolvedChunk`; after the last path it sends
/// `ResolveFinished`. When the op's cancellation flag is set, the
/// forwarder drops the chunk receiver instead — the provider's next
/// send fails and the fetch stops. Queue state is never touched here:
/// mutations happen only when the loop processes the forwarded
/// commands. An immediate `Queue` broadcast (unchanged tracks,
/// `resolving = true`) gives clients instant feedback.
fn start_resolve(&self, kind: ResolveKind, paths: Vec<String>) {
let op = PendingResolve::new(kind);
let cancelled = op.cancel_flag();
let op_id = self.next_op_id.fetch_add(1, Ordering::Relaxed);
{
let Ok(mut pending) = self.pending.lock() else {
error!("pending ops lock poisoned");
return;
};
pending.insert(op_id, op);
}
debug!(op_id, ?paths, "starting queue resolve");
{
let Ok(queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
// Instant feedback: clients see resolving = true before the
// first chunk exists.
self.broadcast_queue(&queue);
}
let provider_tx = self.provider_tx.clone();
let playback_tx = self.playback_tx.clone();
tokio::spawn(
async move {
for path in &paths {
if cancelled.load(Ordering::Relaxed) {
break;
}
let (chunk_tx, chunk_rx) = flume::bounded(4);
let message = ProviderMessage::new(ProviderCommand::ResolveTracks {
path: path.clone(),
chunk_tx,
});
if provider_tx.send_async(message).await.is_err() {
error!("provider channel closed");
break;
}
let mut forwarded = 0usize;
while let Ok(tracks) = chunk_rx.recv_async().await {
// On cancellation this loop exits and drops
// chunk_rx; the provider's next send fails and the
// fetch stops.
if cancelled.load(Ordering::Relaxed) {
break;
}
forwarded += tracks.len();
let apply = PlaybackCommand::ApplyResolvedChunk { op_id, tracks };
if playback_tx
.send_async(PlaybackMessage::new(apply))
.await
.is_err()
{
error!("playback channel closed");
return;
}
}
if forwarded == 0 && !cancelled.load(Ordering::Relaxed) {
warn!(path, "path resolved to no playable tracks");
}
}
// Always reported — also for cancelled or empty ops — so
// the pending map can never leak a stuck resolving flag.
let finished = PlaybackCommand::ResolveFinished { op_id };
let _ = playback_tx.send_async(PlaybackMessage::new(finished)).await;
}
.in_current_span(),
);
}
/// Applies one chunk to the queue for pending op `op_id`, broadcasts
/// the grown queue, and starts playback when the chunk made a track
/// current. Chunks for an unknown op id (finished or cancelled) are
/// dropped silently.
async fn apply_resolved_chunk(&self, op_id: u64, tracks: Vec<Track>) {
let track = {
let Ok(mut queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
let track = {
let Ok(mut pending) = self.pending.lock() else {
error!("pending ops lock poisoned");
return;
};
let Some(op) = pending.get_mut(&op_id) else {
trace!(op_id, "dropping chunk for a finished or cancelled op");
return;
};
op.apply_chunk(&mut queue, &tracks)
};
self.broadcast_queue(&queue);
track
};
self.play_if_some(track).await;
}
/// Removes the finished op and broadcasts the final `Queue` snapshot
/// (clearing `resolving` once no ops remain). An op already removed by
/// cancellation needs no broadcast — the cancelling command mutates
/// the queue and broadcasts itself.
fn finish_resolve(&self, op_id: u64) {
let removed = {
let Ok(mut pending) = self.pending.lock() else {
error!("pending ops lock poisoned");
return;
};
pending.remove(&op_id)
};
let Some(op) = removed else {
trace!(op_id, "resolve finished for a cancelled op");
return;
};
debug!(op_id, tracks = op.applied(), "queue resolve finished");
let Ok(queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
self.broadcast_queue(&queue);
}
/// Cancels every in-flight resolve op (used by `Replace` and `Clear`).
/// The forwarders see the flag, drop their chunk receivers (stopping
/// the fetches) and still report `ResolveFinished`, which is dropped as
/// unknown here.
fn cancel_pending_resolves(&self) {
let Ok(mut pending) = self.pending.lock() else {
error!("pending ops lock poisoned");
return;
};
for (op_id, op) in pending.drain() {
debug!(op_id, "cancelling in-flight resolve");
op.cancel();
}
}
/// A wire snapshot of the queue with the `resolving` flag set from the
/// pending-op map. Callers must not hold the `pending` lock (`queue` is
/// fine — the lock order is queue, then pending).
fn queue_snapshot(&self, queue: &QueueManager) -> ProtoQueue {
let resolving = self
.pending
.lock()
.map(|pending| !pending.is_empty())
.unwrap_or(false);
let mut snapshot: ProtoQueue = queue.clone().into();
snapshot.resolving = resolving;
snapshot
}
/// Broadcasts the current queue snapshot. All queue broadcasts go
/// through here so the `resolving` flag can never be forgotten — and
/// every queue-content change reaches the persister the same way.
fn broadcast_queue(&self, queue: &QueueManager) {
self.send_persist_snapshot(queue);
self.broadcast(StreamUpdate::Queue(self.queue_snapshot(queue)));
}
/// Hands the queue's persistable state to the persister task; a no-op
/// when persistence is disabled. Latest snapshot wins, so calling this
/// on every mutation is free of backpressure (the persister skips
/// writes for unchanged snapshots).
fn send_persist_snapshot(&self, queue: &QueueManager) {
if self.store.is_none() {
return;
}
let proto: ProtoQueue = queue.clone().into();
self.persist_tx.send_replace(Some(QueueSnapshot {
tracks: proto.tracks,
current_position: proto.current_position,
repeat: queue.repeat,
shuffle: queue.shuffle,
}));
}
#[instrument(skip(self))]
async fn get_urls_for_track(&self, path: &str) -> Result<Vec<String>, ProviderError> {
let (result_tx, result_rx) = flume::bounded(1);
let message = ProviderMessage::new(ProviderCommand::GetTrackUrls {
path: path.to_string(),
result_tx,
});
self.provider_tx
.send_async(message)
.await
.map_err(|_| ProviderError::InternalError)?;
result_rx
.recv_async()
.await
.map_err(|_| ProviderError::InternalError)?
}
async fn stop_player(&self) {
if let Err(err) = self.player.stop().await {
debug!("stop had no effect: {err:?}");
}
}
/// Plays the given track if there is one, otherwise stops the player.
#[instrument(skip(self, track), fields(track = track.as_ref().map(|t| t.path.as_str())))]
async fn play_or_stop(&self, track: Option<Track>) {
if track.is_some() {
self.play(track).await;
} else {
self.stop_player().await;
}
}
/// Plays the given track if there is one; does nothing otherwise.
async fn play_if_some(&self, track: Option<Track>) {
if track.is_some() {
self.play(track).await;
}
}
/// Finds the stream URLs of the first playable track, starting at
/// `track` and advancing the queue past unplayable ones. Tracks marked
/// `is_skipped` (captures recorded their source as uncapturable) are
/// skipped without a provider round trip; tracks whose stream URLs
/// fail to resolve are skipped with a warning. Bounded by the queue
/// length at entry — one full pass at most — so an all-skipped queue
/// with repeat on returns `None` instead of spinning
/// (architecture/incremental-captures.md D3).
async fn next_playable_urls(&self, mut track: Track) -> Option<Vec<String>> {
let mut attempts_left = {
let Ok(queue) = self.queue.lock() else {
error!("queue lock poisoned");
return None;
};
queue.len()
};
loop {
let path = track.path.as_str();
if track.is_skipped {
debug!(path, "track is marked skipped, skipping");
} else {
match self.get_urls_for_track(path).await {
Ok(urls) if !urls.is_empty() => return Some(urls),
Ok(_) => warn!(path, "provider returned no stream urls, skipping track"),
Err(err) => warn!(path, "failed to fetch stream urls ({err}), skipping track"),
}
}
attempts_left = attempts_left.saturating_sub(1);
if attempts_left == 0 {
warn!("no playable track in the queue after a full pass");
return None;
}
let next = {
let Ok(mut queue) = self.queue.lock() else {
error!("queue lock poisoned");
return None;
};
queue.next_track()
};
match next {
Some(next_track) => track = next_track,
None => {
debug!("reached the end of the queue without a playable track");
return None;
}
}
}
}
/// Starts playback of the given track, skipping past unplayable ones
/// (see [`Self::next_playable_urls`]); stops the player when nothing
/// in the queue is playable.
#[instrument(skip(self, track), fields(track = track.as_ref().map(|t| t.path.as_str())))]
async fn play(&self, track: Option<Track>) {
let Some(track) = track else {
debug!("nothing to play");
return;
};
let Some(urls) = self.next_playable_urls(track).await else {
self.stop_player().await;
return;
};
{
let Ok(queue) = self.queue.lock() else {
error!("queue lock poisoned");
return;
};
// Current-track moves (Next/Prev/SetCurrent/skips) change the
// persisted position without a queue broadcast.
self.send_persist_snapshot(&queue);
self.broadcast(StreamUpdate::QueueTrack(QueueTrack {
queue_position: queue.current_position() as u32,
track: queue.current_track(),
}));
}
debug!(url_count = urls.len(), "starting player");
if let Err(err) = self.player.play(&urls[0]).await {
error!("player failed to start track: {err:?}");
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
fn track(i: usize) -> Track {
Track {
path: format!("/tidal/playlists/p/{i}"),
artist: "artist".to_string(),
title: format!("track {i}"),
duration: None,
album: None,
is_skipped: false,
}
}
async fn store_in(dir: &TempDir) -> Arc<QueueStore> {
Arc::new(
QueueStore::open(dir.path().join("queues"))
.await
.expect("open store"),
)
}
fn playback_with(store: Option<Arc<QueueStore>>) -> Playback {
let (update_tx, _) = tokio::sync::broadcast::channel(64);
let (provider_tx, _provider_rx) = flume::bounded(16);
Playback::new(update_tx, provider_tx, store)
}
fn fill_queue(playback: &Playback, n: usize) {
let tracks: Vec<Track> = (0..n).map(track).collect();
let mut queue = playback.queue.lock().expect("queue lock");
let _ = queue.replace_with_tracks(&tracks);
}
#[tokio::test]
async fn restore_fills_the_queue_without_starting_playback() {
let dir = TempDir::new().expect("tempdir");
let store = store_in(&dir).await;
store
.persist_current(&QueueSnapshot {
tracks: (0..3).map(track).collect(),
current_position: 1,
repeat: true,
shuffle: false,
})
.await
.expect("persist");
let playback = playback_with(Some(store));
playback.restore_current().await;
let queue = playback.queue.lock().expect("queue lock");
let snapshot: ProtoQueue = queue.clone().into();
assert_eq!(snapshot.tracks.len(), 3);
assert_eq!(queue.current_position(), 1);
assert!(queue.repeat);
// A restarted server stays silent: restoring must not play.
assert_eq!(
*playback.state.lock().expect("state lock"),
PlayState::Stopped
);
}
#[tokio::test]
async fn restore_survives_an_out_of_range_position() {
let dir = TempDir::new().expect("tempdir");
let store = store_in(&dir).await;
store
.persist_current(&QueueSnapshot {
tracks: vec![track(0)],
current_position: 99, // hand-edited folder
repeat: false,
shuffle: false,
})
.await
.expect("persist");
let playback = playback_with(Some(store));
playback.restore_current().await;
let queue = playback.queue.lock().expect("queue lock");
assert_eq!(queue.current_position(), 0);
}
#[tokio::test]
async fn save_queue_command_snapshots_the_live_queue() {
let dir = TempDir::new().expect("tempdir");
let store = store_in(&dir).await;
let playback = playback_with(Some(Arc::clone(&store)));
fill_queue(&playback, 2);
let (result_tx, result_rx) = flume::bounded(1);
playback
.handle_command(PlaybackCommand::SaveQueue {
name: "road trip".to_string(),
result_tx,
})
.await;
result_rx
.recv_async()
.await
.expect("reply")
.expect("save succeeds");
let entries = std::fs::read_dir(store.dir().join("road trip"))
.expect("saved queue folder")
.filter(|e| {
!e.as_ref()
.expect("entry")
.file_name()
.to_string_lossy()
.starts_with('.')
})
.count();
assert_eq!(entries, 2);
}
#[tokio::test]
async fn save_queue_rejects_an_empty_queue() {
let dir = TempDir::new().expect("tempdir");
let playback = playback_with(Some(store_in(&dir).await));
let (result_tx, result_rx) = flume::bounded(1);
playback
.handle_command(PlaybackCommand::SaveQueue {
name: "empty".to_string(),
result_tx,
})
.await;
let result = result_rx.recv_async().await.expect("reply");
assert!(matches!(result, Err(SaveQueueError::EmptyQueue)));
}
#[tokio::test]
async fn skipped_tracks_are_skipped_with_a_bounded_pass() {
// All tracks are marked skipped and repeat is on: `next_track`
// cycles forever, so only the one-full-pass bound ends the loop
// with `None`. The marked tracks are skipped without any provider
// round trip (the provider channel is closed — a call would fail,
// not hang).
let playback = playback_with(None);
let tracks: Vec<Track> = (0..3)
.map(|i| Track {
is_skipped: true,
..track(i)
})
.collect();
let first = {
let mut queue = playback.queue.lock().expect("queue lock");
let first = queue.replace_with_tracks(&tracks);
queue.repeat = true;
first
};
let urls = playback
.next_playable_urls(first.expect("first track"))
.await;
assert!(urls.is_none(), "an all-skipped queue has nothing playable");
}
#[tokio::test]
async fn queue_mutations_reach_the_persist_channel() {
let dir = TempDir::new().expect("tempdir");
let playback = playback_with(Some(store_in(&dir).await));
fill_queue(&playback, 2);
let rx = playback.persist_tx.subscribe();
playback
.handle_command(PlaybackCommand::Remove { positions: vec![1] })
.await;
let snapshot = rx.borrow().clone().expect("snapshot sent");
assert_eq!(snapshot.tracks.len(), 1);
}
}