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, provider_tx: flume::Sender, pub playback_tx: flume::Sender, playback_rx: flume::Receiver, queue: Mutex, state: Mutex, /// 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>, next_op_id: AtomicU64, /// `None` when queue persistence is disabled (no usable queues /// directory) — the queue then lives in memory only. store: Option>, /// 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>, pub player: Player, } impl Playback { pub fn new( update_tx: tokio::sync::broadcast::Sender, provider_tx: flume::Sender, store: Option>, ) -> 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 => { // FIXME: implement mute in the player engine debug!("toggle mute requested (not implemented)"); } 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) { 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) { 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, 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) { 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) { 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> { 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) { 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 { Arc::new( QueueStore::open(dir.path().join("queues")) .await .expect("open store"), ) } fn playback_with(store: Option>) -> 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 = (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 = (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); } }