mirror of
https://github.com/kilyabin/psysonic.git
synced 2026-07-22 14:35:41 +00:00
refactor(audio): extract source selection from audio_play
Pull the ~300 LOC URL → PlayInput dispatch out of audio_play into a new audio/play_input.rs: - PlayInput enum (Bytes / SeekableMedia / Streaming) - PlayInputContext struct holding the precomputed inputs (url, gen, duration_hint, format/cache hints, optional reused chained bytes) - select_play_input() async — handles all four branches: reused chained bytes, psysonic-local://, ranged HTTP with format sniff, and the legacy non-seekable streaming fallback - url_format_hint() — the conservative URL→extension allowlist helper, used to be inline in audio_play audio_play is now focused on the orchestration above the source layer: ghost-command guard, gapless pre-chain detection, bumping generation, loudness/replay-gain math, fade-in setup, building the actual rodio Source pipeline, sink swap with crossfade, spawning progress task. audio/commands.rs: 980 → 676 LOC. play_input.rs: 385 LOC. Behaviour preservation hinges on identical analysis-seed spawning, format-hint fallback chain, range-detection logic, and gen-cancel checks — all moved verbatim. Smoke test focus: psysonic-local playback, manual skip during ranged HTTP, format-hint-less servers (legacy fallback path), preload cache hit replay, manual skip onto pre-chained track.
This commit is contained in:
+20
-324
@@ -1,30 +1,23 @@
|
|||||||
//! Tauri commands: audio_play / chain_preload / preload + the shared
|
//! Tauri commands: audio_play / chain_preload / preload + the shared
|
||||||
//! spawn_progress_task helper. Transport (pause/resume/stop/seek), device,
|
//! spawn_progress_task helper. Transport (pause/resume/stop/seek), device,
|
||||||
//! radio, mix-mode and AutoEQ commands live in sibling modules.
|
//! radio, mix-mode and AutoEQ commands live in sibling modules.
|
||||||
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
|
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
||||||
use std::sync::{Arc, Mutex};
|
use std::sync::Arc;
|
||||||
use std::time::{Duration, Instant};
|
use std::time::{Duration, Instant};
|
||||||
|
|
||||||
use futures_util::StreamExt;
|
use futures_util::StreamExt;
|
||||||
use ringbuf::{HeapCons, HeapRb};
|
|
||||||
use ringbuf::traits::Split;
|
|
||||||
use rodio::Player;
|
use rodio::Player;
|
||||||
use rodio::Source;
|
use rodio::Source;
|
||||||
use symphonia::core::io::MediaSource;
|
use tauri::{AppHandle, Emitter, State};
|
||||||
use tauri::{AppHandle, Emitter, Manager, State};
|
|
||||||
|
|
||||||
use super::decode::{build_source, build_streaming_source, SizedDecoder};
|
use super::decode::{build_source, build_streaming_source, SizedDecoder};
|
||||||
use super::engine::{audio_http_client, AudioEngine};
|
use super::engine::{audio_http_client, AudioEngine};
|
||||||
use super::helpers::*;
|
use super::helpers::*;
|
||||||
use super::ipc::{maybe_emit_normalization_state, NormalizationStatePayload};
|
use super::ipc::{maybe_emit_normalization_state, NormalizationStatePayload};
|
||||||
|
use super::play_input::{select_play_input, url_format_hint, PlayInput, PlayInputContext};
|
||||||
use super::preview::preview_clear_for_new_main_playback;
|
use super::preview::preview_clear_for_new_main_playback;
|
||||||
use super::progress_task::spawn_progress_task;
|
use super::progress_task::spawn_progress_task;
|
||||||
use super::state::{ChainedInfo, PreloadedTrack};
|
use super::state::{ChainedInfo, PreloadedTrack};
|
||||||
use super::stream::{
|
|
||||||
ranged_download_task, track_download_task, AudioStreamReader,
|
|
||||||
LocalFileSource, RangedHttpSource, RADIO_READ_TIMEOUT_SECS,
|
|
||||||
TRACK_STREAM_MAX_BUF_CAPACITY, TRACK_STREAM_MIN_BUF_CAPACITY, LOCAL_FILE_PLAYBACK_SEED_MAX_BYTES,
|
|
||||||
};
|
|
||||||
|
|
||||||
// ─── Commands ─────────────────────────────────────────────────────────────────
|
// ─── Commands ─────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
@@ -135,323 +128,26 @@ pub async fn audio_play(
|
|||||||
*state.current_analysis_track_id.lock().unwrap() = logical_trim.clone();
|
*state.current_analysis_track_id.lock().unwrap() = logical_trim.clone();
|
||||||
let cache_id_for_tasks = analysis_cache_track_id(logical_trim.as_deref(), &url);
|
let cache_id_for_tasks = analysis_cache_track_id(logical_trim.as_deref(), &url);
|
||||||
|
|
||||||
// Extract format hint from URL for better symphonia probing. Strip the
|
let format_hint = url_format_hint(&url);
|
||||||
// query string first so Subsonic-style URLs (`stream.view?...&v=1.16.1&...`)
|
|
||||||
// don't latch onto random query-param substrings; only accept short
|
|
||||||
// alphanumeric tails that look like an actual audio extension.
|
|
||||||
let format_hint = url
|
|
||||||
.split('?').next()
|
|
||||||
.and_then(|path| path.rsplit('.').next())
|
|
||||||
.filter(|ext| {
|
|
||||||
(1..=5).contains(&ext.len())
|
|
||||||
&& ext.chars().all(|c| c.is_ascii_alphanumeric())
|
|
||||||
&& matches!(
|
|
||||||
ext.to_ascii_lowercase().as_str(),
|
|
||||||
"mp3" | "flac" | "ogg" | "oga" | "opus" | "m4a" | "mp4"
|
|
||||||
| "aac" | "wav" | "wave" | "ape" | "wv" | "webm" | "mka"
|
|
||||||
)
|
|
||||||
})
|
|
||||||
.map(|s| s.to_lowercase());
|
|
||||||
|
|
||||||
enum PlayInput {
|
let play_input = match select_play_input(
|
||||||
Bytes(Vec<u8>),
|
PlayInputContext {
|
||||||
/// Seekable on-demand source — `RangedHttpSource` for HTTP streams,
|
url: &url,
|
||||||
/// `LocalFileSource` for `psysonic-local://` files. Goes through
|
|
||||||
/// `build_streaming_source` (no iTunSMPB scan, since we don't have the
|
|
||||||
/// bytes in memory; chained-track gapless trim still applies via the
|
|
||||||
/// re-played `Bytes` path on the next start).
|
|
||||||
SeekableMedia {
|
|
||||||
reader: Box<dyn MediaSource>,
|
|
||||||
format_hint: Option<String>,
|
|
||||||
tag: &'static str,
|
|
||||||
},
|
|
||||||
Streaming {
|
|
||||||
reader: AudioStreamReader,
|
|
||||||
format_hint: Option<String>,
|
|
||||||
},
|
|
||||||
}
|
|
||||||
|
|
||||||
// Data source selection:
|
|
||||||
// 1) Reused chained bytes (manual skip onto pre-chained track)
|
|
||||||
// 2) `psysonic-local://` (offline / hot cache hit) → LocalFileSource (instant)
|
|
||||||
// 3) Manual uncached remote start:
|
|
||||||
// a) Server supports Range + Content-Length → seekable RangedHttpSource
|
|
||||||
// b) Server does not → legacy non-seekable AudioStreamReader fallback
|
|
||||||
// 4) Preloaded/streamed-cache hit → in-memory bytes via fetch_data
|
|
||||||
let play_input = if let Some(d) = reuse_chained_bytes {
|
|
||||||
spawn_analysis_seed_from_in_memory_bytes(
|
|
||||||
&app,
|
|
||||||
cache_id_for_tasks.as_deref(),
|
|
||||||
gen,
|
gen,
|
||||||
&state.generation,
|
duration_hint,
|
||||||
&d,
|
stream_format_suffix: stream_format_suffix.as_deref(),
|
||||||
);
|
format_hint: format_hint.as_deref(),
|
||||||
PlayInput::Bytes(d)
|
cache_id_for_tasks: cache_id_for_tasks.as_deref(),
|
||||||
} else {
|
reuse_chained_bytes,
|
||||||
let stream_cache_hit = {
|
},
|
||||||
let streamed = state.stream_completed_cache.lock().unwrap();
|
&state,
|
||||||
streamed
|
&app,
|
||||||
.as_ref()
|
).await? {
|
||||||
.is_some_and(|p| same_playback_target(&p.url, &url))
|
Some(input) => input,
|
||||||
};
|
None => return Ok(()), // superseded — bail
|
||||||
let preloaded_hit = {
|
|
||||||
let preloaded = state.preloaded.lock().unwrap();
|
|
||||||
preloaded
|
|
||||||
.as_ref()
|
|
||||||
.is_some_and(|p| same_playback_target(&p.url, &url))
|
|
||||||
};
|
|
||||||
let is_local = url.starts_with("psysonic-local://");
|
|
||||||
|
|
||||||
if is_local && !stream_cache_hit && !preloaded_hit {
|
|
||||||
let path = url.strip_prefix("psysonic-local://").unwrap();
|
|
||||||
let file = std::fs::File::open(path).map_err(|e| e.to_string())?;
|
|
||||||
let len = file.metadata().map(|m| m.len()).unwrap_or(0);
|
|
||||||
let local_hint = std::path::Path::new(path)
|
|
||||||
.extension()
|
|
||||||
.and_then(|e| e.to_str())
|
|
||||||
.map(|s| s.to_lowercase());
|
|
||||||
crate::app_deprintln!(
|
|
||||||
"[stream] LocalFileSource selected — size={} KB, hint={:?}",
|
|
||||||
len / 1024,
|
|
||||||
local_hint
|
|
||||||
);
|
|
||||||
if let Some(ref seed_id) = cache_id_for_tasks {
|
|
||||||
let skip_cpu_seed = app
|
|
||||||
.try_state::<crate::analysis_cache::AnalysisCache>()
|
|
||||||
.map(|c| c.cpu_seed_redundant_for_track(seed_id).unwrap_or(false))
|
|
||||||
.unwrap_or(false);
|
|
||||||
if !skip_cpu_seed {
|
|
||||||
let path_owned = std::path::PathBuf::from(path);
|
|
||||||
let app_seed = app.clone();
|
|
||||||
let gen_seed = gen;
|
|
||||||
let gen_arc_seed = state.generation.clone();
|
|
||||||
let seed_id = seed_id.clone();
|
|
||||||
tokio::spawn(async move {
|
|
||||||
if gen_arc_seed.load(Ordering::SeqCst) != gen_seed {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
let data = match tokio::fs::read(&path_owned).await {
|
|
||||||
Ok(d) => d,
|
|
||||||
Err(_) => return,
|
|
||||||
};
|
|
||||||
if gen_arc_seed.load(Ordering::SeqCst) != gen_seed {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
if data.is_empty() || data.len() > LOCAL_FILE_PLAYBACK_SEED_MAX_BYTES {
|
|
||||||
crate::app_deprintln!(
|
|
||||||
"[stream] psysonic-local: skip analysis seed track_id={} bytes={} (over {} MiB cap)",
|
|
||||||
seed_id,
|
|
||||||
data.len(),
|
|
||||||
LOCAL_FILE_PLAYBACK_SEED_MAX_BYTES / (1024 * 1024)
|
|
||||||
);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
crate::app_deprintln!(
|
|
||||||
"[stream] psysonic-local: file read complete track_id={} size_mib={:.2} — full-track analysis (cpu-seed queue)",
|
|
||||||
seed_id,
|
|
||||||
data.len() as f64 / (1024.0 * 1024.0)
|
|
||||||
);
|
|
||||||
let high = crate::audio::engine::analysis_seed_high_priority_for_track(&app_seed, &seed_id);
|
|
||||||
if let Err(e) =
|
|
||||||
crate::submit_analysis_cpu_seed(app_seed.clone(), seed_id.clone(), data, high).await
|
|
||||||
{
|
|
||||||
crate::app_eprintln!(
|
|
||||||
"[analysis] local-file seed failed for {}: {}",
|
|
||||||
seed_id,
|
|
||||||
e
|
|
||||||
);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
let reader = LocalFileSource { file, len };
|
|
||||||
PlayInput::SeekableMedia {
|
|
||||||
reader: Box::new(reader),
|
|
||||||
format_hint: local_hint,
|
|
||||||
tag: "local-file",
|
|
||||||
}
|
|
||||||
} else if !stream_cache_hit && !preloaded_hit && !is_local {
|
|
||||||
// `manual` must NOT gate this branch: natural track end calls `next(false)`,
|
|
||||||
// so auto-advance must still use ranged HTTP when available. Gating used to
|
|
||||||
// force every auto-start through `fetch_data` (full RAM buffer + extra fetch log)
|
|
||||||
// instead of `RangedHttpSource`.
|
|
||||||
let response = audio_http_client(&state).get(&url).send().await.map_err(|e| e.to_string())?;
|
|
||||||
if !response.status().is_success() {
|
|
||||||
if state.generation.load(Ordering::SeqCst) != gen {
|
|
||||||
return Ok(()); // superseded
|
|
||||||
}
|
|
||||||
let status = response.status().as_u16();
|
|
||||||
let msg = format!("HTTP {status}");
|
|
||||||
app.emit("audio:error", &msg).ok();
|
|
||||||
return Err(msg);
|
|
||||||
}
|
|
||||||
|
|
||||||
let mut stream_hint = content_type_to_hint(
|
|
||||||
response
|
|
||||||
.headers()
|
|
||||||
.get(reqwest::header::CONTENT_TYPE)
|
|
||||||
.and_then(|v| v.to_str().ok())
|
|
||||||
.unwrap_or(""),
|
|
||||||
)
|
|
||||||
.or_else(|| {
|
|
||||||
response
|
|
||||||
.headers()
|
|
||||||
.get(reqwest::header::CONTENT_DISPOSITION)
|
|
||||||
.and_then(|v| v.to_str().ok())
|
|
||||||
.and_then(|cd| format_hint_from_content_disposition(cd))
|
|
||||||
})
|
|
||||||
.or_else(|| normalize_stream_suffix_for_hint(stream_format_suffix.as_deref()))
|
|
||||||
.or_else(|| format_hint.clone());
|
|
||||||
|
|
||||||
let supports_range = response.headers()
|
|
||||||
.get(reqwest::header::ACCEPT_RANGES)
|
|
||||||
.and_then(|v| v.to_str().ok())
|
|
||||||
.is_some_and(|v| v.to_ascii_lowercase().contains("bytes"));
|
|
||||||
let total_size = response.content_length();
|
|
||||||
|
|
||||||
if stream_hint.is_none() && supports_range {
|
|
||||||
if let Some(total_u64) = total_size.filter(|&t| t > 0) {
|
|
||||||
let last = total_u64
|
|
||||||
.saturating_sub(1)
|
|
||||||
.min((STREAM_FORMAT_SNIFF_PROBE_BYTES - 1) as u64);
|
|
||||||
if let Ok(pr) = audio_http_client(&state)
|
|
||||||
.get(&url)
|
|
||||||
.header(reqwest::header::RANGE, format!("bytes=0-{last}"))
|
|
||||||
.send()
|
|
||||||
.await
|
|
||||||
{
|
|
||||||
let stat = pr.status();
|
|
||||||
let ok = stat == reqwest::StatusCode::PARTIAL_CONTENT
|
|
||||||
|| stat == reqwest::StatusCode::OK;
|
|
||||||
if ok {
|
|
||||||
match pr.bytes().await {
|
|
||||||
Ok(bytes) if !bytes.is_empty() => {
|
|
||||||
stream_hint = sniff_stream_format_extension(&bytes).or(stream_hint);
|
|
||||||
if stream_hint.is_some() {
|
|
||||||
crate::app_deprintln!(
|
|
||||||
"[stream] ranged: format sniff from {} B prefix → hint={:?}",
|
|
||||||
bytes.len(),
|
|
||||||
stream_hint
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
_ => {}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Guardrail: when format/container hint is unknown, some demuxers may
|
|
||||||
// seek near EOF during probe. With a progressively downloaded ranged
|
|
||||||
// source that can delay first audible samples until most/all bytes are
|
|
||||||
// fetched. Prefer sequential streaming in that case for faster start.
|
|
||||||
if let (true, Some(total), true) = (supports_range, total_size, stream_hint.is_some()) {
|
|
||||||
let total_usize = total as usize;
|
|
||||||
crate::app_deprintln!(
|
|
||||||
"[stream] RangedHttpSource selected — total={} KB, hint={:?}",
|
|
||||||
total_usize / 1024,
|
|
||||||
stream_hint
|
|
||||||
);
|
|
||||||
let buf = Arc::new(Mutex::new(vec![0u8; total_usize]));
|
|
||||||
let downloaded_to = Arc::new(AtomicUsize::new(0));
|
|
||||||
let done = Arc::new(AtomicBool::new(false));
|
|
||||||
let loudness_hold_for_defer = (total_usize <= crate::audio::stream::TRACK_STREAM_PROMOTE_MAX_BYTES)
|
|
||||||
.then_some(state.ranged_loudness_seed_hold.clone());
|
|
||||||
tokio::spawn(ranged_download_task(
|
|
||||||
gen,
|
|
||||||
state.generation.clone(),
|
|
||||||
audio_http_client(&state),
|
|
||||||
app.clone(),
|
|
||||||
duration_hint,
|
|
||||||
url.clone(),
|
|
||||||
response,
|
|
||||||
buf.clone(),
|
|
||||||
downloaded_to.clone(),
|
|
||||||
done.clone(),
|
|
||||||
state.stream_completed_cache.clone(),
|
|
||||||
state.normalization_engine.clone(),
|
|
||||||
state.normalization_target_lufs.clone(),
|
|
||||||
state.loudness_pre_analysis_attenuation_db.clone(),
|
|
||||||
cache_id_for_tasks.clone(),
|
|
||||||
loudness_hold_for_defer,
|
|
||||||
));
|
|
||||||
let reader = RangedHttpSource {
|
|
||||||
buf,
|
|
||||||
downloaded_to,
|
|
||||||
total_size: total,
|
|
||||||
pos: 0,
|
|
||||||
done,
|
|
||||||
gen_arc: state.generation.clone(),
|
|
||||||
gen,
|
|
||||||
};
|
|
||||||
PlayInput::SeekableMedia {
|
|
||||||
reader: Box::new(reader),
|
|
||||||
format_hint: stream_hint,
|
|
||||||
tag: "ranged-stream",
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
crate::app_deprintln!(
|
|
||||||
"[stream] legacy AudioStreamReader (non-seekable) — accept-ranges={}, content-length={:?}, hint={:?}",
|
|
||||||
supports_range, total_size, stream_hint
|
|
||||||
);
|
|
||||||
let buffer_cap = total_size
|
|
||||||
.map(|n| n as usize)
|
|
||||||
.unwrap_or(TRACK_STREAM_MIN_BUF_CAPACITY)
|
|
||||||
.clamp(TRACK_STREAM_MIN_BUF_CAPACITY, TRACK_STREAM_MAX_BUF_CAPACITY);
|
|
||||||
let rb = HeapRb::<u8>::new(buffer_cap);
|
|
||||||
let (prod, cons) = rb.split();
|
|
||||||
let done = Arc::new(AtomicBool::new(false));
|
|
||||||
tokio::spawn(track_download_task(
|
|
||||||
gen,
|
|
||||||
state.generation.clone(),
|
|
||||||
audio_http_client(&state),
|
|
||||||
app.clone(),
|
|
||||||
url.clone(),
|
|
||||||
response,
|
|
||||||
prod,
|
|
||||||
done.clone(),
|
|
||||||
state.stream_completed_cache.clone(),
|
|
||||||
state.normalization_engine.clone(),
|
|
||||||
state.normalization_target_lufs.clone(),
|
|
||||||
state.loudness_pre_analysis_attenuation_db.clone(),
|
|
||||||
cache_id_for_tasks.clone(),
|
|
||||||
));
|
|
||||||
|
|
||||||
let (_new_cons_tx, new_cons_rx) = std::sync::mpsc::channel::<HeapCons<u8>>();
|
|
||||||
let reader = AudioStreamReader {
|
|
||||||
cons: Mutex::new(cons),
|
|
||||||
new_cons_rx: Mutex::new(new_cons_rx),
|
|
||||||
deadline: std::time::Instant::now()
|
|
||||||
+ Duration::from_secs(RADIO_READ_TIMEOUT_SECS),
|
|
||||||
gen_arc: state.generation.clone(),
|
|
||||||
gen,
|
|
||||||
source_tag: "track-stream",
|
|
||||||
eof_when_empty: Some(done),
|
|
||||||
pos: 0,
|
|
||||||
};
|
|
||||||
PlayInput::Streaming {
|
|
||||||
reader,
|
|
||||||
format_hint: stream_hint,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
let data = fetch_data(&url, &state, gen, &app).await?;
|
|
||||||
let data = match data {
|
|
||||||
Some(d) => d,
|
|
||||||
None => return Ok(()), // superseded while downloading
|
|
||||||
};
|
|
||||||
spawn_analysis_seed_from_in_memory_bytes(
|
|
||||||
&app,
|
|
||||||
cache_id_for_tasks.as_deref(),
|
|
||||||
gen,
|
|
||||||
&state.generation,
|
|
||||||
&data,
|
|
||||||
);
|
|
||||||
PlayInput::Bytes(data)
|
|
||||||
}
|
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|
||||||
if state.generation.load(Ordering::SeqCst) != gen {
|
if state.generation.load(Ordering::SeqCst) != gen {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ mod decode;
|
|||||||
mod dev_io;
|
mod dev_io;
|
||||||
pub mod device_commands;
|
pub mod device_commands;
|
||||||
pub mod mix_commands;
|
pub mod mix_commands;
|
||||||
|
mod play_input;
|
||||||
pub mod preload_commands;
|
pub mod preload_commands;
|
||||||
pub(crate) mod progress_task;
|
pub(crate) mod progress_task;
|
||||||
pub mod radio_commands;
|
pub mod radio_commands;
|
||||||
|
|||||||
@@ -0,0 +1,385 @@
|
|||||||
|
//! Source-selection logic for `audio_play`: given a URL + various caches +
|
||||||
|
//! Subsonic hints, decide whether to play from in-memory bytes, a seekable
|
||||||
|
//! local file, a seekable RangedHttpSource, or a non-seekable streaming reader.
|
||||||
|
|
||||||
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||||
|
use std::sync::{Arc, Mutex};
|
||||||
|
use std::time::Duration;
|
||||||
|
|
||||||
|
use ringbuf::traits::Split;
|
||||||
|
use ringbuf::{HeapCons, HeapRb};
|
||||||
|
use symphonia::core::io::MediaSource;
|
||||||
|
use tauri::{AppHandle, Emitter, Manager, State};
|
||||||
|
|
||||||
|
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, sniff_stream_format_extension,
|
||||||
|
spawn_analysis_seed_from_in_memory_bytes, same_playback_target,
|
||||||
|
STREAM_FORMAT_SNIFF_PROBE_BYTES,
|
||||||
|
};
|
||||||
|
use super::stream::{
|
||||||
|
ranged_download_task, track_download_task, AudioStreamReader,
|
||||||
|
LocalFileSource, RangedHttpSource, LOCAL_FILE_PLAYBACK_SEED_MAX_BYTES,
|
||||||
|
RADIO_READ_TIMEOUT_SECS, TRACK_STREAM_MAX_BUF_CAPACITY, TRACK_STREAM_MIN_BUF_CAPACITY,
|
||||||
|
};
|
||||||
|
|
||||||
|
/// What `audio_play` will hand to `build_source` / `build_streaming_source`.
|
||||||
|
pub(super) enum PlayInput {
|
||||||
|
Bytes(Vec<u8>),
|
||||||
|
/// Seekable on-demand source — `RangedHttpSource` for HTTP streams,
|
||||||
|
/// `LocalFileSource` for `psysonic-local://` files. Goes through
|
||||||
|
/// `build_streaming_source` (no iTunSMPB scan, since we don't have the
|
||||||
|
/// bytes in memory; chained-track gapless trim still applies via the
|
||||||
|
/// re-played `Bytes` path on the next start).
|
||||||
|
SeekableMedia {
|
||||||
|
reader: Box<dyn MediaSource>,
|
||||||
|
format_hint: Option<String>,
|
||||||
|
tag: &'static str,
|
||||||
|
},
|
||||||
|
Streaming {
|
||||||
|
reader: AudioStreamReader,
|
||||||
|
format_hint: Option<String>,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Inputs `audio_play` has already computed before source selection.
|
||||||
|
pub(super) struct PlayInputContext<'a> {
|
||||||
|
pub url: &'a str,
|
||||||
|
pub gen: u64,
|
||||||
|
pub duration_hint: f64,
|
||||||
|
pub stream_format_suffix: Option<&'a str>,
|
||||||
|
pub format_hint: Option<&'a str>,
|
||||||
|
pub cache_id_for_tasks: Option<&'a str>,
|
||||||
|
/// `Some(bytes)` when manual-skip onto a pre-chained track reuses bytes
|
||||||
|
/// from the chained-info block.
|
||||||
|
pub reuse_chained_bytes: Option<Vec<u8>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Resolves the play input for `audio_play` honouring (in priority order):
|
||||||
|
/// 1. Reused chained bytes — manual skip onto pre-chained track.
|
||||||
|
/// 2. `psysonic-local://` files — open as seekable LocalFileSource.
|
||||||
|
/// 3. Remote HTTP without preload/stream-cache hit — try ranged HTTP, fall
|
||||||
|
/// back to non-seekable AudioStreamReader.
|
||||||
|
/// 4. Preload/stream-cache hit — replay in-memory bytes via `fetch_data`.
|
||||||
|
///
|
||||||
|
/// Returns `Ok(None)` when the operation was superseded by a later
|
||||||
|
/// `audio_play` call (generation bump) — caller should bail out silently.
|
||||||
|
pub(super) async fn select_play_input(
|
||||||
|
ctx: PlayInputContext<'_>,
|
||||||
|
state: &State<'_, AudioEngine>,
|
||||||
|
app: &AppHandle,
|
||||||
|
) -> Result<Option<PlayInput>, String> {
|
||||||
|
if let Some(d) = ctx.reuse_chained_bytes {
|
||||||
|
spawn_analysis_seed_from_in_memory_bytes(
|
||||||
|
app,
|
||||||
|
ctx.cache_id_for_tasks,
|
||||||
|
ctx.gen,
|
||||||
|
&state.generation,
|
||||||
|
&d,
|
||||||
|
);
|
||||||
|
return Ok(Some(PlayInput::Bytes(d)));
|
||||||
|
}
|
||||||
|
|
||||||
|
let stream_cache_hit = {
|
||||||
|
let streamed = state.stream_completed_cache.lock().unwrap();
|
||||||
|
streamed
|
||||||
|
.as_ref()
|
||||||
|
.is_some_and(|p| same_playback_target(&p.url, ctx.url))
|
||||||
|
};
|
||||||
|
let preloaded_hit = {
|
||||||
|
let preloaded = state.preloaded.lock().unwrap();
|
||||||
|
preloaded
|
||||||
|
.as_ref()
|
||||||
|
.is_some_and(|p| same_playback_target(&p.url, ctx.url))
|
||||||
|
};
|
||||||
|
let is_local = ctx.url.starts_with("psysonic-local://");
|
||||||
|
|
||||||
|
if is_local && !stream_cache_hit && !preloaded_hit {
|
||||||
|
return Ok(Some(open_local_file_input(&ctx, state, app)?));
|
||||||
|
}
|
||||||
|
if !stream_cache_hit && !preloaded_hit && !is_local {
|
||||||
|
return open_ranged_or_streaming_input(&ctx, state, app).await;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Preloaded or stream-cache hit → replay in-memory bytes.
|
||||||
|
let data = match fetch_data(ctx.url, state, ctx.gen, app).await? {
|
||||||
|
Some(d) => d,
|
||||||
|
None => return Ok(None), // superseded while downloading
|
||||||
|
};
|
||||||
|
spawn_analysis_seed_from_in_memory_bytes(
|
||||||
|
app,
|
||||||
|
ctx.cache_id_for_tasks,
|
||||||
|
ctx.gen,
|
||||||
|
&state.generation,
|
||||||
|
&data,
|
||||||
|
);
|
||||||
|
Ok(Some(PlayInput::Bytes(data)))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// `psysonic-local://<path>` → seekable `LocalFileSource`. Spawns a
|
||||||
|
/// background CPU-seed for the analysis cache when the file is small
|
||||||
|
/// enough (skipped if the cache already has a row for this track).
|
||||||
|
fn open_local_file_input(
|
||||||
|
ctx: &PlayInputContext<'_>,
|
||||||
|
state: &State<'_, AudioEngine>,
|
||||||
|
app: &AppHandle,
|
||||||
|
) -> Result<PlayInput, String> {
|
||||||
|
let path = ctx.url.strip_prefix("psysonic-local://").unwrap();
|
||||||
|
let file = std::fs::File::open(path).map_err(|e| e.to_string())?;
|
||||||
|
let len = file.metadata().map(|m| m.len()).unwrap_or(0);
|
||||||
|
let local_hint = std::path::Path::new(path)
|
||||||
|
.extension()
|
||||||
|
.and_then(|e| e.to_str())
|
||||||
|
.map(|s| s.to_lowercase());
|
||||||
|
crate::app_deprintln!(
|
||||||
|
"[stream] LocalFileSource selected — size={} KB, hint={:?}",
|
||||||
|
len / 1024,
|
||||||
|
local_hint
|
||||||
|
);
|
||||||
|
if let Some(seed_id) = ctx.cache_id_for_tasks {
|
||||||
|
let skip_cpu_seed = app
|
||||||
|
.try_state::<crate::analysis_cache::AnalysisCache>()
|
||||||
|
.map(|c| c.cpu_seed_redundant_for_track(seed_id).unwrap_or(false))
|
||||||
|
.unwrap_or(false);
|
||||||
|
if !skip_cpu_seed {
|
||||||
|
let path_owned = std::path::PathBuf::from(path);
|
||||||
|
let app_seed = app.clone();
|
||||||
|
let gen_seed = ctx.gen;
|
||||||
|
let gen_arc_seed = state.generation.clone();
|
||||||
|
let seed_id = seed_id.to_string();
|
||||||
|
tokio::spawn(async move {
|
||||||
|
if gen_arc_seed.load(Ordering::SeqCst) != gen_seed {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let data = match tokio::fs::read(&path_owned).await {
|
||||||
|
Ok(d) => d,
|
||||||
|
Err(_) => return,
|
||||||
|
};
|
||||||
|
if gen_arc_seed.load(Ordering::SeqCst) != gen_seed {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if data.is_empty() || data.len() > LOCAL_FILE_PLAYBACK_SEED_MAX_BYTES {
|
||||||
|
crate::app_deprintln!(
|
||||||
|
"[stream] psysonic-local: skip analysis seed track_id={} bytes={} (over {} MiB cap)",
|
||||||
|
seed_id,
|
||||||
|
data.len(),
|
||||||
|
LOCAL_FILE_PLAYBACK_SEED_MAX_BYTES / (1024 * 1024)
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
crate::app_deprintln!(
|
||||||
|
"[stream] psysonic-local: file read complete track_id={} size_mib={:.2} — full-track analysis (cpu-seed queue)",
|
||||||
|
seed_id,
|
||||||
|
data.len() as f64 / (1024.0 * 1024.0)
|
||||||
|
);
|
||||||
|
let high = crate::audio::engine::analysis_seed_high_priority_for_track(&app_seed, &seed_id);
|
||||||
|
if let Err(e) =
|
||||||
|
crate::submit_analysis_cpu_seed(app_seed.clone(), seed_id.clone(), data, high).await
|
||||||
|
{
|
||||||
|
crate::app_eprintln!(
|
||||||
|
"[analysis] local-file seed failed for {}: {}",
|
||||||
|
seed_id,
|
||||||
|
e
|
||||||
|
);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
let reader = LocalFileSource { file, len };
|
||||||
|
Ok(PlayInput::SeekableMedia {
|
||||||
|
reader: Box::new(reader),
|
||||||
|
format_hint: local_hint,
|
||||||
|
tag: "local-file",
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Manual or auto-advance starts that aren't already cached: try ranged HTTP
|
||||||
|
/// (seekable) first, fall back to a non-seekable `AudioStreamReader` if the
|
||||||
|
/// server doesn't advertise byte-range support or a length.
|
||||||
|
async fn open_ranged_or_streaming_input(
|
||||||
|
ctx: &PlayInputContext<'_>,
|
||||||
|
state: &State<'_, AudioEngine>,
|
||||||
|
app: &AppHandle,
|
||||||
|
) -> Result<Option<PlayInput>, String> {
|
||||||
|
let response = audio_http_client(state).get(ctx.url).send().await.map_err(|e| e.to_string())?;
|
||||||
|
if !response.status().is_success() {
|
||||||
|
if state.generation.load(Ordering::SeqCst) != ctx.gen {
|
||||||
|
return Ok(None); // superseded
|
||||||
|
}
|
||||||
|
let status = response.status().as_u16();
|
||||||
|
let msg = format!("HTTP {status}");
|
||||||
|
app.emit("audio:error", &msg).ok();
|
||||||
|
return Err(msg);
|
||||||
|
}
|
||||||
|
|
||||||
|
let mut stream_hint = content_type_to_hint(
|
||||||
|
response
|
||||||
|
.headers()
|
||||||
|
.get(reqwest::header::CONTENT_TYPE)
|
||||||
|
.and_then(|v| v.to_str().ok())
|
||||||
|
.unwrap_or(""),
|
||||||
|
)
|
||||||
|
.or_else(|| {
|
||||||
|
response
|
||||||
|
.headers()
|
||||||
|
.get(reqwest::header::CONTENT_DISPOSITION)
|
||||||
|
.and_then(|v| v.to_str().ok())
|
||||||
|
.and_then(format_hint_from_content_disposition)
|
||||||
|
})
|
||||||
|
.or_else(|| normalize_stream_suffix_for_hint(ctx.stream_format_suffix))
|
||||||
|
.or_else(|| ctx.format_hint.map(|s| s.to_string()));
|
||||||
|
|
||||||
|
let supports_range = response.headers()
|
||||||
|
.get(reqwest::header::ACCEPT_RANGES)
|
||||||
|
.and_then(|v| v.to_str().ok())
|
||||||
|
.is_some_and(|v| v.to_ascii_lowercase().contains("bytes"));
|
||||||
|
let total_size = response.content_length();
|
||||||
|
|
||||||
|
if stream_hint.is_none() && supports_range {
|
||||||
|
if let Some(total_u64) = total_size.filter(|&t| t > 0) {
|
||||||
|
let last = total_u64
|
||||||
|
.saturating_sub(1)
|
||||||
|
.min((STREAM_FORMAT_SNIFF_PROBE_BYTES - 1) as u64);
|
||||||
|
if let Ok(pr) = audio_http_client(state)
|
||||||
|
.get(ctx.url)
|
||||||
|
.header(reqwest::header::RANGE, format!("bytes=0-{last}"))
|
||||||
|
.send()
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
let stat = pr.status();
|
||||||
|
let ok = stat == reqwest::StatusCode::PARTIAL_CONTENT
|
||||||
|
|| stat == reqwest::StatusCode::OK;
|
||||||
|
if ok {
|
||||||
|
match pr.bytes().await {
|
||||||
|
Ok(bytes) if !bytes.is_empty() => {
|
||||||
|
stream_hint = sniff_stream_format_extension(&bytes).or(stream_hint);
|
||||||
|
if stream_hint.is_some() {
|
||||||
|
crate::app_deprintln!(
|
||||||
|
"[stream] ranged: format sniff from {} B prefix → hint={:?}",
|
||||||
|
bytes.len(),
|
||||||
|
stream_hint
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
_ => {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Guardrail: when format/container hint is unknown, some demuxers may
|
||||||
|
// seek near EOF during probe. With a progressively downloaded ranged
|
||||||
|
// source that can delay first audible samples until most/all bytes are
|
||||||
|
// fetched. Prefer sequential streaming in that case for faster start.
|
||||||
|
if let (true, Some(total), true) = (supports_range, total_size, stream_hint.is_some()) {
|
||||||
|
let total_usize = total as usize;
|
||||||
|
crate::app_deprintln!(
|
||||||
|
"[stream] RangedHttpSource selected — total={} KB, hint={:?}",
|
||||||
|
total_usize / 1024,
|
||||||
|
stream_hint
|
||||||
|
);
|
||||||
|
let buf = Arc::new(Mutex::new(vec![0u8; total_usize]));
|
||||||
|
let downloaded_to = Arc::new(AtomicUsize::new(0));
|
||||||
|
let done = Arc::new(AtomicBool::new(false));
|
||||||
|
let loudness_hold_for_defer = (total_usize <= super::stream::TRACK_STREAM_PROMOTE_MAX_BYTES)
|
||||||
|
.then_some(state.ranged_loudness_seed_hold.clone());
|
||||||
|
tokio::spawn(ranged_download_task(
|
||||||
|
ctx.gen,
|
||||||
|
state.generation.clone(),
|
||||||
|
audio_http_client(state),
|
||||||
|
app.clone(),
|
||||||
|
ctx.duration_hint,
|
||||||
|
ctx.url.to_string(),
|
||||||
|
response,
|
||||||
|
buf.clone(),
|
||||||
|
downloaded_to.clone(),
|
||||||
|
done.clone(),
|
||||||
|
state.stream_completed_cache.clone(),
|
||||||
|
state.normalization_engine.clone(),
|
||||||
|
state.normalization_target_lufs.clone(),
|
||||||
|
state.loudness_pre_analysis_attenuation_db.clone(),
|
||||||
|
ctx.cache_id_for_tasks.map(|s| s.to_string()),
|
||||||
|
loudness_hold_for_defer,
|
||||||
|
));
|
||||||
|
let reader = RangedHttpSource {
|
||||||
|
buf,
|
||||||
|
downloaded_to,
|
||||||
|
total_size: total,
|
||||||
|
pos: 0,
|
||||||
|
done,
|
||||||
|
gen_arc: state.generation.clone(),
|
||||||
|
gen: ctx.gen,
|
||||||
|
};
|
||||||
|
return Ok(Some(PlayInput::SeekableMedia {
|
||||||
|
reader: Box::new(reader),
|
||||||
|
format_hint: stream_hint,
|
||||||
|
tag: "ranged-stream",
|
||||||
|
}));
|
||||||
|
}
|
||||||
|
|
||||||
|
// Legacy non-seekable streaming reader fallback.
|
||||||
|
crate::app_deprintln!(
|
||||||
|
"[stream] legacy AudioStreamReader (non-seekable) — accept-ranges={}, content-length={:?}, hint={:?}",
|
||||||
|
supports_range, total_size, stream_hint
|
||||||
|
);
|
||||||
|
let buffer_cap = total_size
|
||||||
|
.map(|n| n as usize)
|
||||||
|
.unwrap_or(TRACK_STREAM_MIN_BUF_CAPACITY)
|
||||||
|
.clamp(TRACK_STREAM_MIN_BUF_CAPACITY, TRACK_STREAM_MAX_BUF_CAPACITY);
|
||||||
|
let rb = HeapRb::<u8>::new(buffer_cap);
|
||||||
|
let (prod, cons) = rb.split();
|
||||||
|
let done = Arc::new(AtomicBool::new(false));
|
||||||
|
tokio::spawn(track_download_task(
|
||||||
|
ctx.gen,
|
||||||
|
state.generation.clone(),
|
||||||
|
audio_http_client(state),
|
||||||
|
app.clone(),
|
||||||
|
ctx.url.to_string(),
|
||||||
|
response,
|
||||||
|
prod,
|
||||||
|
done.clone(),
|
||||||
|
state.stream_completed_cache.clone(),
|
||||||
|
state.normalization_engine.clone(),
|
||||||
|
state.normalization_target_lufs.clone(),
|
||||||
|
state.loudness_pre_analysis_attenuation_db.clone(),
|
||||||
|
ctx.cache_id_for_tasks.map(|s| s.to_string()),
|
||||||
|
));
|
||||||
|
|
||||||
|
let (_new_cons_tx, new_cons_rx) = std::sync::mpsc::channel::<HeapCons<u8>>();
|
||||||
|
let reader = AudioStreamReader {
|
||||||
|
cons: Mutex::new(cons),
|
||||||
|
new_cons_rx: Mutex::new(new_cons_rx),
|
||||||
|
deadline: std::time::Instant::now()
|
||||||
|
+ Duration::from_secs(RADIO_READ_TIMEOUT_SECS),
|
||||||
|
gen_arc: state.generation.clone(),
|
||||||
|
gen: ctx.gen,
|
||||||
|
source_tag: "track-stream",
|
||||||
|
eof_when_empty: Some(done),
|
||||||
|
pos: 0,
|
||||||
|
};
|
||||||
|
Ok(Some(PlayInput::Streaming {
|
||||||
|
reader,
|
||||||
|
format_hint: stream_hint,
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Pulled out of the format_hint extraction block in `audio_play` — strip the
|
||||||
|
/// query string first so Subsonic-style URLs (`stream.view?...&v=1.16.1&...`)
|
||||||
|
/// don't latch onto random query-param substrings; only accept short
|
||||||
|
/// alphanumeric tails that look like an actual audio extension.
|
||||||
|
pub(super) fn url_format_hint(url: &str) -> Option<String> {
|
||||||
|
url.split('?').next()
|
||||||
|
.and_then(|path| path.rsplit('.').next())
|
||||||
|
.filter(|ext| {
|
||||||
|
(1..=5).contains(&ext.len())
|
||||||
|
&& ext.chars().all(|c| c.is_ascii_alphanumeric())
|
||||||
|
&& matches!(
|
||||||
|
ext.to_ascii_lowercase().as_str(),
|
||||||
|
"mp3" | "flac" | "ogg" | "oga" | "opus" | "m4a" | "mp4"
|
||||||
|
| "aac" | "wav" | "wave" | "ape" | "wv" | "webm" | "mka"
|
||||||
|
)
|
||||||
|
})
|
||||||
|
.map(|s| s.to_lowercase())
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user