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