diff --git a/docs/features/playback.md b/docs/features/playback.md index 9c39511f..6effde20 100644 --- a/docs/features/playback.md +++ b/docs/features/playback.md @@ -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. diff --git a/src-tauri/crates/app/src/audio/analytics.rs b/src-tauri/crates/app/src/audio/analytics.rs index 66eff010..5f26b7f6 100644 --- a/src-tauri/crates/app/src/audio/analytics.rs +++ b/src-tauri/crates/app/src/audio/analytics.rs @@ -62,6 +62,14 @@ pub enum AnalyticsMsg { source_type: String, source_id: Option, }, + /// 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), /// 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 @@ -113,6 +121,11 @@ async fn handle_message( cmd_tx: &CrossbeamSender, app: &AppHandle, ) -> Result<(), String> { + if let AnalyticsMsg::Flush(done) = msg { + done.notify_one(); + return Ok(()); + } + let state = app.state::(); // A remote-queue track has no library row: it writes no play_event and @@ -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 diff --git a/src-tauri/crates/app/src/audio/decoder.rs b/src-tauri/crates/app/src/audio/decoder.rs index 407ccc61..543ece03 100644 --- a/src-tauri/crates/app/src/audio/decoder.rs +++ b/src-tauri/crates/app/src/audio/decoder.rs @@ -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(( @@ -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(( @@ -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(( @@ -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(( @@ -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(( @@ -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, +) -> 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, diff --git a/src-tauri/crates/app/src/audio/engine.rs b/src-tauri/crates/app/src/audio/engine.rs index 2b9b557d..5d1b9292 100644 --- a/src-tauri/crates/app/src/audio/engine.rs +++ b/src-tauri/crates/app/src/audio/engine.rs @@ -459,6 +459,9 @@ pub struct AudioEngine { pub(crate) shared: Arc, output: Mutex, decoder: Mutex>>, + /// The exit path's way to the analytics task: see + /// [`AudioEngine::shut_down_and_flush`]. + analytics_tx: tokio::sync::mpsc::UnboundedSender, /// AppHandle clone so we can rebuild the cpal output thread from /// `set_output_device` without plumbing the handle through every /// Tauri command call site. @@ -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) => { @@ -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), @@ -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). /// diff --git a/src-tauri/crates/app/src/lib.rs b/src-tauri/crates/app/src/lib.rs index f30b452e..9882bd1d 100644 --- a/src-tauri/crates/app/src/lib.rs +++ b/src-tauri/crates/app/src/lib.rs @@ -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::>(); + // 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; + // 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; + } + }); } }); }