Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions docs/features/playback.md
Original file line number Diff line number Diff line change
Expand Up @@ -658,3 +658,9 @@ Three sites write it, and none coordinates with the others, because writing the
Only a real library track is worth saving: radio and the remote queue use negative ids and nothing loaded is `0`, and `restore_state` has no `track` row to find for either, so the write is skipped there.

Two details keep concurrent writers honest. Each write is **pinned to the profile that was playing** — the id is captured before the await, without waiting for the lock, and a lock that can't be taken means a switch is under way, so the write is dropped rather than landing this track's position in the profile being switched to (#485). And `persist_resume_point` sets both rows **in one transaction**, so two writers can't leave one track's id beside another's position. The ticker only remembers a point that actually landed, so a write lost to a busy database is retried on the next tick instead of being skipped as unchanged.

### The last play at exit

A play is credited when it ends (`play_event`, then the scrobble queue), and quitting ends it: the decoder answers `Shutdown` like a skip, crediting 15 s or more of listening through the analytics task. Until 1.8.2 the process exited with that message still in the channel, and the window close was the only path that sent `Shutdown` at all, so the last track listened to before quitting never reached the history, the stats or Last.fm.

`RunEvent::Exit` silences the output, writes the resume point, then calls `AudioEngine::shut_down_and_flush`. It sends `Shutdown`, waits for the decoder thread to finish (each `Shutdown` return re-queues the command so the decoder loop ends once the credit is sent, instead of waiting for the next command), then queues an `AnalyticsMsg::Flush` behind the credit and waits for its answer: the channel is FIFO, so the answer means the credit was handled — written, or failed and logged, since there is no retrying a write once the process is on its way out. The whole wait is bounded at 2 s; a decoder stuck in a slow read loses that play, as before, rather than holding the process open. Only the decoder credits, so a close that already sent `Shutdown` is not counted twice.
15 changes: 14 additions & 1 deletion src-tauri/crates/app/src/audio/analytics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,14 @@ pub enum AnalyticsMsg {
source_type: String,
source_id: Option<i64>,
},
/// Answered once every message sent before it has been handled. The
/// channel is FIFO and the task handles one message at a time, so the
/// notification means the writes those messages asked for have been
/// attempted: done, or failed and logged (no pool, a failed insert).
/// Sent by the exit path, which would otherwise end the process with
/// the last play still in the channel, and which could not retry a
/// failed write anyway.
Flush(Arc<tokio::sync::Notify>),
/// Sent by the decoder when it's approaching the end of the
/// current track and crossfade is enabled. Triggers a
/// `peek_next` and a `SetNextTrack` reply so the decoder can
Expand Down Expand Up @@ -113,6 +121,11 @@ async fn handle_message(
cmd_tx: &CrossbeamSender<AudioCmd>,
app: &AppHandle,
) -> Result<(), String> {
if let AnalyticsMsg::Flush(done) = msg {
done.notify_one();
return Ok(());
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

let state = app.state::<AppState>();

// A remote-queue track has no library row: it writes no play_event and
Expand Down Expand Up @@ -243,7 +256,7 @@ async fn handle_message(
}
// Handled before the pool acquisition above (no play_event, no
// pool needed).
AnalyticsMsg::RemoteTrackEnded { .. } => {}
AnalyticsMsg::RemoteTrackEnded { .. } | AnalyticsMsg::Flush(_) => {}
AnalyticsMsg::PrefetchNext => {
// Look up what would be played next without bumping the
// cursor (the cursor is bumped only when the crossfade
Expand Down
56 changes: 26 additions & 30 deletions src-tauri/crates/app/src/audio/decoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1211,12 +1211,7 @@ fn play_dop_track(
ControlFlow::Continue => {}
ControlFlow::Break => break 'pkt,
ControlFlow::Shutdown => {
transition_state(shared, app, PlayerState::Idle, Some(stream.track_id));
return Ok((
PlaybackEnd::Interrupted,
shared.session_listened_ms(),
finished_from(&stream),
));
return shutdown_outcome(&stream, shared, app, pending_cmd);
}
ControlFlow::LoadNext => {
return Ok((
Expand Down Expand Up @@ -1289,12 +1284,7 @@ fn play_dop_track(
PushOutcome::Ok => {}
PushOutcome::Stop => break 'pkt,
PushOutcome::Shutdown => {
transition_state(shared, app, PlayerState::Idle, Some(stream.track_id));
return Ok((
PlaybackEnd::Interrupted,
shared.session_listened_ms(),
finished_from(&stream),
));
return shutdown_outcome(&stream, shared, app, pending_cmd);
}
PushOutcome::LoadNext => {
return Ok((
Expand Down Expand Up @@ -1441,12 +1431,7 @@ fn play_track(
ControlFlow::Continue => {}
ControlFlow::Break => break 'pkt,
ControlFlow::Shutdown => {
transition_state(shared, app, PlayerState::Idle, Some(stream.track_id));
return Ok((
PlaybackEnd::Interrupted,
shared.session_listened_ms(),
finished_from(&stream),
));
return shutdown_outcome(&stream, shared, app, pending_cmd);
}
ControlFlow::LoadNext => {
return Ok((
Expand Down Expand Up @@ -1705,12 +1690,7 @@ fn play_track(
PushOutcome::Ok => {}
PushOutcome::Stop => break 'pkt,
PushOutcome::Shutdown => {
transition_state(shared, app, PlayerState::Idle, Some(stream.track_id));
return Ok((
PlaybackEnd::Interrupted,
shared.session_listened_ms(),
finished_from(&stream),
));
return shutdown_outcome(&stream, shared, app, pending_cmd);
}
PushOutcome::LoadNext => {
return Ok((
Expand Down Expand Up @@ -1896,12 +1876,7 @@ fn play_track(
PushOutcome::Ok => primary_resampled.clear(),
PushOutcome::Stop => break 'pkt,
PushOutcome::Shutdown => {
transition_state(shared, app, PlayerState::Idle, Some(stream.track_id));
return Ok((
PlaybackEnd::Interrupted,
shared.session_listened_ms(),
finished_from(&stream),
));
return shutdown_outcome(&stream, shared, app, pending_cmd);
}
PushOutcome::LoadNext => {
return Ok((
Expand Down Expand Up @@ -1958,6 +1933,27 @@ fn play_track(
}
}

/// The outcome of a play cut short by `Shutdown`: credited like any other
/// interruption, with the `Shutdown` it consumed queued again so the outer
/// loop ends the thread once that credit is sent. The exit path waits for
/// the thread to finish, then for the analytics task to write what it was
/// sent; without the re-queue the thread went back to waiting for a command
/// and the exit could not tell when the credit had left.
fn shutdown_outcome(
stream: &ActiveStream,
shared: &SharedPlayback,
app: &AppHandle,
pending_cmd: &mut Option<AudioCmd>,
) -> Result<(PlaybackEnd, u64, FinishedTrack), String> {
transition_state(shared, app, PlayerState::Idle, Some(stream.track_id));
*pending_cmd = Some(AudioCmd::Shutdown);
Ok((
PlaybackEnd::Interrupted,
shared.session_listened_ms(),
finished_from(stream),
))
}

fn finished_from(stream: &ActiveStream) -> FinishedTrack {
FinishedTrack {
track_id: stream.track_id,
Expand Down
55 changes: 54 additions & 1 deletion src-tauri/crates/app/src/audio/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -459,6 +459,9 @@ pub struct AudioEngine {
pub(crate) shared: Arc<SharedPlayback>,
output: Mutex<OutputSlot>,
decoder: Mutex<Option<JoinHandle<()>>>,
/// The exit path's way to the analytics task: see
/// [`AudioEngine::shut_down_and_flush`].
analytics_tx: tokio::sync::mpsc::UnboundedSender<AnalyticsMsg>,
/// AppHandle clone so we can rebuild the cpal output thread from
/// `set_output_device` without plumbing the handle through every
/// Tauri command call site.
Expand Down Expand Up @@ -1094,7 +1097,7 @@ impl AudioEngine {
producer,
shared.clone(),
app.clone(),
analytics_tx,
analytics_tx.clone(),
) {
Ok(join) => (Some(handle), Some(join), active),
Err(err) => {
Expand All @@ -1121,6 +1124,7 @@ impl AudioEngine {
pinned,
}),
decoder: Mutex::new(decoder),
analytics_tx,
app,
exclusive_output: std::sync::atomic::AtomicBool::new(exclusive_output),
exclusive_output_active: std::sync::atomic::AtomicBool::new(exclusive_output_active),
Expand Down Expand Up @@ -1203,6 +1207,55 @@ impl AudioEngine {
self.shared.try_claim_load(intent.get())
}

/// Stop the decoder and wait, at most `budget`, for the play it was in
/// the middle of to reach the database.
///
/// A play is credited when it ends, and quitting ends it: the decoder
/// answers `Shutdown` with the same credit a skip gets (15 s or more),
/// sent to the analytics task. The process used to exit with that
/// message still in the channel, so the last thing listened to before
/// quitting never reached the history or the stats. Waiting for the
/// thread to finish proves the message is sent; the [`AnalyticsMsg::Flush`]
/// queued behind it proves it was handled (written, or failed and
/// logged by the analytics task).
///
/// Bounded because a decoder stuck in a slow read must not hold the
/// process open: past the budget the play is lost, as before.
pub async fn shut_down_and_flush(&self, budget: Duration) {
let deadline = tokio::time::Instant::now() + budget;
// A window close may have sent it already, and then the channel is
// closed: the thread is finishing or finished either way.
let _ = self.cmd_tx.send(AudioCmd::Shutdown);
let decoder = self
.decoder
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.take();
if let Some(decoder) = decoder {
while !decoder.is_finished() {
if tokio::time::Instant::now() >= deadline {
tracing::warn!("decoder still running at exit; last play not credited");
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
let done = Arc::new(tokio::sync::Notify::new());
if self
.analytics_tx
.send(AnalyticsMsg::Flush(done.clone()))
.is_err()
{
return;
}
if tokio::time::timeout_at(deadline, done.notified())
.await
.is_err()
{
tracing::warn!("analytics not drained at exit; last play may be lost");
}
}

/// Send a command to the decoder. Returns `AppError::Audio` if the
/// channel is disconnected (decoder thread has exited).
///
Expand Down
15 changes: 14 additions & 1 deletion src-tauri/crates/app/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1299,7 +1299,20 @@ pub fn run() {
// triggers `RunEvent::ExitRequested` / `Exit` instead, which is
// why the resume point was never written (#624).
if matches!(event, tauri::RunEvent::Exit) {
tauri::async_runtime::block_on(write_resume_point(app));
tauri::async_runtime::block_on(async {
let engine = app.try_state::<Arc<AudioEngine>>();
// Silent first, as the close path makes it: the ring still
// holds audio that would play through both writes below.
if let Some(engine) = &engine {
engine.shared().paused_output.store(true, Ordering::Release);
}
write_resume_point(app).await;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
// After the resume point, which reads the position the
// shutdown is about to leave behind.
if let Some(engine) = engine {
engine.shut_down_and_flush(Duration::from_secs(2)).await;
}
});
}
});
}
Expand Down
Loading