From 6a44fcb6603008511033d8e01cd7462cb2300364 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Gurruchaga?= Date: Sun, 9 Aug 2026 07:53:28 -0400 Subject: [PATCH] fix(audio): use monotonic clock for stall detection --- src/audio/player.rs | 92 ++++++++++++++++++--------- src/audio/stream.rs | 147 ++++++++++++++++++++++++++++++++++---------- 2 files changed, 180 insertions(+), 59 deletions(-) diff --git a/src/audio/player.rs b/src/audio/player.rs index 90e79a1..8ddd520 100644 --- a/src/audio/player.rs +++ b/src/audio/player.rs @@ -10,7 +10,7 @@ use tokio::sync::{mpsc, watch}; use tracing::{error, info, warn}; use crate::audio::meter::rms_to_db; -use crate::audio::stream::{FileBackedDownload, FileBackedReader, StreamReader}; +use crate::audio::stream::{FileBackedDownload, FileBackedReader, ProgressClock, StreamReader}; use crate::station::Station; pub enum PlayerCommand { @@ -388,7 +388,7 @@ struct AudioLoopState { reconnect_count: u32, stream_retry_at: Option<(u32, std::time::Instant)>, od: OnDemandTracker, - stream_last_chunk: Option>, + stream_progress: Option, stream_download_done: Option>, download_complete: bool, file_backed: bool, @@ -416,7 +416,7 @@ impl AudioLoopState { reconnect_count: 0, stream_retry_at: None, od: OnDemandTracker::inactive(), - stream_last_chunk: None, + stream_progress: None, stream_download_done: None, download_complete: false, file_backed: false, @@ -433,7 +433,7 @@ struct StreamConnection { player: Player, duration_secs: Option, title_rx: std_mpsc::Receiver, - last_chunk: Arc, + progress_clock: ProgressClock, download_done: Arc, file_backed: bool, written: Option>, @@ -583,7 +583,7 @@ fn cancel_file_backed_download( .file_download .take() .and_then(|download| download.cancel(handle)); - st.stream_last_chunk = None; + st.stream_progress = None; st.stream_download_done = None; st.file_backed = false; st.stream_written = None; @@ -601,7 +601,7 @@ fn cancel_file_backed_download( fn check_download_done(st: &mut AudioLoopState) -> bool { if let Some(ref arc) = st.stream_download_done { if arc.load(Ordering::Acquire) { - st.stream_last_chunk = None; + st.stream_progress = None; st.stream_download_done = None; st.download_complete = true; st.file_download = None; @@ -649,29 +649,28 @@ fn check_on_demand_finished(st: &mut AudioLoopState, state_tx: &watch::Sender stall_threshold * 1000 { + if stalled_for > std::time::Duration::from_secs(stall_threshold) { let delay = backoff_duration( st.reconnect_count, BASE_RECONNECT_DELAY_SECS, @@ -683,8 +682,8 @@ fn check_stream_stall(st: &mut AudioLoopState) { delay.as_secs_f32(), st.reconnect_count + 1 ); - st.stream_last_chunk = None; - st.reconnect_at = Some(std::time::Instant::now() + delay); + st.stream_progress = None; + st.reconnect_at = Some(now + delay); st.reconnect_count += 1; } } @@ -711,7 +710,7 @@ fn open_stream( station.custom_headers.clone(), handle.clone(), ); - let last_chunk = stream_reader.last_chunk_arc(); + let progress_clock = stream_reader.progress_clock(); let download_done = stream_reader.download_done_arc(); let dead_url_arc = stream_reader.dead_url_arc(); @@ -745,7 +744,7 @@ fn open_stream( player, duration_secs, title_rx, - last_chunk, + progress_clock, download_done, file_backed: false, written: None, @@ -804,7 +803,7 @@ fn open_file_backed_stream( Err(e) => return Err(Some(format!("cache file: {e}"))), }; - let last_chunk = reader.last_chunk_arc(); + let progress_clock = reader.progress_clock(); let download_done = reader.download_done_arc(); let dead_url_arc = reader.dead_url_arc(); let written = reader.written_arc(); @@ -836,7 +835,7 @@ fn open_file_backed_stream( player, duration_secs, title_rx, - last_chunk, + progress_clock, download_done, file_backed: true, written: Some(written), @@ -947,7 +946,7 @@ fn handle_play_cmd( Ok(conn) => { st.stream_retry_at = None; st.title_rx = Some(conn.title_rx); - st.stream_last_chunk = Some(conn.last_chunk); + st.stream_progress = Some(conn.progress_clock); st.stream_download_done = Some(conn.download_done); st.file_backed = conn.file_backed; st.stream_written = conn.written; @@ -1002,7 +1001,7 @@ fn handle_crossfade_cmd( Ok(conn) => { st.stream_retry_at = None; st.title_rx = Some(conn.title_rx); - st.stream_last_chunk = Some(conn.last_chunk); + st.stream_progress = Some(conn.progress_clock); st.stream_download_done = Some(conn.download_done); st.file_backed = conn.file_backed; st.stream_written = conn.written; @@ -1172,7 +1171,7 @@ fn handle_seek_cmd( let (mut stream_reader, new_title_rx) = StreamReader::connect(url, byte_offset, 4096, custom_headers, handle.clone()); st.title_rx = Some(new_title_rx); - st.stream_last_chunk = Some(stream_reader.last_chunk_arc()); + st.stream_progress = Some(stream_reader.progress_clock()); st.stream_download_done = Some(stream_reader.download_done_arc()); st.download_complete = false; let prebuffer_secs = st @@ -1259,7 +1258,7 @@ fn handle_device_change( stop_crossfade_out(st); cancel_file_backed_download(st, handle, true); st.title_rx = None; - st.stream_last_chunk = None; + st.stream_progress = None; st.stream_download_done = None; st.download_complete = false; st.file_backed = false; @@ -1554,7 +1553,7 @@ fn audio_loop( st.reconnect_at = None; st.reconnect_count = 0; st.stream_retry_at = None; - st.stream_last_chunk = None; + st.stream_progress = None; st.stream_download_done = None; st.download_complete = false; st.file_backed = false; @@ -1612,6 +1611,43 @@ mod tests { assert_eq!(on_demand_byte_offset(10.0, &station), 160_000); } + #[test] + fn live_stream_stall_reconnects_only_after_thirty_seconds() { + let started = std::time::Instant::now(); + let mut st = AudioLoopState::new(); + st.current_station = Some(station_with_bitrate(None)); + st.stream_progress = Some(ProgressClock::with_last_progress_at(started)); + + check_stream_stall_at(&mut st, started + std::time::Duration::from_secs(30)); + assert!(st.reconnect_at.is_none()); + + let stalled_at = started + std::time::Duration::from_millis(30_001); + check_stream_stall_at(&mut st, stalled_at); + + assert!(st.reconnect_at.is_some_and(|at| at > stalled_at)); + assert!(st.stream_progress.is_none()); + assert_eq!(st.reconnect_count, 1); + } + + #[test] + fn on_demand_stall_reconnects_only_after_sixty_seconds() { + let started = std::time::Instant::now(); + let mut st = AudioLoopState::new(); + st.current_station = Some(station_with_bitrate(None)); + st.od.active = true; + st.stream_progress = Some(ProgressClock::with_last_progress_at(started)); + + check_stream_stall_at(&mut st, started + std::time::Duration::from_secs(60)); + assert!(st.reconnect_at.is_none()); + + let stalled_at = started + std::time::Duration::from_millis(60_001); + check_stream_stall_at(&mut st, stalled_at); + + assert!(st.reconnect_at.is_some_and(|at| at > stalled_at)); + assert!(st.stream_progress.is_none()); + assert_eq!(st.reconnect_count, 1); + } + #[test] fn conditional_preview_stop_ignores_stale_preview_timer() { let mut st = AudioLoopState::new(); @@ -1680,7 +1716,7 @@ mod tests { FileBackedDownload::from_parts_for_test(task, path.clone(), Arc::clone(&done)); let mut st = AudioLoopState::new(); st.file_backed = true; - st.stream_last_chunk = Some(Arc::new(AtomicU64::new(1))); + st.stream_progress = Some(ProgressClock::default()); st.stream_download_done = Some(done); st.stream_written = Some(Arc::new(AtomicU64::new(1))); st.stream_total = Some(Arc::new(AtomicU64::new(10))); @@ -1690,7 +1726,7 @@ mod tests { assert!(st.file_download.is_none()); assert!(!st.file_backed); - assert!(st.stream_last_chunk.is_none()); + assert!(st.stream_progress.is_none()); assert!(st.stream_download_done.is_none()); assert!(st.stream_written.is_none()); assert!(st.stream_total.is_none()); diff --git a/src/audio/stream.rs b/src/audio/stream.rs index ddd9880..4d04d03 100644 --- a/src/audio/stream.rs +++ b/src/audio/stream.rs @@ -5,15 +5,54 @@ use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::Arc; use std::sync::{mpsc, Mutex}; +use std::time::{Duration, Instant}; use crate::metadata::parse_icy_title; +#[derive(Clone, Default)] +pub(crate) struct ProgressClock { + last_progress: Arc>>, +} + +impl ProgressClock { + pub(crate) fn record_progress(&self) { + self.record_progress_at(Instant::now()); + } + + fn record_progress_at(&self, now: Instant) { + let mut last_progress = self + .last_progress + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + *last_progress = Some(now); + } + + pub(crate) fn elapsed(&self) -> Option { + self.elapsed_at(Instant::now()) + } + + pub(crate) fn elapsed_at(&self, now: Instant) -> Option { + let last_progress = *self + .last_progress + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + last_progress.map(|last| now.saturating_duration_since(last)) + } + + #[cfg(test)] + pub(crate) fn with_last_progress_at(last_progress: Instant) -> Self { + Self { + last_progress: Arc::new(Mutex::new(Some(last_progress))), + } + } +} + pub struct StreamReader { rx: Mutex>, chunks: VecDeque, offset: usize, buffered: usize, - last_chunk_at: Arc, + progress_clock: ProgressClock, download_done: Arc, dead_url: Arc, } @@ -23,16 +62,12 @@ impl StreamReader { if chunk.is_empty() { return; } - let now = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .map(|d| d.as_millis() as u64) - .unwrap_or(0); - self.last_chunk_at.store(now, Ordering::Release); + self.progress_clock.record_progress(); self.buffered += chunk.len(); self.chunks.push_back(chunk); } - pub fn last_chunk_arc(&self) -> Arc { - Arc::clone(&self.last_chunk_at) + pub(crate) fn progress_clock(&self) -> ProgressClock { + self.progress_clock.clone() } pub fn connect( @@ -66,13 +101,13 @@ impl StreamReader { } }); - let last_chunk_at = Arc::new(AtomicU64::new(0)); + let progress_clock = ProgressClock::default(); let reader = Self { rx: Mutex::new(audio_rx), chunks: VecDeque::new(), offset: 0, buffered: 0, - last_chunk_at, + progress_clock, download_done, dead_url, }; @@ -96,13 +131,13 @@ impl StreamReader { } }); - let last_chunk_at = Arc::new(AtomicU64::new(0)); + let progress_clock = ProgressClock::default(); Self { rx: Mutex::new(audio_rx), chunks: VecDeque::new(), offset: 0, buffered: 0, - last_chunk_at, + progress_clock, download_done: Arc::new(AtomicBool::new(false)), dead_url: Arc::new(AtomicBool::new(false)), } @@ -186,13 +221,10 @@ impl Seek for StreamReader { const MAX_DOWNLOAD_RETRIES: u32 = 5; const DOWNLOAD_RETRY_PAUSE: std::time::Duration = std::time::Duration::from_secs(2); const FILE_READ_STARVATION_SECS: u64 = 30; -const FILE_READ_STALL_MS: u64 = 12_000; +const FILE_READ_STALL_SECS: u64 = 12; -fn unix_now_ms() -> u64 { - std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .map(|d| d.as_millis() as u64) - .unwrap_or(0) +fn file_read_is_stalled(stalled_for: Option) -> bool { + stalled_for.is_some_and(|elapsed| elapsed >= Duration::from_secs(FILE_READ_STALL_SECS)) } pub fn youtube_cache_dir() -> std::path::PathBuf { @@ -213,7 +245,7 @@ pub struct FileBackedReader { pos: u64, written: Arc, total_len: Arc, - last_chunk_at: Arc, + progress_clock: ProgressClock, download_done: Arc, dead_url: Arc, } @@ -302,13 +334,13 @@ impl FileBackedReader { let (title_tx, title_rx) = mpsc::sync_channel::(1); let written = Arc::new(AtomicU64::new(0)); let total_len = Arc::new(AtomicU64::new(0)); - let last_chunk_at = Arc::new(AtomicU64::new(0)); + let progress_clock = ProgressClock::default(); let download_done = Arc::new(AtomicBool::new(false)); let dead_url = Arc::new(AtomicBool::new(false)); let written_task = Arc::clone(&written); let total_task = Arc::clone(&total_len); - let last_chunk_task = Arc::clone(&last_chunk_at); + let progress_clock_task = progress_clock.clone(); let done_task = Arc::clone(&download_done); let dead_task = Arc::clone(&dead_url); let task = handle.spawn(async move { @@ -318,7 +350,7 @@ impl FileBackedReader { write_file, written_task, total_task, - last_chunk_task, + progress_clock_task, done_task, dead_task, ) @@ -340,7 +372,7 @@ impl FileBackedReader { pos: 0, written, total_len, - last_chunk_at, + progress_clock, download_done, dead_url, }, @@ -349,8 +381,8 @@ impl FileBackedReader { )) } - pub fn last_chunk_arc(&self) -> Arc { - Arc::clone(&self.last_chunk_at) + pub(crate) fn progress_clock(&self) -> ProgressClock { + self.progress_clock.clone() } pub fn download_done_arc(&self) -> Arc { @@ -415,16 +447,15 @@ impl Read for FileBackedReader { ); return Ok(0); } - let last_chunk = self.last_chunk_at.load(Ordering::Acquire); - let stalled_ms = unix_now_ms().saturating_sub(last_chunk); - let stalled = last_chunk != 0 && stalled_ms >= FILE_READ_STALL_MS; + let stalled_for = self.progress_clock.elapsed(); + let stalled = file_read_is_stalled(stalled_for); if stalled || wait_start.elapsed() > std::time::Duration::from_secs(FILE_READ_STARVATION_SECS) { tracing::warn!( pos = self.pos, written, - stalled_ms, + stalled_ms = stalled_for.map(|elapsed| elapsed.as_millis()), "File-backed read underrun: download not progressing, signaling EOF" ); return Ok(0); @@ -468,7 +499,7 @@ async fn download_to_file( mut file: std::fs::File, written: Arc, total_len: Arc, - last_chunk_at: Arc, + progress_clock: ProgressClock, download_done: Arc, dead_url: Arc, ) -> Result<(), reqwest::Error> { @@ -567,7 +598,7 @@ async fn download_to_file( } current_offset += bytes.len() as u64; written.store(current_offset, Ordering::Release); - last_chunk_at.store(unix_now_ms(), Ordering::Release); + progress_clock.record_progress(); retry_count = 0; } Ok(Some(Err(e))) => { @@ -1103,6 +1134,60 @@ mod tests { use super::*; use std::sync::mpsc; + #[test] + fn progress_clock_starts_without_progress() { + let clock = ProgressClock::default(); + + assert_eq!(clock.elapsed_at(Instant::now()), None); + } + + #[test] + fn progress_clock_measures_elapsed_time_monotonically() { + let started = Instant::now(); + let clock = ProgressClock::with_last_progress_at(started); + + assert_eq!( + clock.elapsed_at(started + Duration::from_secs(11)), + Some(Duration::from_secs(11)) + ); + assert_eq!( + clock.elapsed_at(started + Duration::from_secs(FILE_READ_STALL_SECS)), + Some(Duration::from_secs(FILE_READ_STALL_SECS)) + ); + + clock.record_progress_at(started + Duration::from_secs(20)); + assert_eq!( + clock.elapsed_at(started + Duration::from_secs(25)), + Some(Duration::from_secs(5)) + ); + } + + #[test] + fn file_read_stall_uses_the_twelve_second_threshold() { + assert!(!file_read_is_stalled(None)); + assert!(!file_read_is_stalled(Some(Duration::from_secs(11)))); + assert!(file_read_is_stalled(Some(Duration::from_secs(12)))); + } + + #[test] + fn progress_clock_recovers_from_a_poisoned_lock() { + let clock = ProgressClock::default(); + let poisoned_clock = clock.clone(); + let _ = std::thread::spawn(move || { + let _guard = poisoned_clock + .last_progress + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + panic!("poison progress clock for test"); + }) + .join(); + let now = Instant::now(); + + clock.record_progress_at(now); + + assert_eq!(clock.elapsed_at(now), Some(Duration::ZERO)); + } + fn make_stripper(metaint: usize) -> (IcyStripper, mpsc::Receiver) { let (tx, rx) = mpsc::sync_channel(8); (IcyStripper::new(metaint, tx), rx)