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
92 changes: 64 additions & 28 deletions src/audio/player.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -388,7 +388,7 @@ struct AudioLoopState {
reconnect_count: u32,
stream_retry_at: Option<(u32, std::time::Instant)>,
od: OnDemandTracker,
stream_last_chunk: Option<Arc<AtomicU64>>,
stream_progress: Option<ProgressClock>,
stream_download_done: Option<Arc<AtomicBool>>,
download_complete: bool,
file_backed: bool,
Expand Down Expand Up @@ -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,
Expand All @@ -433,7 +433,7 @@ struct StreamConnection {
player: Player,
duration_secs: Option<f32>,
title_rx: std_mpsc::Receiver<String>,
last_chunk: Arc<AtomicU64>,
progress_clock: ProgressClock,
download_done: Arc<AtomicBool>,
file_backed: bool,
written: Option<Arc<AtomicU64>>,
Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -649,29 +649,28 @@ fn check_on_demand_finished(st: &mut AudioLoopState, state_tx: &watch::Sender<Pl
}

fn check_stream_stall(st: &mut AudioLoopState) {
check_stream_stall_at(st, std::time::Instant::now());
}

fn check_stream_stall_at(st: &mut AudioLoopState, now: std::time::Instant) {
if st.reconnect_at.is_some() || st.stream_retry_at.is_some() {
return;
}
let Some(ref arc) = st.stream_last_chunk else {
let Some(ref progress_clock) = st.stream_progress else {
return;
};
if st.current_station.is_none() {
return;
}
let last_ms = arc.load(Ordering::Acquire);
if last_ms == 0 {
let Some(stalled_for) = progress_clock.elapsed_at(now) else {
return;
}
};
let stall_threshold = if st.od.active {
STALL_SECS_ON_DEMAND
} else {
STALL_SECS_LIVE
};
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
if now_ms.saturating_sub(last_ms) > stall_threshold * 1000 {
if stalled_for > std::time::Duration::from_secs(stall_threshold) {
let delay = backoff_duration(
st.reconnect_count,
BASE_RECONNECT_DELAY_SECS,
Expand All @@ -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;
}
}
Expand All @@ -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();

Expand Down Expand Up @@ -745,7 +744,7 @@ fn open_stream(
player,
duration_secs,
title_rx,
last_chunk,
progress_clock,
download_done,
file_backed: false,
written: None,
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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)));
Expand All @@ -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());
Expand Down
Loading