From abc2c0b57966c19fcff65c3cccbfee540f540c3e Mon Sep 17 00:00:00 2001 From: cucadmuh <49571317+cucadmuh@users.noreply.github.com> Date: Fri, 12 Jun 2026 19:41:17 +0300 Subject: [PATCH] refactor(audio): split source-build pipeline into source_build module (#1074) Move the source-building pipeline out of play_input.rs into a focused source_build.rs: BuildSourceArgs/PlaybackSource, the ranged-stream probe fallback (build_playback_source_with_probe_fallback) and its private helpers (build_source_from_play_input, wait_or_fetch_bytes_for_stream_fallback, etc.). play_input.rs now only handles source *selection* and drops from 799 to 420 lines; helpers that became module-internal lose their pub(super) visibility. No behavior change. --- .../crates/psysonic-audio/src/commands.rs | 6 +- .../psysonic-audio/src/device_resume.rs | 6 +- src-tauri/crates/psysonic-audio/src/lib.rs | 1 + .../crates/psysonic-audio/src/play_input.rs | 381 +---------------- .../crates/psysonic-audio/src/source_build.rs | 398 ++++++++++++++++++ 5 files changed, 405 insertions(+), 387 deletions(-) create mode 100644 src-tauri/crates/psysonic-audio/src/source_build.rs diff --git a/src-tauri/crates/psysonic-audio/src/commands.rs b/src-tauri/crates/psysonic-audio/src/commands.rs index 598dddec..585e073d 100644 --- a/src-tauri/crates/psysonic-audio/src/commands.rs +++ b/src-tauri/crates/psysonic-audio/src/commands.rs @@ -14,10 +14,8 @@ use super::decode::build_source; use super::engine::{audio_http_client, AudioEngine}; use super::helpers::*; use super::ipc::{maybe_emit_normalization_state, NormalizationStatePayload}; -use super::play_input::{ - build_playback_source_with_probe_fallback, select_play_input, url_format_hint, BuildSourceArgs, - PlayInputContext, -}; +use super::play_input::{select_play_input, url_format_hint, PlayInputContext}; +use super::source_build::{build_playback_source_with_probe_fallback, BuildSourceArgs}; use super::sink_swap::{ spawn_legacy_stream_start_when_armed, swap_in_new_sink, LegacyStreamStartWhenArmed, SinkSwapInputs, diff --git a/src-tauri/crates/psysonic-audio/src/device_resume.rs b/src-tauri/crates/psysonic-audio/src/device_resume.rs index 8e1c52a3..1dd8e7b0 100644 --- a/src-tauri/crates/psysonic-audio/src/device_resume.rs +++ b/src-tauri/crates/psysonic-audio/src/device_resume.rs @@ -23,9 +23,9 @@ use tauri::Emitter; use tauri::Manager; use super::engine::AudioEngine; -use super::play_input::{ - build_playback_source_with_probe_fallback, url_format_hint, BuildSourceArgs, PlayInput, - PlaybackSource, +use super::play_input::{url_format_hint, PlayInput}; +use super::source_build::{ + build_playback_source_with_probe_fallback, BuildSourceArgs, PlaybackSource, }; use super::sink_swap::{swap_in_new_sink, SinkSwapInputs}; use super::progress_task::spawn_progress_task; diff --git a/src-tauri/crates/psysonic-audio/src/lib.rs b/src-tauri/crates/psysonic-audio/src/lib.rs index 22d06ff1..81957f97 100644 --- a/src-tauri/crates/psysonic-audio/src/lib.rs +++ b/src-tauri/crates/psysonic-audio/src/lib.rs @@ -17,6 +17,7 @@ pub mod mix_commands; mod play_input; pub mod playback_rate; mod sink_swap; +mod source_build; mod preserve_worker; pub mod preload_commands; pub(crate) mod progress_task; diff --git a/src-tauri/crates/psysonic-audio/src/play_input.rs b/src-tauri/crates/psysonic-audio/src/play_input.rs index b87b0f8d..4ed0382d 100644 --- a/src-tauri/crates/psysonic-audio/src/play_input.rs +++ b/src-tauri/crates/psysonic-audio/src/play_input.rs @@ -15,11 +15,10 @@ use super::analysis_dispatch::{ prepare_playback_analysis, spawn_track_analysis_bytes, spawn_track_analysis_file, TrackAnalysisOrigin, }; -use super::decode::{build_source, build_streaming_source, BuiltSource, SizedDecoder}; use super::engine::{audio_http_client, AudioEngine}; use super::helpers::{ content_type_to_hint, fetch_data, format_hint_from_content_disposition, - normalize_stream_suffix_for_hint, resolve_playback_format_hint, sniff_stream_format_extension, + normalize_stream_suffix_for_hint, sniff_stream_format_extension, same_playback_target, STREAM_FORMAT_SNIFF_PROBE_BYTES, }; @@ -419,381 +418,3 @@ pub(crate) fn url_format_hint(url: &str) -> Option { }) .map(|s| s.to_lowercase()) } - -/// Arguments forwarded from `audio_play` into the source-build pipeline. -/// Bundles the format-hint inputs, playback-shaping parameters and the shared -/// done flag so that `build_playback_source_with_probe_fallback` stays below -/// the `clippy::too_many_arguments` threshold. -pub(crate) struct BuildSourceArgs<'a> { - pub url: &'a str, - pub gen: u64, - pub cache_id_for_tasks: Option<&'a str>, - pub server_id: Option<&'a str>, - pub url_format_hint: Option<&'a str>, - pub stream_format_suffix: Option<&'a str>, - pub done_flag: Arc, - pub fade_in_dur: Duration, - pub hi_res_enabled: bool, - pub duration_hint: f64, -} - -/// Output of `build_source_from_play_input`: the wrapped rodio source plus -/// whether the chosen source path is seekable (only the Streaming variant -/// is not). -pub(crate) struct PlaybackSource { - pub(crate) built: BuiltSource, - pub(crate) is_seekable: bool, -} - -fn play_media_format_hint(input: &PlayInput) -> Option { - match input { - PlayInput::SeekableMedia { format_hint, .. } | PlayInput::Streaming { format_hint, .. } => { - format_hint.clone() - } - PlayInput::Bytes(_) => None, - } -} - -/// Ranged HTTP probe/decode failed in a way that may succeed after the -/// background download finishes (moov-at-end, demuxer EOF during partial buffer). -fn is_ranged_stream_probe_failure(err: &str) -> bool { - err.contains("ranged-stream") - && (err.contains("format probe failed") - || err.contains("moov metadata") - || err.contains("end of stream")) -} - -/// Completed ranged download or spill file for `url`, if ready. -async fn try_take_completed_stream_bytes( - url: &str, - state: &State<'_, AudioEngine>, -) -> Option> { - if let Some(data) = super::helpers::take_stream_completed_for_url(state, url) { - return Some(data); - } - let spill_path = { - let guard = state.stream_completed_spill.lock().unwrap(); - guard - .as_ref() - .filter(|p| same_playback_target(&p.url, url)) - .map(|p| p.path.clone()) - }; - if let Some(path) = spill_path { - let data = tokio::fs::read(&path).await.ok()?; - if !data.is_empty() { - return Some(data); - } - } - None -} - -/// Ranged assembly can be byte-complete but missing `moov` (holes) or non-audio HTTP body. -async fn prefer_clean_http_bytes_for_fallback( - url: &str, - gen: u64, - state: &State<'_, AudioEngine>, - app: &AppHandle, - ranged_data: Vec, - format_hint: Option<&str>, - label: &str, -) -> Result>, String> { - let is_mp4 = super::stream::container_hint_is_mp4(format_hint); - if is_mp4 { - super::stream::log_isobmff_buffer_diagnostic(&ranged_data, format_hint, label); - if !super::stream::isobmff_buffer_looks_complete(&ranged_data) - || super::stream::mp4_suspect_zero_holes(&ranged_data) - { - crate::app_deprintln!( - "[stream] ranged buffer looks incomplete or holey — refetching via sequential HTTP" - ); - if let Some(fresh) = fetch_data(url, state, gen, app).await? { - if super::stream::isobmff_buffer_looks_complete(&fresh) { - return Ok(Some(fresh)); - } - super::stream::log_isobmff_buffer_diagnostic(&fresh, format_hint, "http-refetch"); - } - } - } - Ok(Some(ranged_data)) -} - -/// Wait for the in-flight ranged download to finish, then HTTP-fetch if needed. -pub(super) async fn wait_or_fetch_bytes_for_stream_fallback( - url: &str, - gen: u64, - state: &State<'_, AudioEngine>, - app: &AppHandle, - format_hint: Option<&str>, -) -> Result>, String> { - use std::time::{Duration, Instant}; - - let deadline = Instant::now() + Duration::from_secs(TRACK_READ_TIMEOUT_SECS); - loop { - if state.generation.load(Ordering::SeqCst) != gen { - return Ok(None); - } - if let Some(data) = try_take_completed_stream_bytes(url, state).await { - crate::app_deprintln!( - "[stream] full-buffer fallback: using completed download ({} KiB)", - data.len() / 1024 - ); - return prefer_clean_http_bytes_for_fallback( - url, - gen, - state, - app, - data, - format_hint, - "ranged-cache", - ) - .await; - } - if Instant::now() >= deadline { - break; - } - tokio::time::sleep(Duration::from_millis(50)).await; - } - crate::app_deprintln!( - "[stream] full-buffer fallback: download still in progress after {}s — HTTP fetch", - TRACK_READ_TIMEOUT_SECS - ); - fetch_data(url, state, gen, app).await -} - -fn is_in_memory_probe_failure(err: &str) -> bool { - err.contains("format probe failed") - || err.contains("could not open audio stream") - || err.contains("no playable audio track") -} - -/// Like [`build_source_from_play_input`], but on ranged-stream probe failure waits -/// for a full download (or fetches it) and retries from in-memory bytes. -pub(crate) async fn build_playback_source_with_probe_fallback( - play_input: PlayInput, - args: BuildSourceArgs<'_>, - state: &State<'_, AudioEngine>, - app: &AppHandle, -) -> Result { - let BuildSourceArgs { - url, - gen, - cache_id_for_tasks, - server_id, - url_format_hint, - stream_format_suffix, - done_flag, - fade_in_dur, - hi_res_enabled, - duration_hint, - } = args; - let media_hint = play_media_format_hint(&play_input); - let effective_hint = resolve_playback_format_hint( - url_format_hint, - stream_format_suffix, - media_hint.as_deref(), - None, - ); - if let Some(ref h) = effective_hint { - crate::app_deprintln!("[stream] playback format hint: {h}"); - } - - match build_source_from_play_input( - play_input, - state, - effective_hint.as_deref(), - done_flag.clone(), - fade_in_dur, - hi_res_enabled, - duration_hint, - ) - .await - { - Ok(p) => Ok(p), - Err(e) if is_ranged_stream_probe_failure(&e) => { - crate::app_deprintln!( - "[stream] ranged-stream probe failed — trying full-buffer fallback: {}", - e - ); - let data = match wait_or_fetch_bytes_for_stream_fallback( - url, - gen, - state, - app, - effective_hint.as_deref(), - ) - .await? - { - Some(d) => d, - None => return Err(e), - }; - if state.generation.load(Ordering::SeqCst) != gen { - return Err("ranged-stream: superseded during full-buffer fallback".into()); - } - let bytes_hint = resolve_playback_format_hint( - url_format_hint, - stream_format_suffix, - media_hint.as_deref(), - Some(&data), - ); - if bytes_hint.as_ref() != effective_hint.as_ref() { - crate::app_deprintln!( - "[stream] full-buffer fallback: resolved hint {:?} (was {:?})", - bytes_hint, - effective_hint - ); - } - if let Some(track_id) = cache_id_for_tasks - .map(str::trim) - .filter(|s| !s.is_empty()) - { - let (sid, high) = - prepare_playback_analysis(app, state, server_id, track_id, None); - spawn_track_analysis_bytes( - app.clone(), - TrackAnalysisOrigin::StreamDownloadComplete, - sid, - track_id.to_string(), - data.clone(), - high, - Some((gen, state.generation.clone())), - ); - } - match build_source_from_play_input( - PlayInput::Bytes(data.clone()), - state, - bytes_hint.as_deref(), - done_flag.clone(), - fade_in_dur, - hi_res_enabled, - duration_hint, - ) - .await - { - Ok(p) => Ok(p), - Err(pe) if is_in_memory_probe_failure(&pe) => { - if super::stream::container_hint_is_mp4(bytes_hint.as_deref()) { - super::stream::log_isobmff_buffer_diagnostic( - &data, - bytes_hint.as_deref(), - "ranged-cache-probe-fail", - ); - } - crate::app_deprintln!( - "[stream] in-memory probe failed — sequential HTTP refetch: {}", - pe - ); - let fresh = match fetch_data(url, state, gen, app).await? { - Some(d) => d, - None => return Err(pe), - }; - if super::stream::container_hint_is_mp4(bytes_hint.as_deref()) { - super::stream::log_isobmff_buffer_diagnostic( - &fresh, - bytes_hint.as_deref(), - "http-refetch-after-probe-fail", - ); - } - build_source_from_play_input( - PlayInput::Bytes(fresh), - state, - bytes_hint.as_deref(), - done_flag, - fade_in_dur, - hi_res_enabled, - duration_hint, - ) - .await - } - Err(pe) => Err(pe), - } - } - Err(e) => Err(e), - } -} - -/// Dispatch [`PlayInput`] → fully wrapped rodio source. For Bytes the full -/// in-memory pipeline (incl. iTunSMPB scan); for SeekableMedia / Streaming -/// the streaming variant runs the decoder build on a blocking thread. -pub(super) async fn build_source_from_play_input( - play_input: PlayInput, - state: &State<'_, AudioEngine>, - format_hint: Option<&str>, - done_flag: Arc, - fade_in_dur: Duration, - hi_res_enabled: bool, - duration_hint: f64, -) -> Result { - // Always 0 — no application-level resampling. Rodio handles conversion to - // the output device rate internally; we let every track play at its native rate. - let target_rate: u32 = 0; - let mut is_seekable = true; - let built = match play_input { - PlayInput::Bytes(data) => build_source( - data, - duration_hint, - state.eq_gains.clone(), - state.eq_enabled.clone(), - state.eq_pre_gain.clone(), - state.playback_rate.clone(), - done_flag, - fade_in_dur, - state.samples_played.clone(), - target_rate, - format_hint, - hi_res_enabled, - ), - PlayInput::SeekableMedia { - reader, - format_hint: media_hint, - tag, - mp4_probe_gate, - } => { - if let Some(gate) = mp4_probe_gate.as_ref() { - super::stream::wait_for_ranged_mp4_probe_ready(gate).await?; - if gate.gen_arc.load(Ordering::SeqCst) != gate.gen { - return Err("ranged-stream: superseded before moov metadata ready".into()); - } - } - let decoder = tokio::task::spawn_blocking(move || { - SizedDecoder::new_streaming(reader, media_hint.as_deref(), tag) - }) - .await - .map_err(|e| e.to_string())??; - build_streaming_source( - decoder, - duration_hint, - state.eq_gains.clone(), - state.eq_enabled.clone(), - state.eq_pre_gain.clone(), - state.playback_rate.clone(), - done_flag, - fade_in_dur, - state.samples_played.clone(), - target_rate, - None, - ) - } - PlayInput::Streaming { reader, format_hint: stream_hint } => { - is_seekable = false; - let decoder = tokio::task::spawn_blocking(move || { - SizedDecoder::new_streaming(Box::new(reader), stream_hint.as_deref(), "track-stream") - }) - .await - .map_err(|e| e.to_string())??; - build_streaming_source( - decoder, - duration_hint, - state.eq_gains.clone(), - state.eq_enabled.clone(), - state.eq_pre_gain.clone(), - state.playback_rate.clone(), - done_flag, - fade_in_dur, - state.samples_played.clone(), - target_rate, - Some(state.stream_playback_armed.clone()), - ) - } - }?; - Ok(PlaybackSource { built, is_seekable }) -} diff --git a/src-tauri/crates/psysonic-audio/src/source_build.rs b/src-tauri/crates/psysonic-audio/src/source_build.rs new file mode 100644 index 00000000..0f6d7f0a --- /dev/null +++ b/src-tauri/crates/psysonic-audio/src/source_build.rs @@ -0,0 +1,398 @@ +//! Source-building pipeline for `audio_play`: turn a resolved [`PlayInput`] +//! into a fully wrapped rodio source, including the ranged-stream probe +//! fallback (wait for / fetch a full download and retry from in-memory bytes +//! when a partial ranged buffer can't be probed yet). Split out of +//! `play_input.rs` so source *selection* stays separate from source *building*. + +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; +use std::time::Duration; + +use tauri::{AppHandle, State}; + +use super::analysis_dispatch::{ + prepare_playback_analysis, spawn_track_analysis_bytes, TrackAnalysisOrigin, +}; +use super::decode::{build_source, build_streaming_source, BuiltSource, SizedDecoder}; +use super::engine::AudioEngine; +use super::helpers::{fetch_data, resolve_playback_format_hint, same_playback_target}; +use super::play_input::PlayInput; +use super::stream::TRACK_READ_TIMEOUT_SECS; + +/// Arguments forwarded from `audio_play` into the source-build pipeline. +/// Bundles the format-hint inputs, playback-shaping parameters and the shared +/// done flag so that `build_playback_source_with_probe_fallback` stays below +/// the `clippy::too_many_arguments` threshold. +pub(crate) struct BuildSourceArgs<'a> { + pub url: &'a str, + pub gen: u64, + pub cache_id_for_tasks: Option<&'a str>, + pub server_id: Option<&'a str>, + pub url_format_hint: Option<&'a str>, + pub stream_format_suffix: Option<&'a str>, + pub done_flag: Arc, + pub fade_in_dur: Duration, + pub hi_res_enabled: bool, + pub duration_hint: f64, +} + +/// Output of `build_source_from_play_input`: the wrapped rodio source plus +/// whether the chosen source path is seekable (only the Streaming variant +/// is not). +pub(crate) struct PlaybackSource { + pub(crate) built: BuiltSource, + pub(crate) is_seekable: bool, +} + +fn play_media_format_hint(input: &PlayInput) -> Option { + match input { + PlayInput::SeekableMedia { format_hint, .. } | PlayInput::Streaming { format_hint, .. } => { + format_hint.clone() + } + PlayInput::Bytes(_) => None, + } +} + +/// Ranged HTTP probe/decode failed in a way that may succeed after the +/// background download finishes (moov-at-end, demuxer EOF during partial buffer). +fn is_ranged_stream_probe_failure(err: &str) -> bool { + err.contains("ranged-stream") + && (err.contains("format probe failed") + || err.contains("moov metadata") + || err.contains("end of stream")) +} + +/// Completed ranged download or spill file for `url`, if ready. +async fn try_take_completed_stream_bytes( + url: &str, + state: &State<'_, AudioEngine>, +) -> Option> { + if let Some(data) = super::helpers::take_stream_completed_for_url(state, url) { + return Some(data); + } + let spill_path = { + let guard = state.stream_completed_spill.lock().unwrap(); + guard + .as_ref() + .filter(|p| same_playback_target(&p.url, url)) + .map(|p| p.path.clone()) + }; + if let Some(path) = spill_path { + let data = tokio::fs::read(&path).await.ok()?; + if !data.is_empty() { + return Some(data); + } + } + None +} + +/// Ranged assembly can be byte-complete but missing `moov` (holes) or non-audio HTTP body. +async fn prefer_clean_http_bytes_for_fallback( + url: &str, + gen: u64, + state: &State<'_, AudioEngine>, + app: &AppHandle, + ranged_data: Vec, + format_hint: Option<&str>, + label: &str, +) -> Result>, String> { + let is_mp4 = super::stream::container_hint_is_mp4(format_hint); + if is_mp4 { + super::stream::log_isobmff_buffer_diagnostic(&ranged_data, format_hint, label); + if !super::stream::isobmff_buffer_looks_complete(&ranged_data) + || super::stream::mp4_suspect_zero_holes(&ranged_data) + { + crate::app_deprintln!( + "[stream] ranged buffer looks incomplete or holey — refetching via sequential HTTP" + ); + if let Some(fresh) = fetch_data(url, state, gen, app).await? { + if super::stream::isobmff_buffer_looks_complete(&fresh) { + return Ok(Some(fresh)); + } + super::stream::log_isobmff_buffer_diagnostic(&fresh, format_hint, "http-refetch"); + } + } + } + Ok(Some(ranged_data)) +} + +/// Wait for the in-flight ranged download to finish, then HTTP-fetch if needed. +async fn wait_or_fetch_bytes_for_stream_fallback( + url: &str, + gen: u64, + state: &State<'_, AudioEngine>, + app: &AppHandle, + format_hint: Option<&str>, +) -> Result>, String> { + use std::time::Instant; + + let deadline = Instant::now() + Duration::from_secs(TRACK_READ_TIMEOUT_SECS); + loop { + if state.generation.load(Ordering::SeqCst) != gen { + return Ok(None); + } + if let Some(data) = try_take_completed_stream_bytes(url, state).await { + crate::app_deprintln!( + "[stream] full-buffer fallback: using completed download ({} KiB)", + data.len() / 1024 + ); + return prefer_clean_http_bytes_for_fallback( + url, + gen, + state, + app, + data, + format_hint, + "ranged-cache", + ) + .await; + } + if Instant::now() >= deadline { + break; + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + crate::app_deprintln!( + "[stream] full-buffer fallback: download still in progress after {}s — HTTP fetch", + TRACK_READ_TIMEOUT_SECS + ); + fetch_data(url, state, gen, app).await +} + +fn is_in_memory_probe_failure(err: &str) -> bool { + err.contains("format probe failed") + || err.contains("could not open audio stream") + || err.contains("no playable audio track") +} + +/// Like [`build_source_from_play_input`], but on ranged-stream probe failure waits +/// for a full download (or fetches it) and retries from in-memory bytes. +pub(crate) async fn build_playback_source_with_probe_fallback( + play_input: PlayInput, + args: BuildSourceArgs<'_>, + state: &State<'_, AudioEngine>, + app: &AppHandle, +) -> Result { + let BuildSourceArgs { + url, + gen, + cache_id_for_tasks, + server_id, + url_format_hint, + stream_format_suffix, + done_flag, + fade_in_dur, + hi_res_enabled, + duration_hint, + } = args; + let media_hint = play_media_format_hint(&play_input); + let effective_hint = resolve_playback_format_hint( + url_format_hint, + stream_format_suffix, + media_hint.as_deref(), + None, + ); + if let Some(ref h) = effective_hint { + crate::app_deprintln!("[stream] playback format hint: {h}"); + } + + match build_source_from_play_input( + play_input, + state, + effective_hint.as_deref(), + done_flag.clone(), + fade_in_dur, + hi_res_enabled, + duration_hint, + ) + .await + { + Ok(p) => Ok(p), + Err(e) if is_ranged_stream_probe_failure(&e) => { + crate::app_deprintln!( + "[stream] ranged-stream probe failed — trying full-buffer fallback: {}", + e + ); + let data = match wait_or_fetch_bytes_for_stream_fallback( + url, + gen, + state, + app, + effective_hint.as_deref(), + ) + .await? + { + Some(d) => d, + None => return Err(e), + }; + if state.generation.load(Ordering::SeqCst) != gen { + return Err("ranged-stream: superseded during full-buffer fallback".into()); + } + let bytes_hint = resolve_playback_format_hint( + url_format_hint, + stream_format_suffix, + media_hint.as_deref(), + Some(&data), + ); + if bytes_hint.as_ref() != effective_hint.as_ref() { + crate::app_deprintln!( + "[stream] full-buffer fallback: resolved hint {:?} (was {:?})", + bytes_hint, + effective_hint + ); + } + if let Some(track_id) = cache_id_for_tasks + .map(str::trim) + .filter(|s| !s.is_empty()) + { + let (sid, high) = + prepare_playback_analysis(app, state, server_id, track_id, None); + spawn_track_analysis_bytes( + app.clone(), + TrackAnalysisOrigin::StreamDownloadComplete, + sid, + track_id.to_string(), + data.clone(), + high, + Some((gen, state.generation.clone())), + ); + } + match build_source_from_play_input( + PlayInput::Bytes(data.clone()), + state, + bytes_hint.as_deref(), + done_flag.clone(), + fade_in_dur, + hi_res_enabled, + duration_hint, + ) + .await + { + Ok(p) => Ok(p), + Err(pe) if is_in_memory_probe_failure(&pe) => { + if super::stream::container_hint_is_mp4(bytes_hint.as_deref()) { + super::stream::log_isobmff_buffer_diagnostic( + &data, + bytes_hint.as_deref(), + "ranged-cache-probe-fail", + ); + } + crate::app_deprintln!( + "[stream] in-memory probe failed — sequential HTTP refetch: {}", + pe + ); + let fresh = match fetch_data(url, state, gen, app).await? { + Some(d) => d, + None => return Err(pe), + }; + if super::stream::container_hint_is_mp4(bytes_hint.as_deref()) { + super::stream::log_isobmff_buffer_diagnostic( + &fresh, + bytes_hint.as_deref(), + "http-refetch-after-probe-fail", + ); + } + build_source_from_play_input( + PlayInput::Bytes(fresh), + state, + bytes_hint.as_deref(), + done_flag, + fade_in_dur, + hi_res_enabled, + duration_hint, + ) + .await + } + Err(pe) => Err(pe), + } + } + Err(e) => Err(e), + } +} + +/// Dispatch [`PlayInput`] → fully wrapped rodio source. For Bytes the full +/// in-memory pipeline (incl. iTunSMPB scan); for SeekableMedia / Streaming +/// the streaming variant runs the decoder build on a blocking thread. +async fn build_source_from_play_input( + play_input: PlayInput, + state: &State<'_, AudioEngine>, + format_hint: Option<&str>, + done_flag: Arc, + fade_in_dur: Duration, + hi_res_enabled: bool, + duration_hint: f64, +) -> Result { + // Always 0 — no application-level resampling. Rodio handles conversion to + // the output device rate internally; we let every track play at its native rate. + let target_rate: u32 = 0; + let mut is_seekable = true; + let built = match play_input { + PlayInput::Bytes(data) => build_source( + data, + duration_hint, + state.eq_gains.clone(), + state.eq_enabled.clone(), + state.eq_pre_gain.clone(), + state.playback_rate.clone(), + done_flag, + fade_in_dur, + state.samples_played.clone(), + target_rate, + format_hint, + hi_res_enabled, + ), + PlayInput::SeekableMedia { + reader, + format_hint: media_hint, + tag, + mp4_probe_gate, + } => { + if let Some(gate) = mp4_probe_gate.as_ref() { + super::stream::wait_for_ranged_mp4_probe_ready(gate).await?; + if gate.gen_arc.load(Ordering::SeqCst) != gate.gen { + return Err("ranged-stream: superseded before moov metadata ready".into()); + } + } + let decoder = tokio::task::spawn_blocking(move || { + SizedDecoder::new_streaming(reader, media_hint.as_deref(), tag) + }) + .await + .map_err(|e| e.to_string())??; + build_streaming_source( + decoder, + duration_hint, + state.eq_gains.clone(), + state.eq_enabled.clone(), + state.eq_pre_gain.clone(), + state.playback_rate.clone(), + done_flag, + fade_in_dur, + state.samples_played.clone(), + target_rate, + None, + ) + } + PlayInput::Streaming { reader, format_hint: stream_hint } => { + is_seekable = false; + let decoder = tokio::task::spawn_blocking(move || { + SizedDecoder::new_streaming(Box::new(reader), stream_hint.as_deref(), "track-stream") + }) + .await + .map_err(|e| e.to_string())??; + build_streaming_source( + decoder, + duration_hint, + state.eq_gains.clone(), + state.eq_enabled.clone(), + state.eq_pre_gain.clone(), + state.playback_rate.clone(), + done_flag, + fade_in_dur, + state.samples_played.clone(), + target_rate, + Some(state.stream_playback_armed.clone()), + ) + } + }?; + Ok(PlaybackSource { built, is_seekable }) +}