Skip to main content

koan_core/player/
mod.rs

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