Skip to main content

koan_core/audio/
buffer.rs

1use std::fs::File;
2use std::path::{Path, PathBuf};
3use std::sync::Arc;
4use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
5use std::thread;
6
7use symphonia::core::codecs::audio::well_known::{
8    CODEC_ID_AAC, CODEC_ID_ALAC, CODEC_ID_FLAC, CODEC_ID_MP3, CODEC_ID_OPUS, CODEC_ID_PCM_F32LE,
9    CODEC_ID_PCM_S16LE, CODEC_ID_PCM_S24LE, CODEC_ID_PCM_S32LE, CODEC_ID_VORBIS,
10};
11use symphonia::core::codecs::audio::{AudioCodecId, AudioCodecParameters, AudioDecoderOptions};
12use symphonia::core::formats::probe::Hint;
13use symphonia::core::formats::{FormatOptions, FormatReader, SeekMode, SeekTo, Track, TrackType};
14use symphonia::core::io::MediaSourceStream;
15use symphonia::core::meta::MetadataOptions;
16use symphonia::core::units::{Duration, Time, TimeBase, Timestamp};
17use thiserror::Error;
18
19use crate::audio::dsp::{Chain, Setup};
20use crate::audio::opus::OpusBridge;
21use crate::audio::viz::VizBuffer;
22use crate::config::ReplayGainMode;
23use crate::player::state::{QueueItemId, RendererClock};
24
25#[derive(Debug, Error)]
26pub enum DecodeError {
27    #[error("failed to open file: {0}")]
28    Io(#[from] std::io::Error),
29    #[error("no supported audio track found")]
30    NoTrack,
31    #[error("unsupported codec")]
32    UnsupportedCodec,
33    #[error("decode error: {0}")]
34    Decode(String),
35}
36
37/// Info about the decoded audio stream, extracted before decoding starts.
38#[derive(Debug, Clone)]
39pub struct StreamInfo {
40    pub codec: String,
41    pub sample_rate: u32,
42    pub channels: u16,
43    pub bit_depth: Option<u16>,
44    /// Bitrate in kbps. Meaningful for lossy codecs, None for lossless.
45    pub bitrate_kbps: Option<u32>,
46    pub duration_ms: u64,
47}
48
49/// Handle to a running decode thread. Drop to stop it.
50pub struct DecodeHandle {
51    stop: Arc<AtomicBool>,
52    thread: Option<thread::JoinHandle<()>>,
53}
54
55impl DecodeHandle {
56    /// Signal the decode thread to stop without waiting for it to exit.
57    /// Unparks it too: a full ring has it parked for as long as the output
58    /// takes to drain half of it.
59    pub fn signal_stop(&self) {
60        self.stop.store(true, Ordering::Relaxed);
61        if let Some(handle) = &self.thread {
62            handle.thread().unpark();
63        }
64    }
65
66    /// The decode thread, to wake when its ring has room. An output that
67    /// drains the ring faster than it plays, as an encoder for a renderer
68    /// does, would otherwise wait out the decoder's park on a full ring.
69    pub fn thread(&self) -> Option<thread::Thread> {
70        self.thread.as_ref().map(|t| t.thread().clone())
71    }
72
73    /// Create a DecodeHandle with no real thread (for tests only).
74    #[cfg(test)]
75    pub fn new_for_test(stop: Arc<AtomicBool>) -> Self {
76        Self { stop, thread: None }
77    }
78
79    /// Signal the decode thread to stop and wait for it.
80    pub fn stop(&mut self) {
81        self.signal_stop();
82        if let Some(handle) = self.thread.take()
83            && let Err(payload) = handle.join()
84        {
85            let msg = payload
86                .downcast_ref::<String>()
87                .map(|s| s.as_str())
88                .or_else(|| payload.downcast_ref::<&str>().copied())
89                .unwrap_or("unknown");
90            log::error!("decode thread panicked: {}", msg);
91        }
92    }
93}
94
95impl Drop for DecodeHandle {
96    fn drop(&mut self) {
97        self.stop();
98    }
99}
100
101// --- Playback timeline: the source of truth for "what's playing" ---
102
103/// A track boundary in the playback stream. At `sample_offset` cumulative
104/// samples written to the ring buffer, this track starts.
105#[derive(Debug, Clone)]
106pub struct TrackBoundary {
107    pub id: QueueItemId,
108    pub path: PathBuf,
109    pub info: StreamInfo,
110    /// Cumulative interleaved samples written to the ring buffer when this
111    /// track's first sample was pushed. For the first track this is 0
112    /// (or seek_samples if seeking).
113    pub sample_offset: u64,
114    /// Samples of this track's audio written to ring buffer so far.
115    /// Updated as decode progresses. At EOF, equals total decoded samples.
116    pub samples_written: u64,
117    /// The seek offset in samples for this track (non-zero only if user seeked).
118    pub seek_samples: u64,
119    /// The rate the ring buffer's samples play at — the source's, unless DSP
120    /// resampled it to reach an impulse response.
121    pub output_rate: u32,
122}
123
124impl TrackBoundary {
125    /// Where in this track the playhead is when `played` samples have played.
126    fn position_ms(&self, played: u64) -> Option<u64> {
127        let ch = self.info.channels as u64;
128        let rate = self.info.sample_rate as u64;
129        let out = self.output_rate as u64;
130        if ch == 0 || rate == 0 || out == 0 {
131            return None;
132        }
133        // Add seek offset since that's where playback started within the track.
134        let track_samples = played.saturating_sub(self.sample_offset);
135        Some((track_samples / ch) * 1000 / out + (self.seek_samples / ch) * 1000 / rate)
136    }
137}
138
139/// Where the playhead is. A play is its boundary, not its item: an item
140/// repeated gaplessly has a boundary per pass, and each is a play of its own.
141#[derive(Debug, Clone, Copy, PartialEq, Eq)]
142pub struct Playhead {
143    /// The index of the boundary the playhead is past, in this session.
144    pub boundary: usize,
145    pub id: QueueItemId,
146    pub position_ms: u64,
147}
148
149/// Shared timeline that the decode thread writes and the UI reads.
150/// The decode thread appends boundaries; the UI reads them + samples_played
151/// to derive current track and position.
152pub struct PlaybackTimeline {
153    boundaries: parking_lot::RwLock<Vec<TrackBoundary>>,
154    /// Total interleaved samples written to the ring buffer across all tracks.
155    samples_written: AtomicU64,
156    /// Total interleaved samples consumed (played) by the audio engine.
157    /// Written by the audio render callback, read by UI.
158    pub samples_played: Arc<AtomicU64>,
159    /// Told when the decoder queues another track, which is when the moment
160    /// the playhead reaches it becomes known.
161    queued: parking_lot::Mutex<Option<Box<dyn Fn() + Send + Sync>>>,
162    /// Where a renderer playing this session's stream is, in milliseconds
163    /// of the stream. Set, it is the playhead in place of `samples_played`:
164    /// what has been encoded runs seconds ahead of what is heard.
165    clock: parking_lot::Mutex<Option<RendererClock>>,
166}
167
168impl std::fmt::Debug for PlaybackTimeline {
169    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
170        f.debug_struct("PlaybackTimeline")
171            .field("samples_played", &self.samples_played)
172            .finish_non_exhaustive()
173    }
174}
175
176impl PlaybackTimeline {
177    pub fn new() -> Arc<Self> {
178        Arc::new(Self {
179            boundaries: parking_lot::RwLock::new(Vec::new()),
180            samples_written: AtomicU64::new(0),
181            samples_played: Arc::new(AtomicU64::new(0)),
182            queued: parking_lot::Mutex::new(None),
183            clock: parking_lot::Mutex::new(None),
184        })
185    }
186
187    /// Samples played: the engine's count, or the renderer's clock.
188    fn played(&self, bounds: &[TrackBoundary]) -> u64 {
189        let Some(clock) = *self.clock.lock() else {
190            return self.samples_played.load(Ordering::Acquire);
191        };
192        let Some(first) = bounds.first() else {
193            return 0;
194        };
195        let per_second = first.output_rate as u64 * first.info.channels as u64;
196        let at = clock.now_ms() * per_second / 1000;
197        // Rounded down to a whole frame, and never past what was written.
198        (at - at % (first.info.channels as u64).max(1))
199            .min(self.samples_written.load(Ordering::Acquire))
200    }
201
202    /// The renderer's clock, in milliseconds of the stream.
203    pub fn clock(&self) -> Option<RendererClock> {
204        *self.clock.lock()
205    }
206
207    pub fn set_clock(&self, clock: Option<RendererClock>) {
208        *self.clock.lock() = clock;
209    }
210
211    /// The track at `samples` into the session.
212    pub fn track_at(&self, samples: u64) -> Option<QueueItemId> {
213        let bounds = self.boundaries.read();
214        let idx = bounds.partition_point(|b| b.sample_offset <= samples);
215        Some(bounds.get(idx.checked_sub(1)?)?.id)
216    }
217
218    /// How long until the playhead reaches the end of what has been written.
219    pub fn until_end(&self) -> Option<std::time::Duration> {
220        let bounds = self.boundaries.read();
221        let played = self.played(&bounds);
222        let last = bounds.last()?;
223        let per_second = last.output_rate as u64 * last.info.channels as u64;
224        if per_second == 0 {
225            return None;
226        }
227        let left = self
228            .samples_written
229            .load(Ordering::Acquire)
230            .saturating_sub(played);
231        Some(std::time::Duration::from_micros(
232            left.saturating_mul(1_000_000) / per_second,
233        ))
234    }
235
236    /// Call `f` whenever the decoder queues another track.
237    pub fn on_queued(&self, f: impl Fn() + Send + Sync + 'static) {
238        *self.queued.lock() = Some(Box::new(f));
239    }
240
241    /// The play under the playhead and how far into it, without the clones
242    /// `current_playback` makes.
243    pub fn playhead(&self) -> Option<Playhead> {
244        let bounds = self.boundaries.read();
245        let played = self.played(&bounds);
246        let boundary = bounds
247            .partition_point(|b| b.sample_offset <= played)
248            .checked_sub(1)?;
249        let current = bounds.get(boundary)?;
250        Some(Playhead {
251            boundary,
252            id: current.id,
253            position_ms: current.position_ms(played)?,
254        })
255    }
256
257    /// How many boundaries the session has: the index the next one takes.
258    pub fn boundary_count(&self) -> usize {
259        self.boundaries.read().len()
260    }
261
262    /// How far into the play at `boundary` playback has got, stopping at its
263    /// end once the playhead has moved on to the next.
264    pub fn position_in(&self, boundary: usize) -> Option<u64> {
265        let bounds = self.boundaries.read();
266        let played = self.played(&bounds);
267        let at = bounds
268            .get(boundary + 1)
269            .map_or(played, |next| played.min(next.sample_offset));
270        bounds.get(boundary)?.position_ms(at)
271    }
272
273    /// How long until the playhead reaches the track queued after the one
274    /// playing, at the rate this one plays. None with nothing queued.
275    pub fn until_next_track(&self) -> Option<std::time::Duration> {
276        let bounds = self.boundaries.read();
277        let played = self.played(&bounds);
278        let idx = bounds.partition_point(|b| b.sample_offset <= played);
279        let next = bounds.get(idx)?;
280        let current = bounds.get(idx.checked_sub(1)?)?;
281        let per_second = current.output_rate as u64 * current.info.channels as u64;
282        if per_second == 0 {
283            return None;
284        }
285        let left = next.sample_offset - played;
286        Some(std::time::Duration::from_micros(
287            left.saturating_mul(1_000_000) / per_second,
288        ))
289    }
290
291    /// The tracks the decoder has queued after the one under the playhead, in
292    /// order. Committed: the ring cannot be truncated, so only a new session
293    /// takes them back.
294    pub fn queued_after_playhead(&self) -> Vec<QueueItemId> {
295        let bounds = self.boundaries.read();
296        let played = self.played(&bounds);
297        let idx = bounds.partition_point(|b| b.sample_offset <= played);
298        bounds[idx..].iter().map(|b| b.id).collect()
299    }
300
301    /// The decode thread's write access.
302    pub fn writer(&self) -> TimelineWriter<'_> {
303        TimelineWriter { timeline: self }
304    }
305
306    /// Reset for a new playback session.
307    pub fn reset(&self) {
308        self.boundaries.write().clear();
309        self.samples_written.store(0, Ordering::Relaxed);
310        self.samples_played.store(0, Ordering::Relaxed);
311        *self.clock.lock() = None;
312    }
313
314    /// Get a clone of the samples_played Arc for the audio engine.
315    pub fn samples_played_counter(&self) -> Arc<AtomicU64> {
316        self.samples_played.clone()
317    }
318
319    /// Derive current track info and position from the playback head.
320    /// Returns (id, path, stream_info, position_ms).
321    ///
322    /// Acquires the boundaries read lock BEFORE reading `samples_played` so
323    /// channels/sample_rate/boundaries are all from a consistent snapshot.
324    /// Without this ordering, a track transition could update the atomics
325    /// after we read `samples_played` but before we read the boundary list.
326    pub fn current_playback(&self) -> Option<(QueueItemId, PathBuf, StreamInfo, u64)> {
327        // Lock first — ensures we see boundaries consistent with the atomic read.
328        let bounds = self.boundaries.read();
329
330        if bounds.is_empty() {
331            return None;
332        }
333
334        // Read samples_played while holding the lock. This guarantees we
335        // never observe a stale boundary list with a newer samples_played
336        // (or vice versa).
337        let played = self.played(&bounds);
338
339        // Find which track the playback head is in via binary search.
340        // partition_point returns first index where offset > played;
341        // the track we want is one before that.
342        let idx = bounds.partition_point(|b| b.sample_offset <= played);
343        let current = if idx > 0 {
344            &bounds[idx - 1]
345        } else {
346            return None;
347        };
348
349        let position_ms = current.position_ms(played)?;
350
351        Some((
352            current.id,
353            current.path.clone(),
354            current.info.clone(),
355            position_ms,
356        ))
357    }
358}
359
360/// One decode session's write access to the timeline.
361///
362/// Sessions never overlap: the player joins the outgoing decode thread before
363/// it resets the timeline for the next, so nothing here has to tell an old
364/// session's writes from a new one's.
365pub struct TimelineWriter<'a> {
366    timeline: &'a PlaybackTimeline,
367}
368
369impl TimelineWriter<'_> {
370    /// Cumulative samples written to the ring buffer this session.
371    fn samples_written(&self) -> u64 {
372        self.timeline.samples_written.load(Ordering::Relaxed)
373    }
374
375    /// Called by decode thread when starting a new track.
376    fn push_boundary(&self, boundary: TrackBoundary) {
377        self.timeline.boundaries.write().push(boundary);
378        if let Some(queued) = self.timeline.queued.lock().as_ref() {
379            queued();
380        }
381    }
382
383    /// Called by decode thread after pushing samples.
384    fn add_written(&self, count: u64) {
385        let mut bounds = self.timeline.boundaries.write();
386        self.timeline
387            .samples_written
388            .fetch_add(count, Ordering::Relaxed);
389        // Also update the last boundary's samples_written.
390        if let Some(last) = bounds.last_mut() {
391            last.samples_written += count;
392        }
393    }
394}
395
396// ---------------------------------------------------------------------------
397// Source abstraction
398// ---------------------------------------------------------------------------
399
400/// A source entry for the generic decode queue.
401///
402/// Each entry provides an ID, a display path (for logging/timeline),
403/// a format hint, and a factory that constructs a fresh `MediaSourceStream`.
404pub struct SourceEntry {
405    pub id: QueueItemId,
406    /// Path used for logging and `TrackBoundary`. Need not be a real FS path.
407    pub path: PathBuf,
408    /// Format hint for Symphonia (e.g. file extension).
409    pub hint: Hint,
410    /// Factory that creates the `MediaSourceStream`. Called exactly once per track.
411    pub make_mss: Box<dyn FnOnce() -> std::io::Result<MediaSourceStream<'static>> + Send>,
412}
413
414impl SourceEntry {
415    /// Convenience: build a `SourceEntry` from a local file path.
416    pub fn from_file(id: QueueItemId, path: PathBuf) -> Self {
417        let ext = path
418            .extension()
419            .and_then(|e| e.to_str())
420            .unwrap_or("")
421            .to_string();
422        let path_clone = path.clone();
423        let mut hint = Hint::new();
424        if !ext.is_empty() {
425            hint.with_extension(&ext);
426        }
427        Self {
428            id,
429            path,
430            hint,
431            make_mss: Box::new(move || {
432                let file = File::open(&path_clone)?;
433                Ok(MediaSourceStream::new(Box::new(file), Default::default()))
434            }),
435        }
436    }
437}
438
439// ---------------------------------------------------------------------------
440// Probe API
441// ---------------------------------------------------------------------------
442
443/// Probe a `MediaSourceStream` (with hint) and return stream info without decoding.
444pub fn probe_source(mss: MediaSourceStream<'_>, hint: &Hint) -> Result<StreamInfo, DecodeError> {
445    probe_mss(mss, hint)
446}
447
448/// Probe a file and return stream info without decoding.
449pub fn probe_file(path: &Path) -> Result<StreamInfo, DecodeError> {
450    let file_size = std::fs::metadata(path).ok().map(|m| m.len());
451    let file = File::open(path)?;
452    let mss = MediaSourceStream::new(Box::new(file), Default::default());
453    let mut hint = Hint::new();
454    if let Some(ext) = path.extension().and_then(|e| e.to_str()) {
455        hint.with_extension(ext);
456    }
457    let mut info = probe_mss(mss, &hint)?;
458    // For Opus (and other lossy codecs where symphonia couldn't give us a
459    // bitrate), estimate from file size / duration when both are available.
460    if info.bitrate_kbps.is_none()
461        && info.bit_depth.is_none()
462        && let Some(size) = file_size
463        && info.duration_ms > 0
464    {
465        info.bitrate_kbps = Some((size * 8 / info.duration_ms) as u32);
466    }
467    Ok(info)
468}
469
470/// Probe a `MediaSourceStream` with a hint.
471fn probe_mss(mss: MediaSourceStream<'_>, hint: &Hint) -> Result<StreamInfo, DecodeError> {
472    let reader = symphonia::default::get_probe()
473        .probe(
474            hint,
475            mss,
476            FormatOptions::default(),
477            MetadataOptions::default(),
478        )
479        .map_err(|e| match e {
480            // Kept whole rather than flattened to a string: a probe against a
481            // partial file fails by running out of bytes, and the caller has to
482            // tell that apart from a file it cannot make sense of.
483            symphonia::core::errors::Error::IoError(io) => DecodeError::Io(io),
484            other => DecodeError::Decode(other.to_string()),
485        })?;
486
487    let track = reader
488        .default_track(TrackType::Audio)
489        .ok_or(DecodeError::NoTrack)?;
490    let codec_params = track
491        .codec_params
492        .as_ref()
493        .and_then(|p| p.audio())
494        .ok_or(DecodeError::NoTrack)?;
495    let is_opus = codec_params.codec == CODEC_ID_OPUS;
496    // Opus always decodes to 48 kHz regardless of the input sample rate.
497    let sample_rate = if is_opus {
498        48000
499    } else {
500        codec_params.sample_rate.unwrap_or(44100)
501    };
502    let channels = codec_params
503        .channels
504        .as_ref()
505        .map(|c| c.count() as u16)
506        .unwrap_or(2);
507    let bit_depth = if is_opus {
508        None
509    } else {
510        Some(codec_params.bits_per_sample.unwrap_or(16) as u16)
511    };
512    let duration_ms = track_duration_ms(&*reader, track, sample_rate);
513    let codec = codec_name(codec_params.codec);
514
515    // Symphonia doesn't expose bitrate directly. For lossy codecs we can
516    // estimate from bits_per_coded_sample when the demuxer provides it.
517    // Opus estimation from file size is handled in probe_file() where we
518    // have the path; here we only have a MediaSourceStream.
519    let bitrate_kbps = estimate_bitrate_from_codec_params(codec_params);
520
521    Ok(StreamInfo {
522        codec,
523        sample_rate,
524        channels,
525        bit_depth,
526        bitrate_kbps,
527        duration_ms,
528    })
529}
530
531// ---------------------------------------------------------------------------
532// Generic decode API (SourceEntry-based)
533// ---------------------------------------------------------------------------
534
535/// What is done to the samples between the decoder and the ring buffer.
536#[derive(Clone)]
537pub struct Processing {
538    pub rg_mode: ReplayGainMode,
539    pub pre_amp_db: f64,
540    /// The output device's DSP profile. `None` leaves the samples untouched.
541    pub dsp: Option<Arc<Setup>>,
542}
543
544impl Default for Processing {
545    fn default() -> Self {
546        Self {
547            rg_mode: ReplayGainMode::Off,
548            pre_amp_db: 0.0,
549            dsp: None,
550        }
551    }
552}
553
554impl Processing {
555    /// The rate a source at `source` plays at.
556    pub fn output_rate(&self, source: u32) -> u32 {
557        self.dsp.as_ref().map_or(source, |d| d.output_rate(source))
558    }
559}
560
561/// Start decoding from a `SourceEntry` into the ring buffer.
562///
563/// `first`      — the first track's source entry.
564/// `seek_ms`    — if > 0, seek to this position before decoding the first track.
565/// `next_track` — closure returning the next `SourceEntry` for gapless playback.
566///                Called on EOF. Returns None when the playlist is exhausted.
567#[allow(clippy::too_many_arguments)]
568pub fn start_decode<N, F>(
569    first: SourceEntry,
570    producer: rtrb::Producer<f32>,
571    seek_ms: u64,
572    next_track: N,
573    timeline: Arc<PlaybackTimeline>,
574    viz_buffer: Option<Arc<VizBuffer>>,
575    processing: Processing,
576    on_finished: F,
577) -> Result<DecodeHandle, DecodeError>
578where
579    N: Fn() -> Option<SourceEntry> + Send + 'static,
580    F: FnOnce() + Send + 'static,
581{
582    let stop = Arc::new(AtomicBool::new(false));
583    let stop_clone = stop.clone();
584    let thread = thread::Builder::new()
585        .name("koan-decode".into())
586        .spawn(move || {
587            decode_queue_loop(
588                first,
589                producer,
590                &stop_clone,
591                seek_ms,
592                &next_track,
593                &timeline.writer(),
594                viz_buffer.as_deref(),
595                &processing,
596            );
597            // Notify the player that the decode loop finished (playlist
598            // exhausted or error). Only fire if we weren't explicitly stopped
599            // (i.e. this is a natural end, not a seek/skip teardown).
600            if !stop_clone.load(Ordering::Relaxed) {
601                on_finished();
602            }
603        })
604        .map_err(DecodeError::Io)?;
605
606    Ok(DecodeHandle {
607        stop,
608        thread: Some(thread),
609    })
610}
611
612// ---------------------------------------------------------------------------
613// Internal decode loop
614// ---------------------------------------------------------------------------
615
616/// Sources that may fail to open or decode in a row before the session is
617/// abandoned. One bad file must not end the queue, but a queue of nothing but
618/// bad files still has to terminate.
619const MAX_CONSECUTIVE_FAILURES: u32 = 32;
620
621/// Gapless decode loop: decode first entry, then call next_track on EOF.
622///
623/// Every track in a session shares one ring buffer, and therefore the audio
624/// engine configured for it. A track whose PCM format differs from the first
625/// ends the session rather than being written at the wrong format; the player
626/// restarts it on a correctly configured engine.
627///
628/// An unreadable source is skipped, not fatal — the decode head runs up to a
629/// full ring buffer ahead of the DAC, so tearing down here would truncate the
630/// track still being heard as well as dropping the rest of the queue.
631#[allow(clippy::too_many_arguments)]
632fn decode_queue_loop<N>(
633    first: SourceEntry,
634    mut producer: rtrb::Producer<f32>,
635    stop: &AtomicBool,
636    initial_seek_ms: u64,
637    next_track: &N,
638    timeline: &TimelineWriter<'_>,
639    viz_buffer: Option<&VizBuffer>,
640    processing: &Processing,
641) where
642    N: Fn() -> Option<SourceEntry>,
643{
644    // The delay line is indexed against the engine's played counter, which the
645    // player has just reset for this session.
646    if let Some(viz) = viz_buffer {
647        viz.reset();
648    }
649
650    let mut pending = Some(first);
651    let mut seek_ms = initial_seek_ms;
652    let mut format: Option<PcmFormat> = None;
653    let mut failures: u32 = 0;
654    let mut chain: Option<Chain> = None;
655
656    while let Some(entry) = pending.take() {
657        if stop.load(Ordering::Relaxed) {
658            break;
659        }
660
661        let SourceEntry {
662            id,
663            path,
664            hint,
665            make_mss,
666        } = entry;
667
668        let outcome = make_mss().map_err(DecodeError::Io).and_then(|mss| {
669            decode_single(
670                id,
671                &path,
672                &hint,
673                mss,
674                &mut producer,
675                stop,
676                seek_ms,
677                timeline,
678                viz_buffer,
679                processing,
680                &mut chain,
681                format,
682            )
683        });
684
685        match outcome {
686            Ok(Decoded::Complete(decoded_format)) => {
687                format = Some(decoded_format);
688                failures = 0;
689            }
690            Ok(Decoded::FormatMismatch) => break,
691            Err(e) => {
692                if stop.load(Ordering::Relaxed) {
693                    break;
694                }
695                failures += 1;
696                log::error!("skipping {}: {}", path.display(), e);
697                if failures >= MAX_CONSECUTIVE_FAILURES {
698                    log::error!(
699                        "{} sources failed in a row, decode thread giving up",
700                        failures
701                    );
702                    break;
703                }
704            }
705        }
706        // A stop ends the session wherever it lands; looking ahead would peek
707        // the playlist and log a transition that never happens.
708        if stop.load(Ordering::Relaxed) {
709            break;
710        }
711
712        seek_ms = 0;
713        pending = (next_track)();
714        match pending {
715            Some(ref next) => log::info!("gapless transition → {}", next.path.display()),
716            None => log::info!("playlist exhausted, decode thread finishing"),
717        }
718    }
719
720    // What the chain still holds is the end of the last track.
721    if let (Some(chain), Some(format)) = (chain.as_mut(), format)
722        && !stop.load(Ordering::Relaxed)
723    {
724        write_ring(&mut producer, chain.flush(), stop, viz_buffer, format);
725    }
726
727    wait_for_drain(&producer, stop, format);
728}
729
730/// Write `samples` into the ring, blocking while it is full. False once the
731/// session has been stopped.
732///
733/// The viz buffer is fed inside the loop so it receives samples at the rate
734/// the output consumes them, not in packet-sized bursts. Without this, FLAC
735/// packets (~93ms each at 44.1kHz) would update it only ~11 times a second,
736/// making waveform modes visibly choppy.
737fn write_ring(
738    producer: &mut rtrb::Producer<f32>,
739    samples: &[f32],
740    stop: &AtomicBool,
741    viz_buffer: Option<&VizBuffer>,
742    (sample_rate, channels): PcmFormat,
743) -> bool {
744    let mut offset = 0;
745    while offset < samples.len() {
746        if stop.load(Ordering::Relaxed) {
747            return false;
748        }
749
750        let slots = producer.slots();
751        if slots == 0 {
752            // A full ring is the steady state, so this is the wait playback
753            // spends nearly all its time in: until the output has played
754            // half of it, which leaves the other half — a second or more at
755            // any rate koan plays — to refill against. The viz delay line
756            // is read at the playhead, so a refill arriving in one burst
757            // reads the same as one trickling in. A stop unparks it.
758            let half = producer.buffer().capacity() / 2;
759            thread::park_timeout(time_to_play(half, Some((sample_rate, channels))));
760            continue;
761        }
762
763        let chunk_size = slots.min(samples.len() - offset);
764        if let Ok(mut chunk) = producer.write_chunk_uninit(chunk_size) {
765            let to_write = &samples[offset..offset + chunk_size];
766            let (first, second) = chunk.as_mut_slices();
767            let first_len = first.len().min(to_write.len());
768            for (slot, &val) in first.iter_mut().zip(&to_write[..first_len]) {
769                slot.write(val);
770            }
771            if first_len < to_write.len() {
772                for (slot, &val) in second.iter_mut().zip(&to_write[first_len..]) {
773                    slot.write(val);
774                }
775            }
776            // SAFETY: All slots in the chunk have been initialized by the
777            // two loops above — first.len() + second.len() == chunk_size,
778            // and every slot is written via MaybeUninit::write().
779            unsafe { chunk.commit_all() };
780
781            if let Some(viz) = viz_buffer {
782                viz.push_samples(to_write, channels, sample_rate);
783            }
784
785            offset += chunk_size;
786        }
787    }
788    true
789}
790
791/// How long the output takes to play `samples` interleaved samples. Without a
792/// format to go by, a step short enough to poll with.
793fn time_to_play(samples: usize, format: Option<PcmFormat>) -> std::time::Duration {
794    match format {
795        Some((rate, channels)) if rate > 0 && channels > 0 => std::time::Duration::from_micros(
796            samples as u64 * 1_000_000 / (rate as u64 * channels as u64),
797        ),
798        _ => std::time::Duration::from_millis(2),
799    }
800}
801
802/// Block until the audio engine has consumed everything in the ring buffer.
803///
804/// A session ends only once its audio has been heard, so the player can tear
805/// the engine down without clipping the tail of the last track decoded.
806/// Returns early if playback is torn down underneath us.
807fn wait_for_drain(producer: &rtrb::Producer<f32>, stop: &AtomicBool, format: Option<PcmFormat>) {
808    let capacity = producer.buffer().capacity();
809    while !stop.load(Ordering::Relaxed) && !producer.is_abandoned() {
810        let left = capacity.saturating_sub(producer.slots());
811        if left == 0 {
812            return;
813        }
814        thread::park_timeout(time_to_play(left, format));
815    }
816}
817
818// ---------------------------------------------------------------------------
819// Core decode single track
820// ---------------------------------------------------------------------------
821
822/// The PCM format of a decoded stream: sample rate in Hz and channel count.
823/// The audio engine is configured from this, so the ring buffer may only ever
824/// hold samples of one such format at a time.
825type PcmFormat = (u32, u16);
826
827/// Outcome of decoding one source.
828enum Decoded {
829    /// Decoded to EOF, in the given format.
830    Complete(PcmFormat),
831    /// The source's format differs from the stream already in the ring buffer.
832    /// Nothing further was written — the engine must be reconfigured first.
833    FormatMismatch,
834}
835
836/// Decode a single source into the producer. Returns on clean EOF.
837///
838/// `expected` — the format already in the ring buffer, if any. A source that
839/// does not match it is rejected without writing samples or pushing a boundary.
840#[allow(clippy::too_many_arguments)]
841fn decode_single(
842    queue_item_id: QueueItemId,
843    path: &Path,
844    hint: &Hint,
845    mss: MediaSourceStream<'_>,
846    producer: &mut rtrb::Producer<f32>,
847    stop: &AtomicBool,
848    seek_ms: u64,
849    timeline: &TimelineWriter<'_>,
850    viz_buffer: Option<&VizBuffer>,
851    processing: &Processing,
852    chain: &mut Option<Chain>,
853    expected: Option<PcmFormat>,
854) -> Result<Decoded, DecodeError> {
855    let mut reader = symphonia::default::get_probe()
856        .probe(
857            hint,
858            mss,
859            FormatOptions::default(),
860            MetadataOptions::default(),
861        )
862        .map_err(|e| DecodeError::Decode(e.to_string()))?;
863
864    let track = reader
865        .default_track(TrackType::Audio)
866        .ok_or(DecodeError::NoTrack)?;
867    let track_id = track.id;
868    let time_base = track.time_base;
869    let codec_params = track
870        .codec_params
871        .as_ref()
872        .and_then(|p| p.audio())
873        .ok_or(DecodeError::NoTrack)?;
874    let is_opus_codec = codec_params.codec == CODEC_ID_OPUS;
875
876    // Opus always decodes to 48 kHz regardless of the internal rate.
877    let sample_rate = if is_opus_codec {
878        48000
879    } else {
880        codec_params.sample_rate.unwrap_or(44100)
881    };
882    let channels = codec_params
883        .channels
884        .as_ref()
885        .map(|c| c.count() as u16)
886        .unwrap_or(2);
887
888    let duration_ms = track_duration_ms(&*reader, track, sample_rate);
889
890    // Try codec_params first; fall back to file-size estimation for Opus/lossy.
891    let mut bitrate_kbps = estimate_bitrate_from_codec_params(codec_params);
892    if bitrate_kbps.is_none()
893        && is_opus_codec
894        && let Ok(meta) = std::fs::metadata(path)
895        && duration_ms > 0
896    {
897        bitrate_kbps = Some((meta.len() * 8 / duration_ms) as u32);
898    }
899
900    let info = StreamInfo {
901        codec: codec_name(codec_params.codec),
902        sample_rate,
903        channels,
904        bit_depth: if is_opus_codec {
905            None
906        } else {
907            Some(codec_params.bits_per_sample.unwrap_or(16) as u16)
908        },
909        bitrate_kbps,
910        duration_ms,
911    };
912
913    // The ring holds the output format, which DSP may have moved off the
914    // source's: two sources resampled to one impulse response's rate play
915    // gaplessly on one engine.
916    let out_rate = processing.output_rate(sample_rate);
917    let format = (out_rate, channels);
918    if let Some(expected) = expected
919        && expected != format
920    {
921        log::info!(
922            "format change at {}: {}Hz/{}ch → {}Hz/{}ch, restarting audio engine",
923            path.display(),
924            expected.0,
925            expected.1,
926            out_rate,
927            channels
928        );
929        return Ok(Decoded::FormatMismatch);
930    }
931
932    // Build either a Symphonia decoder or our Opus bridge.
933    let mut symphonia_decoder = if is_opus_codec {
934        None
935    } else {
936        Some(
937            symphonia::default::get_codecs()
938                .make_audio_decoder(codec_params, &AudioDecoderOptions::default())
939                .map_err(|_| DecodeError::UnsupportedCodec)?,
940        )
941    };
942    let mut opus_bridge = if is_opus_codec {
943        Some(OpusBridge::new(codec_params).map_err(|e| DecodeError::Decode(e.to_string()))?)
944    } else {
945        None
946    };
947
948    // Seek if requested (only for the first track usually).
949    //
950    // Accurate rather than coarse: a coarse seek picks a byte offset by
951    // interpolating linearly over the file and then derives its reported
952    // timestamp from that same guess, so on VBR MP3 it lands seconds from the
953    // request *and* reports a position it never reached. Accurate mode walks
954    // frame headers and is truthful about both, for 1.5-3ms on files up to
955    // 79MB. The timeline records where playback actually resumed, so the
956    // transport shows the position being heard.
957    let mut seek_samples = 0;
958    if seek_ms > 0 {
959        let seeked = reader
960            .seek(
961                SeekMode::Accurate,
962                SeekTo::Time {
963                    time: Time::from_millis_u64(seek_ms),
964                    track_id: Some(track_id),
965                },
966            )
967            .map_err(|e| DecodeError::Decode(format!("seek failed: {}", e)))?;
968        seek_samples = landing_samples(time_base, seeked.actual_ts, sample_rate, channels)
969            .unwrap_or(seek_ms * sample_rate as u64 * channels as u64 / 1000);
970        if let Some(ref mut dec) = symphonia_decoder {
971            dec.reset();
972        }
973        if let Some(ref mut opus) = opus_bridge {
974            opus.reset();
975        }
976    }
977
978    // The chain carries across a gapless boundary: resetting its filters would
979    // put a transient at every track change. A resampler for the outgoing rate
980    // gives up its tail first, which belongs to the track before.
981    if let Some(setup) = &processing.dsp {
982        match chain {
983            Some(c) => {
984                let tail = c.set_source_rate(sample_rate);
985                if !write_ring(producer, tail, stop, viz_buffer, format) {
986                    return Ok(Decoded::Complete(format));
987                }
988            }
989            None => *chain = Some(Chain::new(setup, sample_rate, channels)),
990        }
991    }
992
993    // Record this track's boundary in the timeline.
994    let write_offset = timeline.samples_written();
995    timeline.push_boundary(TrackBoundary {
996        id: queue_item_id,
997        path: path.to_path_buf(),
998        info,
999        sample_offset: write_offset,
1000        samples_written: 0,
1001        seek_samples,
1002        output_rate: out_rate,
1003    });
1004
1005    // Read ReplayGain tags and select the active gain for this track.
1006    let rg_mode = processing.rg_mode;
1007    let pre_amp_db = processing.pre_amp_db;
1008    let rg_gain = if rg_mode != ReplayGainMode::Off {
1009        match crate::audio::replaygain::read_tags(path) {
1010            Ok(rg_info) => {
1011                let selected = crate::audio::replaygain::select_gain(&rg_info, rg_mode);
1012                if let Some((gain_db, _)) = selected {
1013                    log::info!(
1014                        "replaygain: applying {:.2} dB ({:?}) to {}",
1015                        gain_db,
1016                        rg_mode,
1017                        path.display()
1018                    );
1019                }
1020                selected
1021            }
1022            Err(e) => {
1023                log::debug!("replaygain: no tags for {}: {}", path.display(), e);
1024                None
1025            }
1026        }
1027    } else {
1028        None
1029    };
1030    let mut rg_scratch: Vec<f32> = Vec::new();
1031
1032    let mut sample_buf: Vec<f32> = Vec::new();
1033
1034    loop {
1035        if stop.load(Ordering::Relaxed) {
1036            return Ok(Decoded::Complete(format));
1037        }
1038
1039        let packet = match reader.next_packet() {
1040            Ok(Some(p)) => p,
1041            Ok(None) => return Ok(Decoded::Complete(format)),
1042            Err(e) => return Err(DecodeError::Decode(e.to_string())),
1043        };
1044
1045        if packet.track_id != track_id {
1046            continue;
1047        }
1048
1049        // Decode the packet — either via Opus bridge or Symphonia codec.
1050        let samples: &[f32] = if let Some(ref mut opus) = opus_bridge {
1051            match opus.decode_packet(&packet.data) {
1052                Ok(s) => s,
1053                Err(e) => {
1054                    log::warn!("opus decode error (skipping packet): {}", e);
1055                    continue;
1056                }
1057            }
1058        } else {
1059            let decoder = symphonia_decoder.as_mut().unwrap();
1060            let decoded = match decoder.decode(&packet) {
1061                Ok(d) => d,
1062                Err(symphonia::core::errors::Error::DecodeError(e)) => {
1063                    log::warn!("decode error (skipping packet): {}", e);
1064                    continue;
1065                }
1066                Err(e) => return Err(DecodeError::Decode(e.to_string())),
1067            };
1068
1069            let spec = decoded.spec();
1070            let (decoded_rate, decoded_channels) = (spec.rate(), spec.channels().count() as u16);
1071            // The engine is configured from the probed format. PCM that
1072            // disagrees with it would play at the wrong speed, so end the
1073            // session instead and let the player reconfigure.
1074            if (decoded_rate, decoded_channels) != (sample_rate, channels) {
1075                log::warn!(
1076                    "{}: decoded {}Hz/{}ch but stream declares {}Hz/{}ch, restarting audio engine",
1077                    path.display(),
1078                    decoded_rate,
1079                    decoded_channels,
1080                    sample_rate,
1081                    channels
1082                );
1083                return Ok(Decoded::FormatMismatch);
1084            }
1085            decoded.copy_to_vec_interleaved(&mut sample_buf);
1086            &sample_buf[..]
1087        };
1088
1089        if samples.is_empty() {
1090            continue;
1091        }
1092
1093        // Apply ReplayGain if active. Uses a reusable scratch buffer to avoid
1094        // allocating per packet. Zero overhead when RG is off.
1095        let samples = if let Some((gain_db, peak)) = rg_gain {
1096            rg_scratch.clear();
1097            rg_scratch.extend_from_slice(samples);
1098            crate::audio::replaygain::apply_gain(&mut rg_scratch, gain_db, peak, pre_amp_db);
1099            &rg_scratch[..]
1100        } else {
1101            samples
1102        };
1103
1104        let (samples, length) = match chain.as_mut() {
1105            Some(c) => c.process(samples),
1106            None => (samples, samples.len() as u64),
1107        };
1108        if !write_ring(producer, samples, stop, viz_buffer, format) {
1109            return Ok(Decoded::Complete(format));
1110        }
1111        timeline.add_written(length);
1112    }
1113}
1114
1115/// Interleaved sample offset of a seek's landing point.
1116///
1117/// `actual_ts` is in the track's timebase, which is not always the reciprocal
1118/// of the sample rate (Matroska ticks in milliseconds), so it is converted
1119/// through `Time` rather than assumed to be a frame count.
1120fn landing_samples(
1121    time_base: Option<TimeBase>,
1122    actual_ts: Timestamp,
1123    sample_rate: u32,
1124    channels: u16,
1125) -> Option<u64> {
1126    let (seconds, nanos) = time_base?.calc_time(actual_ts)?.parts();
1127    let rate = sample_rate as u64;
1128    let frames = seconds.max(0) as u64 * rate + (nanos as u64 * rate) / 1_000_000_000;
1129    Some(frames * channels as u64)
1130}
1131
1132/// Duration of a track in milliseconds.
1133///
1134/// The container's stated duration is authoritative because a track's timebase
1135/// is not always the reciprocal of the sample rate — Matroska ticks in
1136/// milliseconds, and states its duration at media level rather than per track.
1137/// Falls back to the playable frame count when no duration is stated at all.
1138pub(crate) fn track_duration_ms(
1139    reader: &(impl FormatReader + ?Sized),
1140    track: &Track,
1141    sample_rate: u32,
1142) -> u64 {
1143    fn to_ms(time_base: Option<TimeBase>, duration: Option<Duration>) -> Option<u64> {
1144        let time = time_base?.calc_duration(duration?)?;
1145        Some(time.as_millis().max(0) as u64)
1146    }
1147
1148    let media = reader.media_info();
1149    to_ms(track.time_base, track.duration)
1150        .or_else(|| to_ms(media.time_base, media.duration))
1151        .or_else(|| {
1152            track
1153                .num_frames
1154                .map(|frames| frames * 1000 / sample_rate as u64)
1155        })
1156        .unwrap_or(0)
1157}
1158
1159/// Estimate bitrate (kbps) from Symphonia codec parameters.
1160///
1161/// Symphonia doesn't expose a `bit_rate` field. For lossy codecs like MP3/AAC
1162/// we can derive it from `bits_per_coded_sample` when the demuxer populates it.
1163/// Returns `None` for lossless codecs or when the info isn't available.
1164fn estimate_bitrate_from_codec_params(params: &AudioCodecParameters) -> Option<u32> {
1165    let is_lossy = matches!(
1166        params.codec,
1167        CODEC_ID_MP3 | CODEC_ID_AAC | CODEC_ID_VORBIS | CODEC_ID_OPUS
1168    );
1169    if !is_lossy {
1170        return None;
1171    }
1172
1173    // bits_per_coded_sample * sample_rate / 1000 gives kbps for CBR streams.
1174    // Few demuxers fill this in, but it's our best shot without file size.
1175    let bpcs = params.bits_per_coded_sample?;
1176    let sr = params.sample_rate?;
1177    let channels = params
1178        .channels
1179        .as_ref()
1180        .map(|c| c.count() as u32)
1181        .unwrap_or(2);
1182    Some(bpcs * sr * channels / 1000)
1183}
1184
1185pub fn codec_name(codec: AudioCodecId) -> String {
1186    match codec {
1187        CODEC_ID_FLAC => "FLAC",
1188        CODEC_ID_MP3 => "MP3",
1189        CODEC_ID_AAC => "AAC",
1190        CODEC_ID_VORBIS => "Vorbis",
1191        CODEC_ID_OPUS => "Opus",
1192        CODEC_ID_ALAC => "ALAC",
1193        CODEC_ID_PCM_S16LE => "PCM/16",
1194        CODEC_ID_PCM_S24LE => "PCM/24",
1195        CODEC_ID_PCM_S32LE => "PCM/32",
1196        CODEC_ID_PCM_F32LE => "PCM/f32",
1197        other => return format!("Unknown({:?})", other),
1198    }
1199    .to_string()
1200}
1201
1202#[cfg(test)]
1203mod tests {
1204    use std::path::PathBuf;
1205    use std::sync::atomic::Ordering;
1206
1207    use super::*;
1208    use crate::player::state::QueueItemId;
1209
1210    fn make_info(sample_rate: u32, channels: u16) -> StreamInfo {
1211        StreamInfo {
1212            codec: "FLAC".to_string(),
1213            sample_rate,
1214            channels,
1215            bit_depth: Some(16),
1216            bitrate_kbps: None,
1217            duration_ms: 10_000,
1218        }
1219    }
1220
1221    fn make_boundary(
1222        id: QueueItemId,
1223        sample_offset: u64,
1224        seek_samples: u64,
1225        channels: u16,
1226        sample_rate: u32,
1227    ) -> TrackBoundary {
1228        TrackBoundary {
1229            id,
1230            path: PathBuf::from("/music/track.flac"),
1231            info: make_info(sample_rate, channels),
1232            sample_offset,
1233            samples_written: 0,
1234            seek_samples,
1235            output_rate: sample_rate,
1236        }
1237    }
1238
1239    fn writer(timeline: &PlaybackTimeline) -> TimelineWriter<'_> {
1240        timeline.writer()
1241    }
1242
1243    // --- PlaybackTimeline tests ---
1244
1245    #[test]
1246    fn test_timeline_single_track() {
1247        // Push one boundary at offset 0 with stereo 44100 Hz audio.
1248        // After simulating 44100 frames (88200 interleaved samples) played,
1249        // current_playback() should report track index 0 at position 1000 ms.
1250        let timeline = PlaybackTimeline::new();
1251        let tl = writer(&timeline);
1252        let id = QueueItemId::new();
1253        // sample_offset=0, seek_samples=0, channels=2, sample_rate=44100
1254        tl.push_boundary(make_boundary(id, 0, 0, 2, 44100));
1255        tl.add_written(88200); // 1 second of audio
1256
1257        // Simulate 1 second played: 44100 frames * 2 channels = 88200 interleaved samples
1258        timeline.samples_played.store(88200, Ordering::Relaxed);
1259
1260        let result = timeline.current_playback();
1261        assert!(
1262            result.is_some(),
1263            "expected Some for single track with samples played"
1264        );
1265        let (result_id, _path, _info, position_ms) = result.unwrap();
1266        assert_eq!(result_id, id);
1267        assert_eq!(
1268            position_ms, 1000,
1269            "1 second of 44100 Hz stereo should be 1000 ms"
1270        );
1271    }
1272
1273    #[test]
1274    fn test_timeline_gapless_transition() {
1275        // Two tracks in gapless sequence. Track 1 ends at sample 88200 (1 sec stereo 44100 Hz).
1276        // Track 2 begins at sample_offset 88200. When playback head is at 100000 (past the boundary),
1277        // current_playback() should report track 2.
1278        let timeline = PlaybackTimeline::new();
1279        let tl = writer(&timeline);
1280        let id1 = QueueItemId::new();
1281        let id2 = QueueItemId::new();
1282
1283        // Track 1: starts at offset 0
1284        tl.push_boundary(make_boundary(id1, 0, 0, 2, 44100));
1285        tl.add_written(88200);
1286
1287        // Track 2: starts at offset 88200 (immediately after track 1's samples)
1288        tl.push_boundary(make_boundary(id2, 88200, 0, 2, 44100));
1289        tl.add_written(44100); // half a second of track 2
1290
1291        // Set playback head past the track 1/2 boundary
1292        timeline.samples_played.store(90000, Ordering::Relaxed);
1293
1294        let result = timeline.current_playback();
1295        assert!(result.is_some());
1296        let (result_id, _path, _info, position_ms) = result.unwrap();
1297        assert_eq!(
1298            result_id, id2,
1299            "playback head past boundary should report second track"
1300        );
1301        // (90000 - 88200) / 2 channels * 1000 / 44100 = 900 / 44100 ≈ 20 ms
1302        assert_eq!(position_ms, 20, "position within track 2 should be ~20 ms");
1303    }
1304
1305    #[test]
1306    fn test_timeline_zero_samples() {
1307        // With 0 samples played and a boundary at offset 0, current_playback() should
1308        // still return the first track at position 0 ms.
1309        let timeline = PlaybackTimeline::new();
1310        let tl = writer(&timeline);
1311        let id = QueueItemId::new();
1312        tl.push_boundary(make_boundary(id, 0, 0, 2, 44100));
1313        tl.add_written(1000);
1314        timeline.samples_played.store(0, Ordering::Relaxed);
1315
1316        let result = timeline.current_playback();
1317        assert!(
1318            result.is_some(),
1319            "expected Some at 0 samples played with a boundary at offset 0"
1320        );
1321        let (result_id, _path, _info, position_ms) = result.unwrap();
1322        assert_eq!(result_id, id);
1323        assert_eq!(position_ms, 0);
1324    }
1325
1326    #[test]
1327    fn test_timeline_past_all_boundaries() {
1328        // When samples_played exceeds all boundaries, the last track should be reported.
1329        // The binary search finds the last boundary whose sample_offset <= played.
1330        let timeline = PlaybackTimeline::new();
1331        let tl = writer(&timeline);
1332        let id1 = QueueItemId::new();
1333        let id2 = QueueItemId::new();
1334
1335        tl.push_boundary(make_boundary(id1, 0, 0, 2, 44100));
1336        tl.add_written(88200);
1337        tl.push_boundary(make_boundary(id2, 88200, 0, 2, 44100));
1338        tl.add_written(88200);
1339
1340        // Simulate playback far past both tracks
1341        timeline
1342            .samples_played
1343            .store(999_999_999, Ordering::Relaxed);
1344
1345        let result = timeline.current_playback();
1346        assert!(result.is_some());
1347        let (result_id, _path, _info, _position_ms) = result.unwrap();
1348        assert_eq!(
1349            result_id, id2,
1350            "samples past all boundaries should report the last track"
1351        );
1352    }
1353
1354    #[test]
1355    fn test_timeline_seek_offset() {
1356        // When a seek offset is set, position_ms should include the seek position.
1357        // seek_samples = 88200 means playback started 1 second into the track.
1358        // With 0 additional samples played past the boundary, position should be 1000 ms.
1359        let timeline = PlaybackTimeline::new();
1360        let tl = writer(&timeline);
1361        let id = QueueItemId::new();
1362        let seek_samples = 88200u64; // 1 second at 44100 Hz stereo
1363        tl.push_boundary(make_boundary(id, 0, seek_samples, 2, 44100));
1364        tl.add_written(44100); // half a second written so far
1365        // samples_played at the track boundary (0 frames past the track start)
1366        timeline.samples_played.store(0, Ordering::Relaxed);
1367
1368        let result = timeline.current_playback();
1369        assert!(result.is_some());
1370        let (_result_id, _path, _info, position_ms) = result.unwrap();
1371        // track_samples = 0 - 0 = 0; seek contribution = (88200/2)*1000/44100 = 1000 ms
1372        assert_eq!(
1373            position_ms, 1000,
1374            "position should include seek offset of 1000 ms"
1375        );
1376    }
1377
1378    #[test]
1379    fn the_next_track_is_due_when_the_current_one_has_played_out() {
1380        let timeline = PlaybackTimeline::new();
1381        let tl = writer(&timeline);
1382        let (a, b) = (QueueItemId::new(), QueueItemId::new());
1383        // One second of 44.1kHz stereo before b starts.
1384        tl.push_boundary(make_boundary(a, 0, 0, 2, 44100));
1385        assert_eq!(timeline.until_next_track(), None, "nothing queued after it");
1386        tl.push_boundary(make_boundary(b, 88200, 0, 2, 44100));
1387
1388        timeline.samples_played.store(44100, Ordering::Relaxed);
1389        assert_eq!(
1390            timeline.until_next_track(),
1391            Some(std::time::Duration::from_millis(500))
1392        );
1393        timeline.samples_played.store(88200, Ordering::Relaxed);
1394        assert_eq!(timeline.until_next_track(), None, "b is playing now");
1395    }
1396
1397    #[test]
1398    fn the_queued_callback_hears_each_track() {
1399        let timeline = PlaybackTimeline::new();
1400        let heard = Arc::new(std::sync::atomic::AtomicUsize::new(0));
1401        let counter = heard.clone();
1402        timeline.on_queued(move || {
1403            counter.fetch_add(1, Ordering::Relaxed);
1404        });
1405        let tl = writer(&timeline);
1406        tl.push_boundary(make_boundary(QueueItemId::new(), 0, 0, 2, 44100));
1407        tl.push_boundary(make_boundary(QueueItemId::new(), 88200, 0, 2, 44100));
1408        assert_eq!(heard.load(Ordering::Relaxed), 2);
1409    }
1410
1411    #[test]
1412    fn a_track_played_past_reads_as_played_to_its_end() {
1413        let timeline = PlaybackTimeline::new();
1414        let tl = writer(&timeline);
1415        let (a, b) = (QueueItemId::new(), QueueItemId::new());
1416        tl.push_boundary(make_boundary(a, 0, 0, 2, 44100));
1417        tl.push_boundary(make_boundary(b, 88200, 0, 2, 44100));
1418        timeline
1419            .samples_played
1420            .store(88200 + 44100, Ordering::Relaxed);
1421
1422        assert_eq!(timeline.position_in(0), Some(1000));
1423        assert_eq!(timeline.position_in(1), Some(500));
1424        assert_eq!(
1425            timeline.playhead(),
1426            Some(Playhead {
1427                boundary: 1,
1428                id: b,
1429                position_ms: 500
1430            })
1431        );
1432    }
1433
1434    #[test]
1435    fn an_item_queued_twice_is_two_plays() {
1436        let timeline = PlaybackTimeline::new();
1437        let tl = writer(&timeline);
1438        let a = QueueItemId::new();
1439        tl.push_boundary(make_boundary(a, 0, 0, 2, 44100));
1440        tl.push_boundary(make_boundary(a, 88200, 0, 2, 44100));
1441        timeline
1442            .samples_played
1443            .store(88200 + 44100, Ordering::Relaxed);
1444
1445        assert_eq!(timeline.position_in(0), Some(1000), "the first pass ended");
1446        assert_eq!(timeline.position_in(1), Some(500));
1447        assert_eq!(timeline.playhead().map(|p| p.boundary), Some(1));
1448    }
1449
1450    #[test]
1451    fn test_timeline_reset() {
1452        // After reset(), current_playback() returns None and all counters are cleared.
1453        let timeline = PlaybackTimeline::new();
1454        let tl = writer(&timeline);
1455        let id = QueueItemId::new();
1456        tl.push_boundary(make_boundary(id, 0, 0, 2, 44100));
1457        tl.add_written(88200);
1458        timeline.samples_played.store(44100, Ordering::Relaxed);
1459
1460        // Sanity check: playback is live before reset
1461        assert!(timeline.current_playback().is_some());
1462
1463        timeline.reset();
1464
1465        assert!(
1466            timeline.current_playback().is_none(),
1467            "after reset, current_playback should return None"
1468        );
1469        assert_eq!(
1470            timeline.samples_played.load(Ordering::Relaxed),
1471            0,
1472            "samples_played should be 0 after reset"
1473        );
1474        assert_eq!(
1475            timeline.samples_written.load(Ordering::Relaxed),
1476            0,
1477            "samples_written should be 0 after reset"
1478        );
1479    }
1480
1481    // --- Probe and decode integration tests ---
1482
1483    #[test]
1484    fn probe_file_extracts_stream_info() {
1485        let dir = tempfile::tempdir().unwrap();
1486        let wav_path = dir.path().join("probe_test.wav");
1487        crate::test_utils::generate_wav(&wav_path, 44100, 2, 1.0, 16);
1488
1489        let info = probe_file(&wav_path).expect("probe_file should succeed on a valid WAV");
1490        assert_eq!(info.sample_rate, 44100, "sample rate mismatch");
1491        assert_eq!(info.channels, 2, "channel count mismatch");
1492        assert_eq!(info.bit_depth, Some(16), "bit depth mismatch");
1493        assert!(
1494            info.duration_ms > 900 && info.duration_ms < 1100,
1495            "duration should be ~1000ms, got {}",
1496            info.duration_ms
1497        );
1498        assert!(
1499            info.codec.contains("PCM"),
1500            "codec should be PCM variant, got {}",
1501            info.codec
1502        );
1503    }
1504
1505    #[test]
1506    fn decode_single_produces_samples() {
1507        let dir = tempfile::tempdir().unwrap();
1508        let wav_path = dir.path().join("tone.wav");
1509        // 440 Hz sine, mono, 0.1s — enough to verify non-zero decode output.
1510        crate::test_utils::generate_wav_tone(&wav_path, 44100, 440.0, 0.1);
1511
1512        // Set up rtrb ring buffer.
1513        let (mut producer, mut consumer) = rtrb::RingBuffer::new(44100 * 2);
1514
1515        let timeline = PlaybackTimeline::new();
1516        let tl = writer(&timeline);
1517        let stop = Arc::new(AtomicBool::new(false));
1518
1519        let id = QueueItemId::new();
1520        let entry = SourceEntry::from_file(id, wav_path.clone());
1521        let hint = entry.hint.clone();
1522        let mss = (entry.make_mss)().expect("should open WAV file");
1523
1524        let result = decode_single(
1525            id,
1526            &wav_path,
1527            &hint,
1528            mss,
1529            &mut producer,
1530            &stop,
1531            0,
1532            &tl,
1533            None,
1534            &Processing::default(),
1535            &mut None,
1536            None,
1537        );
1538        assert!(
1539            matches!(result, Ok(Decoded::Complete((44100, 1)))),
1540            "decode_single should complete at the source format"
1541        );
1542
1543        // Read samples from the consumer side.
1544        let available = consumer.slots();
1545        assert!(available > 0, "expected samples in ring buffer, got 0");
1546
1547        // Verify at least some samples are non-zero (it's a sine wave, not silence).
1548        let mut found_nonzero = false;
1549        while consumer.slots() > 0 {
1550            if let Ok(chunk) = consumer.read_chunk(consumer.slots().min(1024)) {
1551                let (first, second) = chunk.as_slices();
1552                for &s in first.iter().chain(second.iter()) {
1553                    if s.abs() > 0.001 {
1554                        found_nonzero = true;
1555                        break;
1556                    }
1557                }
1558                chunk.commit_all();
1559            }
1560            if found_nonzero {
1561                break;
1562            }
1563        }
1564        assert!(
1565            found_nonzero,
1566            "expected non-zero samples from 440Hz sine decode"
1567        );
1568    }
1569
1570    // --- Ring buffer format contract ---
1571
1572    /// Decode a queue of files through `decode_queue_loop` with a consumer
1573    /// draining in the background. Returns the boundaries the decode thread
1574    /// pushed onto the timeline.
1575    fn run_queue(paths: &[PathBuf]) -> Vec<TrackBoundary> {
1576        run_queue_with(paths, &Processing::default()).0
1577    }
1578
1579    /// `run_queue`, through `processing`. Also returns how many samples the
1580    /// consumer read and how many the timeline counted.
1581    fn run_queue_with(
1582        paths: &[PathBuf],
1583        processing: &Processing,
1584    ) -> (Vec<TrackBoundary>, u64, u64) {
1585        let (producer, mut consumer) = rtrb::RingBuffer::new(1 << 16);
1586        let timeline = PlaybackTimeline::new();
1587        let tl = writer(&timeline);
1588        let stop = Arc::new(AtomicBool::new(false));
1589
1590        let drain_stop = Arc::new(AtomicBool::new(false));
1591        let drain_flag = drain_stop.clone();
1592        let drainer = std::thread::spawn(move || {
1593            let mut read = 0u64;
1594            while !drain_flag.load(Ordering::Relaxed) {
1595                let n = consumer.slots();
1596                if n > 0
1597                    && let Ok(chunk) = consumer.read_chunk(n)
1598                {
1599                    read += n as u64;
1600                    chunk.commit_all();
1601                }
1602                std::thread::sleep(std::time::Duration::from_micros(200));
1603            }
1604            read
1605        });
1606
1607        let rest: std::sync::Mutex<Vec<PathBuf>> = std::sync::Mutex::new(paths[1..].to_vec());
1608        let next_track = move || {
1609            let mut rest = rest.lock().ok()?;
1610            if rest.is_empty() {
1611                return None;
1612            }
1613            Some(SourceEntry::from_file(QueueItemId::new(), rest.remove(0)))
1614        };
1615
1616        decode_queue_loop(
1617            SourceEntry::from_file(QueueItemId::new(), paths[0].clone()),
1618            producer,
1619            &stop,
1620            0,
1621            &next_track,
1622            &tl,
1623            None,
1624            processing,
1625        );
1626
1627        drain_stop.store(true, Ordering::Relaxed);
1628        let read = drainer.join().unwrap();
1629
1630        let written = timeline.samples_written.load(Ordering::Relaxed);
1631        (timeline.boundaries.read().clone(), read, written)
1632    }
1633
1634    #[test]
1635    fn convolution_at_one_rate_plays_two_source_rates_gaplessly() {
1636        let dir = tempfile::tempdir().unwrap();
1637        let a = dir.path().join("a.wav");
1638        let b = dir.path().join("b.wav");
1639        crate::test_utils::generate_wav(&a, 44100, 2, 0.5, 16);
1640        crate::test_utils::generate_wav(&b, 96000, 2, 0.5, 16);
1641        let mut ir = vec![0.0; 129];
1642        ir[64] = 1.0;
1643        let processing = Processing {
1644            dsp: Some(Arc::new(Setup::new(
1645                vec![],
1646                vec![crate::audio::dsp::Impulse::from_channels(48000, vec![ir])],
1647            ))),
1648            ..Processing::default()
1649        };
1650
1651        let (bounds, read, written) = run_queue_with(&[a, b], &processing);
1652        assert_eq!(bounds.len(), 2, "both resample to 48 kHz, so one session");
1653        assert!(bounds.iter().all(|b| b.output_rate == 48000));
1654        assert_eq!(bounds[0].info.sample_rate, 44100);
1655        assert_eq!(bounds[1].info.sample_rate, 96000);
1656        // Half a second of each at 48 kHz, stereo — and every sample the
1657        // timeline counted reached the output, delays trimmed and tails flushed.
1658        assert_eq!(written, 2 * 24000 * 2);
1659        assert_eq!(read, written);
1660        assert_eq!(bounds[1].sample_offset, 24000 * 2);
1661    }
1662
1663    #[test]
1664    fn gapless_continues_when_format_matches() {
1665        let dir = tempfile::tempdir().unwrap();
1666        let a = dir.path().join("a.wav");
1667        let b = dir.path().join("b.wav");
1668        crate::test_utils::generate_wav(&a, 44100, 2, 0.1, 16);
1669        crate::test_utils::generate_wav(&b, 44100, 2, 0.1, 16);
1670
1671        let bounds = run_queue(&[a, b]);
1672        assert_eq!(
1673            bounds.len(),
1674            2,
1675            "same-format tracks should decode gaplessly"
1676        );
1677    }
1678
1679    #[test]
1680    fn gapless_stops_at_sample_rate_change() {
1681        let dir = tempfile::tempdir().unwrap();
1682        let a = dir.path().join("a.wav");
1683        let b = dir.path().join("b.wav");
1684        crate::test_utils::generate_wav(&a, 44100, 2, 0.1, 16);
1685        crate::test_utils::generate_wav(&b, 48000, 2, 0.1, 16);
1686
1687        let bounds = run_queue(&[a, b]);
1688        assert_eq!(
1689            bounds.len(),
1690            1,
1691            "a 48kHz track must not join a 44.1kHz ring buffer"
1692        );
1693        assert_eq!(bounds[0].info.sample_rate, 44100);
1694    }
1695
1696    #[test]
1697    fn gapless_stops_at_channel_change() {
1698        let dir = tempfile::tempdir().unwrap();
1699        let a = dir.path().join("a.wav");
1700        let b = dir.path().join("b.wav");
1701        crate::test_utils::generate_wav(&a, 44100, 2, 0.1, 16);
1702        crate::test_utils::generate_wav(&b, 44100, 1, 0.1, 16);
1703
1704        let bounds = run_queue(&[a, b]);
1705        assert_eq!(
1706            bounds.len(),
1707            1,
1708            "a mono track must not join a stereo ring buffer"
1709        );
1710        assert_eq!(bounds[0].info.channels, 2);
1711    }
1712
1713    #[test]
1714    fn drain_waits_for_the_consumer() {
1715        let (mut producer, mut consumer) = rtrb::RingBuffer::new(64);
1716        for _ in 0..64 {
1717            producer.push(0.0).unwrap();
1718        }
1719        let stop = Arc::new(AtomicBool::new(false));
1720
1721        let reader = std::thread::spawn(move || {
1722            std::thread::sleep(std::time::Duration::from_millis(20));
1723            let chunk = consumer.read_chunk(64).unwrap();
1724            chunk.commit_all();
1725            consumer
1726        });
1727
1728        wait_for_drain(&producer, &stop, None);
1729        assert_eq!(producer.slots(), 64, "drain must wait for an empty buffer");
1730        drop(reader.join().unwrap());
1731    }
1732
1733    #[test]
1734    fn drain_returns_when_playback_is_torn_down() {
1735        let (producer, consumer) = rtrb::RingBuffer::<f32>::new(64);
1736        let stop = Arc::new(AtomicBool::new(true));
1737        wait_for_drain(&producer, &stop, None);
1738        drop(consumer);
1739    }
1740
1741    // --- Failure handling in the gapless queue ---
1742
1743    /// A file that exists and has an audio extension but no audio in it.
1744    fn write_garbage(path: &Path) {
1745        std::fs::write(path, b"this is not a wav file").unwrap();
1746    }
1747
1748    #[test]
1749    fn an_unreadable_track_is_skipped_and_the_queue_continues() {
1750        let dir = tempfile::tempdir().unwrap();
1751        let a = dir.path().join("a.wav");
1752        let bad = dir.path().join("bad.wav");
1753        let c = dir.path().join("c.wav");
1754        crate::test_utils::generate_wav(&a, 44100, 2, 0.1, 16);
1755        write_garbage(&bad);
1756        crate::test_utils::generate_wav(&c, 44100, 2, 0.1, 16);
1757
1758        let bounds = run_queue(&[a.clone(), bad, c.clone()]);
1759        let decoded: Vec<_> = bounds.iter().map(|b| b.path.clone()).collect();
1760        assert_eq!(
1761            decoded,
1762            vec![a, c],
1763            "one bad file must not take the rest of the queue with it"
1764        );
1765    }
1766
1767    #[test]
1768    fn a_missing_track_is_skipped_and_the_queue_continues() {
1769        let dir = tempfile::tempdir().unwrap();
1770        let missing = dir.path().join("gone.wav");
1771        let b = dir.path().join("b.wav");
1772        crate::test_utils::generate_wav(&b, 44100, 2, 0.1, 16);
1773
1774        let bounds = run_queue(&[missing, b.clone()]);
1775        assert_eq!(bounds.len(), 1);
1776        assert_eq!(
1777            bounds[0].path, b,
1778            "a bad first track must not end the session"
1779        );
1780    }
1781
1782    #[test]
1783    fn an_entirely_unreadable_queue_terminates() {
1784        let dir = tempfile::tempdir().unwrap();
1785        let bad = dir.path().join("bad.wav");
1786        write_garbage(&bad);
1787
1788        // Nothing decodes, so the ring stays empty; the consumer only has to
1789        // outlive the producer.
1790        let (producer, _consumer) = rtrb::RingBuffer::<f32>::new(1 << 12);
1791        let timeline = PlaybackTimeline::new();
1792        let tl = writer(&timeline);
1793        let stop = Arc::new(AtomicBool::new(false));
1794
1795        let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
1796        let counter = calls.clone();
1797        let bad_path = bad.clone();
1798        let next_track = move || {
1799            counter.fetch_add(1, Ordering::Relaxed);
1800            Some(SourceEntry::from_file(QueueItemId::new(), bad_path.clone()))
1801        };
1802
1803        decode_queue_loop(
1804            SourceEntry::from_file(QueueItemId::new(), bad),
1805            producer,
1806            &stop,
1807            0,
1808            &next_track,
1809            &tl,
1810            None,
1811            &Processing::default(),
1812        );
1813
1814        // The first source plus one per next_track call, capped.
1815        assert_eq!(
1816            calls.load(Ordering::Relaxed) + 1,
1817            MAX_CONSECUTIVE_FAILURES as usize
1818        );
1819        assert!(timeline.boundaries.read().is_empty());
1820    }
1821
1822    // --- Seek landing position ---
1823
1824    #[test]
1825    fn landing_samples_converts_frame_timebases() {
1826        // The usual audio case: one tick per frame.
1827        let tb = TimeBase::try_from_recip(44100).unwrap();
1828        assert_eq!(
1829            landing_samples(Some(tb), Timestamp::from(44100u32), 44100, 2),
1830            Some(88_200)
1831        );
1832    }
1833
1834    #[test]
1835    fn landing_samples_converts_millisecond_timebases() {
1836        // Matroska ticks in milliseconds, not frames.
1837        let tb = TimeBase::try_new(1, 1000).unwrap();
1838        assert_eq!(
1839            landing_samples(Some(tb), Timestamp::from(1500u32), 48000, 2),
1840            Some(48000 * 3 / 2 * 2)
1841        );
1842    }
1843
1844    #[test]
1845    fn landing_samples_needs_a_timebase() {
1846        assert_eq!(
1847            landing_samples(None, Timestamp::from(1000u32), 44100, 2),
1848            None
1849        );
1850    }
1851
1852    /// Build a 300s VBR MP3: 30s of near-silence then a loud tone, which gives
1853    /// lame a wide enough bitrate spread to make coarse seeking miss badly.
1854    #[cfg(test)]
1855    fn make_vbr_mp3(dir: &Path) -> PathBuf {
1856        let wav = dir.join("source.wav");
1857        let mp3 = dir.join("source.mp3");
1858        let ok = std::process::Command::new("sox")
1859            .args(["-n", "-r", "44100", "-c", "2"])
1860            .arg(&wav)
1861            .args([
1862                "synth", "30", "sine", "200", "vol", "0.02", ":", "synth", "270", "sine", "880",
1863                "vol", "0.9",
1864            ])
1865            .status()
1866            .expect("sox not installed")
1867            .success();
1868        assert!(ok, "sox failed");
1869        let ok = std::process::Command::new("lame")
1870            .args(["-V", "2", "--quiet"])
1871            .arg(&wav)
1872            .arg(&mp3)
1873            .status()
1874            .expect("lame not installed")
1875            .success();
1876        assert!(ok, "lame failed");
1877        mp3
1878    }
1879
1880    /// A VBR seek must report where playback actually resumed: the reported
1881    /// start plus the audio that actually followed has to add back up to the
1882    /// file's duration. Under a coarse seek this file lands 3.7s late while
1883    /// reporting 83ms early — a 3.8s lie for the rest of the track.
1884    #[test]
1885    #[ignore = "generates a fixture with sox + lame; run with cargo test -- --ignored"]
1886    fn seek_on_vbr_reports_where_it_landed() {
1887        let dir = tempfile::tempdir().unwrap();
1888        let path = make_vbr_mp3(dir.path());
1889        let info = probe_file(&path).unwrap();
1890        let channels = info.channels as u64;
1891        let rate = info.sample_rate as u64;
1892
1893        let seek_ms = 150_000u64;
1894        let (mut producer, mut consumer) = rtrb::RingBuffer::<f32>::new(1 << 16);
1895        let stop = Arc::new(AtomicBool::new(false));
1896        let timeline = PlaybackTimeline::new();
1897        let tl = writer(&timeline);
1898
1899        let drain_stop = stop.clone();
1900        let drained = std::thread::spawn(move || {
1901            let mut total = 0u64;
1902            while !drain_stop.load(Ordering::Relaxed) {
1903                let slots = consumer.slots();
1904                if slots == 0 {
1905                    std::thread::sleep(std::time::Duration::from_micros(200));
1906                    continue;
1907                }
1908                let chunk = consumer.read_chunk(slots).unwrap();
1909                total += slots as u64;
1910                chunk.commit_all();
1911            }
1912            total
1913        });
1914
1915        let file = File::open(&path).unwrap();
1916        let mss = MediaSourceStream::new(Box::new(file), Default::default());
1917        let mut hint = Hint::new();
1918        hint.with_extension("mp3");
1919        decode_single(
1920            QueueItemId::new(),
1921            &path,
1922            &hint,
1923            mss,
1924            &mut producer,
1925            &stop,
1926            seek_ms,
1927            &tl,
1928            None,
1929            &Processing::default(),
1930            &mut None,
1931            None,
1932        )
1933        .unwrap();
1934
1935        let written = timeline.samples_written.load(Ordering::Relaxed);
1936        stop.store(true, Ordering::Relaxed);
1937        drained.join().unwrap();
1938
1939        let reported_start_ms = {
1940            let bounds = timeline.boundaries.read();
1941            (bounds[0].seek_samples / channels) * 1000 / rate
1942        };
1943        let decoded_ms = (written / channels) * 1000 / rate;
1944
1945        // Where playback started plus how much audio followed is the whole file.
1946        let total_ms = reported_start_ms + decoded_ms;
1947        assert!(
1948            total_ms.abs_diff(info.duration_ms) < 500,
1949            "reported start {}ms + {}ms decoded = {}ms, but the file is {}ms",
1950            reported_start_ms,
1951            decoded_ms,
1952            total_ms,
1953            info.duration_ms
1954        );
1955        // Landing is frame-granular, never sample-exact.
1956        assert!(
1957            reported_start_ms.abs_diff(seek_ms) < 100,
1958            "seek to {}ms reported {}ms",
1959            seek_ms,
1960            reported_start_ms
1961        );
1962    }
1963}