978 lines
39 KiB
Rust
978 lines
39 KiB
Rust
use crate::capture::CaptureError;
|
|
use crate::crabidy_store::{self, CrabidyStore, QueueSnapshot};
|
|
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, VecDeque};
|
|
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 state
|
|
/// directory) — the queue then lives in memory only.
|
|
store: Option<Arc<CrabidyStore>>,
|
|
/// 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<CrabidyStore>>,
|
|
audio_device: Option<String>,
|
|
) -> 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::new(audio_device);
|
|
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 {
|
|
crabidy_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,
|
|
}),
|
|
// Stamped by the RPC handler, which owns the auth
|
|
// switch; the playback loop only knows queue state.
|
|
auth_enabled: false,
|
|
}
|
|
};
|
|
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_snapshot(&name, &snapshot.tracks).await,
|
|
None => Err(CaptureError::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` with an *exponential read-ahead*: a
|
|
/// concurrency window that starts at 1 and doubles (1, 2, 4, 8, 16, then
|
|
/// steady 16) after each path completes. It keeps up to `window` paths
|
|
/// resolving at once (each via `ProviderCommand::ResolveTracks` on its own
|
|
/// bounded chunk channel) but forwards chunks to the playback loop as
|
|
/// `PlaybackCommand::ApplyResolvedChunk` in strict path order, so the queue
|
|
/// keeps the requested order and the first chunk still lands — and starts
|
|
/// playback — as fast as a single resolve. The first track therefore plays
|
|
/// as soon as one resolve yields it, while the window fills the queue far
|
|
/// ahead of playback: a very short or skipped leading track cannot outrun
|
|
/// resolution and cause a silent stall (architecture/progressive-queueing.md D8).
|
|
///
|
|
/// After the last path it sends `ResolveFinished`. When the op's
|
|
/// cancellation flag is set, the forwarder stops launching resolves and
|
|
/// drops its chunk receivers — each 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 {
|
|
// Exponential read-ahead window: 1, 2, 4, 8, 16, then steady.
|
|
const MAX_WINDOW: usize = 16;
|
|
let mut window = 1usize;
|
|
let mut next = 0usize;
|
|
let mut inflight: VecDeque<(String, flume::Receiver<Vec<Track>>)> = VecDeque::new();
|
|
'resolve: loop {
|
|
if cancelled.load(Ordering::Relaxed) {
|
|
break;
|
|
}
|
|
// Top the window up with fresh concurrent resolves. Each
|
|
// resolves in the background into its own bounded channel
|
|
// while we drain the oldest one below.
|
|
while inflight.len() < window && next < paths.len() {
|
|
let path = paths[next].clone();
|
|
next += 1;
|
|
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 'resolve;
|
|
}
|
|
inflight.push_back((path, chunk_rx));
|
|
}
|
|
// Drain the oldest (lowest path index) resolve to
|
|
// completion before the next, so chunks reach the loop in
|
|
// request order even though resolution ran concurrently.
|
|
let Some((path, chunk_rx)) = inflight.pop_front() else {
|
|
break;
|
|
};
|
|
let mut forwarded = 0usize;
|
|
while let Ok(tracks) = chunk_rx.recv_async().await {
|
|
// On cancellation this drops chunk_rx (and, on the next
|
|
// iteration, the rest of `inflight`); each provider's
|
|
// next send then fails and its fetch stops.
|
|
if cancelled.load(Ordering::Relaxed) {
|
|
break 'resolve;
|
|
}
|
|
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");
|
|
}
|
|
window = (window * 2).min(MAX_WINDOW);
|
|
}
|
|
// 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 — or when this op still wants to start but its earlier attempt
|
|
/// ran the player dry on skipped/short head tracks before more resolved.
|
|
/// 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 attempt = {
|
|
let Ok(mut queue) = self.queue.lock() else {
|
|
error!("queue lock poisoned");
|
|
return;
|
|
};
|
|
let (designated, wants_start) = {
|
|
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;
|
|
};
|
|
let designated = op.apply_chunk(&mut queue, &tracks);
|
|
(designated, op.wants_start())
|
|
};
|
|
self.broadcast_queue(&queue);
|
|
// Start on the track this chunk made current; failing that, if the
|
|
// op still wants to start (a prior attempt found nothing playable
|
|
// yet), retry from the current position — `next_playable_urls`
|
|
// advances past skipped/unplayable heads to the first track that
|
|
// has now resolved.
|
|
designated.or_else(|| {
|
|
if wants_start {
|
|
queue.current_track()
|
|
} else {
|
|
None
|
|
}
|
|
})
|
|
};
|
|
// Confirming the start clears the op's `wants_start`, so later chunks
|
|
// only extend the queue and never restart the playing track.
|
|
if self.play(attempt).await {
|
|
if let Ok(mut pending) = self.pending.lock() {
|
|
if let Some(op) = pending.get_mut(&op_id) {
|
|
op.mark_started();
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
/// 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;
|
|
}
|
|
}
|
|
|
|
/// 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.
|
|
///
|
|
/// Returns `true` when a playable track was found and handed to the
|
|
/// player. A device error while starting is logged but still counts as
|
|
/// started — retrying the same track would not help. Returns `false` when
|
|
/// there was nothing to play or nothing playable remained (the player is
|
|
/// then stopped), so the caller can retry once more tracks resolve.
|
|
#[instrument(skip(self, track), fields(track = track.as_ref().map(|t| t.path.as_str())))]
|
|
async fn play(&self, track: Option<Track>) -> bool {
|
|
let Some(track) = track else {
|
|
debug!("nothing to play");
|
|
return false;
|
|
};
|
|
let Some(urls) = self.next_playable_urls(track).await else {
|
|
self.stop_player().await;
|
|
return false;
|
|
};
|
|
{
|
|
let Ok(queue) = self.queue.lock() else {
|
|
error!("queue lock poisoned");
|
|
return false;
|
|
};
|
|
// 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:?}");
|
|
}
|
|
true
|
|
}
|
|
}
|
|
|
|
#[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,
|
|
provider_item_id: String::new(),
|
|
is_captured: false,
|
|
}
|
|
}
|
|
|
|
async fn store_in(dir: &TempDir) -> Arc<CrabidyStore> {
|
|
Arc::new(
|
|
CrabidyStore::open(dir.path().join("state"), dir.path().join("store"))
|
|
.await
|
|
.expect("open store"),
|
|
)
|
|
}
|
|
|
|
fn playback_with(store: Option<Arc<CrabidyStore>>) -> Playback {
|
|
let (update_tx, _) = tokio::sync::broadcast::channel(64);
|
|
let (provider_tx, _provider_rx) = flume::bounded(16);
|
|
Playback::new(update_tx, provider_tx, store, None)
|
|
}
|
|
|
|
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.tree_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(CaptureError::BadSource(_))));
|
|
}
|
|
|
|
#[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);
|
|
}
|
|
}
|