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