Skip to main content

koan_core/player/
mod.rs

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