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
640/// Block until the audio engine has consumed everything in the ring buffer.
641///
642/// A session ends only once its audio has been heard, so the player can tear
643/// the engine down without clipping the tail of the last track decoded.
644/// Returns early if playback is torn down underneath us.
645fn wait_for_drain(producer: &rtrb::Producer<f32>, stop: &AtomicBool) {
646    let capacity = producer.buffer().capacity();
647    while !stop.load(Ordering::Relaxed) && !producer.is_abandoned() {
648        if producer.slots() >= capacity {
649            return;
650        }
651        thread::sleep(std::time::Duration::from_millis(2));
652    }
653}
654
655// ---------------------------------------------------------------------------
656// Core decode single track
657// ---------------------------------------------------------------------------
658
659/// The PCM format of a decoded stream: sample rate in Hz and channel count.
660/// The audio engine is configured from this, so the ring buffer may only ever
661/// hold samples of one such format at a time.
662type PcmFormat = (u32, u16);
663
664/// Outcome of decoding one source.
665enum Decoded {
666    /// Decoded to EOF, in the given format.
667    Complete(PcmFormat),
668    /// The source's format differs from the stream already in the ring buffer.
669    /// Nothing further was written — the engine must be reconfigured first.
670    FormatMismatch,
671}
672
673/// Decode a single source into the producer. Returns on clean EOF.
674///
675/// `expected` — the format already in the ring buffer, if any. A source that
676/// does not match it is rejected without writing samples or pushing a boundary.
677#[allow(clippy::too_many_arguments)]
678fn decode_single(
679    queue_item_id: QueueItemId,
680    path: &Path,
681    hint: &Hint,
682    mss: MediaSourceStream<'_>,
683    producer: &mut rtrb::Producer<f32>,
684    stop: &AtomicBool,
685    seek_ms: u64,
686    timeline: &TimelineWriter<'_>,
687    viz_buffer: Option<&VizBuffer>,
688    rg_mode: ReplayGainMode,
689    pre_amp_db: f64,
690    expected: Option<PcmFormat>,
691) -> Result<Decoded, DecodeError> {
692    let mut reader = symphonia::default::get_probe()
693        .probe(
694            hint,
695            mss,
696            FormatOptions::default(),
697            MetadataOptions::default(),
698        )
699        .map_err(|e| DecodeError::Decode(e.to_string()))?;
700
701    let track = reader
702        .default_track(TrackType::Audio)
703        .ok_or(DecodeError::NoTrack)?;
704    let track_id = track.id;
705    let time_base = track.time_base;
706    let codec_params = track
707        .codec_params
708        .as_ref()
709        .and_then(|p| p.audio())
710        .ok_or(DecodeError::NoTrack)?;
711    let is_opus_codec = codec_params.codec == CODEC_ID_OPUS;
712
713    // Opus always decodes to 48 kHz regardless of the internal rate.
714    let sample_rate = if is_opus_codec {
715        48000
716    } else {
717        codec_params.sample_rate.unwrap_or(44100)
718    };
719    let channels = codec_params
720        .channels
721        .as_ref()
722        .map(|c| c.count() as u16)
723        .unwrap_or(2);
724
725    let duration_ms = track_duration_ms(&*reader, track, sample_rate);
726
727    // Try codec_params first; fall back to file-size estimation for Opus/lossy.
728    let mut bitrate_kbps = estimate_bitrate_from_codec_params(codec_params);
729    if bitrate_kbps.is_none()
730        && is_opus_codec
731        && let Ok(meta) = std::fs::metadata(path)
732        && duration_ms > 0
733    {
734        bitrate_kbps = Some((meta.len() * 8 / duration_ms) as u32);
735    }
736
737    let info = StreamInfo {
738        codec: codec_name(codec_params.codec),
739        sample_rate,
740        channels,
741        bit_depth: if is_opus_codec {
742            None
743        } else {
744            Some(codec_params.bits_per_sample.unwrap_or(16) as u16)
745        },
746        bitrate_kbps,
747        duration_ms,
748    };
749
750    if let Some(expected) = expected
751        && expected != (sample_rate, channels)
752    {
753        log::info!(
754            "format change at {}: {}Hz/{}ch → {}Hz/{}ch, restarting audio engine",
755            path.display(),
756            expected.0,
757            expected.1,
758            sample_rate,
759            channels
760        );
761        return Ok(Decoded::FormatMismatch);
762    }
763
764    // Build either a Symphonia decoder or our Opus bridge.
765    let mut symphonia_decoder = if is_opus_codec {
766        None
767    } else {
768        Some(
769            symphonia::default::get_codecs()
770                .make_audio_decoder(codec_params, &AudioDecoderOptions::default())
771                .map_err(|_| DecodeError::UnsupportedCodec)?,
772        )
773    };
774    let mut opus_bridge = if is_opus_codec {
775        Some(OpusBridge::new(codec_params).map_err(|e| DecodeError::Decode(e.to_string()))?)
776    } else {
777        None
778    };
779
780    // Seek if requested (only for the first track usually).
781    //
782    // Accurate rather than coarse: a coarse seek picks a byte offset by
783    // interpolating linearly over the file and then derives its reported
784    // timestamp from that same guess, so on VBR MP3 it lands seconds from the
785    // request *and* reports a position it never reached. Accurate mode walks
786    // frame headers and is truthful about both, for 1.5-3ms on files up to
787    // 79MB. The timeline records where playback actually resumed, so the
788    // transport shows the position being heard.
789    let mut seek_samples = 0;
790    if seek_ms > 0 {
791        let seeked = reader
792            .seek(
793                SeekMode::Accurate,
794                SeekTo::Time {
795                    time: Time::from_millis_u64(seek_ms),
796                    track_id: Some(track_id),
797                },
798            )
799            .map_err(|e| DecodeError::Decode(format!("seek failed: {}", e)))?;
800        seek_samples = landing_samples(time_base, seeked.actual_ts, sample_rate, channels)
801            .unwrap_or(seek_ms * sample_rate as u64 * channels as u64 / 1000);
802        if let Some(ref mut dec) = symphonia_decoder {
803            dec.reset();
804        }
805        if let Some(ref mut opus) = opus_bridge {
806            opus.reset();
807        }
808    }
809
810    // Record this track's boundary in the timeline.
811    let write_offset = timeline.samples_written();
812    timeline.push_boundary(TrackBoundary {
813        id: queue_item_id,
814        path: path.to_path_buf(),
815        info,
816        sample_offset: write_offset,
817        samples_written: 0,
818        seek_samples,
819    });
820
821    // Read ReplayGain tags and select the active gain for this track.
822    let rg_gain = if rg_mode != ReplayGainMode::Off {
823        match crate::audio::replaygain::read_tags(path) {
824            Ok(rg_info) => {
825                let selected = crate::audio::replaygain::select_gain(&rg_info, rg_mode);
826                if let Some((gain_db, _)) = selected {
827                    log::info!(
828                        "replaygain: applying {:.2} dB ({:?}) to {}",
829                        gain_db,
830                        rg_mode,
831                        path.display()
832                    );
833                }
834                selected
835            }
836            Err(e) => {
837                log::debug!("replaygain: no tags for {}: {}", path.display(), e);
838                None
839            }
840        }
841    } else {
842        None
843    };
844    let mut rg_scratch: Vec<f32> = Vec::new();
845
846    let mut sample_buf: Vec<f32> = Vec::new();
847
848    loop {
849        if stop.load(Ordering::Relaxed) || !timeline.is_current() {
850            return Ok(Decoded::Complete((sample_rate, channels)));
851        }
852
853        let packet = match reader.next_packet() {
854            Ok(Some(p)) => p,
855            Ok(None) => return Ok(Decoded::Complete((sample_rate, channels))),
856            Err(e) => return Err(DecodeError::Decode(e.to_string())),
857        };
858
859        if packet.track_id != track_id {
860            continue;
861        }
862
863        // Decode the packet — either via Opus bridge or Symphonia codec.
864        let samples: &[f32] = if let Some(ref mut opus) = opus_bridge {
865            match opus.decode_packet(&packet.data) {
866                Ok(s) => s,
867                Err(e) => {
868                    log::warn!("opus decode error (skipping packet): {}", e);
869                    continue;
870                }
871            }
872        } else {
873            let decoder = symphonia_decoder.as_mut().unwrap();
874            let decoded = match decoder.decode(&packet) {
875                Ok(d) => d,
876                Err(symphonia::core::errors::Error::DecodeError(e)) => {
877                    log::warn!("decode error (skipping packet): {}", e);
878                    continue;
879                }
880                Err(e) => return Err(DecodeError::Decode(e.to_string())),
881            };
882
883            let spec = decoded.spec();
884            let (decoded_rate, decoded_channels) = (spec.rate(), spec.channels().count() as u16);
885            // The engine is configured from the probed format. PCM that
886            // disagrees with it would play at the wrong speed, so end the
887            // session instead and let the player reconfigure.
888            if (decoded_rate, decoded_channels) != (sample_rate, channels) {
889                log::warn!(
890                    "{}: decoded {}Hz/{}ch but stream declares {}Hz/{}ch, restarting audio engine",
891                    path.display(),
892                    decoded_rate,
893                    decoded_channels,
894                    sample_rate,
895                    channels
896                );
897                return Ok(Decoded::FormatMismatch);
898            }
899            decoded.copy_to_vec_interleaved(&mut sample_buf);
900            &sample_buf[..]
901        };
902
903        if samples.is_empty() {
904            continue;
905        }
906
907        // Apply ReplayGain if active. Uses a reusable scratch buffer to avoid
908        // allocating per packet. Zero overhead when RG is off.
909        let samples = if let Some((gain_db, peak)) = rg_gain {
910            rg_scratch.clear();
911            rg_scratch.extend_from_slice(samples);
912            crate::audio::replaygain::apply_gain(&mut rg_scratch, gain_db, peak, pre_amp_db);
913            &rg_scratch[..]
914        } else {
915            samples
916        };
917
918        // Push samples into ring buffer, blocking if full.
919        // VizBuffer is updated incrementally inside this loop so it receives
920        // samples at the real-time audio consumption rate (paced by the audio
921        // callback draining the rtrb consumer), not in packet-sized bursts.
922        // Without this, FLAC packets (~93ms each at 44.1kHz) would update the
923        // viz buffer only ~11 times/sec, making waveform modes visibly choppy.
924        let mut offset = 0;
925        while offset < samples.len() {
926            // Also drops out on a stale generation, which keeps a dying thread
927            // from pushing into the viz delay line the next session just reset.
928            if stop.load(Ordering::Relaxed) || !timeline.is_current() {
929                return Ok(Decoded::Complete((sample_rate, channels)));
930            }
931
932            let slots = producer.slots();
933            if slots == 0 {
934                thread::sleep(std::time::Duration::from_micros(500));
935                continue;
936            }
937
938            let chunk_size = slots.min(samples.len() - offset);
939            if let Ok(mut chunk) = producer.write_chunk_uninit(chunk_size) {
940                let to_write = &samples[offset..offset + chunk_size];
941                let (first, second) = chunk.as_mut_slices();
942                let first_len = first.len().min(to_write.len());
943                for (slot, &val) in first.iter_mut().zip(&to_write[..first_len]) {
944                    slot.write(val);
945                }
946                if first_len < to_write.len() {
947                    for (slot, &val) in second.iter_mut().zip(&to_write[first_len..]) {
948                        slot.write(val);
949                    }
950                }
951                // SAFETY: All slots in the chunk have been initialized by the
952                // two loops above — first.len() + second.len() == chunk_size,
953                // and every slot is written via MaybeUninit::write().
954                unsafe { chunk.commit_all() };
955
956                // Feed viz buffer at the same rate as rtrb consumption.
957                if let Some(viz) = viz_buffer {
958                    viz.push_samples(to_write, channels, sample_rate);
959                }
960
961                offset += chunk_size;
962            }
963        }
964
965        timeline.add_written(samples.len() as u64);
966    }
967}
968
969/// Interleaved sample offset of a seek's landing point.
970///
971/// `actual_ts` is in the track's timebase, which is not always the reciprocal
972/// of the sample rate (Matroska ticks in milliseconds), so it is converted
973/// through `Time` rather than assumed to be a frame count.
974fn landing_samples(
975    time_base: Option<TimeBase>,
976    actual_ts: Timestamp,
977    sample_rate: u32,
978    channels: u16,
979) -> Option<u64> {
980    let (seconds, nanos) = time_base?.calc_time(actual_ts)?.parts();
981    let rate = sample_rate as u64;
982    let frames = seconds.max(0) as u64 * rate + (nanos as u64 * rate) / 1_000_000_000;
983    Some(frames * channels as u64)
984}
985
986/// Duration of a track in milliseconds.
987///
988/// The container's stated duration is authoritative because a track's timebase
989/// is not always the reciprocal of the sample rate — Matroska ticks in
990/// milliseconds, and states its duration at media level rather than per track.
991/// Falls back to the playable frame count when no duration is stated at all.
992pub(crate) fn track_duration_ms(
993    reader: &(impl FormatReader + ?Sized),
994    track: &Track,
995    sample_rate: u32,
996) -> u64 {
997    fn to_ms(time_base: Option<TimeBase>, duration: Option<Duration>) -> Option<u64> {
998        let time = time_base?.calc_duration(duration?)?;
999        Some(time.as_millis().max(0) as u64)
1000    }
1001
1002    let media = reader.media_info();
1003    to_ms(track.time_base, track.duration)
1004        .or_else(|| to_ms(media.time_base, media.duration))
1005        .or_else(|| {
1006            track
1007                .num_frames
1008                .map(|frames| frames * 1000 / sample_rate as u64)
1009        })
1010        .unwrap_or(0)
1011}
1012
1013/// Estimate bitrate (kbps) from Symphonia codec parameters.
1014///
1015/// Symphonia doesn't expose a `bit_rate` field. For lossy codecs like MP3/AAC
1016/// we can derive it from `bits_per_coded_sample` when the demuxer populates it.
1017/// Returns `None` for lossless codecs or when the info isn't available.
1018fn estimate_bitrate_from_codec_params(params: &AudioCodecParameters) -> Option<u32> {
1019    let is_lossy = matches!(
1020        params.codec,
1021        CODEC_ID_MP3 | CODEC_ID_AAC | CODEC_ID_VORBIS | CODEC_ID_OPUS
1022    );
1023    if !is_lossy {
1024        return None;
1025    }
1026
1027    // bits_per_coded_sample * sample_rate / 1000 gives kbps for CBR streams.
1028    // Few demuxers fill this in, but it's our best shot without file size.
1029    let bpcs = params.bits_per_coded_sample?;
1030    let sr = params.sample_rate?;
1031    let channels = params
1032        .channels
1033        .as_ref()
1034        .map(|c| c.count() as u32)
1035        .unwrap_or(2);
1036    Some(bpcs * sr * channels / 1000)
1037}
1038
1039pub fn codec_name(codec: AudioCodecId) -> String {
1040    match codec {
1041        CODEC_ID_FLAC => "FLAC",
1042        CODEC_ID_MP3 => "MP3",
1043        CODEC_ID_AAC => "AAC",
1044        CODEC_ID_VORBIS => "Vorbis",
1045        CODEC_ID_OPUS => "Opus",
1046        CODEC_ID_ALAC => "ALAC",
1047        CODEC_ID_PCM_S16LE => "PCM/16",
1048        CODEC_ID_PCM_S24LE => "PCM/24",
1049        CODEC_ID_PCM_S32LE => "PCM/32",
1050        CODEC_ID_PCM_F32LE => "PCM/f32",
1051        other => return format!("Unknown({:?})", other),
1052    }
1053    .to_string()
1054}
1055
1056#[cfg(test)]
1057mod tests {
1058    use std::path::PathBuf;
1059    use std::sync::atomic::Ordering;
1060
1061    use super::*;
1062    use crate::player::state::QueueItemId;
1063
1064    fn make_info(sample_rate: u32, channels: u16) -> StreamInfo {
1065        StreamInfo {
1066            codec: "FLAC".to_string(),
1067            sample_rate,
1068            channels,
1069            bit_depth: Some(16),
1070            bitrate_kbps: None,
1071            duration_ms: 10_000,
1072        }
1073    }
1074
1075    fn make_boundary(
1076        id: QueueItemId,
1077        sample_offset: u64,
1078        seek_samples: u64,
1079        channels: u16,
1080        sample_rate: u32,
1081    ) -> TrackBoundary {
1082        TrackBoundary {
1083            id,
1084            path: PathBuf::from("/music/track.flac"),
1085            info: make_info(sample_rate, channels),
1086            sample_offset,
1087            samples_written: 0,
1088            seek_samples,
1089        }
1090    }
1091
1092    /// A writer for the timeline's current generation.
1093    fn writer(timeline: &PlaybackTimeline) -> TimelineWriter<'_> {
1094        timeline.writer(timeline.generation())
1095    }
1096
1097    // --- PlaybackTimeline tests ---
1098
1099    #[test]
1100    fn test_timeline_single_track() {
1101        // Push one boundary at offset 0 with stereo 44100 Hz audio.
1102        // After simulating 44100 frames (88200 interleaved samples) played,
1103        // current_playback() should report track index 0 at position 1000 ms.
1104        let timeline = PlaybackTimeline::new();
1105        let tl = writer(&timeline);
1106        let id = QueueItemId::new();
1107        // sample_offset=0, seek_samples=0, channels=2, sample_rate=44100
1108        tl.push_boundary(make_boundary(id, 0, 0, 2, 44100));
1109        tl.add_written(88200); // 1 second of audio
1110
1111        // Simulate 1 second played: 44100 frames * 2 channels = 88200 interleaved samples
1112        timeline.samples_played.store(88200, Ordering::Relaxed);
1113
1114        let result = timeline.current_playback();
1115        assert!(
1116            result.is_some(),
1117            "expected Some for single track with samples played"
1118        );
1119        let (result_id, _path, _info, position_ms) = result.unwrap();
1120        assert_eq!(result_id, id);
1121        assert_eq!(
1122            position_ms, 1000,
1123            "1 second of 44100 Hz stereo should be 1000 ms"
1124        );
1125    }
1126
1127    #[test]
1128    fn test_timeline_gapless_transition() {
1129        // Two tracks in gapless sequence. Track 1 ends at sample 88200 (1 sec stereo 44100 Hz).
1130        // Track 2 begins at sample_offset 88200. When playback head is at 100000 (past the boundary),
1131        // current_playback() should report track 2.
1132        let timeline = PlaybackTimeline::new();
1133        let tl = writer(&timeline);
1134        let id1 = QueueItemId::new();
1135        let id2 = QueueItemId::new();
1136
1137        // Track 1: starts at offset 0
1138        tl.push_boundary(make_boundary(id1, 0, 0, 2, 44100));
1139        tl.add_written(88200);
1140
1141        // Track 2: starts at offset 88200 (immediately after track 1's samples)
1142        tl.push_boundary(make_boundary(id2, 88200, 0, 2, 44100));
1143        tl.add_written(44100); // half a second of track 2
1144
1145        // Set playback head past the track 1/2 boundary
1146        timeline.samples_played.store(90000, Ordering::Relaxed);
1147
1148        let result = timeline.current_playback();
1149        assert!(result.is_some());
1150        let (result_id, _path, _info, position_ms) = result.unwrap();
1151        assert_eq!(
1152            result_id, id2,
1153            "playback head past boundary should report second track"
1154        );
1155        // (90000 - 88200) / 2 channels * 1000 / 44100 = 900 / 44100 ≈ 20 ms
1156        assert_eq!(position_ms, 20, "position within track 2 should be ~20 ms");
1157    }
1158
1159    #[test]
1160    fn test_timeline_zero_samples() {
1161        // With 0 samples played and a boundary at offset 0, current_playback() should
1162        // still return the first track at position 0 ms.
1163        let timeline = PlaybackTimeline::new();
1164        let tl = writer(&timeline);
1165        let id = QueueItemId::new();
1166        tl.push_boundary(make_boundary(id, 0, 0, 2, 44100));
1167        tl.add_written(1000);
1168        timeline.samples_played.store(0, Ordering::Relaxed);
1169
1170        let result = timeline.current_playback();
1171        assert!(
1172            result.is_some(),
1173            "expected Some at 0 samples played with a boundary at offset 0"
1174        );
1175        let (result_id, _path, _info, position_ms) = result.unwrap();
1176        assert_eq!(result_id, id);
1177        assert_eq!(position_ms, 0);
1178    }
1179
1180    #[test]
1181    fn test_timeline_past_all_boundaries() {
1182        // When samples_played exceeds all boundaries, the last track should be reported.
1183        // The binary search finds the last boundary whose sample_offset <= played.
1184        let timeline = PlaybackTimeline::new();
1185        let tl = writer(&timeline);
1186        let id1 = QueueItemId::new();
1187        let id2 = QueueItemId::new();
1188
1189        tl.push_boundary(make_boundary(id1, 0, 0, 2, 44100));
1190        tl.add_written(88200);
1191        tl.push_boundary(make_boundary(id2, 88200, 0, 2, 44100));
1192        tl.add_written(88200);
1193
1194        // Simulate playback far past both tracks
1195        timeline
1196            .samples_played
1197            .store(999_999_999, Ordering::Relaxed);
1198
1199        let result = timeline.current_playback();
1200        assert!(result.is_some());
1201        let (result_id, _path, _info, _position_ms) = result.unwrap();
1202        assert_eq!(
1203            result_id, id2,
1204            "samples past all boundaries should report the last track"
1205        );
1206    }
1207
1208    #[test]
1209    fn test_timeline_seek_offset() {
1210        // When a seek offset is set, position_ms should include the seek position.
1211        // seek_samples = 88200 means playback started 1 second into the track.
1212        // With 0 additional samples played past the boundary, position should be 1000 ms.
1213        let timeline = PlaybackTimeline::new();
1214        let tl = writer(&timeline);
1215        let id = QueueItemId::new();
1216        let seek_samples = 88200u64; // 1 second at 44100 Hz stereo
1217        tl.push_boundary(make_boundary(id, 0, seek_samples, 2, 44100));
1218        tl.add_written(44100); // half a second written so far
1219        // samples_played at the track boundary (0 frames past the track start)
1220        timeline.samples_played.store(0, Ordering::Relaxed);
1221
1222        let result = timeline.current_playback();
1223        assert!(result.is_some());
1224        let (_result_id, _path, _info, position_ms) = result.unwrap();
1225        // track_samples = 0 - 0 = 0; seek contribution = (88200/2)*1000/44100 = 1000 ms
1226        assert_eq!(
1227            position_ms, 1000,
1228            "position should include seek offset of 1000 ms"
1229        );
1230    }
1231
1232    #[test]
1233    fn test_timeline_reset() {
1234        // After reset(), current_playback() returns None and all counters are cleared.
1235        let timeline = PlaybackTimeline::new();
1236        let tl = writer(&timeline);
1237        let id = QueueItemId::new();
1238        tl.push_boundary(make_boundary(id, 0, 0, 2, 44100));
1239        tl.add_written(88200);
1240        timeline.samples_played.store(44100, Ordering::Relaxed);
1241
1242        // Sanity check: playback is live before reset
1243        assert!(timeline.current_playback().is_some());
1244
1245        timeline.reset();
1246
1247        assert!(
1248            timeline.current_playback().is_none(),
1249            "after reset, current_playback should return None"
1250        );
1251        assert_eq!(
1252            timeline.samples_played.load(Ordering::Relaxed),
1253            0,
1254            "samples_played should be 0 after reset"
1255        );
1256        assert_eq!(
1257            timeline.samples_written.load(Ordering::Relaxed),
1258            0,
1259            "samples_written should be 0 after reset"
1260        );
1261    }
1262
1263    // --- Probe and decode integration tests ---
1264
1265    #[test]
1266    fn probe_file_extracts_stream_info() {
1267        let dir = tempfile::tempdir().unwrap();
1268        let wav_path = dir.path().join("probe_test.wav");
1269        crate::test_utils::generate_wav(&wav_path, 44100, 2, 1.0, 16);
1270
1271        let info = probe_file(&wav_path).expect("probe_file should succeed on a valid WAV");
1272        assert_eq!(info.sample_rate, 44100, "sample rate mismatch");
1273        assert_eq!(info.channels, 2, "channel count mismatch");
1274        assert_eq!(info.bit_depth, Some(16), "bit depth mismatch");
1275        assert!(
1276            info.duration_ms > 900 && info.duration_ms < 1100,
1277            "duration should be ~1000ms, got {}",
1278            info.duration_ms
1279        );
1280        assert!(
1281            info.codec.contains("PCM"),
1282            "codec should be PCM variant, got {}",
1283            info.codec
1284        );
1285    }
1286
1287    #[test]
1288    fn decode_single_produces_samples() {
1289        let dir = tempfile::tempdir().unwrap();
1290        let wav_path = dir.path().join("tone.wav");
1291        // 440 Hz sine, mono, 0.1s — enough to verify non-zero decode output.
1292        crate::test_utils::generate_wav_tone(&wav_path, 44100, 440.0, 0.1);
1293
1294        // Set up rtrb ring buffer.
1295        let (mut producer, mut consumer) = rtrb::RingBuffer::new(44100 * 2);
1296
1297        let timeline = PlaybackTimeline::new();
1298        let tl = writer(&timeline);
1299        let stop = Arc::new(AtomicBool::new(false));
1300
1301        let id = QueueItemId::new();
1302        let entry = SourceEntry::from_file(id, wav_path.clone());
1303        let hint = entry.hint.clone();
1304        let mss = (entry.make_mss)().expect("should open WAV file");
1305
1306        let result = decode_single(
1307            id,
1308            &wav_path,
1309            &hint,
1310            mss,
1311            &mut producer,
1312            &stop,
1313            0,
1314            &tl,
1315            None,
1316            crate::config::ReplayGainMode::Off,
1317            0.0,
1318            None,
1319        );
1320        assert!(
1321            matches!(result, Ok(Decoded::Complete((44100, 1)))),
1322            "decode_single should complete at the source format"
1323        );
1324
1325        // Read samples from the consumer side.
1326        let available = consumer.slots();
1327        assert!(available > 0, "expected samples in ring buffer, got 0");
1328
1329        // Verify at least some samples are non-zero (it's a sine wave, not silence).
1330        let mut found_nonzero = false;
1331        while consumer.slots() > 0 {
1332            if let Ok(chunk) = consumer.read_chunk(consumer.slots().min(1024)) {
1333                let (first, second) = chunk.as_slices();
1334                for &s in first.iter().chain(second.iter()) {
1335                    if s.abs() > 0.001 {
1336                        found_nonzero = true;
1337                        break;
1338                    }
1339                }
1340                chunk.commit_all();
1341            }
1342            if found_nonzero {
1343                break;
1344            }
1345        }
1346        assert!(
1347            found_nonzero,
1348            "expected non-zero samples from 440Hz sine decode"
1349        );
1350    }
1351
1352    // --- Ring buffer format contract ---
1353
1354    /// Decode a queue of files through `decode_queue_loop` with a consumer
1355    /// draining in the background. Returns the boundaries the decode thread
1356    /// pushed onto the timeline.
1357    fn run_queue(paths: &[PathBuf]) -> Vec<TrackBoundary> {
1358        let (producer, mut consumer) = rtrb::RingBuffer::new(1 << 16);
1359        let timeline = PlaybackTimeline::new();
1360        let tl = writer(&timeline);
1361        let stop = Arc::new(AtomicBool::new(false));
1362
1363        let drain_stop = Arc::new(AtomicBool::new(false));
1364        let drain_flag = drain_stop.clone();
1365        let drainer = std::thread::spawn(move || {
1366            while !drain_flag.load(Ordering::Relaxed) {
1367                let n = consumer.slots();
1368                if n > 0
1369                    && let Ok(chunk) = consumer.read_chunk(n)
1370                {
1371                    chunk.commit_all();
1372                }
1373                std::thread::sleep(std::time::Duration::from_micros(200));
1374            }
1375        });
1376
1377        let rest: std::sync::Mutex<Vec<PathBuf>> = std::sync::Mutex::new(paths[1..].to_vec());
1378        let next_track = move || {
1379            let mut rest = rest.lock().ok()?;
1380            if rest.is_empty() {
1381                return None;
1382            }
1383            Some(SourceEntry::from_file(QueueItemId::new(), rest.remove(0)))
1384        };
1385
1386        decode_queue_loop(
1387            SourceEntry::from_file(QueueItemId::new(), paths[0].clone()),
1388            producer,
1389            &stop,
1390            0,
1391            &next_track,
1392            &tl,
1393            None,
1394            crate::config::ReplayGainMode::Off,
1395            0.0,
1396        );
1397
1398        drain_stop.store(true, Ordering::Relaxed);
1399        drainer.join().unwrap();
1400
1401        timeline.boundaries.read().clone()
1402    }
1403
1404    #[test]
1405    fn gapless_continues_when_format_matches() {
1406        let dir = tempfile::tempdir().unwrap();
1407        let a = dir.path().join("a.wav");
1408        let b = dir.path().join("b.wav");
1409        crate::test_utils::generate_wav(&a, 44100, 2, 0.1, 16);
1410        crate::test_utils::generate_wav(&b, 44100, 2, 0.1, 16);
1411
1412        let bounds = run_queue(&[a, b]);
1413        assert_eq!(
1414            bounds.len(),
1415            2,
1416            "same-format tracks should decode gaplessly"
1417        );
1418    }
1419
1420    #[test]
1421    fn gapless_stops_at_sample_rate_change() {
1422        let dir = tempfile::tempdir().unwrap();
1423        let a = dir.path().join("a.wav");
1424        let b = dir.path().join("b.wav");
1425        crate::test_utils::generate_wav(&a, 44100, 2, 0.1, 16);
1426        crate::test_utils::generate_wav(&b, 48000, 2, 0.1, 16);
1427
1428        let bounds = run_queue(&[a, b]);
1429        assert_eq!(
1430            bounds.len(),
1431            1,
1432            "a 48kHz track must not join a 44.1kHz ring buffer"
1433        );
1434        assert_eq!(bounds[0].info.sample_rate, 44100);
1435    }
1436
1437    #[test]
1438    fn gapless_stops_at_channel_change() {
1439        let dir = tempfile::tempdir().unwrap();
1440        let a = dir.path().join("a.wav");
1441        let b = dir.path().join("b.wav");
1442        crate::test_utils::generate_wav(&a, 44100, 2, 0.1, 16);
1443        crate::test_utils::generate_wav(&b, 44100, 1, 0.1, 16);
1444
1445        let bounds = run_queue(&[a, b]);
1446        assert_eq!(
1447            bounds.len(),
1448            1,
1449            "a mono track must not join a stereo ring buffer"
1450        );
1451        assert_eq!(bounds[0].info.channels, 2);
1452    }
1453
1454    #[test]
1455    fn drain_waits_for_the_consumer() {
1456        let (mut producer, mut consumer) = rtrb::RingBuffer::new(64);
1457        for _ in 0..64 {
1458            producer.push(0.0).unwrap();
1459        }
1460        let stop = Arc::new(AtomicBool::new(false));
1461
1462        let reader = std::thread::spawn(move || {
1463            std::thread::sleep(std::time::Duration::from_millis(20));
1464            let chunk = consumer.read_chunk(64).unwrap();
1465            chunk.commit_all();
1466            consumer
1467        });
1468
1469        wait_for_drain(&producer, &stop);
1470        assert_eq!(producer.slots(), 64, "drain must wait for an empty buffer");
1471        drop(reader.join().unwrap());
1472    }
1473
1474    #[test]
1475    fn drain_returns_when_playback_is_torn_down() {
1476        let (producer, consumer) = rtrb::RingBuffer::<f32>::new(64);
1477        let stop = Arc::new(AtomicBool::new(true));
1478        wait_for_drain(&producer, &stop);
1479        drop(consumer);
1480    }
1481
1482    // --- Generation guard on timeline writes ---
1483
1484    #[test]
1485    fn a_writer_knows_when_its_session_has_ended() {
1486        let timeline = PlaybackTimeline::new();
1487        let tl = writer(&timeline);
1488        assert!(tl.is_current());
1489
1490        timeline.reset();
1491        assert!(!tl.is_current(), "reset must retire the outgoing writer");
1492        assert!(writer(&timeline).is_current());
1493    }
1494
1495    #[test]
1496    fn stale_writes_are_dropped_after_reset() {
1497        let timeline = PlaybackTimeline::new();
1498        let dying = writer(&timeline);
1499        dying.push_boundary(make_boundary(QueueItemId::new(), 0, 0, 2, 44100));
1500        dying.add_written(88200);
1501
1502        timeline.reset();
1503
1504        // The outgoing decode thread checks `stop` only at the top of its chunk
1505        // loop, so its last packet lands after the reset.
1506        dying.add_written(4608);
1507        dying.push_boundary(make_boundary(QueueItemId::new(), 0, 0, 2, 44100));
1508
1509        assert_eq!(timeline.samples_written.load(Ordering::Relaxed), 0);
1510        assert!(timeline.boundaries.read().is_empty());
1511    }
1512
1513    #[test]
1514    fn a_dying_decode_thread_cannot_blank_the_transport() {
1515        let timeline = PlaybackTimeline::new();
1516        let dying = writer(&timeline);
1517        dying.push_boundary(make_boundary(QueueItemId::new(), 0, 0, 2, 44100));
1518        dying.add_written(88200);
1519
1520        // A skip: reset, then the old thread's final packet, then the new
1521        // session stamps its first boundary at whatever the counter now says.
1522        timeline.reset();
1523        dying.add_written(4608);
1524
1525        let fresh = writer(&timeline);
1526        let id = QueueItemId::new();
1527        let write_offset = fresh.samples_written();
1528        fresh.push_boundary(make_boundary(id, write_offset, 0, 2, 44100));
1529
1530        assert_eq!(write_offset, 0, "first boundary must start at 0");
1531        let (playing, _, _, position_ms) = timeline
1532            .current_playback()
1533            .expect("transport must not go blank at samples_played = 0");
1534        assert_eq!(playing, id);
1535        assert_eq!(position_ms, 0);
1536    }
1537
1538    // --- Failure handling in the gapless queue ---
1539
1540    /// A file that exists and has an audio extension but no audio in it.
1541    fn write_garbage(path: &Path) {
1542        std::fs::write(path, b"this is not a wav file").unwrap();
1543    }
1544
1545    #[test]
1546    fn an_unreadable_track_is_skipped_and_the_queue_continues() {
1547        let dir = tempfile::tempdir().unwrap();
1548        let a = dir.path().join("a.wav");
1549        let bad = dir.path().join("bad.wav");
1550        let c = dir.path().join("c.wav");
1551        crate::test_utils::generate_wav(&a, 44100, 2, 0.1, 16);
1552        write_garbage(&bad);
1553        crate::test_utils::generate_wav(&c, 44100, 2, 0.1, 16);
1554
1555        let bounds = run_queue(&[a.clone(), bad, c.clone()]);
1556        let decoded: Vec<_> = bounds.iter().map(|b| b.path.clone()).collect();
1557        assert_eq!(
1558            decoded,
1559            vec![a, c],
1560            "one bad file must not take the rest of the queue with it"
1561        );
1562    }
1563
1564    #[test]
1565    fn a_missing_track_is_skipped_and_the_queue_continues() {
1566        let dir = tempfile::tempdir().unwrap();
1567        let missing = dir.path().join("gone.wav");
1568        let b = dir.path().join("b.wav");
1569        crate::test_utils::generate_wav(&b, 44100, 2, 0.1, 16);
1570
1571        let bounds = run_queue(&[missing, b.clone()]);
1572        assert_eq!(bounds.len(), 1);
1573        assert_eq!(
1574            bounds[0].path, b,
1575            "a bad first track must not end the session"
1576        );
1577    }
1578
1579    #[test]
1580    fn an_entirely_unreadable_queue_terminates() {
1581        let dir = tempfile::tempdir().unwrap();
1582        let bad = dir.path().join("bad.wav");
1583        write_garbage(&bad);
1584
1585        // Nothing decodes, so the ring stays empty; the consumer only has to
1586        // outlive the producer.
1587        let (producer, _consumer) = rtrb::RingBuffer::<f32>::new(1 << 12);
1588        let timeline = PlaybackTimeline::new();
1589        let tl = writer(&timeline);
1590        let stop = Arc::new(AtomicBool::new(false));
1591
1592        let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
1593        let counter = calls.clone();
1594        let bad_path = bad.clone();
1595        let next_track = move || {
1596            counter.fetch_add(1, Ordering::Relaxed);
1597            Some(SourceEntry::from_file(QueueItemId::new(), bad_path.clone()))
1598        };
1599
1600        decode_queue_loop(
1601            SourceEntry::from_file(QueueItemId::new(), bad),
1602            producer,
1603            &stop,
1604            0,
1605            &next_track,
1606            &tl,
1607            None,
1608            crate::config::ReplayGainMode::Off,
1609            0.0,
1610        );
1611
1612        // The first source plus one per next_track call, capped.
1613        assert_eq!(
1614            calls.load(Ordering::Relaxed) + 1,
1615            MAX_CONSECUTIVE_FAILURES as usize
1616        );
1617        assert!(timeline.boundaries.read().is_empty());
1618    }
1619
1620    // --- Seek landing position ---
1621
1622    #[test]
1623    fn landing_samples_converts_frame_timebases() {
1624        // The usual audio case: one tick per frame.
1625        let tb = TimeBase::try_from_recip(44100).unwrap();
1626        assert_eq!(
1627            landing_samples(Some(tb), Timestamp::from(44100u32), 44100, 2),
1628            Some(88_200)
1629        );
1630    }
1631
1632    #[test]
1633    fn landing_samples_converts_millisecond_timebases() {
1634        // Matroska ticks in milliseconds, not frames.
1635        let tb = TimeBase::try_new(1, 1000).unwrap();
1636        assert_eq!(
1637            landing_samples(Some(tb), Timestamp::from(1500u32), 48000, 2),
1638            Some(48000 * 3 / 2 * 2)
1639        );
1640    }
1641
1642    #[test]
1643    fn landing_samples_needs_a_timebase() {
1644        assert_eq!(
1645            landing_samples(None, Timestamp::from(1000u32), 44100, 2),
1646            None
1647        );
1648    }
1649
1650    /// Build a 300s VBR MP3: 30s of near-silence then a loud tone, which gives
1651    /// lame a wide enough bitrate spread to make coarse seeking miss badly.
1652    #[cfg(test)]
1653    fn make_vbr_mp3(dir: &Path) -> PathBuf {
1654        let wav = dir.join("source.wav");
1655        let mp3 = dir.join("source.mp3");
1656        let ok = std::process::Command::new("sox")
1657            .args(["-n", "-r", "44100", "-c", "2"])
1658            .arg(&wav)
1659            .args([
1660                "synth", "30", "sine", "200", "vol", "0.02", ":", "synth", "270", "sine", "880",
1661                "vol", "0.9",
1662            ])
1663            .status()
1664            .expect("sox not installed")
1665            .success();
1666        assert!(ok, "sox failed");
1667        let ok = std::process::Command::new("lame")
1668            .args(["-V", "2", "--quiet"])
1669            .arg(&wav)
1670            .arg(&mp3)
1671            .status()
1672            .expect("lame not installed")
1673            .success();
1674        assert!(ok, "lame failed");
1675        mp3
1676    }
1677
1678    /// A VBR seek must report where playback actually resumed: the reported
1679    /// start plus the audio that actually followed has to add back up to the
1680    /// file's duration. Under a coarse seek this file lands 3.7s late while
1681    /// reporting 83ms early — a 3.8s lie for the rest of the track.
1682    #[test]
1683    #[ignore = "generates a fixture with sox + lame; run with cargo test -- --ignored"]
1684    fn seek_on_vbr_reports_where_it_landed() {
1685        let dir = tempfile::tempdir().unwrap();
1686        let path = make_vbr_mp3(dir.path());
1687        let info = probe_file(&path).unwrap();
1688        let channels = info.channels as u64;
1689        let rate = info.sample_rate as u64;
1690
1691        let seek_ms = 150_000u64;
1692        let (mut producer, mut consumer) = rtrb::RingBuffer::<f32>::new(1 << 16);
1693        let stop = Arc::new(AtomicBool::new(false));
1694        let timeline = PlaybackTimeline::new();
1695        let tl = writer(&timeline);
1696
1697        let drain_stop = stop.clone();
1698        let drained = std::thread::spawn(move || {
1699            let mut total = 0u64;
1700            while !drain_stop.load(Ordering::Relaxed) {
1701                let slots = consumer.slots();
1702                if slots == 0 {
1703                    std::thread::sleep(std::time::Duration::from_micros(200));
1704                    continue;
1705                }
1706                let chunk = consumer.read_chunk(slots).unwrap();
1707                total += slots as u64;
1708                chunk.commit_all();
1709            }
1710            total
1711        });
1712
1713        let file = File::open(&path).unwrap();
1714        let mss = MediaSourceStream::new(Box::new(file), Default::default());
1715        let mut hint = Hint::new();
1716        hint.with_extension("mp3");
1717        decode_single(
1718            QueueItemId::new(),
1719            &path,
1720            &hint,
1721            mss,
1722            &mut producer,
1723            &stop,
1724            seek_ms,
1725            &tl,
1726            None,
1727            ReplayGainMode::Off,
1728            0.0,
1729            None,
1730        )
1731        .unwrap();
1732
1733        let written = timeline.samples_written.load(Ordering::Relaxed);
1734        stop.store(true, Ordering::Relaxed);
1735        drained.join().unwrap();
1736
1737        let reported_start_ms = {
1738            let bounds = timeline.boundaries.read();
1739            (bounds[0].seek_samples / channels) * 1000 / rate
1740        };
1741        let decoded_ms = (written / channels) * 1000 / rate;
1742
1743        // Where playback started plus how much audio followed is the whole file.
1744        let total_ms = reported_start_ms + decoded_ms;
1745        assert!(
1746            total_ms.abs_diff(info.duration_ms) < 500,
1747            "reported start {}ms + {}ms decoded = {}ms, but the file is {}ms",
1748            reported_start_ms,
1749            decoded_ms,
1750            total_ms,
1751            info.duration_ms
1752        );
1753        // Landing is frame-granular, never sample-exact.
1754        assert!(
1755            reported_start_ms.abs_diff(seek_ms) < 100,
1756            "seek to {}ms reported {}ms",
1757            seek_ms,
1758            reported_start_ms
1759        );
1760    }
1761}