Skip to main content

koan_core/player/
mod.rs

1pub mod commands;
2pub mod history;
3pub mod state;
4pub mod undo;
5
6use std::path::{Path, PathBuf};
7use std::sync::Arc;
8use std::thread;
9
10use thiserror::Error;
11
12use crate::audio::{
13    analyzer::VizAnalyzer,
14    backend::{self, AudioBackend, AudioEngineHandle, BackendError, SampleRateWatch},
15    buffer, streaming,
16    viz::{VizBuffer, VizSnapshot},
17};
18use crate::remote::client::PlaybackReportState;
19use buffer::PlaybackTimeline;
20use commands::{CommandChannel, PlayerCommand};
21use history::{InFlight, PlayEvent, PlayRecorder, PlaybackReport};
22use state::{
23    ItemState, LoadState, PlaybackSource, PlaybackState, QueueItemId, SharedPlayerState, TrackInfo,
24};
25use undo::{UndoEntry, UndoStack};
26
27/// Ring buffer size in samples. ~1s at 192kHz stereo.
28pub(crate) const RING_BUFFER_SIZE: usize = 192_000 * 2;
29
30/// Kept back from the end of a track when seeking, so dragging the thumb all
31/// the way over lands in the last moment of it rather than in the next track.
32const SEEK_END_GUARD_MS: u64 = 500;
33
34#[derive(Debug, Error)]
35pub enum PlayerError {
36    #[error("backend error: {0}")]
37    Backend(#[from] BackendError),
38    #[error("decode error: {0}")]
39    Decode(#[from] buffer::DecodeError),
40}
41
42/// Everything needed to read a track that is still downloading: where it is,
43/// how far the transfer has got, and how the container has to be opened.
44#[derive(Clone)]
45struct StreamSource {
46    path: PathBuf,
47    bytes_written: Arc<crate::remote::downloads::ByteFeed>,
48    total: u64,
49    mode: streaming::ProbeMode,
50}
51
52/// A file's own extension, lowercased. A download in progress is named
53/// `track.m4a.part`, and its extension is the track's, not `part`.
54fn media_extension(path: &Path) -> Option<String> {
55    crate::remote::download::strip_part_suffix(path)
56        .extension()
57        .and_then(|e| e.to_str())
58        .map(str::to_ascii_lowercase)
59}
60
61/// Symphonia's format hint for a path — its extension, where it has one.
62fn hint_for(path: &Path) -> symphonia::core::formats::probe::Hint {
63    let mut hint = symphonia::core::formats::probe::Hint::new();
64    if let Some(ext) = media_extension(path) {
65        hint.with_extension(&ext);
66    }
67    hint
68}
69
70/// How to open a partial file once the whole description has failed.
71///
72/// Most containers describe their frames from the front and need an end they
73/// can actually reach: FLAC bisects towards the end it is given. Two need the
74/// whole file's end instead. Ogg takes the end it is handed as the end of the
75/// stream, so a file ending at the write head is a track already over. MP4
76/// bounds its top-level boxes by the end, so a `moov` larger than what has
77/// arrived overruns it and the file will not open until it has all landed.
78/// See `ProbeMode`.
79fn lengthless_mode_for(path: &Path) -> streaming::ProbeMode {
80    match media_extension(path).as_deref() {
81        Some("ogg" | "oga" | "opus" | "spx" | "m4a" | "m4b" | "mp4" | "mov") => {
82            streaming::ProbeMode::LengthlessWholeEnd
83        }
84        _ => streaming::ProbeMode::Lengthless,
85    }
86}
87
88/// The player controller. Owns the audio pipeline and processes commands.
89pub struct Player {
90    shared_state: Arc<SharedPlayerState>,
91    commands: CommandChannel,
92    active_playback: Option<ActivePlayback>,
93    timeline: Arc<PlaybackTimeline>,
94    viz_buffer: Arc<VizBuffer>,
95    viz_snapshot: Arc<VizSnapshot>,
96    /// Background FFT analysis thread. Held for its lifetime; dropped on Player drop.
97    _viz_analyzer: VizAnalyzer,
98    undo_stack: UndoStack,
99    /// When Some, undo entries are collected into this buffer instead of pushed
100    /// directly onto the undo stack. Flushed on EndUndoBatch.
101    batch_buffer: Option<Vec<UndoEntry>>,
102    /// Configured output device name. None = system default.
103    output_device_name: Option<String>,
104    /// Platform audio backend (CoreAudio on macOS and iOS, cpal on Linux).
105    backend: Box<dyn AudioBackend>,
106    /// Debounce: timestamp of last NextTrack/PrevTrack to suppress key repeat.
107    last_skip: std::time::Instant,
108    /// How the file currently streaming had to be opened. A seek reopens it and
109    /// must not undo what the probe settled on.
110    stream_mode: streaming::ProbeMode,
111    /// Writes plays away from this thread. None when there is no database to
112    /// write to, and in tests, which must not touch the real library.
113    history: Option<PlayRecorder>,
114    /// How much of the current track has been heard so far.
115    in_flight: Option<InFlight>,
116    /// Playback sessions started — lets tests assert how many engine restarts
117    /// an operation costs.
118    #[cfg(test)]
119    playback_starts: usize,
120}
121
122/// Holds the resources for an active playback session.
123struct ActivePlayback {
124    engine: Box<dyn AudioEngineHandle>,
125    decode_handle: buffer::DecodeHandle,
126    /// Set when the decoder is reading a download as it arrives.
127    stream: Option<LiveStream>,
128    /// Keeps the device rate subscription alive for as long as this engine is
129    /// the one feeding the DAC. Dropped with it.
130    _rate_watch: Option<Box<dyn SampleRateWatch>>,
131}
132
133/// A download being decoded as it lands. The reader may be parked at the write
134/// head waiting for bytes, so stopping has to tell it to give up and wake it,
135/// or the join waits on the network.
136struct LiveStream {
137    feed: Arc<crate::remote::downloads::ByteFeed>,
138    abandoned: Arc<std::sync::atomic::AtomicBool>,
139}
140
141impl LiveStream {
142    fn abandon(&self) {
143        self.abandoned
144            .store(true, std::sync::atomic::Ordering::Release);
145        self.feed.done();
146    }
147}
148
149impl Default for Player {
150    fn default() -> Self {
151        Self::new()
152    }
153}
154
155impl Player {
156    pub fn new() -> Self {
157        let viz_buffer = VizBuffer::new();
158        let viz_snapshot = VizSnapshot::new();
159        let timeline = PlaybackTimeline::new();
160        let cfg = crate::config::Config::load_or_default();
161        let viz_analyzer = VizAnalyzer::spawn_with_snapshot(
162            Arc::clone(&viz_buffer),
163            &cfg.visualizer,
164            Arc::clone(&viz_snapshot),
165            timeline.samples_played_counter(),
166        );
167
168        Self {
169            shared_state: SharedPlayerState::new(),
170            commands: CommandChannel::new(),
171            active_playback: None,
172            timeline,
173            viz_buffer,
174            viz_snapshot,
175            _viz_analyzer: viz_analyzer,
176            undo_stack: UndoStack::new(),
177            batch_buffer: None,
178            output_device_name: cfg.playback.output_device,
179            backend: crate::audio::platform_backend(),
180            last_skip: std::time::Instant::now(),
181            stream_mode: streaming::ProbeMode::Full,
182            history: None,
183            in_flight: None,
184            #[cfg(test)]
185            playback_starts: 0,
186        }
187    }
188
189    /// Get a clone of the shared state for UI reads.
190    pub fn shared_state(&self) -> Arc<SharedPlayerState> {
191        self.shared_state.clone()
192    }
193
194    /// Get the playback timeline for UI reads.
195    pub fn timeline(&self) -> Arc<PlaybackTimeline> {
196        self.timeline.clone()
197    }
198
199    /// Get the visualization buffer for the TUI.
200    pub fn viz_buffer(&self) -> Arc<VizBuffer> {
201        self.viz_buffer.clone()
202    }
203
204    /// Get the shared analysis snapshot for the TUI.
205    /// The analysis thread writes here; the UI thread reads a clone each frame.
206    pub fn viz_snapshot(&self) -> Arc<VizSnapshot> {
207        self.viz_snapshot.clone()
208    }
209
210    /// Access undo stack (for tests and UI state queries).
211    pub fn undo_stack(&self) -> &UndoStack {
212        &self.undo_stack
213    }
214
215    /// Create an audio engine for a stream, switching the output device to the
216    /// source rate first so output is bit-perfect.
217    ///
218    /// The engine is always configured with the source's own rate and channel
219    /// count — the format the decode thread writes into the ring buffer. A
220    /// device that cannot take the requested rate (MPEG-2/2.5 MP3 rates are
221    /// commonly refused) resamples instead of playing at the wrong speed.
222    #[allow(clippy::type_complexity)]
223    fn create_engine_for(
224        &self,
225        info: &buffer::StreamInfo,
226        consumer: rtrb::Consumer<f32>,
227    ) -> Result<(Box<dyn AudioEngineHandle>, Option<Box<dyn SampleRateWatch>>), PlayerError> {
228        let device = self.resolve_device()?;
229        let device_rate = self.backend.get_device_sample_rate(&device)?;
230        let source_rate = info.sample_rate as f64;
231
232        // The track info is already published, so anything read between here
233        // and the switch landing would pair this track with the last one's
234        // output rate — and a rate switch is not instant. Say nothing instead.
235        self.shared_state.clear_output_sample_rate();
236
237        let settled = if (device_rate - source_rate).abs() > 0.1 {
238            log::info!(
239                "switching device sample rate: {}Hz → {}Hz",
240                device_rate,
241                source_rate
242            );
243            match self.backend.set_device_sample_rate(&device, source_rate) {
244                Ok(rate) => rate,
245                Err(e) => {
246                    log::warn!("failed to set device sample rate: {}", e);
247                    device_rate
248                }
249            }
250        } else {
251            device_rate
252        };
253
254        if (settled - source_rate).abs() > 0.1 {
255            log::warn!(
256                "device stayed at {}Hz (wanted {}Hz) — output is resampled, not bit-perfect",
257                settled,
258                source_rate
259            );
260        }
261
262        // The front ends compare this against the source rate to say whether
263        // anything had to resample.
264        self.shared_state
265            .set_output_sample_rate(settled.round() as u32);
266
267        // koan is not the only client of this device. Subscribe so the front
268        // ends learn about a rate someone else moved instead of trusting the
269        // reading above until the next track happens to build an engine.
270        let watch_state = self.shared_state.clone();
271        let watch_name = device.name.clone();
272        let rate_watch = self.backend.watch_device_sample_rate(
273            &device,
274            Box::new(move |rate| {
275                log::info!("device sample rate changed externally: {rate}Hz on '{watch_name}'");
276                watch_state.set_output_sample_rate(rate.round() as u32);
277            }),
278        );
279
280        let engine = self.backend.create_engine(
281            &device,
282            source_rate,
283            info.channels as u32,
284            consumer,
285            self.timeline.samples_played_counter(),
286        )?;
287
288        Ok((engine, rate_watch))
289    }
290
291    /// Resolve the output device: use configured device name if set,
292    /// falling back to system default if not set or if the named device is unavailable.
293    fn resolve_device(&self) -> Result<backend::DeviceInfo, PlayerError> {
294        if let Some(ref name) = self.output_device_name {
295            match self.backend.list_devices() {
296                Ok(devices) => {
297                    if let Some(dev) = devices.into_iter().find(|d| d.name == *name) {
298                        return Ok(dev);
299                    }
300                    log::warn!(
301                        "configured output device '{}' not found, falling back to default",
302                        name,
303                    );
304                }
305                Err(e) => {
306                    log::warn!("failed to list devices while resolving '{}': {}", name, e);
307                }
308            }
309        }
310        Ok(self.backend.default_device()?)
311    }
312
313    /// Switch the output device. Persists to config and restarts the engine
314    /// on the current track if playing.
315    pub fn set_output_device(&mut self, name: String) {
316        log::info!("switching output device to: {}", name);
317        self.output_device_name = Some(name.clone());
318
319        if let Err(e) = crate::config::Config::persist(|cfg| {
320            cfg.playback.output_device = Some(name);
321        }) {
322            log::error!("failed to save output device config: {}", e);
323        }
324
325        self.restart_on_current_track();
326    }
327
328    /// Clear the configured output device, reverting to system default.
329    pub fn clear_output_device(&mut self) {
330        log::info!("reverting to system default output device");
331        self.output_device_name = None;
332
333        if let Err(e) = crate::config::Config::persist(|cfg| {
334            cfg.playback.output_device = None;
335        }) {
336            log::error!("failed to save output device config: {}", e);
337        }
338
339        self.restart_on_current_track();
340    }
341
342    /// If a track is currently playing or paused, restart playback at the
343    /// current position (e.g. after switching output devices). Preserves pause state.
344    fn restart_on_current_track(&mut self) {
345        if let Some(info) = self.shared_state.track_info() {
346            let position_ms = self.shared_state.position_ms();
347            if let Err(e) = self.restart_current(&info, position_ms) {
348                log::error!("failed to restart playback on device switch: {}", e);
349            }
350        }
351    }
352
353    /// Get the current output device name (if configured).
354    pub fn output_device_name(&self) -> Option<&str> {
355        self.output_device_name.as_deref()
356    }
357
358    /// Get a command sender for the UI layer.
359    pub fn command_sender(&self) -> crossbeam_channel::Sender<PlayerCommand> {
360        self.commands.tx.clone()
361    }
362
363    /// Play a specific item in the playlist by ID.
364    /// Sets cursor, starts playback if Ready or streaming-ready, otherwise waits for TrackReady.
365    pub fn play(&mut self, id: QueueItemId) {
366        self.shared_state.set_cursor(Some(id));
367
368        match self.shared_state.item_playback_source(id) {
369            Some(PlaybackSource::Ready(path)) => {
370                if let Err(e) = self.start_playback(id, &path, 0) {
371                    log::error!("play failed: {}", e);
372                }
373            }
374            Some(PlaybackSource::Streaming {
375                path,
376                bytes_written,
377                total,
378            }) => {
379                // Stop what is playing and park here. The probe answers on its
380                // own thread; if it cannot, TrackReady starts the track once
381                // the whole file has landed.
382                self.report(PlaybackReportState::Stopped);
383                self.stop_engine();
384                self.shared_state.set_playback_state(PlaybackState::Stopped);
385                self.probe_stream_for_playback(id, &path, bytes_written, total);
386            }
387            None => {
388                // Item not ready — stop current playback, wait for TrackReady.
389                self.report(PlaybackReportState::Stopped);
390                self.stop_engine();
391                self.shared_state.set_playback_state(PlaybackState::Stopped);
392                log::info!("play: item {:?} not ready, waiting for TrackReady", id);
393            }
394        }
395    }
396
397    /// Start playback of a file.
398    ///
399    /// A failure leaves the player cleanly stopped. Displaying a track that no
400    /// engine is playing freezes the position and makes the transport lie.
401    fn start_playback(
402        &mut self,
403        id: QueueItemId,
404        path: &Path,
405        seek_ms: u64,
406    ) -> Result<(), PlayerError> {
407        #[cfg(test)]
408        {
409            self.playback_starts += 1;
410        }
411        let result = self.open_playback(id, path, seek_ms);
412        if result.is_err() {
413            self.stop_playback_and_clear_state();
414        }
415        self.wake_analyzer();
416        result
417    }
418
419    fn open_playback(
420        &mut self,
421        id: QueueItemId,
422        path: &Path,
423        seek_ms: u64,
424    ) -> Result<(), PlayerError> {
425        self.stop_engine();
426
427        let info = buffer::probe_file(path)?;
428
429        // Set track_info + position immediately so the UI never sees a gap.
430        // For seeks, this keeps the bar at the target position instead of
431        // flashing to 0 while the new timeline spins up.
432        self.shared_state.set_track_info(Some(TrackInfo {
433            id,
434            path: path.to_path_buf(),
435            codec: info.codec.clone(),
436            sample_rate: info.sample_rate,
437            bit_depth: info.bit_depth,
438            bitrate_kbps: info.bitrate_kbps,
439            channels: info.channels,
440            duration_ms: info.duration_ms,
441        }));
442        self.shared_state.set_position_ms(seek_ms);
443        self.on_track_changed(id, seek_ms);
444        log::info!(
445            "playing: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
446            path.display(),
447            id,
448            info.codec,
449            info.sample_rate,
450            info.channels,
451            info.duration_ms,
452            if seek_ms > 0 {
453                format!(" @{}ms", seek_ms)
454            } else {
455                String::new()
456            }
457        );
458
459        let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
460
461        self.timeline.reset();
462
463        let next_track = self.decode_cursor(id);
464
465        // Load ReplayGain config for this playback session.
466        let cfg = crate::config::Config::load_or_default();
467        let rg_mode = cfg.playback.replaygain;
468        let pre_amp_db = cfg.playback.pre_amp_db;
469
470        let finish_tx = self.commands.tx.clone();
471        let (_stream_info, decode_handle) = buffer::start_decode_file(
472            id,
473            path,
474            producer,
475            seek_ms,
476            next_track,
477            self.timeline.clone(),
478            Some(self.viz_buffer.clone()),
479            rg_mode,
480            pre_amp_db,
481            move || {
482                finish_tx.send(PlayerCommand::DecodeFinished).ok();
483            },
484        )?;
485
486        let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
487        engine.start()?;
488
489        self.shared_state.set_playback_state(PlaybackState::Playing);
490
491        self.active_playback = Some(ActivePlayback {
492            engine,
493            decode_handle,
494            stream: None,
495            _rate_watch: rate_watch,
496        });
497
498        Ok(())
499    }
500
501    /// Gapless lookahead: the decode thread keeps its own cursor, separate
502    /// from the UI cursor, so it can look ahead through the playlist without
503    /// moving what the UI shows as now playing.
504    fn decode_cursor(
505        &self,
506        id: QueueItemId,
507    ) -> impl Fn() -> Option<(QueueItemId, PathBuf)> + Send + 'static {
508        let state = self.shared_state.clone();
509        let cursor = parking_lot::Mutex::new(Some(id));
510        move || {
511            let current = cursor.lock().take()?;
512            let next = state.peek_next_ready_after(current);
513            if let Some((next_id, _)) = &next {
514                *cursor.lock() = Some(*next_id);
515            }
516            next
517        }
518    }
519
520    /// Probe a partially-downloaded file on its own thread, and start it when
521    /// the answer comes back.
522    ///
523    /// Nothing here waits. Probing reads as much of the container as it takes
524    /// to describe itself — for Ogg, its last page, which means the whole
525    /// remaining download — and this is the thread that answers play, pause and
526    /// seek. So the probe goes elsewhere and its result returns as a command.
527    ///
528    /// A format that describes itself up front (FLAC, MP3) comes back in
529    /// milliseconds and starts early, which is the point of streaming. One that
530    /// does not comes back whenever it comes back, by which time the download
531    /// has usually landed and `TrackReady` has started the track from disk —
532    /// and the late answer is simply dropped. Either way the player kept
533    /// answering commands throughout.
534    fn probe_stream_for_playback(
535        &self,
536        id: QueueItemId,
537        path: &Path,
538        bytes_written: Arc<crate::remote::downloads::ByteFeed>,
539        total: u64,
540    ) {
541        let path = path.to_path_buf();
542        let tx = self.commands.tx.clone();
543        let hint = hint_for(&path);
544
545        // Abandon the moment the track stops being the one wanted. A probe of
546        // a container that needs its tail otherwise reads to the end of a
547        // download nobody is waiting for any more, and skipping through a
548        // queue that is still caching would leave one doing so per skip.
549        let status = {
550            let downloading = self.stream_status_fn(id);
551            let state = self.shared_state.clone();
552            Arc::new(move || {
553                if state.is_cursor(id) {
554                    downloading()
555                } else {
556                    streaming::StreamStatus::Failed
557                }
558            }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
559        };
560
561        let spawned = thread::Builder::new()
562            .name("koan-stream-probe".into())
563            .spawn(move || {
564                // `wait` says whether a read may sit at the write head for
565                // more of the download. The first attempt must not: a
566                // container that goes looking for its tail would wait for the
567                // whole transfer, and failing at once is how that is detected.
568                // The second has no length to go looking with, so whatever it
569                // still wants is in front of it and worth waiting for.
570                let attempt = |mode, wait: bool| {
571                    let open = if wait {
572                        streaming::PartialFileSource::open(
573                            &path,
574                            bytes_written.clone(),
575                            total,
576                            status.clone(),
577                            mode,
578                        )
579                    } else {
580                        streaming::PartialFileSource::open_for_probe(
581                            &path,
582                            bytes_written.clone(),
583                            total,
584                            status.clone(),
585                            mode,
586                        )
587                    };
588                    open.map_err(buffer::DecodeError::Io).and_then(|source| {
589                        let mss = symphonia::core::io::MediaSourceStream::new(
590                            Box::new(source),
591                            Default::default(),
592                        );
593                        buffer::probe_source(mss, &hint)
594                    })
595                };
596
597                // Ask for the whole description first. Neither attempt waits at
598                // the write head, so a container that needs bytes which have
599                // not arrived fails here rather than reading the transfer out.
600                let info = match attempt(streaming::ProbeMode::Full, false) {
601                    Ok(info) => Some((info, streaming::ProbeMode::Full)),
602                    Err(e) => {
603                        // Try again claiming no length. Ogg goes looking for its
604                        // last page only when told there is one to find; without
605                        // it the track opens now and plays, at the price of
606                        // seeking and of the duration that page carries. Both
607                        // come back when the download lands.
608                        log::info!(
609                            "stream probe: {} needs more than has arrived ({}), opening without a length",
610                            path.display(),
611                            e
612                        );
613                        let lengthless = lengthless_mode_for(&path);
614                        attempt(lengthless, true)
615                            .ok()
616                            .map(|info| (info, lengthless))
617                    }
618                };
619
620                match info {
621                    Some((info, mode)) => {
622                        tx.send(PlayerCommand::StreamProbed {
623                            id,
624                            info: Box::new(info),
625                            mode,
626                        })
627                        .ok();
628                    }
629                    // Not a failure of the track: it plays from disk once the
630                    // download lands, and the cursor is still parked on it.
631                    None => log::info!(
632                        "stream probe: {} cannot start early, waiting for the download",
633                        path.display()
634                    ),
635                }
636            });
637
638        if let Err(e) = spawned {
639            log::warn!("stream probe: could not spawn for {:?}: {}", id, e);
640        }
641    }
642
643    /// A probe finished. Start the track if it is still the one wanted and
644    /// nothing has started it in the meantime.
645    fn stream_probed(
646        &mut self,
647        id: QueueItemId,
648        info: buffer::StreamInfo,
649        mode: streaming::ProbeMode,
650    ) {
651        if !self.shared_state.is_cursor(id) {
652            return; // Moved on.
653        }
654        if self.shared_state.playback_state() != PlaybackState::Stopped {
655            return; // Already playing — the download landed first, or the user did.
656        }
657
658        match self.shared_state.item_playback_source(id) {
659            // The download landed while probing: play it as an ordinary file.
660            Some(PlaybackSource::Ready(path)) => {
661                if let Err(e) = self.start_playback(id, &path, 0) {
662                    log::error!("stream probe: playback failed: {}", e);
663                }
664            }
665            Some(PlaybackSource::Streaming {
666                path,
667                bytes_written,
668                total,
669            }) => {
670                let source = StreamSource {
671                    path,
672                    bytes_written,
673                    total,
674                    mode,
675                };
676                if let Err(e) = self.start_streaming_playback(id, source, 0, info) {
677                    log::error!("stream probe: streaming playback failed: {}", e);
678                }
679            }
680            None => {}
681        }
682    }
683
684    /// What the streaming source asks per read to know whether the download is
685    /// still going. Asked each time rather than passed once: a transfer can
686    /// land, or die, at any point during playback.
687    fn stream_status_fn(
688        &self,
689        id: QueueItemId,
690    ) -> Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync> {
691        let state = self.shared_state.clone();
692        Arc::new(move || match state.item_load_state(id) {
693            Some(LoadState::Ready) => streaming::StreamStatus::Complete,
694            Some(LoadState::Failed(_)) => streaming::StreamStatus::Failed,
695            _ => streaming::StreamStatus::Downloading,
696        })
697    }
698
699    /// Start streaming playback from a partially-downloaded file.
700    ///
701    /// The decoder reads the `.part` file straight off disk through a
702    /// `PartialFileSource`, which blocks when it reaches the write head. The
703    /// download's final rename does not disturb an already-open descriptor, so
704    /// a transfer landing mid-track needs no handover.
705    ///
706    /// `info` is already known — from the off-thread probe when starting, or
707    /// from what is playing when seeking. Nothing here probes.
708    fn start_streaming_playback(
709        &mut self,
710        id: QueueItemId,
711        source: StreamSource,
712        seek_ms: u64,
713        info: buffer::StreamInfo,
714    ) -> Result<(), PlayerError> {
715        let result = self.open_streaming_playback(id, source, seek_ms, info);
716        if result.is_err() {
717            self.stop_playback_and_clear_state();
718        }
719        result
720    }
721
722    fn open_streaming_playback(
723        &mut self,
724        id: QueueItemId,
725        source: StreamSource,
726        seek_ms: u64,
727        info: buffer::StreamInfo,
728    ) -> Result<(), PlayerError> {
729        self.stop_engine();
730        // Held so a seek can reopen the same way without probing again.
731        self.stream_mode = source.mode;
732        let path = source.path.as_path();
733
734        let live = LiveStream {
735            feed: source.bytes_written.clone(),
736            abandoned: Default::default(),
737        };
738        let status = {
739            let downloading = self.stream_status_fn(id);
740            let abandoned = live.abandoned.clone();
741            Arc::new(move || {
742                if abandoned.load(std::sync::atomic::Ordering::Acquire) {
743                    streaming::StreamStatus::Failed
744                } else {
745                    downloading()
746                }
747            }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
748        };
749        let open_source = {
750            let StreamSource {
751                path,
752                bytes_written,
753                total,
754                mode,
755            } = source.clone();
756            let status = status.clone();
757            move || {
758                streaming::PartialFileSource::open(
759                    &path,
760                    bytes_written.clone(),
761                    total,
762                    status.clone(),
763                    mode,
764                )
765            }
766        };
767
768        self.shared_state.set_track_info(Some(TrackInfo {
769            id,
770            path: path.to_path_buf(),
771            codec: info.codec.clone(),
772            sample_rate: info.sample_rate,
773            bit_depth: info.bit_depth,
774            bitrate_kbps: info.bitrate_kbps,
775            channels: info.channels,
776            duration_ms: info.duration_ms,
777        }));
778        self.shared_state.set_position_ms(seek_ms);
779        self.on_track_changed(id, seek_ms);
780        log::info!(
781            "streaming: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
782            path.display(),
783            id,
784            info.codec,
785            info.sample_rate,
786            info.channels,
787            info.duration_ms,
788            if seek_ms > 0 {
789                format!(" @{}ms", seek_ms)
790            } else {
791                String::new()
792            },
793        );
794
795        let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
796
797        self.timeline.reset();
798
799        // Gapless lookahead after streaming: next track uses normal file path.
800        let next_track = self.decode_cursor(id);
801
802        let first = buffer::SourceEntry {
803            id,
804            path: path.to_path_buf(),
805            hint: hint_for(path),
806            make_mss: Box::new(move || {
807                Ok(symphonia::core::io::MediaSourceStream::new(
808                    Box::new(open_source()?),
809                    Default::default(),
810                ))
811            }),
812        };
813
814        // Load ReplayGain config for this streaming session.
815        let cfg = crate::config::Config::load_or_default();
816        let rg_mode = cfg.playback.replaygain;
817        let pre_amp_db = cfg.playback.pre_amp_db;
818
819        let finish_tx = self.commands.tx.clone();
820        let (_stream_info, decode_handle) = buffer::start_decode(
821            first,
822            producer,
823            seek_ms,
824            move || {
825                let (next_id, next_path) = next_track()?;
826                Some(buffer::SourceEntry::from_file(next_id, next_path))
827            },
828            self.timeline.clone(),
829            Some(self.viz_buffer.clone()),
830            rg_mode,
831            pre_amp_db,
832            move || {
833                finish_tx.send(PlayerCommand::DecodeFinished).ok();
834            },
835        )?;
836
837        let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
838        engine.start()?;
839
840        self.shared_state.set_playback_state(PlaybackState::Playing);
841
842        self.active_playback = Some(ActivePlayback {
843            engine,
844            decode_handle,
845            stream: Some(live),
846            _rate_watch: rate_watch,
847        });
848
849        Ok(())
850    }
851
852    /// Seek within the current track, preserving pause state.
853    ///
854    /// A track still downloading is seekable only as far as its bytes reach, so
855    /// the target is clamped to `seekable_ms` and the restart goes back through
856    /// the streaming path — reopening a partial file as a plain file would
857    /// decode whatever happens to be on disk and end the track early.
858    pub fn seek(&mut self, position_ms: u64) {
859        let Some(info) = self.shared_state.track_info() else {
860            return;
861        };
862        // Stop just short of the end rather than falling into the next track.
863        let seekable = self.shared_state.seekable_ms();
864        if seekable == 0 {
865            // Nothing of this track can be reached yet — a partial container
866            // that has not said what it is. Restarting it at zero is not what
867            // anyone asked for, so the seek is simply declined.
868            log::debug!("seek declined: {:?} is not seekable yet", info.id);
869            return;
870        }
871        let ceiling = seekable.min(
872            self.shared_state
873                .duration_ms()
874                .saturating_sub(SEEK_END_GUARD_MS),
875        );
876        let clamped = position_ms.min(ceiling);
877
878        if let Err(e) = self.restart_current(&info, clamped) {
879            log::error!("seek failed: {}", e);
880        }
881    }
882
883    /// Restart what is playing at `position_ms`, preserving pause state.
884    ///
885    /// What a seek does, and what switching output device does, and what going
886    /// back from the first track does. All three restart the same track, so all
887    /// three resolve the source the same way: from the queue item, never from
888    /// `info.path`, which names the `.part` file for a track that was still
889    /// downloading when it started and is not renamed when the download lands.
890    fn restart_current(&mut self, info: &TrackInfo, position_ms: u64) -> Result<(), PlayerError> {
891        let was_paused = self.shared_state.playback_state() == PlaybackState::Paused;
892
893        match self.shared_state.item_playback_source(info.id) {
894            Some(PlaybackSource::Streaming {
895                path,
896                bytes_written,
897                total,
898            }) => {
899                // No probe: what is playing already said what this is, and
900                // reading an Ogg's last page to learn it again would mean
901                // waiting for the rest of the download.
902                let known = buffer::StreamInfo {
903                    codec: info.codec.clone(),
904                    sample_rate: info.sample_rate,
905                    channels: info.channels,
906                    bit_depth: info.bit_depth,
907                    bitrate_kbps: info.bitrate_kbps,
908                    duration_ms: info.duration_ms,
909                };
910                let source = StreamSource {
911                    path,
912                    bytes_written,
913                    total,
914                    mode: self.stream_mode,
915                };
916                self.start_streaming_playback(info.id, source, position_ms, known)?;
917            }
918            Some(PlaybackSource::Ready(path)) => {
919                self.start_playback(info.id, &path, position_ms)?;
920            }
921            None => return Ok(()),
922        }
923
924        if was_paused {
925            self.pause_now();
926        }
927        self.report(if was_paused {
928            PlaybackReportState::Paused
929        } else {
930            PlaybackReportState::Playing
931        });
932        Ok(())
933    }
934
935    /// Skip to next track in playlist.
936    pub fn next_track(&mut self) {
937        match self.shared_state.advance_cursor_loadable() {
938            Some(id) => self.play(id),
939            None => {
940                log::info!("no more tracks in playlist");
941                self.stop_playback_and_clear_state();
942            }
943        }
944    }
945
946    /// Go back to previous track.
947    pub fn prev_track(&mut self) {
948        match self.shared_state.retreat_cursor() {
949            Some((id, _)) => self.play(id),
950            None => {
951                // No previous track — restart current from the beginning.
952                if let Some(info) = self.shared_state.track_info()
953                    && let Err(e) = self.restart_current(&info, 0)
954                {
955                    log::error!("restart failed: {}", e);
956                }
957            }
958        }
959    }
960
961    /// Pause playback, fading out if the config asks for it.
962    ///
963    /// A fade leaves the unit running until it reaches silence;
964    /// `update_playback_state` stops it from there.
965    pub fn pause(&mut self) {
966        let Some(ref playback) = self.active_playback else {
967            return;
968        };
969        if crate::config::Config::load_or_default()
970            .playback
971            .fade_on_pause
972        {
973            playback.engine.fade_out();
974            self.shared_state.set_playback_state(PlaybackState::Paused);
975        } else {
976            self.pause_now();
977        }
978        self.report(PlaybackReportState::Paused);
979    }
980
981    /// Pause without a fade — for a restart that should come back paused,
982    /// where there is nothing audible to fade.
983    fn pause_now(&mut self) {
984        if let Some(ref playback) = self.active_playback {
985            if let Err(e) = playback.engine.stop() {
986                log::error!("pause failed: {}", e);
987                return;
988            }
989            self.shared_state.set_playback_state(PlaybackState::Paused);
990        }
991    }
992
993    /// Resume playback. Fades back in if the pause faded out.
994    pub fn resume(&mut self) {
995        if let Some(ref playback) = self.active_playback {
996            let engine = &playback.engine;
997            let resumed = if engine.is_running() || engine.is_silent() {
998                engine.fade_in()
999            } else {
1000                engine.start()
1001            };
1002            if let Err(e) = resumed {
1003                log::error!("resume failed: {}", e);
1004                return;
1005            }
1006            self.shared_state.set_playback_state(PlaybackState::Playing);
1007            self.wake_analyzer();
1008            self.report(PlaybackReportState::Playing);
1009        }
1010    }
1011
1012    /// Tell the analyser there is about to be something to hear.
1013    ///
1014    /// It parks when nothing is playing and nothing is reading, and the one
1015    /// thing it cannot be signalled from is the play head — that counter is
1016    /// written by the audio render callback, which may never take a lock. So
1017    /// the player says so instead, on the two edges where silence ends.
1018    fn wake_analyzer(&self) {
1019        self.viz_snapshot.wake();
1020    }
1021
1022    /// Stop playback and clear playlist.
1023    pub fn stop(&mut self) {
1024        self.shared_state.clear_playlist();
1025        self.stop_playback_and_clear_state();
1026    }
1027
1028    /// Stop the audio engine and decode thread without touching shared state.
1029    ///
1030    /// Output stops first, then the decode thread is joined, then the engine
1031    /// drops: tearing CoreAudio down under a live producer is the end-of-queue
1032    /// crash (#89).
1033    fn stop_engine(&mut self) {
1034        let Some(playback) = self.active_playback.take() else {
1035            return;
1036        };
1037        let ActivePlayback {
1038            engine,
1039            mut decode_handle,
1040            stream,
1041            _rate_watch,
1042        } = playback;
1043
1044        let _ = engine.stop();
1045        // Stop first, so the failed read the abandon causes reads as a stop
1046        // rather than a bad source to skip past.
1047        decode_handle.signal_stop();
1048        if let Some(stream) = stream {
1049            stream.abandon();
1050        }
1051        decode_handle.stop();
1052        drop(engine);
1053    }
1054
1055    /// Full stop: tear down engine + clear all display state.
1056    fn stop_playback_and_clear_state(&mut self) {
1057        self.report(PlaybackReportState::Stopped);
1058        self.finish_play();
1059        self.stop_engine();
1060        self.timeline.reset();
1061        self.shared_state.set_playback_state(PlaybackState::Stopped);
1062        self.shared_state.set_position_ms(0);
1063        self.shared_state.set_track_info(None);
1064    }
1065
1066    /// Remove a track from the playlist. If it was the cursor, resume at the
1067    /// track that followed it.
1068    ///
1069    /// `remove_item` clears the cursor, and an unset cursor means "start from the
1070    /// top" — so the successor is pinned down by parking the cursor on the removed
1071    /// track's predecessor first. `None` is correct only when it was the first item.
1072    pub fn remove_from_playlist(&mut self, id: QueueItemId) {
1073        let was_cursor = self.shared_state.is_cursor(id);
1074        let resume_after = was_cursor
1075            .then(|| self.shared_state.item_before(id))
1076            .flatten();
1077        self.shared_state.remove_item(id);
1078        if was_cursor {
1079            self.shared_state.set_cursor(resume_after);
1080            self.next_track();
1081        }
1082    }
1083
1084    /// A download finished — if cursor is waiting on this item, start playback.
1085    /// If already streaming this item, re-read its metadata from the complete file.
1086    pub fn track_ready(&mut self, id: QueueItemId) {
1087        // Mark as Ready (download thread already did this, but be safe).
1088        self.shared_state.update_item_state(id, ItemState::Ready);
1089
1090        if !self.shared_state.is_cursor(id) {
1091            return;
1092        }
1093
1094        let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
1095        let current_track_id = self.shared_state.track_info().map(|t| t.id);
1096
1097        if is_playing && current_track_id == Some(id) {
1098            // Already streaming this track — download just finished.
1099            // Re-read the full tags now that the whole file is here.
1100            log::info!(
1101                "track_ready: download complete while streaming {:?}, refreshing metadata",
1102                id
1103            );
1104            self.refresh_track_metadata(id);
1105            return;
1106        }
1107
1108        // Cursor is on this item but not yet playing — start playback now.
1109        if !is_playing && let Some(path) = self.shared_state.item_path_if_ready(id) {
1110            log::info!("track_ready: starting playback for {:?}", id);
1111            if let Err(e) = self.start_playback(id, &path, 0) {
1112                log::error!("track_ready playback failed: {}", e);
1113            }
1114        }
1115    }
1116
1117    /// Called when enough data has been buffered for streaming playback.
1118    /// If the cursor is waiting on this track and nothing is playing, start streaming.
1119    pub fn track_stream_ready(&mut self, id: QueueItemId) {
1120        if !self.shared_state.is_cursor(id) {
1121            return;
1122        }
1123
1124        let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
1125        if is_playing {
1126            return; // Already playing something — don't interrupt.
1127        }
1128
1129        match self.shared_state.item_playback_source(id) {
1130            Some(PlaybackSource::Streaming {
1131                path,
1132                bytes_written,
1133                total,
1134            }) => {
1135                log::info!("track_stream_ready: probing partial file for {:?}", id);
1136                self.probe_stream_for_playback(id, &path, bytes_written, total);
1137            }
1138            Some(PlaybackSource::Ready(path)) => {
1139                // Download finished between threshold and now — just play normally.
1140                log::info!(
1141                    "track_stream_ready: track already ready, starting normal playback for {:?}",
1142                    id
1143                );
1144                if let Err(e) = self.start_playback(id, &path, 0) {
1145                    log::error!("track_stream_ready playback failed: {}", e);
1146                }
1147            }
1148            None => {} // Not enough data yet — wait.
1149        }
1150    }
1151
1152    /// Re-read full lofty metadata for a track after its download completes.
1153    /// Called from track_ready() when a streaming track finishes downloading.
1154    /// What the item takes from it is `update_item_metadata`'s call.
1155    fn refresh_track_metadata(&mut self, id: QueueItemId) {
1156        use crate::index::metadata;
1157
1158        let path = match self.shared_state.item_path_if_ready(id) {
1159            Some(p) => p,
1160            None => return,
1161        };
1162
1163        match metadata::read_metadata(&path) {
1164            Ok(meta) => {
1165                self.shared_state.update_item_metadata(
1166                    id,
1167                    meta.title,
1168                    meta.artist,
1169                    meta.album_artist.unwrap_or_default(),
1170                    meta.album,
1171                    meta.duration_ms.map(|d| d as u64),
1172                );
1173
1174                // Re-probe the complete file for accurate duration + stream info.
1175                // The initial probe was done on partial streaming data and may have
1176                // underestimated duration, causing premature seek clamping or wrong
1177                // progress bar display.
1178                //
1179                // The path is taken over at the same time. Playback started
1180                // against the `.part` file and the download's last act is to
1181                // rename it, so what `track_info` holds now names nothing.
1182                if let Some(current) = self.shared_state.track_info()
1183                    && current.id == id
1184                {
1185                    let probed = buffer::probe_file(&path).ok();
1186                    let duration_ms = probed
1187                        .as_ref()
1188                        .map(|s| s.duration_ms)
1189                        .filter(|d| *d > current.duration_ms)
1190                        .unwrap_or(current.duration_ms);
1191                    if duration_ms != current.duration_ms {
1192                        log::info!(
1193                            "track_ready: duration corrected {}ms → {}ms",
1194                            current.duration_ms,
1195                            duration_ms
1196                        );
1197                    }
1198                    self.shared_state.set_track_info(Some(TrackInfo {
1199                        duration_ms,
1200                        path: path.clone(),
1201                        ..current
1202                    }));
1203                }
1204
1205                // Signal UI to re-read cover art and update souvlaki media controls.
1206                self.shared_state.signal_metadata_refresh();
1207                log::info!("track_ready: metadata refreshed for {:?}", id);
1208            }
1209            Err(e) => {
1210                log::warn!("track_ready: metadata refresh failed for {:?}: {}", id, e);
1211            }
1212        }
1213    }
1214
1215    /// The needle has moved to `id`. Close out the outgoing track and write
1216    /// the new one to history straight away, so history reads in play order
1217    /// even for a track that is skipped a moment later.
1218    ///
1219    /// A seek restarts playback of the same track, so identity is checked
1220    /// rather than closing unconditionally — otherwise scrubbing around a
1221    /// track would enter it into history once per seek.
1222    fn on_track_changed(&mut self, id: QueueItemId, position_ms: u64) {
1223        if self.in_flight.as_ref().is_some_and(|f| f.item == id) {
1224            return;
1225        }
1226        self.finish_play();
1227        let track_id = self.shared_state.item_db_id(id);
1228        self.in_flight = Some(InFlight::new(id, track_id));
1229        if let (Some(track_id), Some(recorder)) = (track_id, self.history.as_ref()) {
1230            recorder.record(PlayEvent::Started {
1231                track_id,
1232                position_ms,
1233            });
1234        }
1235    }
1236
1237    /// Tell the remote server where the track playing now stands. A track
1238    /// starting is reported by `on_track_changed`; this covers what happens
1239    /// to it afterwards.
1240    fn report(&self, state: PlaybackReportState) {
1241        let Some(track_id) = self.in_flight.as_ref().and_then(InFlight::track_id) else {
1242            return;
1243        };
1244        if let Some(recorder) = self.history.as_ref() {
1245            recorder.record(PlayEvent::Playback(PlaybackReport {
1246                track_id,
1247                state,
1248                position_ms: self.shared_state.position_ms(),
1249            }));
1250        }
1251    }
1252
1253    /// Tell history how long the current track was heard for. Returns what was
1254    /// reported, which is how the tests see it.
1255    fn finish_play(&mut self) -> Option<PlayEvent> {
1256        let flight = self.in_flight.take()?;
1257        let event = PlayEvent::Finished {
1258            track_id: flight.track_id()?,
1259            listened_ms: flight.listened_ms(),
1260        };
1261        if let Some(recorder) = self.history.as_ref() {
1262            recorder.record(event);
1263        }
1264        Some(event)
1265    }
1266
1267    /// Poll the timeline and update shared state with current track/position.
1268    /// Called from the command loop on each tick.
1269    pub fn update_playback_state(&mut self) {
1270        let Some(playback) = self.active_playback.as_ref() else {
1271            return;
1272        };
1273
1274        if self.shared_state.playback_state() == PlaybackState::Paused
1275            && playback.engine.is_running()
1276            && playback.engine.is_silent()
1277            && let Err(e) = playback.engine.stop()
1278        {
1279            log::error!("stopping after fade failed: {}", e);
1280        }
1281
1282        if let Some((id, path, info, position_ms)) = self.timeline.current_playback() {
1283            self.shared_state.set_position_ms(position_ms);
1284
1285            // A gapless transition moves the needle without anything on this
1286            // thread having asked it to, so the play is banked from here.
1287            self.on_track_changed(id, position_ms);
1288            if let Some(f) = self.in_flight.as_mut() {
1289                f.advance(position_ms);
1290            }
1291
1292            // Update track_info + cursor if the timeline shows a different track
1293            // (gapless transition happened).
1294            let current_id = self.shared_state.track_info().map(|t| t.id);
1295            if current_id != Some(id) {
1296                log::info!("timeline: now playing {:?}", id);
1297                self.shared_state.set_track_info(Some(TrackInfo {
1298                    id,
1299                    path,
1300                    codec: info.codec,
1301                    sample_rate: info.sample_rate,
1302                    bit_depth: info.bit_depth,
1303                    bitrate_kbps: info.bitrate_kbps,
1304                    channels: info.channels,
1305                    duration_ms: info.duration_ms,
1306                }));
1307                self.shared_state.set_cursor(Some(id));
1308            }
1309        }
1310    }
1311
1312    /// A download the cursor is parked on will never land.
1313    ///
1314    /// `play()` leaves the cursor on an item that is not yet Ready and stops,
1315    /// waiting for `TrackReady`. When the download fails instead, that wait has
1316    /// no end — so walk on to the next item that can still load, or stop
1317    /// cleanly if there is none.
1318    pub fn track_failed(&mut self, id: QueueItemId) {
1319        if !self.shared_state.is_cursor(id) {
1320            return;
1321        }
1322        // Only a parked cursor is waiting on this. Playing means it is being
1323        // streamed from the partial file — the pump sees the failure and ends
1324        // the decode, which advances the queue — and paused is the user's.
1325        if self.shared_state.playback_state() != PlaybackState::Stopped {
1326            return;
1327        }
1328        log::info!("track {:?} cannot load, moving on", id);
1329        self.next_track();
1330    }
1331
1332    /// Decode thread naturally finished (playlist exhausted or error).
1333    /// Advance to the next playable track; otherwise stop cleanly.
1334    ///
1335    /// A track that has not finished downloading parks the cursor on it, so its
1336    /// `TrackReady`/`TrackStreamReady` resumes the queue instead of being
1337    /// discarded as "not the cursor".
1338    fn on_decode_finished(&mut self) {
1339        log::info!("decode finished, checking for next track");
1340        match self.shared_state.advance_cursor_loadable() {
1341            Some(id) => self.play(id),
1342            None => {
1343                log::info!("no more tracks — stopping");
1344                self.stop_playback_and_clear_state();
1345            }
1346        }
1347    }
1348
1349    /// Snapshot items with their predecessors for an undo of "these were removed".
1350    /// In playlist order, so undo re-inserts each item after a predecessor that
1351    /// is already back in place.
1352    fn snapshot_for_undo(
1353        &self,
1354        ids: &[QueueItemId],
1355    ) -> Vec<(Box<state::PlaylistItem>, Option<QueueItemId>)> {
1356        self.shared_state
1357            .items_before(ids)
1358            .into_iter()
1359            .filter_map(|(id, after)| Some((Box::new(self.shared_state.get_item(id)?), after)))
1360            .collect()
1361    }
1362
1363    /// Route an undo entry to the batch buffer (if batching) or the undo stack.
1364    fn push_undo(&mut self, entry: UndoEntry) {
1365        if let Some(ref mut batch) = self.batch_buffer {
1366            batch.push(entry);
1367        } else {
1368            self.undo_stack.push(entry);
1369        }
1370    }
1371
1372    /// Process a single command.
1373    pub fn process_command(&mut self, cmd: PlayerCommand) {
1374        match cmd {
1375            PlayerCommand::Play(id) => self.play(id),
1376            PlayerCommand::Pause => self.pause(),
1377            PlayerCommand::Resume => self.resume(),
1378            PlayerCommand::Stop => self.stop(),
1379            PlayerCommand::Seek(pos) => self.seek(pos),
1380            PlayerCommand::NextTrack => {
1381                // Debounce: suppress key repeat from terminal (150ms window).
1382                let now = std::time::Instant::now();
1383                if now.duration_since(self.last_skip).as_millis() >= 150 {
1384                    self.last_skip = now;
1385                    self.next_track();
1386                }
1387            }
1388            PlayerCommand::PrevTrack => {
1389                let now = std::time::Instant::now();
1390                if now.duration_since(self.last_skip).as_millis() >= 150 {
1391                    self.last_skip = now;
1392                    self.prev_track();
1393                }
1394            }
1395            PlayerCommand::AddToPlaylist(items) => {
1396                let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1397                self.shared_state.add_items(items);
1398                self.push_undo(UndoEntry::Added { ids });
1399            }
1400            PlayerCommand::UpdatePaths(updates) => {
1401                self.shared_state.update_paths(&updates);
1402                if let Some(info) = self.shared_state.track_info()
1403                    && let Some((_, new_path)) = updates.iter().find(|(id, _)| *id == info.id)
1404                {
1405                    self.shared_state.set_track_info(Some(TrackInfo {
1406                        path: new_path.clone(),
1407                        ..info
1408                    }));
1409                }
1410            }
1411            PlayerCommand::InsertInPlaylist { items, after } => {
1412                let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1413                self.shared_state.insert_items_after(items, after);
1414                self.push_undo(UndoEntry::Inserted { ids });
1415            }
1416            PlayerCommand::ClearPlaylist => {
1417                // Stop engine + clear display state WITHOUT touching the playlist,
1418                // then snapshot, then clear. This avoids the race where stop()
1419                // would clear the playlist before we capture it for undo.
1420                self.stop_playback_and_clear_state();
1421                let (items, cursor) = self.shared_state.snapshot_playlist();
1422                self.shared_state.clear_playlist();
1423                self.push_undo(UndoEntry::Replaced { items, cursor });
1424            }
1425            PlayerCommand::ReplacePlaylist { items, start } => {
1426                // Same order as ClearPlaylist: stop and clear display state
1427                // before snapshotting, or the snapshot captures an already
1428                // emptied playlist and undo restores nothing.
1429                self.stop_playback_and_clear_state();
1430                let (old_items, cursor) = self.shared_state.snapshot_playlist();
1431                self.shared_state.clear_playlist();
1432                self.push_undo(UndoEntry::Replaced {
1433                    items: old_items,
1434                    cursor,
1435                });
1436
1437                if items.is_empty() {
1438                    return;
1439                }
1440                let start_id = items.get(start).unwrap_or(&items[0]).id;
1441                self.shared_state.add_items(items);
1442                self.play(start_id);
1443            }
1444            PlayerCommand::RemoveFromPlaylist(id) => {
1445                let item = self.shared_state.get_item(id);
1446                let after = self.shared_state.item_before(id);
1447                self.remove_from_playlist(id);
1448                if let Some(item) = item {
1449                    self.push_undo(UndoEntry::Removed {
1450                        items: vec![(Box::new(item), after)],
1451                    });
1452                }
1453            }
1454            PlayerCommand::RemoveFromPlaylistBatch(ids) => {
1455                // Snapshot before removing anything, and resolve the resume point
1456                // once: removing one at a time would restart the engine for every
1457                // deleted track that the cursor lands on along the way.
1458                let items_with_pos = self.snapshot_for_undo(&ids);
1459                let resume_after = match self.shared_state.cursor() {
1460                    Some(cursor) if ids.contains(&cursor) => {
1461                        Some(self.shared_state.surviving_item_before(cursor, &ids))
1462                    }
1463                    _ => None,
1464                };
1465
1466                self.shared_state.remove_items(&ids);
1467
1468                if let Some(resume_after) = resume_after {
1469                    self.shared_state.set_cursor(resume_after);
1470                    self.next_track();
1471                }
1472
1473                if !items_with_pos.is_empty() {
1474                    self.push_undo(UndoEntry::Removed {
1475                        items: items_with_pos,
1476                    });
1477                }
1478            }
1479            PlayerCommand::MoveInPlaylist { id, target, after } => {
1480                let was_after = self.shared_state.item_before(id);
1481                self.shared_state.move_item(id, target, after);
1482                self.push_undo(UndoEntry::Moved { id, was_after });
1483            }
1484            PlayerCommand::MoveItemsInPlaylist { ids, target, after } => {
1485                let entries = self.shared_state.items_before(&ids);
1486                self.shared_state.move_items(&ids, target, after);
1487                self.push_undo(UndoEntry::MovedBatch { entries });
1488            }
1489            PlayerCommand::ReorderPlaylist(order) => {
1490                // Undoable like any other move: undoing it puts the queue back
1491                // and, by doing so, ends the lock — which is the honest result
1492                // of having rearranged the queue by hand.
1493                let entries = self.shared_state.items_before(&order);
1494                self.shared_state.reorder_to(&order);
1495                self.push_undo(UndoEntry::MovedBatch { entries });
1496            }
1497            PlayerCommand::TrackReady(id) => self.track_ready(id),
1498            PlayerCommand::DecodeFinished => self.on_decode_finished(),
1499            PlayerCommand::TrackStreamReady(id) => self.track_stream_ready(id),
1500            PlayerCommand::StreamProbed { id, info, mode } => self.stream_probed(id, *info, mode),
1501            PlayerCommand::TrackFailed(id) => self.track_failed(id),
1502            PlayerCommand::Undo => self.execute_undo(),
1503            PlayerCommand::Redo => self.execute_redo(),
1504            PlayerCommand::BeginUndoBatch => {
1505                self.batch_buffer = Some(Vec::new());
1506            }
1507            PlayerCommand::EndUndoBatch => {
1508                if let Some(entries) = self.batch_buffer.take() {
1509                    if entries.len() == 1 {
1510                        // Single entry — push directly, no wrapping.
1511                        self.undo_stack.push(entries.into_iter().next().unwrap());
1512                    } else if !entries.is_empty() {
1513                        self.undo_stack.push(UndoEntry::Batch(entries));
1514                    }
1515                }
1516            }
1517            PlayerCommand::SetOutputDevice(name) => self.set_output_device(name),
1518            PlayerCommand::RestartOutput => {
1519                log::info!("restarting audio output");
1520                self.restart_on_current_track();
1521            }
1522            PlayerCommand::ClearOutputDevice => self.clear_output_device(),
1523        }
1524    }
1525
1526    /// Apply an undo/redo entry: mutate the playlist and return the inverse entry.
1527    fn apply_entry(&mut self, entry: UndoEntry) -> UndoEntry {
1528        match entry {
1529            UndoEntry::Added { ids } | UndoEntry::Inserted { ids } => {
1530                // Undo of "items were added": snapshot them with positions, then remove.
1531                let items_with_pos = self.snapshot_for_undo(&ids);
1532                self.shared_state.remove_items(&ids);
1533                UndoEntry::Removed {
1534                    items: items_with_pos,
1535                }
1536            }
1537            UndoEntry::Removed { items } => {
1538                // Undo of "items were removed": re-insert each at its position.
1539                let mut ids = Vec::with_capacity(items.len());
1540                for (item, after) in items {
1541                    ids.push(item.id);
1542                    self.shared_state.insert_item_at(*item, after);
1543                }
1544                UndoEntry::Added { ids }
1545            }
1546            UndoEntry::Moved { id, was_after } => {
1547                let current_after = self.shared_state.item_before(id);
1548                self.shared_state.move_item_to(id, was_after);
1549                UndoEntry::Moved {
1550                    id,
1551                    was_after: current_after,
1552                }
1553            }
1554            UndoEntry::MovedBatch { entries } => {
1555                let ids: Vec<QueueItemId> = entries.iter().map(|(id, _)| *id).collect();
1556                let current_positions = self.shared_state.items_before(&ids);
1557                self.shared_state.move_items_to(&entries);
1558                UndoEntry::MovedBatch {
1559                    entries: current_positions,
1560                }
1561            }
1562            UndoEntry::Replaced { items, cursor } => {
1563                let (current_items, current_cursor) = self.shared_state.snapshot_playlist();
1564                self.shared_state.restore_playlist(items, cursor);
1565                UndoEntry::Replaced {
1566                    items: current_items,
1567                    cursor: current_cursor,
1568                }
1569            }
1570            UndoEntry::Batch(entries) => {
1571                // Apply entries in reverse order, collect inverses.
1572                let mut inverses: Vec<_> = entries
1573                    .into_iter()
1574                    .rev()
1575                    .map(|e| self.apply_entry(e))
1576                    .collect();
1577                inverses.reverse();
1578                UndoEntry::Batch(inverses)
1579            }
1580        }
1581    }
1582
1583    /// Put playback back in agreement with the playlist.
1584    ///
1585    /// The engine keeps decoding whatever it was on while the playlist changes
1586    /// underneath it, which an undo can turn into a lie: undoing a replace
1587    /// restores the queue but leaves the engine playing a track that queue does
1588    /// not contain. The transport then describes an item nothing can select,
1589    /// and the decode lookahead — which finds the next track by locating the
1590    /// current one — has nothing to follow, so the queue ends at the end of the
1591    /// track instead of carrying on.
1592    ///
1593    /// Done once, after the entry is applied, rather than inside each variant:
1594    /// any undo that takes items away can orphan the engine, not only
1595    /// `Replaced`.
1596    fn reconcile_playback(&mut self) {
1597        let Some(playing) = self.shared_state.track_info().map(|t| t.id) else {
1598            return;
1599        };
1600        if self.shared_state.get_item(playing).is_some() {
1601            return;
1602        }
1603        // Pick the restored queue back up where its cursor says it was, but
1604        // only if something was already playing — an undo is not a reason to
1605        // start the music, and the position is not part of what was snapshotted
1606        // so the track begins again.
1607        let resume = (self.shared_state.playback_state() == PlaybackState::Playing)
1608            .then(|| self.shared_state.cursor())
1609            .flatten();
1610        self.stop_playback_and_clear_state();
1611        if let Some(id) = resume {
1612            self.play(id);
1613        }
1614    }
1615
1616    /// Execute an undo operation, pushing the inverse onto the redo stack.
1617    fn execute_undo(&mut self) {
1618        let Some(entry) = self.undo_stack.pop_undo() else {
1619            return;
1620        };
1621        let inverse = self.apply_entry(entry);
1622        self.undo_stack.push_redo(inverse);
1623        self.reconcile_playback();
1624    }
1625
1626    /// Execute a redo operation, pushing the inverse onto the undo stack.
1627    fn execute_redo(&mut self) {
1628        let Some(entry) = self.undo_stack.pop_redo() else {
1629            return;
1630        };
1631        let inverse = self.apply_entry(entry);
1632        self.undo_stack.push_undo_keep_redo(inverse);
1633        self.reconcile_playback();
1634    }
1635
1636    /// Run the command loop. Blocks until the sender is dropped.
1637    pub fn run(&mut self) {
1638        use std::time::Duration;
1639
1640        let rx = self.commands.rx.clone();
1641        loop {
1642            // Poll with timeout so we update position even without commands.
1643            match rx.recv_timeout(Duration::from_millis(50)) {
1644                Ok(cmd) => self.process_command(cmd),
1645                Err(crossbeam_channel::RecvTimeoutError::Timeout) => {}
1646                Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break,
1647            }
1648            self.update_playback_state();
1649        }
1650        self.stop();
1651    }
1652
1653    /// Spawn the player on a background thread, returning the shared state,
1654    /// timeline, visualization snapshot, and command sender.
1655    pub fn spawn() -> (
1656        Arc<SharedPlayerState>,
1657        Arc<PlaybackTimeline>,
1658        Arc<VizSnapshot>,
1659        crossbeam_channel::Sender<PlayerCommand>,
1660    ) {
1661        let mut player = Self::new();
1662        player.history = PlayRecorder::spawn();
1663        let state = player.shared_state();
1664        let timeline = player.timeline();
1665        let viz_snapshot = player.viz_snapshot();
1666        let tx = player.command_sender();
1667
1668        thread::Builder::new()
1669            .name("koan-player".into())
1670            .spawn(move || player.run())
1671            .expect("failed to spawn player thread");
1672
1673        (state, timeline, viz_snapshot, tx)
1674    }
1675}
1676
1677#[cfg(test)]
1678mod tests {
1679    #[test]
1680    fn a_download_in_progress_is_known_by_its_own_extension() {
1681        use std::path::Path;
1682        // `.part` is the transfer's, not the track's.
1683        assert_eq!(
1684            lengthless_mode_for(Path::new("/c/t.m4a.part")),
1685            streaming::ProbeMode::LengthlessWholeEnd
1686        );
1687        assert_eq!(
1688            lengthless_mode_for(Path::new("/c/t.OPUS.part")),
1689            streaming::ProbeMode::LengthlessWholeEnd
1690        );
1691        assert_eq!(
1692            lengthless_mode_for(Path::new("/c/t.flac.part")),
1693            streaming::ProbeMode::Lengthless
1694        );
1695        assert_eq!(
1696            media_extension(Path::new("/c/t.m4a.part")).as_deref(),
1697            Some("m4a")
1698        );
1699        assert_eq!(
1700            media_extension(Path::new("/c/t.mp3")).as_deref(),
1701            Some("mp3")
1702        );
1703    }
1704
1705    use super::*;
1706    use state::PlaylistItem;
1707    use std::path::PathBuf;
1708    use std::sync::atomic::AtomicU64;
1709
1710    fn make_item(title: &str) -> PlaylistItem {
1711        PlaylistItem {
1712            playlist_entry_id: None,
1713            id: QueueItemId::new(),
1714            db_id: None,
1715            path: PathBuf::from(format!("/music/{title}.flac")),
1716            title: title.to_string(),
1717            artist: String::new(),
1718            album_artist: String::new(),
1719            album: String::new(),
1720            year: None,
1721            codec: None,
1722            track_number: None,
1723            disc: None,
1724            duration_ms: None,
1725            state: ItemState::Ready,
1726        }
1727    }
1728
1729    fn playlist_ids(player: &Player) -> Vec<QueueItemId> {
1730        let (items, _) = player.shared_state.snapshot_playlist();
1731        items.iter().map(|i| i.id).collect()
1732    }
1733
1734    fn playlist_titles(player: &Player) -> Vec<String> {
1735        let (items, _) = player.shared_state.snapshot_playlist();
1736        items.iter().map(|i| i.title.clone()).collect()
1737    }
1738
1739    fn pending_item(title: &str) -> PlaylistItem {
1740        PlaylistItem {
1741            playlist_entry_id: None,
1742            state: ItemState::Pending,
1743            ..make_item(title)
1744        }
1745    }
1746
1747    /// Stand in for an engine that is playing `id`. The test items have no
1748    /// files behind them, so `start_playback` can never get far enough to leave
1749    /// this state on its own.
1750    fn pretend_playing(player: &mut Player, id: QueueItemId) {
1751        let item = player
1752            .shared_state
1753            .get_item(id)
1754            .expect("item is in the queue");
1755        player.shared_state.set_track_info(Some(TrackInfo {
1756            id,
1757            path: item.path,
1758            codec: String::new(),
1759            sample_rate: 44_100,
1760            bit_depth: None,
1761            bitrate_kbps: None,
1762            channels: 2,
1763            duration_ms: 1_000,
1764        }));
1765        player
1766            .shared_state
1767            .set_playback_state(PlaybackState::Playing);
1768    }
1769
1770    fn playing_id(player: &Player) -> Option<QueueItemId> {
1771        player.shared_state.track_info().map(|t| t.id)
1772    }
1773
1774    /// Build `n` ready items, add them, and return their IDs.
1775    fn seed(player: &mut Player, n: usize) -> Vec<QueueItemId> {
1776        let items: Vec<_> = (0..n).map(|i| make_item(&format!("t{i}"))).collect();
1777        let ids = items.iter().map(|i| i.id).collect();
1778        player.process_command(PlayerCommand::AddToPlaylist(items));
1779        ids
1780    }
1781
1782    // --- cursor transitions ---
1783
1784    /// Feed the player a track's worth of playback ticks, as the 50ms poll would.
1785    fn listen(player: &mut Player, from_ms: u64, to_ms: u64) {
1786        let mut at = from_ms;
1787        if let Some(f) = player.in_flight.as_mut() {
1788            f.advance(at); // the position the needle landed on
1789        }
1790        while at < to_ms {
1791            at = (at + 50).min(to_ms);
1792            if let Some(f) = player.in_flight.as_mut() {
1793                f.advance(at);
1794            }
1795        }
1796    }
1797
1798    fn start(player: &mut Player, track_id: i64) -> QueueItemId {
1799        let id = QueueItemId::new();
1800        player.on_track_changed(id, 0);
1801        // The item is not in a playlist here, so there is no db_id to find.
1802        player
1803            .in_flight
1804            .as_mut()
1805            .unwrap()
1806            .track_id_for_test(track_id);
1807        id
1808    }
1809
1810    #[test]
1811    fn a_gapless_transition_closes_the_outgoing_track_and_opens_the_next() {
1812        let mut player = Player::new();
1813        start(&mut player, 11);
1814        listen(&mut player, 0, 200_000);
1815
1816        let b = QueueItemId::new();
1817        player.on_track_changed(b, 0);
1818        let f = player
1819            .in_flight
1820            .as_ref()
1821            .expect("the next track is counting");
1822        assert_eq!(f.item, b);
1823        assert_eq!(f.listened_ms(), 0, "and starts from nothing");
1824    }
1825
1826    #[test]
1827    fn a_track_skipped_seconds_in_is_still_history() {
1828        let mut player = Player::new();
1829        start(&mut player, 7);
1830        listen(&mut player, 0, 2_000);
1831
1832        let event = player
1833            .finish_play()
1834            .expect("putting something on is a thing you did, however briefly");
1835        assert!(matches!(
1836            event,
1837            history::PlayEvent::Finished {
1838                track_id: 7,
1839                listened_ms: 2_000
1840            }
1841        ));
1842    }
1843
1844    #[test]
1845    fn a_track_is_closed_out_once() {
1846        let mut player = Player::new();
1847        start(&mut player, 7);
1848        listen(&mut player, 0, 200_000);
1849
1850        assert!(player.finish_play().is_some());
1851        assert!(player.finish_play().is_none());
1852    }
1853
1854    #[test]
1855    fn seeking_around_a_track_does_not_enter_it_twice() {
1856        let mut player = Player::new();
1857        let id = start(&mut player, 7);
1858        listen(&mut player, 0, 120_000);
1859
1860        // A seek restarts playback of the same item.
1861        player.on_track_changed(id, 30_000);
1862        assert_eq!(
1863            player.in_flight.as_ref().unwrap().listened_ms(),
1864            120_000,
1865            "the seek kept the count rather than restarting it"
1866        );
1867        listen(&mut player, 30_000, 40_000);
1868
1869        let Some(history::PlayEvent::Finished { listened_ms, .. }) = player.finish_play() else {
1870            panic!("still one play");
1871        };
1872        assert_eq!(listened_ms, 130_000);
1873        assert!(player.finish_play().is_none());
1874    }
1875
1876    #[test]
1877    fn a_track_that_is_not_in_the_library_is_not_recorded() {
1878        let mut player = Player::new();
1879        let id = QueueItemId::new();
1880        player.on_track_changed(id, 0);
1881        listen(&mut player, 0, 200_000);
1882        assert!(player.finish_play().is_none());
1883    }
1884
1885    #[test]
1886    fn stopping_closes_out_what_was_heard() {
1887        let mut player = Player::new();
1888        start(&mut player, 7);
1889        listen(&mut player, 0, 150_000);
1890
1891        player.stop_playback_and_clear_state();
1892        assert!(player.in_flight.is_none(), "the stop consumed it");
1893    }
1894
1895    #[test]
1896    fn the_server_hears_each_turn_playback_takes() {
1897        use PlaybackReportState::{Paused, Playing, Stopped};
1898        use history::PlaybackReport;
1899
1900        let dir = tempfile::tempdir().unwrap();
1901        let path = dir.path().join("t.wav");
1902        crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
1903
1904        let mut player = Player::new();
1905        player.backend = Box::new(StuckBackend {
1906            rate: 8_000.0,
1907            asked: Default::default(),
1908        });
1909        let (recorder, events) = PlayRecorder::capture();
1910        player.history = Some(recorder);
1911
1912        let item = PlaylistItem {
1913            db_id: Some(5),
1914            path,
1915            ..make_item("t")
1916        };
1917        let id = item.id;
1918        player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
1919        player.process_command(PlayerCommand::Play(id));
1920        player.process_command(PlayerCommand::Pause);
1921        player.process_command(PlayerCommand::Seek(4_000));
1922        player.process_command(PlayerCommand::Resume);
1923        player.process_command(PlayerCommand::Stop);
1924
1925        let report = |state, position_ms| {
1926            PlayEvent::Playback(PlaybackReport {
1927                track_id: 5,
1928                state,
1929                position_ms,
1930            })
1931        };
1932        assert_eq!(
1933            events.try_iter().collect::<Vec<_>>(),
1934            vec![
1935                PlayEvent::Started {
1936                    track_id: 5,
1937                    position_ms: 0
1938                },
1939                report(Paused, 0),
1940                report(Paused, 4_000),
1941                report(Playing, 4_000),
1942                report(Stopped, 4_000),
1943                PlayEvent::Finished {
1944                    track_id: 5,
1945                    listened_ms: 0
1946                },
1947            ]
1948        );
1949    }
1950
1951    #[test]
1952    fn removing_the_playing_track_resumes_at_its_successor() {
1953        let mut player = Player::new();
1954        let ids = seed(&mut player, 5);
1955        player.shared_state.set_cursor(Some(ids[2]));
1956
1957        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
1958
1959        assert_eq!(
1960            player.shared_state.cursor(),
1961            Some(ids[3]),
1962            "playback must continue at the next track, not restart the queue"
1963        );
1964        assert_eq!(player.playback_starts, 1);
1965    }
1966
1967    #[test]
1968    fn removing_the_first_playing_track_resumes_at_the_new_first() {
1969        let mut player = Player::new();
1970        let ids = seed(&mut player, 3);
1971        player.shared_state.set_cursor(Some(ids[0]));
1972
1973        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[0]));
1974
1975        assert_eq!(player.shared_state.cursor(), Some(ids[1]));
1976    }
1977
1978    #[test]
1979    fn next_track_parks_on_a_track_that_has_not_downloaded_yet() {
1980        let mut player = Player::new();
1981        let playing = make_item("playing");
1982        let waiting = pending_item("waiting");
1983        let later = make_item("later");
1984        let (playing_id, waiting_id) = (playing.id, waiting.id);
1985        player.process_command(PlayerCommand::AddToPlaylist(vec![playing, waiting, later]));
1986        player.shared_state.set_cursor(Some(playing_id));
1987
1988        player.process_command(PlayerCommand::DecodeFinished);
1989
1990        assert_eq!(
1991            player.shared_state.cursor(),
1992            Some(waiting_id),
1993            "the cursor parks on the track being fetched"
1994        );
1995        assert_eq!(
1996            player.playback_starts, 0,
1997            "nothing to play until its bytes land"
1998        );
1999
2000        // The download completes. Because the cursor is parked here, the
2001        // TrackReady actually reaches the player and the queue resumes.
2002        player
2003            .shared_state
2004            .update_item_state(waiting_id, ItemState::Ready);
2005        player.process_command(PlayerCommand::TrackReady(waiting_id));
2006
2007        assert_eq!(player.playback_starts, 1);
2008        assert_eq!(player.shared_state.cursor(), Some(waiting_id));
2009    }
2010
2011    #[test]
2012    fn a_download_that_cannot_land_moves_the_cursor_on() {
2013        let mut player = Player::new();
2014        let waiting = pending_item("waiting");
2015        let later = make_item("later");
2016        let (waiting_id, later_id) = (waiting.id, later.id);
2017        player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, later]));
2018
2019        player.process_command(PlayerCommand::Play(waiting_id));
2020        assert_eq!(player.playback_starts, 0, "nothing to play yet");
2021
2022        // The download gives up. Ready will never come.
2023        player
2024            .shared_state
2025            .update_item_state(waiting_id, ItemState::Failed("remote unavailable".into()));
2026        player.process_command(PlayerCommand::TrackFailed(waiting_id));
2027
2028        assert_eq!(
2029            player.shared_state.cursor(),
2030            Some(later_id),
2031            "the queue moves past a track that can never load"
2032        );
2033        assert_eq!(player.playback_starts, 1);
2034    }
2035
2036    #[test]
2037    fn a_queue_that_can_never_load_stops_rather_than_waiting() {
2038        let mut player = Player::new();
2039        let first = pending_item("first");
2040        let second = pending_item("second");
2041        let (first_id, second_id) = (first.id, second.id);
2042        player.process_command(PlayerCommand::AddToPlaylist(vec![first, second]));
2043
2044        player.process_command(PlayerCommand::Play(first_id));
2045        for id in [first_id, second_id] {
2046            player
2047                .shared_state
2048                .update_item_state(id, ItemState::Failed("remote unavailable".into()));
2049            player.process_command(PlayerCommand::TrackFailed(id));
2050        }
2051
2052        assert_eq!(player.playback_starts, 0);
2053        assert_eq!(
2054            player.shared_state.playback_state(),
2055            PlaybackState::Stopped,
2056            "a stop the UI can see, not an indefinite wait for TrackReady"
2057        );
2058    }
2059
2060    #[test]
2061    fn a_failure_elsewhere_in_the_queue_leaves_the_cursor_alone() {
2062        let mut player = Player::new();
2063        let waiting = pending_item("waiting");
2064        let other = pending_item("other");
2065        let (waiting_id, other_id) = (waiting.id, other.id);
2066        player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, other]));
2067        player.process_command(PlayerCommand::Play(waiting_id));
2068
2069        player
2070            .shared_state
2071            .update_item_state(other_id, ItemState::Failed("remote unavailable".into()));
2072        player.process_command(PlayerCommand::TrackFailed(other_id));
2073
2074        assert_eq!(
2075            player.shared_state.cursor(),
2076            Some(waiting_id),
2077            "a track still downloading keeps the cursor"
2078        );
2079    }
2080
2081    #[test]
2082    fn batch_delete_containing_the_cursor_restarts_the_engine_once() {
2083        let mut player = Player::new();
2084        let ids = seed(&mut player, 5);
2085        player.shared_state.set_cursor(Some(ids[2]));
2086
2087        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![
2088            ids[1], ids[2], ids[3],
2089        ]));
2090
2091        assert_eq!(playlist_titles(&player), vec!["t0", "t4"]);
2092        assert_eq!(player.shared_state.cursor(), Some(ids[4]));
2093        assert_eq!(
2094            player.playback_starts, 1,
2095            "one resume for the whole selection, not one per deleted track"
2096        );
2097    }
2098
2099    #[test]
2100    fn batch_delete_below_the_cursor_leaves_playback_alone() {
2101        let mut player = Player::new();
2102        let ids = seed(&mut player, 4);
2103        player.shared_state.set_cursor(Some(ids[0]));
2104
2105        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![ids[2], ids[3]]));
2106
2107        assert_eq!(player.shared_state.cursor(), Some(ids[0]));
2108        assert_eq!(player.playback_starts, 0);
2109    }
2110
2111    #[test]
2112    fn undo_of_a_batch_delete_restores_the_original_order() {
2113        // The TUI collects a selection from a HashSet, so the IDs arrive in
2114        // arbitrary order — scrambled here so a snapshot that trusts that order
2115        // re-inserts C before B and lands it at the end of the playlist.
2116        let mut player = Player::new();
2117        let items = vec![
2118            make_item("A"),
2119            make_item("B"),
2120            make_item("C"),
2121            make_item("D"),
2122        ];
2123        let (b_id, c_id) = (items[1].id, items[2].id);
2124        player.process_command(PlayerCommand::AddToPlaylist(items));
2125
2126        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![c_id, b_id]));
2127        assert_eq!(playlist_titles(&player), vec!["A", "D"]);
2128
2129        player.process_command(PlayerCommand::Undo);
2130        assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2131    }
2132
2133    // --- AddToPlaylist undo/redo ---
2134
2135    #[test]
2136    fn undo_add_removes_items() {
2137        let mut player = Player::new();
2138        let items = vec![make_item("A"), make_item("B")];
2139        let ids: Vec<_> = items.iter().map(|i| i.id).collect();
2140
2141        player.process_command(PlayerCommand::AddToPlaylist(items));
2142        assert_eq!(playlist_ids(&player), ids);
2143        assert!(player.undo_stack().can_undo());
2144
2145        player.process_command(PlayerCommand::Undo);
2146        assert!(playlist_ids(&player).is_empty());
2147        assert!(player.undo_stack().can_redo());
2148    }
2149
2150    #[test]
2151    fn redo_add_restores_items() {
2152        let mut player = Player::new();
2153        let items = vec![make_item("A"), make_item("B")];
2154
2155        player.process_command(PlayerCommand::AddToPlaylist(items));
2156        player.process_command(PlayerCommand::Undo);
2157        assert!(playlist_ids(&player).is_empty());
2158
2159        player.process_command(PlayerCommand::Redo);
2160        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2161    }
2162
2163    // --- RemoveFromPlaylist undo/redo ---
2164
2165    #[test]
2166    fn undo_remove_restores_item_at_position() {
2167        let mut player = Player::new();
2168        let items = vec![make_item("A"), make_item("B"), make_item("C")];
2169        let b_id = items[1].id;
2170
2171        player.process_command(PlayerCommand::AddToPlaylist(items));
2172        player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2173        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2174
2175        player.process_command(PlayerCommand::Undo);
2176        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2177    }
2178
2179    #[test]
2180    fn undo_remove_first_item() {
2181        let mut player = Player::new();
2182        let items = vec![make_item("A"), make_item("B")];
2183        let a_id = items[0].id;
2184
2185        player.process_command(PlayerCommand::AddToPlaylist(items));
2186        player.process_command(PlayerCommand::RemoveFromPlaylist(a_id));
2187        assert_eq!(playlist_titles(&player), vec!["B"]);
2188
2189        player.process_command(PlayerCommand::Undo);
2190        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2191    }
2192
2193    #[test]
2194    fn undo_batch_remove_restores_all() {
2195        let mut player = Player::new();
2196        let items = vec![
2197            make_item("A"),
2198            make_item("B"),
2199            make_item("C"),
2200            make_item("D"),
2201        ];
2202        let b_id = items[1].id;
2203        let c_id = items[2].id;
2204
2205        player.process_command(PlayerCommand::AddToPlaylist(items));
2206        let version_before = player.shared_state.playlist_version();
2207        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![b_id, c_id]));
2208        assert_eq!(playlist_titles(&player), vec!["A", "D"]);
2209        // One bump for the whole batch. Bumping per item is what made clearing
2210        // a large queue crawl, and every bump wakes every client watching.
2211        assert_eq!(
2212            player.shared_state.playlist_version(),
2213            version_before + 1,
2214            "batch removal must bump the playlist version exactly once"
2215        );
2216
2217        // Single undo restores both
2218        player.process_command(PlayerCommand::Undo);
2219        assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2220    }
2221
2222    #[test]
2223    fn redo_batch_remove() {
2224        let mut player = Player::new();
2225        let items = vec![make_item("A"), make_item("B"), make_item("C")];
2226        let a_id = items[0].id;
2227        let b_id = items[1].id;
2228
2229        player.process_command(PlayerCommand::AddToPlaylist(items));
2230        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![a_id, b_id]));
2231        player.process_command(PlayerCommand::Undo);
2232        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2233
2234        player.process_command(PlayerCommand::Redo);
2235        assert_eq!(playlist_titles(&player), vec!["C"]);
2236    }
2237
2238    #[test]
2239    fn redo_remove() {
2240        let mut player = Player::new();
2241        let items = vec![make_item("A"), make_item("B"), make_item("C")];
2242        let b_id = items[1].id;
2243
2244        player.process_command(PlayerCommand::AddToPlaylist(items));
2245        player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2246        player.process_command(PlayerCommand::Undo);
2247        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2248
2249        player.process_command(PlayerCommand::Redo);
2250        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2251    }
2252
2253    // --- InsertInPlaylist undo/redo ---
2254
2255    #[test]
2256    fn undo_insert_removes_inserted_items() {
2257        let mut player = Player::new();
2258        let items = vec![make_item("A"), make_item("C")];
2259        let a_id = items[0].id;
2260
2261        player.process_command(PlayerCommand::AddToPlaylist(items));
2262
2263        let inserted = vec![make_item("B")];
2264        player.process_command(PlayerCommand::InsertInPlaylist {
2265            items: inserted,
2266            after: a_id,
2267        });
2268        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2269
2270        player.process_command(PlayerCommand::Undo);
2271        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2272    }
2273
2274    // --- MoveInPlaylist undo/redo ---
2275
2276    #[test]
2277    fn undo_move_restores_position() {
2278        let mut player = Player::new();
2279        let items = vec![make_item("A"), make_item("B"), make_item("C")];
2280        let a_id = items[0].id;
2281        let c_id = items[2].id;
2282
2283        player.process_command(PlayerCommand::AddToPlaylist(items));
2284
2285        // Move A after C: [B, C, A]
2286        player.process_command(PlayerCommand::MoveInPlaylist {
2287            id: a_id,
2288            target: c_id,
2289            after: true,
2290        });
2291        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2292
2293        player.process_command(PlayerCommand::Undo);
2294        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2295    }
2296
2297    #[test]
2298    fn redo_move() {
2299        let mut player = Player::new();
2300        let items = vec![make_item("A"), make_item("B"), make_item("C")];
2301        let a_id = items[0].id;
2302        let c_id = items[2].id;
2303
2304        player.process_command(PlayerCommand::AddToPlaylist(items));
2305        player.process_command(PlayerCommand::MoveInPlaylist {
2306            id: a_id,
2307            target: c_id,
2308            after: true,
2309        });
2310        player.process_command(PlayerCommand::Undo);
2311        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2312
2313        player.process_command(PlayerCommand::Redo);
2314        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2315    }
2316
2317    // --- MoveItemsInPlaylist (batch) undo/redo ---
2318
2319    #[test]
2320    fn undo_batch_move() {
2321        let mut player = Player::new();
2322        let items = vec![
2323            make_item("A"),
2324            make_item("B"),
2325            make_item("C"),
2326            make_item("D"),
2327        ];
2328        let a_id = items[0].id;
2329        let b_id = items[1].id;
2330        let d_id = items[3].id;
2331
2332        player.process_command(PlayerCommand::AddToPlaylist(items));
2333
2334        // Move A,B after D: [C, D, A, B]
2335        player.process_command(PlayerCommand::MoveItemsInPlaylist {
2336            ids: vec![a_id, b_id],
2337            target: d_id,
2338            after: true,
2339        });
2340        assert_eq!(playlist_titles(&player), vec!["C", "D", "A", "B"]);
2341
2342        player.process_command(PlayerCommand::Undo);
2343        assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2344    }
2345
2346    // --- ClearPlaylist undo/redo ---
2347
2348    #[test]
2349    fn undo_clear_restores_playlist() {
2350        let mut player = Player::new();
2351        let items = vec![make_item("A"), make_item("B"), make_item("C")];
2352
2353        player.process_command(PlayerCommand::AddToPlaylist(items));
2354        player.process_command(PlayerCommand::ClearPlaylist);
2355        assert!(playlist_ids(&player).is_empty());
2356
2357        player.process_command(PlayerCommand::Undo);
2358        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2359    }
2360
2361    /// The bug: replacing the queue starts the new track, and undoing restored
2362    /// the old queue while leaving the engine on a track that queue no longer
2363    /// contains — a transport describing a row nobody can see, and a decode
2364    /// lookahead with nothing to follow.
2365    #[test]
2366    fn undoing_a_replace_does_not_leave_the_engine_on_an_orphaned_track() {
2367        let mut player = Player::new();
2368        let original = seed(&mut player, 3);
2369        player.shared_state.set_cursor(Some(original[0]));
2370        pretend_playing(&mut player, original[0]);
2371
2372        let replacement = vec![make_item("something else")];
2373        let orphan = replacement[0].id;
2374        player.process_command(PlayerCommand::ReplacePlaylist {
2375            items: replacement,
2376            start: 0,
2377        });
2378        // What `play()` would have left behind if the file existed.
2379        pretend_playing(&mut player, orphan);
2380
2381        player.process_command(PlayerCommand::Undo);
2382
2383        assert_eq!(playlist_ids(&player), original, "the queue comes back");
2384        assert!(
2385            player.shared_state.get_item(orphan).is_none(),
2386            "and the replacement is gone from it"
2387        );
2388        assert!(
2389            playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
2390            "so nothing may still be playing out of it"
2391        );
2392    }
2393
2394    /// The same orphaning, reached by undoing an add rather than a replace.
2395    #[test]
2396    fn undoing_an_add_does_not_leave_the_engine_on_a_removed_track() {
2397        let mut player = Player::new();
2398        seed(&mut player, 2);
2399        let added = seed(&mut player, 1);
2400        pretend_playing(&mut player, added[0]);
2401
2402        player.process_command(PlayerCommand::Undo);
2403
2404        assert!(player.shared_state.get_item(added[0]).is_none());
2405        assert!(
2406            playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
2407            "the engine cannot be left on the item the undo removed"
2408        );
2409    }
2410
2411    /// An undo that leaves the playing item where it is must not restart it.
2412    #[test]
2413    fn undoing_a_move_leaves_playback_alone() {
2414        let mut player = Player::new();
2415        let ids = seed(&mut player, 3);
2416        player.shared_state.set_cursor(Some(ids[0]));
2417        pretend_playing(&mut player, ids[0]);
2418        let starts = player.playback_starts;
2419
2420        player.process_command(PlayerCommand::MoveInPlaylist {
2421            id: ids[2],
2422            target: ids[0],
2423            after: false,
2424        });
2425        player.process_command(PlayerCommand::Undo);
2426
2427        assert_eq!(playlist_ids(&player), ids);
2428        assert_eq!(playing_id(&player), Some(ids[0]), "still on the same track");
2429        assert_eq!(player.playback_starts, starts, "and not restarted");
2430    }
2431
2432    #[test]
2433    fn redo_clear() {
2434        let mut player = Player::new();
2435        let items = vec![make_item("A"), make_item("B")];
2436
2437        player.process_command(PlayerCommand::AddToPlaylist(items));
2438        player.process_command(PlayerCommand::ClearPlaylist);
2439        player.process_command(PlayerCommand::Undo);
2440        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2441
2442        player.process_command(PlayerCommand::Redo);
2443        assert!(playlist_ids(&player).is_empty());
2444    }
2445
2446    // --- Multi-step undo/redo ---
2447
2448    #[test]
2449    fn multiple_undos_in_sequence() {
2450        let mut player = Player::new();
2451
2452        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("A")]));
2453        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
2454        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("C")]));
2455        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2456
2457        player.process_command(PlayerCommand::Undo);
2458        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2459
2460        player.process_command(PlayerCommand::Undo);
2461        assert_eq!(playlist_titles(&player), vec!["A"]);
2462
2463        player.process_command(PlayerCommand::Undo);
2464        assert!(playlist_ids(&player).is_empty());
2465    }
2466
2467    #[test]
2468    fn undo_redo_undo_cycle() {
2469        let mut player = Player::new();
2470        let items = vec![make_item("A"), make_item("B")];
2471
2472        player.process_command(PlayerCommand::AddToPlaylist(items));
2473        player.process_command(PlayerCommand::Undo);
2474        assert!(playlist_ids(&player).is_empty());
2475
2476        player.process_command(PlayerCommand::Redo);
2477        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2478
2479        player.process_command(PlayerCommand::Undo);
2480        assert!(playlist_ids(&player).is_empty());
2481    }
2482
2483    #[test]
2484    fn new_action_clears_redo_stack() {
2485        let mut player = Player::new();
2486        let items = vec![make_item("A")];
2487
2488        player.process_command(PlayerCommand::AddToPlaylist(items));
2489        player.process_command(PlayerCommand::Undo);
2490        assert!(player.undo_stack().can_redo());
2491
2492        // New action should clear redo
2493        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
2494        assert!(!player.undo_stack().can_redo());
2495    }
2496
2497    #[test]
2498    fn undo_on_empty_stack_is_noop() {
2499        let mut player = Player::new();
2500        player.process_command(PlayerCommand::Undo);
2501        assert!(playlist_ids(&player).is_empty());
2502    }
2503
2504    #[test]
2505    fn redo_on_empty_stack_is_noop() {
2506        let mut player = Player::new();
2507        player.process_command(PlayerCommand::Redo);
2508        assert!(playlist_ids(&player).is_empty());
2509    }
2510
2511    // --- Non-undoable commands don't push entries ---
2512
2513    #[test]
2514    fn playback_commands_not_undoable() {
2515        let mut player = Player::new();
2516        player.process_command(PlayerCommand::Pause);
2517        player.process_command(PlayerCommand::Resume);
2518        player.process_command(PlayerCommand::NextTrack);
2519        player.process_command(PlayerCommand::PrevTrack);
2520        assert!(!player.undo_stack().can_undo());
2521    }
2522
2523    #[test]
2524    fn update_paths_not_undoable() {
2525        let mut player = Player::new();
2526        let items = vec![make_item("A")];
2527        let id = items[0].id;
2528        player.process_command(PlayerCommand::AddToPlaylist(items));
2529
2530        let undo_count = player.undo_stack().undo_len();
2531        player.process_command(PlayerCommand::UpdatePaths(vec![(
2532            id,
2533            PathBuf::from("/new/path.flac"),
2534        )]));
2535        assert_eq!(player.undo_stack().undo_len(), undo_count);
2536    }
2537
2538    // --- Complex scenarios ---
2539
2540    #[test]
2541    fn add_remove_undo_undo_produces_original() {
2542        let mut player = Player::new();
2543        let items = vec![make_item("A"), make_item("B"), make_item("C")];
2544        let b_id = items[1].id;
2545        let original_titles = vec!["A", "B", "C"];
2546
2547        player.process_command(PlayerCommand::AddToPlaylist(items));
2548        player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2549        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2550
2551        // Undo remove → back to A, B, C
2552        player.process_command(PlayerCommand::Undo);
2553        assert_eq!(playlist_titles(&player), original_titles);
2554
2555        // Undo add → empty
2556        player.process_command(PlayerCommand::Undo);
2557        assert!(playlist_ids(&player).is_empty());
2558    }
2559
2560    #[test]
2561    fn interleaved_adds_and_moves_undo() {
2562        let mut player = Player::new();
2563        let items = vec![make_item("A"), make_item("B"), make_item("C")];
2564        let a_id = items[0].id;
2565        let c_id = items[2].id;
2566
2567        player.process_command(PlayerCommand::AddToPlaylist(items));
2568
2569        // Move A after C: [B, C, A]
2570        player.process_command(PlayerCommand::MoveInPlaylist {
2571            id: a_id,
2572            target: c_id,
2573            after: true,
2574        });
2575        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2576
2577        // Add D: [B, C, A, D]
2578        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("D")]));
2579        assert_eq!(playlist_titles(&player), vec!["B", "C", "A", "D"]);
2580
2581        // Undo add D: [B, C, A]
2582        player.process_command(PlayerCommand::Undo);
2583        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2584
2585        // Undo move: [A, B, C]
2586        player.process_command(PlayerCommand::Undo);
2587        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2588    }
2589
2590    /// Regression test for GitHub #89: AudioEngine must be dropped synchronously
2591    /// in stop_engine() before the caller changes sample rates. If the engine is
2592    /// dropped on a background thread, CoreAudio's internal buffer list can be
2593    /// freed while AudioUnitUninitialize is still tearing it down → crash.
2594    #[test]
2595    fn stop_engine_drops_engine_synchronously() {
2596        use std::sync::atomic::{AtomicBool, Ordering};
2597
2598        struct MockEngine {
2599            dropped: Arc<AtomicBool>,
2600        }
2601        impl AudioEngineHandle for MockEngine {
2602            fn start(&self) -> Result<(), BackendError> {
2603                Ok(())
2604            }
2605            fn stop(&self) -> Result<(), BackendError> {
2606                Ok(())
2607            }
2608            fn is_running(&self) -> bool {
2609                false
2610            }
2611            fn fade_out(&self) {}
2612            fn fade_in(&self) -> Result<(), BackendError> {
2613                Ok(())
2614            }
2615            fn is_silent(&self) -> bool {
2616                false
2617            }
2618        }
2619        impl Drop for MockEngine {
2620            fn drop(&mut self) {
2621                self.dropped.store(true, Ordering::SeqCst);
2622            }
2623        }
2624
2625        let dropped = Arc::new(AtomicBool::new(false));
2626
2627        // Build a minimal decode handle that won't block.
2628        let stop_flag = Arc::new(AtomicBool::new(false));
2629        let decode_handle = buffer::DecodeHandle::new_for_test(stop_flag);
2630
2631        let mut player = Player::new();
2632        player.active_playback = Some(ActivePlayback {
2633            engine: Box::new(MockEngine {
2634                dropped: dropped.clone(),
2635            }),
2636            decode_handle,
2637            stream: None,
2638            _rate_watch: None,
2639        });
2640
2641        player.stop_engine();
2642
2643        // The engine must already be dropped when stop_engine returns.
2644        // If this fails, the engine was moved to a background thread — the
2645        // exact race condition that causes the #89 crash.
2646        assert!(
2647            dropped.load(Ordering::SeqCst),
2648            "AudioEngine must be dropped synchronously in stop_engine (GitHub #89)"
2649        );
2650    }
2651
2652    #[test]
2653    fn abandoning_a_stream_wakes_a_reader_parked_at_the_write_head() {
2654        let live = LiveStream {
2655            feed: crate::remote::downloads::ByteFeed::new(),
2656            abandoned: Default::default(),
2657        };
2658        let feed = live.feed.clone();
2659        let started = std::time::Instant::now();
2660        let reader = thread::spawn(move || {
2661            feed.wait_past(
2662                0,
2663                std::time::Instant::now() + std::time::Duration::from_secs(30),
2664            )
2665        });
2666        thread::sleep(std::time::Duration::from_millis(50));
2667        live.abandon();
2668        reader.join().unwrap();
2669
2670        assert!(live.abandoned.load(std::sync::atomic::Ordering::Acquire));
2671        assert!(started.elapsed() < std::time::Duration::from_secs(5));
2672    }
2673
2674    // --- Engine format matches the decoded PCM ---
2675
2676    /// Backend pinned to one sample rate that refuses every switch, recording
2677    /// the format the engine is asked for.
2678    struct StuckBackend {
2679        rate: f64,
2680        asked: Arc<std::sync::Mutex<Option<(f64, u32)>>>,
2681    }
2682
2683    struct NullEngine;
2684    impl AudioEngineHandle for NullEngine {
2685        fn start(&self) -> Result<(), BackendError> {
2686            Ok(())
2687        }
2688        fn stop(&self) -> Result<(), BackendError> {
2689            Ok(())
2690        }
2691        fn is_running(&self) -> bool {
2692            false
2693        }
2694        fn fade_out(&self) {}
2695        fn fade_in(&self) -> Result<(), BackendError> {
2696            Ok(())
2697        }
2698        fn is_silent(&self) -> bool {
2699            false
2700        }
2701    }
2702
2703    impl AudioBackend for StuckBackend {
2704        fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2705            Ok(vec![self.default_device()?])
2706        }
2707        fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2708            Ok(backend::DeviceInfo {
2709                name: "Stuck DAC".into(),
2710                sample_rates: vec![self.rate],
2711                platform_id: 0,
2712            })
2713        }
2714        fn supported_sample_rates(
2715            &self,
2716            _device: &backend::DeviceInfo,
2717        ) -> Result<Vec<f64>, BackendError> {
2718            Ok(vec![self.rate])
2719        }
2720        fn get_device_sample_rate(
2721            &self,
2722            _device: &backend::DeviceInfo,
2723        ) -> Result<f64, BackendError> {
2724            Ok(self.rate)
2725        }
2726        fn set_device_sample_rate(
2727            &self,
2728            _device: &backend::DeviceInfo,
2729            rate: f64,
2730        ) -> Result<f64, BackendError> {
2731            Err(BackendError::UnsupportedSampleRate(rate))
2732        }
2733        fn create_engine(
2734            &self,
2735            _device: &backend::DeviceInfo,
2736            sample_rate: f64,
2737            channels: u32,
2738            _consumer: rtrb::Consumer<f32>,
2739            _samples_played: Arc<AtomicU64>,
2740        ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2741            *self.asked.lock().unwrap() = Some((sample_rate, channels));
2742            Ok(Box::new(NullEngine))
2743        }
2744    }
2745
2746    fn engine_format_for(source_rate: u32, channels: u16, device_rate: f64) -> (f64, u32) {
2747        let asked = Arc::new(std::sync::Mutex::new(None));
2748        let mut player = Player::new();
2749        player.backend = Box::new(StuckBackend {
2750            rate: device_rate,
2751            asked: asked.clone(),
2752        });
2753
2754        let info = buffer::StreamInfo {
2755            codec: "MP3".into(),
2756            sample_rate: source_rate,
2757            channels,
2758            bit_depth: Some(16),
2759            bitrate_kbps: None,
2760            duration_ms: 1000,
2761        };
2762        let (_producer, consumer) = rtrb::RingBuffer::new(16);
2763        player
2764            .create_engine_for(&info, consumer)
2765            .expect("engine creation should succeed");
2766        let asked = *asked.lock().unwrap();
2767        asked.expect("engine was never created")
2768    }
2769
2770    #[test]
2771    fn engine_uses_source_rate_when_device_refuses_switch() {
2772        // MPEG-2 MP3 rates are routinely rejected by output devices. The engine
2773        // must still be told the rate the PCM actually is.
2774        assert_eq!(engine_format_for(22050, 2, 48000.0), (22050.0, 2));
2775        assert_eq!(engine_format_for(32000, 2, 44100.0), (32000.0, 2));
2776    }
2777
2778    #[test]
2779    fn engine_uses_source_channel_count() {
2780        assert_eq!(engine_format_for(44100, 1, 44100.0), (44100.0, 1));
2781    }
2782
2783    /// The rate the device settled at, as the front ends read it.
2784    fn settled_rate_for(source_rate: u32, device_rate: f64) -> Option<u32> {
2785        let mut player = Player::new();
2786        player.backend = Box::new(StuckBackend {
2787            rate: device_rate,
2788            asked: Arc::new(std::sync::Mutex::new(None)),
2789        });
2790        let state = player.shared_state.clone();
2791
2792        let info = buffer::StreamInfo {
2793            codec: "MP3".into(),
2794            sample_rate: source_rate,
2795            channels: 2,
2796            bit_depth: Some(16),
2797            bitrate_kbps: None,
2798            duration_ms: 1000,
2799        };
2800        let (_producer, consumer) = rtrb::RingBuffer::new(16);
2801        player
2802            .create_engine_for(&info, consumer)
2803            .expect("engine creation should succeed");
2804        state.output_sample_rate()
2805    }
2806
2807    #[test]
2808    fn settled_device_rate_reaches_the_shared_state() {
2809        // A device that refuses the switch is being fed resampled audio, and
2810        // that is the case the front ends have to be able to see. Before this
2811        // the comparison happened once, in a log line.
2812        assert_eq!(settled_rate_for(22050, 48000.0), Some(48000));
2813        // No switch needed, so nothing resampled: the two rates agree.
2814        assert_eq!(settled_rate_for(44100, 44100.0), Some(44100));
2815    }
2816
2817    /// A device that takes its time reclocking, as real hardware does.
2818    struct SlowBackend {
2819        observed: Arc<std::sync::Mutex<Vec<Option<u32>>>>,
2820        state: Arc<SharedPlayerState>,
2821    }
2822
2823    impl AudioBackend for SlowBackend {
2824        fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2825            Ok(vec![self.default_device()?])
2826        }
2827        fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2828            Ok(backend::DeviceInfo {
2829                name: "Slow DAC".into(),
2830                sample_rates: vec![44100.0, 48000.0],
2831                platform_id: 0,
2832            })
2833        }
2834        fn supported_sample_rates(
2835            &self,
2836            _device: &backend::DeviceInfo,
2837        ) -> Result<Vec<f64>, BackendError> {
2838            Ok(vec![44100.0, 48000.0])
2839        }
2840        fn get_device_sample_rate(
2841            &self,
2842            _device: &backend::DeviceInfo,
2843        ) -> Result<f64, BackendError> {
2844            Ok(48000.0)
2845        }
2846        fn set_device_sample_rate(
2847            &self,
2848            _device: &backend::DeviceInfo,
2849            rate: f64,
2850        ) -> Result<f64, BackendError> {
2851            // What a front end polling mid-switch would see.
2852            self.observed
2853                .lock()
2854                .unwrap()
2855                .push(self.state.output_sample_rate());
2856            Ok(rate)
2857        }
2858        fn create_engine(
2859            &self,
2860            _device: &backend::DeviceInfo,
2861            _sample_rate: f64,
2862            _channels: u32,
2863            _consumer: rtrb::Consumer<f32>,
2864            _samples_played: Arc<AtomicU64>,
2865        ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2866            Ok(Box::new(NullEngine))
2867        }
2868    }
2869
2870    #[test]
2871    fn the_previous_rate_is_not_published_while_the_device_reclocks() {
2872        // A 48 kHz track followed by a 44.1 kHz one: for as long as the switch
2873        // takes — the better part of a second on USB — the new track's info is
2874        // published against the old track's output rate. A front end polling in
2875        // that window used to latch "44.1 → 48" and, since nothing about the
2876        // codec or the source rate changed afterwards, never let go of it.
2877        let mut player = Player::new();
2878        let state = player.shared_state.clone();
2879        state.set_output_sample_rate(48000);
2880
2881        let observed = Arc::new(std::sync::Mutex::new(Vec::new()));
2882        player.backend = Box::new(SlowBackend {
2883            observed: observed.clone(),
2884            state: state.clone(),
2885        });
2886
2887        let info = buffer::StreamInfo {
2888            codec: "FLAC".into(),
2889            sample_rate: 44100,
2890            channels: 2,
2891            bit_depth: Some(16),
2892            bitrate_kbps: None,
2893            duration_ms: 1000,
2894        };
2895        let (_producer, consumer) = rtrb::RingBuffer::new(16);
2896        player
2897            .create_engine_for(&info, consumer)
2898            .expect("engine creation should succeed");
2899
2900        assert_eq!(
2901            *observed.lock().unwrap(),
2902            vec![None],
2903            "mid-switch the output rate must read as unknown, not as the last track's"
2904        );
2905        assert_eq!(state.output_sample_rate(), Some(44100));
2906    }
2907
2908    /// Backend that hands its rate-change callback back to the test.
2909    struct WatchedBackend {
2910        inner: StuckBackend,
2911        #[allow(clippy::type_complexity)]
2912        captured: Arc<std::sync::Mutex<Option<Box<dyn Fn(f64) + Send + Sync>>>>,
2913    }
2914
2915    struct NullWatch;
2916    impl backend::SampleRateWatch for NullWatch {}
2917
2918    impl AudioBackend for WatchedBackend {
2919        fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2920            self.inner.list_devices()
2921        }
2922        fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2923            self.inner.default_device()
2924        }
2925        fn supported_sample_rates(
2926            &self,
2927            device: &backend::DeviceInfo,
2928        ) -> Result<Vec<f64>, BackendError> {
2929            self.inner.supported_sample_rates(device)
2930        }
2931        fn get_device_sample_rate(
2932            &self,
2933            device: &backend::DeviceInfo,
2934        ) -> Result<f64, BackendError> {
2935            self.inner.get_device_sample_rate(device)
2936        }
2937        fn set_device_sample_rate(
2938            &self,
2939            device: &backend::DeviceInfo,
2940            rate: f64,
2941        ) -> Result<f64, BackendError> {
2942            self.inner.set_device_sample_rate(device, rate)
2943        }
2944        fn watch_device_sample_rate(
2945            &self,
2946            _device: &backend::DeviceInfo,
2947            on_change: Box<dyn Fn(f64) + Send + Sync>,
2948        ) -> Option<Box<dyn backend::SampleRateWatch>> {
2949            *self.captured.lock().unwrap() = Some(on_change);
2950            Some(Box::new(NullWatch))
2951        }
2952        fn create_engine(
2953            &self,
2954            device: &backend::DeviceInfo,
2955            sample_rate: f64,
2956            channels: u32,
2957            consumer: rtrb::Consumer<f32>,
2958            samples_played: Arc<AtomicU64>,
2959        ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2960            self.inner
2961                .create_engine(device, sample_rate, channels, consumer, samples_played)
2962        }
2963    }
2964
2965    #[test]
2966    fn external_rate_change_reaches_the_shared_state() {
2967        // The device is shared. Another client moving the rate mid-track used
2968        // to leave the front ends asserting bit-perfection while the HAL
2969        // resampled underneath them.
2970        let captured = Arc::new(std::sync::Mutex::new(None));
2971        let mut player = Player::new();
2972        player.backend = Box::new(WatchedBackend {
2973            inner: StuckBackend {
2974                rate: 44100.0,
2975                asked: Arc::new(std::sync::Mutex::new(None)),
2976            },
2977            captured: captured.clone(),
2978        });
2979        let state = player.shared_state.clone();
2980
2981        let info = buffer::StreamInfo {
2982            codec: "FLAC".into(),
2983            sample_rate: 44100,
2984            channels: 2,
2985            bit_depth: Some(16),
2986            bitrate_kbps: None,
2987            duration_ms: 1000,
2988        };
2989        let (_producer, consumer) = rtrb::RingBuffer::new(16);
2990        player
2991            .create_engine_for(&info, consumer)
2992            .expect("engine creation should succeed");
2993        assert_eq!(state.output_sample_rate(), Some(44100));
2994
2995        let on_change = captured.lock().unwrap().take().expect("watch registered");
2996        on_change(48000.0);
2997        assert_eq!(state.output_sample_rate(), Some(48000));
2998    }
2999}