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