mirror of
https://github.com/kilyabin/psysonic.git
synced 2026-07-22 06:25:41 +00:00
feat(audio): seekable streaming via RangedHttpSource + instant local playback
Reworks the manual streaming-start path so playback is seekable from the first frame and hot-cache hits start without a multi-hundred-MB pre-read. Why - The legacy AudioStreamReader was a non-seekable ring buffer. Forward seeks during streaming let Symphonia consume the buffer through to EOF (next song), and backward seeks were silently broken. - fetch_data() loaded `psysonic-local://` files via tokio::fs::read, blocking playback for seconds on hot-cache or offline files (e.g. a 100+ MB hi-res FLAC) and producing ALSA underruns when the new sink swap arrived too late. What changed - New RangedHttpSource MediaSource: pre-allocates a Vec<u8> in track size, background ranged_download_task fills linearly from offset 0 with HTTP Range reconnects. is_seekable=true; reads block on data, seeks update the cursor only. Picked when the server response carries Accept-Ranges: bytes + Content-Length. - New LocalFileSource MediaSource: thin std::fs::File wrapper used for every `psysonic-local://` URL — Symphonia probes the first ~64 KB and starts playback before the rest is read. - audio_play wires both paths through a generic PlayInput::SeekableMedia variant so the build_streaming_source pipeline stays single-shape. - current_is_seekable AtomicBool on AudioEngine: audio_seek short-circuits with `not seekable` for the legacy fallback so the frontend restart- fallback (playerStore seek catch) engages instead of letting a forward seek consume the ring buffer. - Dev-only [stream] / [seek] eprintln logs (cfg(debug_assertions)) so source selection, codec resolution (codec name + lossless flag + bit depth + rate + channels), download progress and seek targets are visible in the terminal during testing. - content_type_to_hint now covers wav/wave/opus; URL-derived format hints strip the query string first and only accept known audio extensions (Subsonic stream URLs have no extension — would otherwise latch onto query-param fragments). Compatibility - Servers without Accept-Ranges fall back to the existing AudioStreamReader path; seek is rejected up-front so the existing restart-fallback runs. - Radio + audio_chain_preload paths untouched. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
+511
-41
@@ -1,6 +1,6 @@
|
||||
use std::io::{Cursor, Read, Seek, SeekFrom};
|
||||
use std::sync::{Arc, Mutex, OnceLock};
|
||||
use std::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering};
|
||||
use std::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, AtomicUsize, Ordering};
|
||||
use std::time::{Duration, Instant};
|
||||
#[cfg(unix)]
|
||||
use libc;
|
||||
@@ -668,6 +668,167 @@ impl MediaSource for AudioStreamReader {
|
||||
fn byte_len(&self) -> Option<u64> { None }
|
||||
}
|
||||
|
||||
// ── RangedHttpSource — seekable HTTP-backed MediaSource ──────────────────────
|
||||
//
|
||||
// Pre-allocates a Vec<u8> of total track size. A background task fills it
|
||||
// linearly from offset 0 via streaming HTTP. Read blocks (with timeout) until
|
||||
// requested bytes are downloaded; Seek only updates the cursor.
|
||||
//
|
||||
// Reports is_seekable=true so Symphonia performs time-based seeks via the
|
||||
// format reader. Backward seeks: instant (data in buffer). Forward seeks
|
||||
// beyond downloaded_to: Read blocks until the linear download catches up.
|
||||
//
|
||||
// Requires server to have responded with both Content-Length and
|
||||
// `Accept-Ranges: bytes` so reconnects can resume via HTTP Range.
|
||||
|
||||
struct RangedHttpSource {
|
||||
/// Pre-allocated buffer of total size. Filled linearly from offset 0.
|
||||
buf: Arc<Mutex<Vec<u8>>>,
|
||||
/// Bytes contiguously downloaded from offset 0.
|
||||
downloaded_to: Arc<AtomicUsize>,
|
||||
total_size: u64,
|
||||
pos: u64,
|
||||
/// Set when the download task terminates (success or hard error).
|
||||
done: Arc<AtomicBool>,
|
||||
gen_arc: Arc<AtomicU64>,
|
||||
gen: u64,
|
||||
}
|
||||
|
||||
impl Read for RangedHttpSource {
|
||||
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
|
||||
if self.gen_arc.load(Ordering::SeqCst) != self.gen {
|
||||
return Ok(0);
|
||||
}
|
||||
if self.pos >= self.total_size {
|
||||
return Ok(0);
|
||||
}
|
||||
let max_read = ((self.total_size - self.pos) as usize).min(buf.len());
|
||||
if max_read == 0 {
|
||||
return Ok(0);
|
||||
}
|
||||
let target_end = self.pos + max_read as u64;
|
||||
|
||||
let deadline = Instant::now() + Duration::from_secs(RADIO_READ_TIMEOUT_SECS);
|
||||
#[cfg(debug_assertions)]
|
||||
let block_started = Instant::now();
|
||||
#[cfg(debug_assertions)]
|
||||
let mut blocked = false;
|
||||
loop {
|
||||
if self.gen_arc.load(Ordering::SeqCst) != self.gen {
|
||||
return Ok(0);
|
||||
}
|
||||
let dl = self.downloaded_to.load(Ordering::SeqCst) as u64;
|
||||
if dl >= target_end {
|
||||
#[cfg(debug_assertions)]
|
||||
if blocked {
|
||||
eprintln!(
|
||||
"[stream] read unblocked after {:.2}s — pos={} target_end={} dl={}",
|
||||
block_started.elapsed().as_secs_f64(),
|
||||
self.pos,
|
||||
target_end,
|
||||
dl
|
||||
);
|
||||
}
|
||||
break;
|
||||
}
|
||||
#[cfg(debug_assertions)]
|
||||
if !blocked {
|
||||
blocked = true;
|
||||
eprintln!(
|
||||
"[stream] read blocking — pos={} need={} dl_to={} (waiting for {} more bytes)",
|
||||
self.pos,
|
||||
target_end,
|
||||
dl,
|
||||
target_end - dl
|
||||
);
|
||||
}
|
||||
// Download finished but our cursor is past downloaded_to (e.g. seek
|
||||
// beyond a partial download that aborted). Return what we have.
|
||||
if self.done.load(Ordering::SeqCst) {
|
||||
if dl > self.pos {
|
||||
let avail = (dl - self.pos) as usize;
|
||||
let src = self.buf.lock().unwrap();
|
||||
let start = self.pos as usize;
|
||||
buf[..avail].copy_from_slice(&src[start..start + avail]);
|
||||
drop(src);
|
||||
self.pos += avail as u64;
|
||||
return Ok(avail);
|
||||
}
|
||||
return Ok(0);
|
||||
}
|
||||
if Instant::now() >= deadline {
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::TimedOut,
|
||||
"ranged-http: no data within timeout",
|
||||
));
|
||||
}
|
||||
std::thread::sleep(Duration::from_millis(RADIO_YIELD_MS));
|
||||
}
|
||||
|
||||
let src = self.buf.lock().unwrap();
|
||||
let start = self.pos as usize;
|
||||
let end = start + max_read;
|
||||
buf[..max_read].copy_from_slice(&src[start..end]);
|
||||
drop(src);
|
||||
self.pos += max_read as u64;
|
||||
Ok(max_read)
|
||||
}
|
||||
}
|
||||
|
||||
impl Seek for RangedHttpSource {
|
||||
fn seek(&mut self, pos: SeekFrom) -> std::io::Result<u64> {
|
||||
let new_pos: i64 = match pos {
|
||||
SeekFrom::Start(p) => p as i64,
|
||||
SeekFrom::Current(p) => self.pos as i64 + p,
|
||||
SeekFrom::End(p) => self.total_size as i64 + p,
|
||||
};
|
||||
if new_pos < 0 {
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidInput,
|
||||
"ranged-http: seek before start",
|
||||
));
|
||||
}
|
||||
self.pos = (new_pos as u64).min(self.total_size);
|
||||
Ok(self.pos)
|
||||
}
|
||||
}
|
||||
|
||||
impl MediaSource for RangedHttpSource {
|
||||
fn is_seekable(&self) -> bool { true }
|
||||
fn byte_len(&self) -> Option<u64> { Some(self.total_size) }
|
||||
}
|
||||
|
||||
// ── LocalFileSource — seekable file-backed MediaSource ───────────────────────
|
||||
//
|
||||
// Wraps `std::fs::File` so the decoder reads on-demand from disk instead of
|
||||
// pre-loading the whole file into a Vec. Used for `psysonic-local://` URLs
|
||||
// (offline library + hot playback cache hits) — gives instant track-start
|
||||
// because Symphonia only needs to read ~64 KB during probe before playback
|
||||
// can begin, vs the previous behaviour of `tokio::fs::read` which blocked
|
||||
// until the entire file (often 100+ MB for hi-res FLAC) was in RAM.
|
||||
|
||||
struct LocalFileSource {
|
||||
file: std::fs::File,
|
||||
len: u64,
|
||||
}
|
||||
|
||||
impl Read for LocalFileSource {
|
||||
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
|
||||
self.file.read(buf)
|
||||
}
|
||||
}
|
||||
|
||||
impl Seek for LocalFileSource {
|
||||
fn seek(&mut self, pos: SeekFrom) -> std::io::Result<u64> {
|
||||
self.file.seek(pos)
|
||||
}
|
||||
}
|
||||
|
||||
impl MediaSource for LocalFileSource {
|
||||
fn is_seekable(&self) -> bool { true }
|
||||
fn byte_len(&self) -> Option<u64> { Some(self.len) }
|
||||
}
|
||||
|
||||
// ── Pause / Reconnect Coordination ───────────────────────────────────────────
|
||||
|
||||
pub(crate) struct RadioSharedFlags {
|
||||
@@ -1000,12 +1161,152 @@ async fn track_download_task(
|
||||
}
|
||||
}
|
||||
|
||||
/// Linear downloader for `RangedHttpSource`: fills the pre-allocated buffer
|
||||
/// from offset 0 to total_size. Reconnects via HTTP Range from the current
|
||||
/// `downloaded` offset on transient errors. On completion (full track) the
|
||||
/// data is promoted to `stream_completed_cache` for fast replay.
|
||||
async fn ranged_download_task(
|
||||
gen: u64,
|
||||
gen_arc: Arc<AtomicU64>,
|
||||
http_client: reqwest::Client,
|
||||
url: String,
|
||||
initial_response: reqwest::Response,
|
||||
buf: Arc<Mutex<Vec<u8>>>,
|
||||
downloaded_to: Arc<AtomicUsize>,
|
||||
done: Arc<AtomicBool>,
|
||||
promote_cache_slot: Arc<Mutex<Option<PreloadedTrack>>>,
|
||||
) {
|
||||
let total_size = buf.lock().unwrap().len();
|
||||
let mut downloaded: usize = 0;
|
||||
let mut reconnects: u32 = 0;
|
||||
let mut next_response: Option<reqwest::Response> = Some(initial_response);
|
||||
#[cfg(debug_assertions)]
|
||||
let dl_started = Instant::now();
|
||||
#[cfg(debug_assertions)]
|
||||
let mut next_progress_mb: usize = 1;
|
||||
|
||||
'outer: loop {
|
||||
let response = if let Some(r) = next_response.take() {
|
||||
r
|
||||
} else {
|
||||
let mut req = http_client.get(&url);
|
||||
if downloaded > 0 {
|
||||
req = req.header(reqwest::header::RANGE, format!("bytes={downloaded}-"));
|
||||
}
|
||||
match req.send().await {
|
||||
Ok(r) => r,
|
||||
Err(err) => {
|
||||
if reconnects >= TRACK_STREAM_MAX_RECONNECTS {
|
||||
eprintln!(
|
||||
"[audio] ranged reconnect failed after {} attempts: {}",
|
||||
reconnects, err
|
||||
);
|
||||
break 'outer;
|
||||
}
|
||||
reconnects += 1;
|
||||
tokio::time::sleep(Duration::from_millis(200)).await;
|
||||
continue 'outer;
|
||||
}
|
||||
}
|
||||
};
|
||||
if downloaded > 0 && response.status() != reqwest::StatusCode::PARTIAL_CONTENT {
|
||||
eprintln!(
|
||||
"[audio] ranged reconnect returned {}, expected 206",
|
||||
response.status()
|
||||
);
|
||||
break 'outer;
|
||||
}
|
||||
if downloaded == 0 && !response.status().is_success() {
|
||||
eprintln!("[audio] ranged HTTP {}", response.status());
|
||||
break 'outer;
|
||||
}
|
||||
|
||||
let mut byte_stream = response.bytes_stream();
|
||||
while let Some(chunk) = byte_stream.next().await {
|
||||
if gen_arc.load(Ordering::SeqCst) != gen {
|
||||
done.store(true, Ordering::SeqCst);
|
||||
return;
|
||||
}
|
||||
let chunk = match chunk {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
if reconnects >= TRACK_STREAM_MAX_RECONNECTS {
|
||||
eprintln!(
|
||||
"[audio] ranged dl error after {} reconnects: {}",
|
||||
reconnects, e
|
||||
);
|
||||
break 'outer;
|
||||
}
|
||||
reconnects += 1;
|
||||
eprintln!(
|
||||
"[audio] ranged dl error (attempt {}/{}): {} — reconnecting",
|
||||
reconnects, TRACK_STREAM_MAX_RECONNECTS, e
|
||||
);
|
||||
next_response = None;
|
||||
continue 'outer;
|
||||
}
|
||||
};
|
||||
reconnects = 0;
|
||||
let writable = total_size.saturating_sub(downloaded);
|
||||
if writable == 0 {
|
||||
break;
|
||||
}
|
||||
let n = chunk.len().min(writable);
|
||||
{
|
||||
let mut b = buf.lock().unwrap();
|
||||
b[downloaded..downloaded + n].copy_from_slice(&chunk[..n]);
|
||||
}
|
||||
downloaded += n;
|
||||
downloaded_to.store(downloaded, Ordering::SeqCst);
|
||||
#[cfg(debug_assertions)]
|
||||
{
|
||||
let mb = downloaded / (1024 * 1024);
|
||||
if mb >= next_progress_mb {
|
||||
let pct = (downloaded as f64 / total_size as f64 * 100.0) as u32;
|
||||
eprintln!(
|
||||
"[stream] dl progress: {} MB / {} MB ({}%)",
|
||||
mb,
|
||||
total_size / (1024 * 1024),
|
||||
pct
|
||||
);
|
||||
next_progress_mb = mb + 1;
|
||||
}
|
||||
}
|
||||
if downloaded >= total_size {
|
||||
break;
|
||||
}
|
||||
}
|
||||
// Stream ended cleanly (or hit total_size).
|
||||
break 'outer;
|
||||
}
|
||||
|
||||
done.store(true, Ordering::SeqCst);
|
||||
|
||||
#[cfg(debug_assertions)]
|
||||
eprintln!(
|
||||
"[stream] dl done: {} / {} bytes in {:.2}s ({} reconnects)",
|
||||
downloaded,
|
||||
total_size,
|
||||
dl_started.elapsed().as_secs_f64(),
|
||||
reconnects
|
||||
);
|
||||
|
||||
if downloaded == total_size && total_size > 0 && total_size <= TRACK_STREAM_PROMOTE_MAX_BYTES {
|
||||
let data = buf.lock().unwrap().clone();
|
||||
*promote_cache_slot.lock().unwrap() = Some(PreloadedTrack { url, data });
|
||||
#[cfg(debug_assertions)]
|
||||
eprintln!("[stream] promoted to stream_completed_cache for replay");
|
||||
}
|
||||
}
|
||||
|
||||
fn content_type_to_hint(ct: &str) -> Option<String> {
|
||||
let ct = ct.to_ascii_lowercase();
|
||||
if ct.contains("mpeg") || ct.contains("mp3") { Some("mp3".into()) }
|
||||
else if ct.contains("aac") || ct.contains("aacp") { Some("aac".into()) }
|
||||
else if ct.contains("ogg") { Some("ogg".into()) }
|
||||
else if ct.contains("flac") { Some("flac".into()) }
|
||||
else if ct.contains("wav") || ct.contains("wave") { Some("wav".into()) }
|
||||
else if ct.contains("opus") { Some("opus".into()) }
|
||||
else { None }
|
||||
}
|
||||
|
||||
@@ -1050,6 +1351,38 @@ impl MediaSource for SizedCursorSource {
|
||||
// Implements Iterator<Item = i16> + Source — identical interface to
|
||||
// rodio::Decoder, so the rest of the source chain is unchanged.
|
||||
|
||||
/// Dev-build only: format codec parameters into a human-readable line for
|
||||
/// the terminal so you can verify whether playback is genuinely lossless.
|
||||
#[cfg(debug_assertions)]
|
||||
fn log_codec_resolution(
|
||||
tag: &str,
|
||||
params: &symphonia::core::codecs::CodecParameters,
|
||||
container_hint: Option<&str>,
|
||||
) {
|
||||
let codec_name = symphonia::default::get_codecs()
|
||||
.get_codec(params.codec)
|
||||
.map(|d| d.short_name)
|
||||
.unwrap_or("?");
|
||||
let rate = params.sample_rate.map(|r| format!("{} Hz", r)).unwrap_or_else(|| "? Hz".into());
|
||||
let bits = params.bits_per_sample
|
||||
.or(params.bits_per_coded_sample)
|
||||
.map(|b| format!("{}-bit", b))
|
||||
.unwrap_or_else(|| "?-bit".into());
|
||||
let ch = params.channels
|
||||
.map(|c| format!("{}ch", c.count()))
|
||||
.unwrap_or_else(|| "?ch".into());
|
||||
let lossless = codec_name.starts_with("pcm")
|
||||
|| matches!(
|
||||
codec_name,
|
||||
"flac" | "alac" | "wavpack" | "monkeys-audio" | "tta" | "shorten"
|
||||
);
|
||||
let kind = if lossless { "LOSSLESS" } else { "lossy" };
|
||||
eprintln!(
|
||||
"[stream] {tag}: codec={codec_name} ({kind}) {bits} {rate} {ch} container={}",
|
||||
container_hint.unwrap_or("?")
|
||||
);
|
||||
}
|
||||
|
||||
/// Max retries for IO/packet-read errors (fatal — network drop, truncated file).
|
||||
const DECODE_MAX_RETRIES: usize = 3;
|
||||
/// Max *consecutive* DecodeErrors before giving up on a file.
|
||||
@@ -1138,6 +1471,9 @@ impl SizedDecoder {
|
||||
.zip(track.codec_params.n_frames)
|
||||
.map(|(base, frames)| base.calc_time(frames));
|
||||
|
||||
#[cfg(debug_assertions)]
|
||||
log_codec_resolution("bytes", &track.codec_params, format_hint);
|
||||
|
||||
let mut decoder = psysonic_codec_registry()
|
||||
.make(&track.codec_params, &DecoderOptions::default())
|
||||
.map_err(|e| {
|
||||
@@ -1222,6 +1558,8 @@ impl SizedDecoder {
|
||||
.find(|t| t.codec_params.codec != CODEC_TYPE_NULL)
|
||||
.ok_or_else(|| format!("{source_tag}: no audio track found"))?;
|
||||
let track_id = track.id;
|
||||
#[cfg(debug_assertions)]
|
||||
log_codec_resolution(source_tag, &track.codec_params, format_hint);
|
||||
// Live streams have no known total frame count → total_duration = None.
|
||||
let total_duration = None;
|
||||
let mut decoder = try_make_radio_decoder(&track.codec_params, &DecoderOptions::default())
|
||||
@@ -1687,6 +2025,11 @@ pub struct AudioEngine {
|
||||
/// Last fully downloaded manual-stream track bytes (same playback identity),
|
||||
/// used to recover seek/replay without waiting for network again.
|
||||
pub(crate) stream_completed_cache: Arc<Mutex<Option<PreloadedTrack>>>,
|
||||
/// True when the currently playing source supports seeking (in-memory bytes
|
||||
/// or `RangedHttpSource`); false for the legacy non-seekable streaming
|
||||
/// fallback (`AudioStreamReader`). `audio_seek` rejects with a "not
|
||||
/// seekable" error when false so the frontend restart-fallback can engage.
|
||||
pub(crate) current_is_seekable: Arc<AtomicBool>,
|
||||
pub crossfade_enabled: Arc<AtomicBool>,
|
||||
pub crossfade_secs: Arc<AtomicU32>,
|
||||
pub fading_out_sink: Arc<Mutex<Option<Sink>>>,
|
||||
@@ -1945,6 +2288,7 @@ pub fn create_engine() -> (AudioEngine, std::thread::JoinHandle<()>) {
|
||||
eq_pre_gain: Arc::new(AtomicU32::new(0f32.to_bits())),
|
||||
preloaded: Arc::new(Mutex::new(None)),
|
||||
stream_completed_cache: Arc::new(Mutex::new(None)),
|
||||
current_is_seekable: Arc::new(AtomicBool::new(true)),
|
||||
crossfade_enabled: Arc::new(AtomicBool::new(false)),
|
||||
crossfade_secs: Arc::new(AtomicU32::new(3.0f32.to_bits())),
|
||||
fading_out_sink: Arc::new(Mutex::new(None)),
|
||||
@@ -2197,13 +2541,36 @@ pub async fn audio_play(
|
||||
old.stop();
|
||||
}
|
||||
|
||||
// Extract format hint from URL for better symphonia probing.
|
||||
let format_hint = url.rsplit('.').next()
|
||||
.and_then(|ext| ext.split('?').next())
|
||||
// Extract format hint from URL for better symphonia probing. 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.
|
||||
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 {
|
||||
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>,
|
||||
@@ -2212,8 +2579,11 @@ pub async fn audio_play(
|
||||
|
||||
// Data source selection:
|
||||
// 1) Reused chained bytes (manual skip onto pre-chained track)
|
||||
// 2) Manual uncached remote start -> stream immediately (no full prefetch)
|
||||
// 3) Existing byte path (preloaded/offline/full HTTP fetch)
|
||||
// 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 {
|
||||
PlayInput::Bytes(d)
|
||||
} else {
|
||||
@@ -2231,7 +2601,27 @@ pub async fn audio_play(
|
||||
};
|
||||
let is_local = url.starts_with("psysonic-local://");
|
||||
|
||||
if manual && !stream_cache_hit && !preloaded_hit && !is_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());
|
||||
#[cfg(debug_assertions)]
|
||||
eprintln!(
|
||||
"[stream] LocalFileSource selected — size={} KB, hint={:?}",
|
||||
len / 1024,
|
||||
local_hint
|
||||
);
|
||||
let reader = LocalFileSource { file, len };
|
||||
PlayInput::SeekableMedia {
|
||||
reader: Box::new(reader),
|
||||
format_hint: local_hint,
|
||||
tag: "local-file",
|
||||
}
|
||||
} else if manual && !stream_cache_hit && !preloaded_hit && !is_local {
|
||||
let response = state.http_client.get(&url).send().await.map_err(|e| e.to_string())?;
|
||||
if !response.status().is_success() {
|
||||
if state.generation.load(Ordering::SeqCst) != gen {
|
||||
@@ -2251,41 +2641,88 @@ pub async fn audio_play(
|
||||
.unwrap_or(""),
|
||||
).or_else(|| format_hint.clone());
|
||||
|
||||
let buffer_cap = response
|
||||
.content_length()
|
||||
.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(),
|
||||
state.http_client.clone(),
|
||||
url.clone(),
|
||||
response,
|
||||
prod,
|
||||
done.clone(),
|
||||
state.stream_completed_cache.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();
|
||||
|
||||
// Track streaming has no reconnect producer; keep an empty channel.
|
||||
let (_new_cons_tx, new_cons_rx) = std::sync::mpsc::channel::<HeapConsumer<u8>>();
|
||||
let reader = AudioStreamReader {
|
||||
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,
|
||||
if let (true, Some(total)) = (supports_range, total_size) {
|
||||
let total_usize = total as usize;
|
||||
#[cfg(debug_assertions)]
|
||||
eprintln!(
|
||||
"[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));
|
||||
tokio::spawn(ranged_download_task(
|
||||
gen,
|
||||
state.generation.clone(),
|
||||
state.http_client.clone(),
|
||||
url.clone(),
|
||||
response,
|
||||
buf.clone(),
|
||||
downloaded_to.clone(),
|
||||
done.clone(),
|
||||
state.stream_completed_cache.clone(),
|
||||
));
|
||||
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 {
|
||||
#[cfg(debug_assertions)]
|
||||
eprintln!(
|
||||
"[stream] legacy AudioStreamReader (non-seekable) — accept-ranges={}, content-length={:?}",
|
||||
supports_range, total_size
|
||||
);
|
||||
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(),
|
||||
state.http_client.clone(),
|
||||
url.clone(),
|
||||
response,
|
||||
prod,
|
||||
done.clone(),
|
||||
state.stream_completed_cache.clone(),
|
||||
));
|
||||
|
||||
let (_new_cons_tx, new_cons_rx) = std::sync::mpsc::channel::<HeapConsumer<u8>>();
|
||||
let reader = AudioStreamReader {
|
||||
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?;
|
||||
@@ -2337,6 +2774,7 @@ pub async fn audio_play(
|
||||
// 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 new_is_seekable = true;
|
||||
let built = match play_input {
|
||||
PlayInput::Bytes(data) => build_source(
|
||||
data,
|
||||
@@ -2351,7 +2789,27 @@ pub async fn audio_play(
|
||||
format_hint.as_deref(),
|
||||
hi_res_enabled,
|
||||
),
|
||||
PlayInput::SeekableMedia { reader, format_hint, tag } => {
|
||||
let decoder = tokio::task::spawn_blocking(move || {
|
||||
SizedDecoder::new_streaming(reader, format_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(),
|
||||
done_flag.clone(),
|
||||
fade_in_dur,
|
||||
state.samples_played.clone(),
|
||||
target_rate,
|
||||
)
|
||||
}
|
||||
PlayInput::Streaming { reader, format_hint } => {
|
||||
new_is_seekable = false;
|
||||
let decoder = tokio::task::spawn_blocking(move || {
|
||||
SizedDecoder::new_streaming(Box::new(reader), format_hint.as_deref(), "track-stream")
|
||||
})
|
||||
@@ -2371,6 +2829,7 @@ pub async fn audio_play(
|
||||
)
|
||||
}
|
||||
}.map_err(|e| { app.emit("audio:error", &e).ok(); e })?;
|
||||
state.current_is_seekable.store(new_is_seekable, Ordering::SeqCst);
|
||||
let source = built.source;
|
||||
let duration_secs = built.duration_secs;
|
||||
let output_rate = built.output_rate;
|
||||
@@ -2969,6 +3428,17 @@ pub fn audio_seek(seconds: f64, state: State<'_, AudioEngine>) -> Result<(), Str
|
||||
}
|
||||
}
|
||||
|
||||
// Reject seek up-front for non-seekable streaming sources so the frontend's
|
||||
// restart-fallback engages instead of rolling the dice on the format reader
|
||||
// (which can consume the ring buffer to EOF for forward seeks → next song).
|
||||
if !state.current_is_seekable.load(Ordering::SeqCst) {
|
||||
#[cfg(debug_assertions)]
|
||||
eprintln!("[seek] rejected → not-seekable source (legacy stream)");
|
||||
return Err("source is not seekable".into());
|
||||
}
|
||||
#[cfg(debug_assertions)]
|
||||
eprintln!("[seek] target={:.2}s", seconds);
|
||||
|
||||
// Seeking back invalidates any pending gapless chain.
|
||||
let cur_pos = {
|
||||
let cur = state.current.lock().unwrap();
|
||||
|
||||
Reference in New Issue
Block a user