Skip to main content

koan_core/player/
mod.rs

1pub mod commands;
2pub mod history;
3mod renderer;
4pub mod state;
5pub mod undo;
6
7use std::path::{Path, PathBuf};
8use std::sync::Arc;
9use std::thread;
10
11use thiserror::Error;
12
13use crate::audio::{
14    analyzer::VizAnalyzer,
15    backend::{self, AudioBackend, AudioEngineHandle, BackendError, SampleRateWatch},
16    buffer, streaming,
17    viz::{VizBuffer, VizSnapshot},
18};
19use crate::remote::client::PlaybackReportState;
20use buffer::PlaybackTimeline;
21use commands::{CommandChannel, PlayerCommand};
22use history::{InFlight, PlayEvent, PlayRecorder, PlaybackReport};
23use state::{
24    ItemState, PlayMode, PlaybackSource, PlaybackState, QueueItemId, Repeat, SharedPlayerState,
25    TrackInfo,
26};
27use undo::{UndoEntry, UndoStack};
28
29/// Ring buffer size in samples. ~1s at 192kHz stereo.
30pub(crate) const RING_BUFFER_SIZE: usize = 192_000 * 2;
31
32/// Kept back from the end of a track when seeking, so dragging the thumb all
33/// the way over lands in the last moment of it rather than in the next track.
34const SEEK_END_GUARD_MS: u64 = 500;
35
36/// Past the moment the next track should start, so the wake finds the
37/// playhead already in it rather than a hair short.
38const BOUNDARY_SLACK: std::time::Duration = std::time::Duration::from_millis(5);
39
40/// How often to look at a fading pause, for the few checks it takes to reach
41/// silence.
42const FADE_CHECK: std::time::Duration = std::time::Duration::from_millis(50);
43
44#[derive(Debug, Error)]
45pub enum PlayerError {
46    #[error("backend error: {0}")]
47    Backend(#[from] BackendError),
48    #[error("decode error: {0}")]
49    Decode(#[from] buffer::DecodeError),
50    #[error("renderer: {0}")]
51    Renderer(String),
52    /// The renderer cannot play this track. It is marked as such in the
53    /// queue and skipped.
54    #[error("{0}")]
55    Unplayable(String),
56}
57
58/// Everything needed to read a track that is still downloading: where it is,
59/// how far the transfer has got, and how the container has to be opened.
60#[derive(Clone)]
61struct StreamSource {
62    path: PathBuf,
63    bytes_written: Arc<crate::remote::downloads::ByteFeed>,
64    total: u64,
65    mode: streaming::ProbeMode,
66}
67
68/// A file's own extension, lowercased. A download in progress is named
69/// `track.m4a.part`, and its extension is the track's, not `part`.
70fn media_extension(path: &Path) -> Option<String> {
71    crate::remote::download::strip_part_suffix(path)
72        .extension()
73        .and_then(|e| e.to_str())
74        .map(str::to_ascii_lowercase)
75}
76
77/// Symphonia's format hint for a path — its extension, where it has one.
78fn hint_for(path: &Path) -> symphonia::core::formats::probe::Hint {
79    let mut hint = symphonia::core::formats::probe::Hint::new();
80    if let Some(ext) = media_extension(path) {
81        hint.with_extension(&ext);
82    }
83    hint
84}
85
86/// How to open a partial file once the whole description has failed.
87///
88/// Most containers describe their frames from the front and need an end they
89/// can actually reach: FLAC bisects towards the end it is given. Two need the
90/// whole file's end instead. Ogg takes the end it is handed as the end of the
91/// stream, so a file ending at the write head is a track already over. MP4
92/// bounds its top-level boxes by the end, so a `moov` larger than what has
93/// arrived overruns it and the file will not open until it has all landed.
94/// See `ProbeMode`.
95fn lengthless_mode_for(path: &Path) -> streaming::ProbeMode {
96    match media_extension(path).as_deref() {
97        Some("ogg" | "oga" | "opus" | "spx" | "m4a" | "m4b" | "mp4" | "mov") => {
98            streaming::ProbeMode::LengthlessWholeEnd
99        }
100        _ => streaming::ProbeMode::Lengthless,
101    }
102}
103
104/// The player controller. Owns the audio pipeline and processes commands.
105pub struct Player {
106    shared_state: Arc<SharedPlayerState>,
107    commands: CommandChannel,
108    /// What the player is doing with the track under the cursor.
109    transport: Transport,
110    timeline: Arc<PlaybackTimeline>,
111    viz_buffer: Arc<VizBuffer>,
112    viz_snapshot: Arc<VizSnapshot>,
113    /// Background FFT analysis thread. Held for its lifetime; dropped on Player drop.
114    _viz_analyzer: VizAnalyzer,
115    undo_stack: UndoStack,
116    /// When Some, undo entries are collected into this buffer instead of pushed
117    /// directly onto the undo stack. Flushed on EndUndoBatch.
118    batch_buffer: Option<Vec<UndoEntry>>,
119    /// Configured output device name. None = system default.
120    output_device_name: Option<String>,
121    /// Platform audio backend (CoreAudio on macOS and iOS, cpal on Linux).
122    backend: Box<dyn AudioBackend>,
123    /// How the file currently streaming had to be opened. A seek reopens it and
124    /// must not undo what the probe settled on.
125    stream_mode: streaming::ProbeMode,
126    /// Writes plays away from this thread. None when there is no database to
127    /// write to, and in tests, which must not touch the real library.
128    history: Option<PlayRecorder>,
129    /// This player's download queue, which follows its playlist. Held here so
130    /// it lives as long as the player it fetches for. `None` for a player
131    /// made without `spawn`, which fetches nothing.
132    downloads: Option<crate::remote::queue::DownloadQueue>,
133    /// How much of the current track has been heard so far.
134    in_flight: Option<InFlight>,
135    /// When the silence after a rate switch runs out and the track is heard.
136    lead_in_ends: Option<std::time::Instant>,
137    /// Bumped by every session opened and every one torn down, so that what a
138    /// torn-down session reports after the fact is recognised as stale.
139    session: u64,
140    /// Waiting for a pause's fade to reach silence, to hear where it did.
141    silence_waiters: Vec<crossbeam_channel::Sender<u64>>,
142    /// The DSP setup last loaded, and the config and device it was loaded
143    /// for. Reading impulse responses off disk on every seek would be wasted.
144    dsp: Option<DspCache>,
145    /// The UPnP renderer chosen as the output, when one is: every session
146    /// opened while it is set plays there. Not a transport of its own.
147    renderer: Option<renderer::RendererLink>,
148    /// Shuffle and repeat. Published by `publish`, which the queue's own
149    /// reads follow.
150    mode: PlayMode,
151    /// Playback sessions started — lets tests assert how many engine restarts
152    /// an operation costs.
153    #[cfg(test)]
154    playback_starts: usize,
155    /// What `dsp_for` answers in tests, in place of reading the profiles
156    /// from config: config is the process's, and the suite runs in parallel.
157    /// The first is this device's, the second the renderer's.
158    #[cfg(test)]
159    dsp_override: Option<Arc<crate::audio::dsp::Setup>>,
160    #[cfg(test)]
161    renderer_dsp_override: Option<Arc<crate::audio::dsp::Setup>>,
162    /// The renderer used last time is being looked for, and is to be gone
163    /// back to if it turns up before anyone plays or picks an output.
164    resume_renderer: bool,
165}
166
167struct DspCache {
168    config: Arc<crate::config::Config>,
169    device: String,
170    setup: Option<Arc<crate::audio::dsp::Setup>>,
171}
172
173/// Whether a session plays or sits paused, and how a track waited for opens.
174/// A fade out is `Paused` with the engine still running until it is silent.
175#[derive(Clone, Copy, PartialEq, Eq)]
176enum Run {
177    Playing,
178    Paused,
179}
180
181/// What the player is doing. Everything it publishes — the playback state,
182/// the wait, the track — is derived from this in `publish`, so none of them
183/// can disagree with it or with each other.
184enum Transport {
185    /// Nothing loaded and nothing asked for.
186    Idle,
187    /// A track asked for that cannot open yet. Nothing opens a track without
188    /// one of these or a command naming it.
189    Waiting(Waiting),
190    /// A session is open, decoding into the ring.
191    Loaded(Session),
192}
193
194/// A track that cannot open yet, and how it opens once it can: at the start
195/// as soon as enough of it has streamed, or at a later position once the whole
196/// file is on disk. Streaming it would start it before the position can be
197/// reached, and a seek once it can would let the part before it be heard.
198#[derive(Clone, Copy)]
199struct Waiting {
200    id: QueueItemId,
201    position_ms: u64,
202    start: Run,
203}
204
205impl Waiting {
206    fn may_stream(&self, id: QueueItemId) -> bool {
207        self.id == id && self.position_ms == 0
208    }
209}
210
211/// An open playback session: the track under the playhead, whether it plays,
212/// and the output it plays on.
213struct Session {
214    /// The track under the playhead. Opened with the first; moved on by a
215    /// gapless transition, corrected when a download it streams lands.
216    track: TrackInfo,
217    run: Run,
218    /// The steps the decoder has taken through the queue, in order. Always
219    /// empty for a renderer, which has no decoder here to look ahead.
220    lookahead: Arc<parking_lot::Mutex<Vec<state::Lookahead>>>,
221    output: Output,
222}
223
224/// Where a session's sound comes out, and so whether koan decodes it.
225enum Output {
226    /// Decoded here, into the ring, drained by an engine: this device's own
227    /// output.
228    Local(Local),
229    /// A renderer: handed the original file, which it decodes and plays
230    /// gaplessly itself, or a stream decoded and processed here. See
231    /// `renderer`.
232    Renderer(Box<renderer::Play>),
233}
234
235struct Local {
236    engine: Box<dyn AudioEngineHandle>,
237    decode_handle: buffer::DecodeHandle,
238    /// Set when the decoder is reading a download as it arrives.
239    stream: Option<LiveStream>,
240    /// Keeps the device rate subscription alive for as long as this engine is
241    /// the one feeding the DAC. Dropped with it.
242    _rate_watch: Option<Box<dyn SampleRateWatch>>,
243    /// What the output device's profile does to this session, for the badge.
244    dsp: Option<crate::audio::dsp::DspStatus>,
245}
246
247impl Session {
248    /// The engine, for a session on this device.
249    fn engine(&self) -> Option<&dyn AudioEngineHandle> {
250        match &self.output {
251            Output::Local(local) => Some(local.engine.as_ref()),
252            Output::Renderer(_) => None,
253        }
254    }
255}
256
257/// Where a session reads its first track from.
258enum Source {
259    File(PathBuf),
260    Stream(StreamSource),
261}
262
263impl Source {
264    fn path(&self) -> &Path {
265        match self {
266            Source::File(path) => path,
267            Source::Stream(source) => &source.path,
268        }
269    }
270}
271
272/// A download being decoded as it lands. The reader may be parked at the write
273/// head waiting for bytes, so stopping has to tell it to give up and wake it,
274/// or the join waits on the network.
275struct LiveStream {
276    feed: Arc<crate::remote::downloads::ByteFeed>,
277    abandoned: Arc<std::sync::atomic::AtomicBool>,
278}
279
280impl LiveStream {
281    fn abandon(&self) {
282        self.abandoned
283            .store(true, std::sync::atomic::Ordering::Release);
284        self.feed.done();
285    }
286}
287
288impl Default for Player {
289    fn default() -> Self {
290        Self::new()
291    }
292}
293
294impl Player {
295    pub fn new() -> Self {
296        let viz_buffer = VizBuffer::new();
297        let viz_snapshot = VizSnapshot::new();
298        let timeline = PlaybackTimeline::new();
299        let cfg = crate::config::Config::cached();
300        let viz_analyzer = VizAnalyzer::spawn_with_snapshot(
301            Arc::clone(&viz_buffer),
302            &cfg.visualizer,
303            Arc::clone(&viz_snapshot),
304            timeline.samples_played_counter(),
305        );
306
307        let shared_state = SharedPlayerState::new();
308        shared_state.attach_timeline(timeline.clone());
309        let commands = CommandChannel::new();
310        let tx = commands.tx.clone();
311        timeline.on_queued(move || {
312            let _ = tx.try_send(PlayerCommand::TrackQueued);
313        });
314
315        Self {
316            shared_state,
317            commands,
318            transport: Transport::Idle,
319            lead_in_ends: None,
320            dsp: None,
321            session: 0,
322            silence_waiters: Vec::new(),
323            renderer: None,
324            mode: PlayMode::default(),
325            timeline,
326            viz_buffer,
327            viz_snapshot,
328            _viz_analyzer: viz_analyzer,
329            undo_stack: UndoStack::new(),
330            batch_buffer: None,
331            output_device_name: cfg.playback.output_device.clone(),
332            backend: crate::audio::platform_backend(),
333            stream_mode: streaming::ProbeMode::Full,
334            history: None,
335            downloads: None,
336            in_flight: None,
337            #[cfg(test)]
338            playback_starts: 0,
339            #[cfg(test)]
340            dsp_override: None,
341            #[cfg(test)]
342            renderer_dsp_override: None,
343            resume_renderer: false,
344        }
345    }
346
347    /// Get a clone of the shared state for UI reads.
348    pub fn shared_state(&self) -> Arc<SharedPlayerState> {
349        self.shared_state.clone()
350    }
351
352    /// Get the playback timeline for UI reads.
353    pub fn timeline(&self) -> Arc<PlaybackTimeline> {
354        self.timeline.clone()
355    }
356
357    /// Get the visualization buffer for the TUI.
358    pub fn viz_buffer(&self) -> Arc<VizBuffer> {
359        self.viz_buffer.clone()
360    }
361
362    /// Get the shared analysis snapshot for the TUI.
363    /// The analysis thread writes here; the UI thread reads a clone each frame.
364    pub fn viz_snapshot(&self) -> Arc<VizSnapshot> {
365        self.viz_snapshot.clone()
366    }
367
368    /// Access undo stack (for tests and UI state queries).
369    pub fn undo_stack(&self) -> &UndoStack {
370        &self.undo_stack
371    }
372
373    /// Create an audio engine for a stream, switching the output device to the
374    /// source rate first so output is bit-perfect.
375    ///
376    /// The engine is always configured with the source's own rate and channel
377    /// count — the format the decode thread writes into the ring buffer. A
378    /// device that cannot take the requested rate (MPEG-2/2.5 MP3 rates are
379    /// commonly refused) resamples instead of playing at the wrong speed.
380    #[allow(clippy::type_complexity)]
381    fn create_engine_for(
382        &mut self,
383        info: &buffer::StreamInfo,
384        consumer: rtrb::Consumer<f32>,
385    ) -> Result<(Box<dyn AudioEngineHandle>, Option<Box<dyn SampleRateWatch>>), PlayerError> {
386        let device = self.resolve_device()?;
387        let device_rate = self.backend.get_device_sample_rate(&device)?;
388        // The setup the decode thread was just handed, not a fresh read: a
389        // config edited in between would put the engine at another rate.
390        let dsp = match &self.dsp {
391            Some(cache) if cache.device == device.name => cache.setup.clone(),
392            _ => self.dsp_for(&device.name),
393        };
394        // The rate the decode thread writes at: the source's, unless DSP
395        // resamples it to reach an impulse response.
396        let source_rate =
397            dsp.as_ref()
398                .map_or(info.sample_rate, |d| d.output_rate(info.sample_rate)) as f64;
399
400        // Anything read between here and the switch landing would pair this
401        // track with the last one's output rate, and a rate switch is not
402        // instant. Say nothing instead.
403        self.shared_state.clear_output_sample_rate();
404
405        let settled = if (device_rate - source_rate).abs() > 0.1 {
406            log::info!(
407                "switching device sample rate: {}Hz → {}Hz",
408                device_rate,
409                source_rate
410            );
411            match self.backend.set_device_sample_rate(&device, source_rate) {
412                Ok(rate) => rate,
413                Err(e) => {
414                    log::warn!("failed to set device sample rate: {}", e);
415                    device_rate
416                }
417            }
418        } else {
419            device_rate
420        };
421
422        if (settled - source_rate).abs() > 0.1 {
423            log::warn!(
424                "device stayed at {}Hz (wanted {}Hz) — output is resampled, not bit-perfect",
425                settled,
426                source_rate
427            );
428        }
429
430        // The front ends compare this against the source rate to say whether
431        // anything had to resample.
432        self.shared_state
433            .set_output_sample_rate(settled.round() as u32);
434
435        // koan is not the only client of this device. Subscribe so the front
436        // ends learn about a rate someone else moved instead of trusting the
437        // reading above until the next track happens to build an engine.
438        let watch_state = self.shared_state.clone();
439        let watch_name = device.name.clone();
440        let rate_watch = self.backend.watch_device_sample_rate(
441            &device,
442            Box::new(move |rate| {
443                log::info!("device sample rate changed externally: {rate}Hz on '{watch_name}'");
444                watch_state.set_output_sample_rate(rate.round() as u32);
445            }),
446        );
447
448        let engine = self.backend.create_engine(
449            &device,
450            source_rate,
451            info.channels as u32,
452            consumer,
453            self.timeline.samples_played_counter(),
454        )?;
455        // iOS answers 0 for a rate that belongs to the app's session, and no
456        // switch it can make is one to wait out. A new engine inside the
457        // silence, a seek, still has the rest of the relock ahead of it.
458        let now = std::time::Instant::now();
459        let lead_in = if device_rate > 0.0 && (settled - device_rate).abs() > 0.1 {
460            std::time::Duration::from_millis(
461                crate::config::Config::cached()
462                    .playback
463                    .rate_switch_lead_in_ms as u64,
464            )
465        } else {
466            self.lead_in_ends
467                .map(|end| end.saturating_duration_since(now))
468                .unwrap_or_default()
469        };
470        self.lead_in_ends = None;
471        if !lead_in.is_zero() {
472            engine.lead_in((settled * lead_in.as_secs_f64()) as u64);
473            self.lead_in_ends = Some(now + lead_in);
474        }
475
476        Ok((engine, rate_watch))
477    }
478
479    /// The DSP profile for `device`, loaded once per config and device.
480    fn dsp_for(&mut self, device: &str) -> Option<Arc<crate::audio::dsp::Setup>> {
481        let config = crate::config::Config::cached();
482        #[cfg(test)]
483        if let Some(setup) = if self
484            .renderer
485            .as_ref()
486            .is_some_and(|l| l.device_name() == device)
487        {
488            self.renderer_dsp_override.clone()
489        } else {
490            self.dsp_override.clone()
491        } {
492            self.dsp = Some(DspCache {
493                config,
494                device: device.to_string(),
495                setup: Some(setup.clone()),
496            });
497            return Some(setup);
498        }
499        if let Some(cache) = &self.dsp
500            && Arc::ptr_eq(&cache.config, &config)
501            && cache.device == device
502        {
503            return cache.setup.clone();
504        }
505        let setup = config.dsp.profile_for(device).and_then(|profile| {
506            crate::audio::dsp::Setup::load(profile, &crate::config::config_dir())
507                .inspect_err(|e| {
508                    log::error!(
509                        "dsp: profile '{}' not loaded, playing without it: {e}",
510                        profile.name
511                    )
512                })
513                .ok()
514                .flatten()
515                .map(Arc::new)
516        });
517        self.dsp = Some(DspCache {
518            config,
519            device: device.to_string(),
520            setup: setup.clone(),
521        });
522        setup
523    }
524
525    /// Load the profiles again, and restart where playback is if what the
526    /// output in use plays through has changed. The setups are compared as
527    /// loaded, responses included, since the files can change without the
528    /// config doing so. A preset given to another device, from the Play on
529    /// menu, changes nothing here and restarts nothing.
530    fn reload_dsp(&mut self) {
531        let device = match &self.renderer {
532            Some(link) => Ok(link.device_name().to_string()),
533            None => self.resolve_device().map(|d| d.name),
534        };
535        let Ok(device) = device else {
536            self.dsp = None;
537            return;
538        };
539        let was = self
540            .dsp
541            .take()
542            .filter(|c| c.device == device)
543            .map(|c| c.setup);
544        let now = self.dsp_for(&device);
545        let changed = match was {
546            Some(was) => was != now,
547            None => now.is_some(),
548        };
549        if changed {
550            self.restart_on_current_track();
551        }
552    }
553
554    /// ReplayGain and DSP for a session on the output device.
555    fn processing(&mut self) -> buffer::Processing {
556        let cfg = crate::config::Config::cached();
557        let dsp = match self.resolve_device() {
558            Ok(device) => self.dsp_for(&device.name),
559            Err(_) => None,
560        };
561        buffer::Processing {
562            rg_mode: cfg.playback.replaygain,
563            pre_amp_db: cfg.playback.pre_amp_db,
564            dsp,
565        }
566    }
567
568    /// Resolve the output device: use configured device name if set,
569    /// falling back to system default if not set or if the named device is unavailable.
570    fn resolve_device(&self) -> Result<backend::DeviceInfo, PlayerError> {
571        if let Some(ref name) = self.output_device_name {
572            match self.backend.list_devices() {
573                Ok(devices) => {
574                    if let Some(dev) = devices.into_iter().find(|d| d.name == *name) {
575                        return Ok(dev);
576                    }
577                    log::warn!(
578                        "configured output device '{}' not found, falling back to default",
579                        name,
580                    );
581                }
582                Err(e) => {
583                    log::warn!("failed to list devices while resolving '{}': {}", name, e);
584                }
585            }
586        }
587        Ok(self.backend.default_device()?)
588    }
589
590    /// Switch the output device. Persists to config and restarts the engine
591    /// on the current track if playing.
592    pub fn set_output_device(&mut self, name: String) {
593        log::info!("switching output device to: {}", name);
594        self.output_device_name = Some(name.clone());
595
596        if let Err(e) = crate::config::Config::persist(|cfg| {
597            cfg.playback.output_device = Some(name);
598            cfg.playback.renderer = None;
599            cfg.playback.renderer_name = None;
600        }) {
601            log::error!("failed to save output device config: {}", e);
602        }
603
604        if self.renderer.is_some() {
605            self.use_renderer(None);
606        } else {
607            self.restart_on_current_track();
608        }
609    }
610
611    /// Clear the configured output device, reverting to system default.
612    pub fn clear_output_device(&mut self) {
613        log::info!("reverting to system default output device");
614        self.output_device_name = None;
615
616        if let Err(e) = crate::config::Config::persist(|cfg| {
617            cfg.playback.output_device = None;
618            cfg.playback.renderer = None;
619            cfg.playback.renderer_name = None;
620        }) {
621            log::error!("failed to save output device config: {}", e);
622        }
623
624        if self.renderer.is_some() {
625            self.use_renderer(None);
626        } else {
627            self.restart_on_current_track();
628        }
629    }
630
631    /// If a track is currently playing or paused, restart playback at the
632    /// current position (e.g. after switching output devices). Preserves pause state.
633    fn restart_on_current_track(&mut self) {
634        let position_ms = self.shared_state.position_ms();
635        if let Err(e) = self.restart_current(position_ms) {
636            log::error!("failed to restart playback on device switch: {}", e);
637        }
638    }
639
640    /// Get the current output device name (if configured).
641    pub fn output_device_name(&self) -> Option<&str> {
642        self.output_device_name.as_deref()
643    }
644
645    /// Get a command sender for the UI layer.
646    pub fn command_sender(&self) -> crossbeam_channel::Sender<PlayerCommand> {
647        self.commands.tx.clone()
648    }
649
650    fn session(&self) -> Option<&Session> {
651        match &self.transport {
652            Transport::Loaded(session) => Some(session),
653            _ => None,
654        }
655    }
656
657    fn waiting(&self) -> Option<Waiting> {
658        match self.transport {
659            Transport::Waiting(waiting) => Some(waiting),
660            _ => None,
661        }
662    }
663
664    /// Whether the track waited for may open from its download as it lands:
665    /// at the start, and not to a renderer, which is handed whole files only.
666    fn may_stream(&self, waiting: &Waiting, id: QueueItemId) -> bool {
667        waiting.may_stream(id) && self.renderer.is_none()
668    }
669
670    fn forget_waiting(&mut self) {
671        if matches!(self.transport, Transport::Waiting(_)) {
672            self.transport = Transport::Idle;
673        }
674    }
675
676    /// Tell the front ends what the player is doing, after each command and
677    /// each tick. Derived wholly from the transport: nothing else writes the
678    /// playback state, the wait or the track.
679    fn publish(&self) {
680        let state = &self.shared_state;
681        let (playback, waiting, track) = match &self.transport {
682            Transport::Idle => (PlaybackState::Stopped, false, None),
683            Transport::Waiting(waiting) => (
684                match waiting.start {
685                    Run::Playing => PlaybackState::Stopped,
686                    Run::Paused => PlaybackState::Paused,
687                },
688                true,
689                None,
690            ),
691            Transport::Loaded(session) => (
692                match session.run {
693                    Run::Playing => PlaybackState::Playing,
694                    Run::Paused => PlaybackState::Paused,
695                },
696                false,
697                Some(&session.track),
698            ),
699        };
700        if state.track_info().as_ref() != track {
701            state.set_track_info(track.cloned());
702        }
703        state.set_transport(playback, waiting);
704        // Only a session decoded here can have been processed: a renderer
705        // handed the original file plays it as it is.
706        let dsp = match &self.transport {
707            Transport::Loaded(Session {
708                output: Output::Local(local),
709                ..
710            }) => local.dsp.clone(),
711            Transport::Loaded(Session {
712                output: Output::Renderer(play),
713                ..
714            }) => play.dsp(),
715            _ => None,
716        };
717        if state.dsp() != dsp {
718            state.set_dsp(dsp);
719        }
720        state.set_play_mode(self.mode);
721    }
722
723    /// What the listener last asked for: to hear something, to have it
724    /// paused, or nothing at all. An implicit move — the playing track
725    /// removed, a download failed, an undo, the end of the queue's decoding —
726    /// carries it over to the track it moves to.
727    fn intent(&self) -> Option<Run> {
728        match &self.transport {
729            Transport::Idle => None,
730            Transport::Waiting(waiting) => Some(waiting.start),
731            Transport::Loaded(session) => Some(session.run),
732        }
733    }
734
735    /// Move to `to` as the listener's last request would have it: playing
736    /// on, paused, or, with nothing asked for, the cursor alone.
737    fn carry_on(&mut self, to: Option<QueueItemId>, intent: Option<Run>) {
738        match (to, intent) {
739            (Some(id), Some(start)) => self.cue(id, 0, start),
740            (Some(id), None) => {
741                self.stop_playback_and_clear_state();
742                self.shared_state.set_cursor(Some(id));
743            }
744            (None, _) => self.stop_playback_and_clear_state(),
745        }
746    }
747
748    /// Play a specific item in the playlist by ID.
749    /// Sets cursor, starts playback if Ready or streaming-ready, otherwise waits for TrackReady.
750    pub fn play(&mut self, id: QueueItemId) {
751        // Asked for by name, so a new play even of the item playing: from the
752        // top, not a seek within what was heard.
753        if self.in_flight.as_ref().is_some_and(|f| f.item == id) {
754            self.finish_play();
755        }
756        self.forget_waiting();
757        self.shared_state.set_cursor(Some(id));
758
759        match self.shared_state.item_playback_source(id) {
760            Some(PlaybackSource::Ready(path)) => {
761                if let Err(e) = self.start_playback(id, &path, 0, Run::Playing) {
762                    log::error!("play failed: {}", e);
763                }
764            }
765            Some(PlaybackSource::Streaming {
766                path,
767                bytes_written,
768                total,
769            }) => {
770                // Stop what is playing and park here. The probe answers on its
771                // own thread; if it cannot, TrackReady starts the track once
772                // the whole file has landed. A renderer takes whole files
773                // only, so for one there is no probe: TrackReady it is.
774                self.park(id, 0, Run::Playing);
775                if self.renderer.is_none() {
776                    self.probe_stream_for_playback(id, &path, bytes_written, total);
777                }
778            }
779            None => {
780                self.park(id, 0, Run::Playing);
781                log::info!("play: item {:?} not ready, waiting for TrackReady", id);
782            }
783        }
784    }
785
786    /// Load a track at `position_ms`, playing or paused — where a restored
787    /// session or a hand-off picks up.
788    ///
789    /// A track still downloading waits until it can open — see `Waiting`.
790    fn cue(&mut self, id: QueueItemId, position_ms: u64, start: Run) {
791        self.forget_waiting();
792        self.shared_state.set_cursor(Some(id));
793        let Some(PlaybackSource::Ready(path)) = self.shared_state.item_playback_source(id) else {
794            if position_ms == 0 && start == Run::Playing {
795                self.play(id);
796                return;
797            }
798            self.park(id, position_ms, start);
799            log::info!("cue: {id:?} not on disk yet, opening at {position_ms}ms once it is");
800            return;
801        };
802        if let Err(e) = self.start_playback(id, &path, position_ms, start) {
803            log::error!("cue failed: {}", e);
804            return;
805        }
806        self.report(match start {
807            Run::Playing => PlaybackReportState::Playing,
808            Run::Paused => PlaybackReportState::Paused,
809        });
810    }
811
812    /// Stop what is playing and wait for `id` to become playable. A wait that
813    /// is to open paused already reads as paused, so a client offers to play
814    /// rather than showing a track on its way.
815    fn park(&mut self, id: QueueItemId, position_ms: u64, start: Run) {
816        self.report(PlaybackReportState::Stopped);
817        self.finish_play();
818        self.stop_engine();
819        self.timeline.reset();
820        self.shared_state.set_position_ms(position_ms);
821        self.transport = Transport::Waiting(Waiting {
822            id,
823            position_ms,
824            start,
825        });
826    }
827
828    fn start_playback(
829        &mut self,
830        id: QueueItemId,
831        path: &Path,
832        seek_ms: u64,
833        start: Run,
834    ) -> Result<(), PlayerError> {
835        self.open_session(id, Source::File(path.to_path_buf()), None, seek_ms, start)
836    }
837
838    /// Open a session on `id` at `seek_ms`, playing or paused.
839    ///
840    /// `info` is what is already known of the stream: from the off-thread
841    /// probe for a download, or from what is playing for a restart. A file
842    /// opened without it is probed here. A failure leaves the player cleanly
843    /// stopped; displaying a track that no engine is playing freezes the
844    /// position and makes the transport lie.
845    fn open_session(
846        &mut self,
847        id: QueueItemId,
848        source: Source,
849        info: Option<buffer::StreamInfo>,
850        seek_ms: u64,
851        start: Run,
852    ) -> Result<(), PlayerError> {
853        #[cfg(test)]
854        {
855            self.playback_starts += 1;
856        }
857        self.forget_waiting();
858        let mut result = self.try_open_session(id, source, info, seek_ms, start);
859        // A track the renderer cannot play is marked so in the queue and
860        // walked past, opening the next the way this one would have opened.
861        // A loop rather than a call back through `play`: a long run of them,
862        // a DSD library on a renderer without DSD, would grow the stack per
863        // track.
864        let mut skipped = id;
865        while matches!(result, Err(PlayerError::Unplayable(_)))
866            && self.shared_state.is_cursor(skipped)
867        {
868            let Some(next) = self.shared_state.advance_cursor_loadable() else {
869                log::info!("upnp: nothing further the renderer can play");
870                self.stop_playback_and_clear_state();
871                return Ok(());
872            };
873            match self.shared_state.item_playback_source(next) {
874                Some(PlaybackSource::Ready(path)) => {
875                    skipped = next;
876                    result = self.try_open_session(next, Source::File(path), None, 0, start);
877                }
878                // Not on disk yet: waited for, like any track not yet loaded.
879                _ => {
880                    self.cue(next, 0, start);
881                    return Ok(());
882                }
883            }
884        }
885        if result.is_err() {
886            self.stop_playback_and_clear_state();
887        }
888        self.wake_analyzer();
889        result
890    }
891
892    fn try_open_session(
893        &mut self,
894        id: QueueItemId,
895        source: Source,
896        info: Option<buffer::StreamInfo>,
897        seek_ms: u64,
898        start: Run,
899    ) -> Result<(), PlayerError> {
900        if self.renderer.is_some() {
901            return self.try_open_on_renderer(id, source, info, seek_ms, start);
902        }
903        self.stop_engine();
904        let info = match info {
905            Some(info) => info,
906            None => buffer::probe_file(source.path())?,
907        };
908        let path = source.path().to_path_buf();
909        let streaming = matches!(source, Source::Stream(_));
910
911        let (first, stream) = match source {
912            Source::File(path) => (buffer::SourceEntry::from_file(id, path), None),
913            Source::Stream(source) => {
914                // Held so a seek can reopen the same way without probing again.
915                self.stream_mode = source.mode;
916                let live = LiveStream {
917                    feed: source.bytes_written.clone(),
918                    abandoned: Default::default(),
919                };
920                let status = {
921                    let downloading = self.stream_status_fn(id);
922                    let abandoned = live.abandoned.clone();
923                    Arc::new(move || {
924                        if abandoned.load(std::sync::atomic::Ordering::Acquire) {
925                            streaming::StreamStatus::Failed
926                        } else {
927                            downloading()
928                        }
929                    }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
930                };
931                let entry = buffer::SourceEntry {
932                    id,
933                    path: source.path.clone(),
934                    hint: hint_for(&source.path),
935                    make_mss: Box::new(move || {
936                        let partial = streaming::PartialFileSource::open(
937                            &source.path,
938                            source.bytes_written.clone(),
939                            source.total,
940                            status.clone(),
941                            source.mode,
942                        )?;
943                        Ok(symphonia::core::io::MediaSourceStream::new(
944                            Box::new(partial),
945                            Default::default(),
946                        ))
947                    }),
948                };
949                (entry, Some(live))
950            }
951        };
952
953        let track = TrackInfo {
954            id,
955            path: path.clone(),
956            codec: info.codec.clone(),
957            sample_rate: info.sample_rate,
958            bit_depth: info.bit_depth,
959            bitrate_kbps: info.bitrate_kbps,
960            channels: info.channels,
961            duration_ms: info.duration_ms,
962        };
963        // For seeks, this keeps the bar at the target position instead of
964        // flashing to 0 while the new timeline spins up.
965        self.shared_state.set_position_ms(seek_ms);
966        self.on_track_changed(id, seek_ms);
967        log::info!(
968            "{}: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
969            if streaming { "streaming" } else { "playing" },
970            path.display(),
971            id,
972            info.codec,
973            info.sample_rate,
974            info.channels,
975            info.duration_ms,
976            if seek_ms > 0 {
977                format!(" @{}ms", seek_ms)
978            } else {
979                String::new()
980            }
981        );
982
983        let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
984        self.timeline.reset();
985        let lookahead = Arc::new(parking_lot::Mutex::new(Vec::new()));
986        let next_track = self.decode_cursor(id, lookahead.clone());
987        // Files and downloads in progress open here alike, so both get the
988        // output's processing.
989        let processing = self.processing();
990        self.session += 1;
991        let session = self.session;
992        let finish_tx = self.commands.tx.clone();
993        let decode_handle = buffer::start_decode(
994            first,
995            producer,
996            seek_ms,
997            move || {
998                let (next_id, next_path) = next_track()?;
999                Some(buffer::SourceEntry::from_file(next_id, next_path))
1000            },
1001            self.timeline.clone(),
1002            Some(self.viz_buffer.clone()),
1003            processing,
1004            move || {
1005                finish_tx.send(PlayerCommand::DecodeFinished(session)).ok();
1006            },
1007        )?;
1008
1009        let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
1010        // The setup `create_engine_for` chose the rate by.
1011        let dsp = self
1012            .dsp
1013            .as_ref()
1014            .and_then(|c| c.setup.as_ref())
1015            .map(|s| s.status(info.sample_rate));
1016        // A session opened paused leaves the unit stopped: starting it and
1017        // stopping it again lets a moment of the track out.
1018        if start == Run::Playing {
1019            engine.start()?;
1020        }
1021        self.transport = Transport::Loaded(Session {
1022            track,
1023            run: start,
1024            lookahead,
1025            output: Output::Local(Local {
1026                engine,
1027                decode_handle,
1028                stream,
1029                _rate_watch: rate_watch,
1030                dsp,
1031            }),
1032        });
1033        Ok(())
1034    }
1035
1036    /// Gapless lookahead: the decode thread keeps its own cursor, separate
1037    /// from the UI cursor, so it can look ahead through the playlist without
1038    /// moving what the UI shows as now playing.
1039    ///
1040    /// Each step is kept in `steps`, so a queue edit is checked against what
1041    /// the decoder decided rather than against what it would decide now —
1042    /// which differs once a download has landed or a file failed to open. The
1043    /// lock is held across the read and the record: an edit either lands
1044    /// before the read or finds the step recorded.
1045    fn decode_cursor(
1046        &self,
1047        id: QueueItemId,
1048        steps: Arc<parking_lot::Mutex<Vec<state::Lookahead>>>,
1049    ) -> impl Fn() -> Option<(QueueItemId, PathBuf)> + Send + 'static {
1050        let state = self.shared_state.clone();
1051        let timeline = self.timeline.clone();
1052        move || {
1053            let mut steps = steps.lock();
1054            let current = match steps.last() {
1055                None => id,
1056                Some(step) => step.chosen.as_ref()?.0,
1057            };
1058            let mut step = state.lookahead_after(current)?;
1059            step.boundary = timeline.boundary_count();
1060            let next = step.chosen.clone();
1061            steps.push(step);
1062            next
1063        }
1064    }
1065
1066    /// Probe a partially-downloaded file on its own thread, and start it when
1067    /// the answer comes back.
1068    ///
1069    /// Nothing here waits. Probing reads as much of the container as it takes
1070    /// to describe itself — for Ogg, its last page, which means the whole
1071    /// remaining download — and this is the thread that answers play, pause and
1072    /// seek. So the probe goes elsewhere and its result returns as a command.
1073    ///
1074    /// A format that describes itself up front (FLAC, MP3) comes back in
1075    /// milliseconds and starts early, which is the point of streaming. One that
1076    /// does not comes back whenever it comes back, by which time the download
1077    /// has usually landed and `TrackReady` has started the track from disk —
1078    /// and the late answer is simply dropped. Either way the player kept
1079    /// answering commands throughout.
1080    fn probe_stream_for_playback(
1081        &self,
1082        id: QueueItemId,
1083        path: &Path,
1084        bytes_written: Arc<crate::remote::downloads::ByteFeed>,
1085        total: u64,
1086    ) {
1087        let path = path.to_path_buf();
1088        let tx = self.commands.tx.clone();
1089        let hint = hint_for(&path);
1090
1091        // Abandon the moment the track stops being the one wanted. A probe of
1092        // a container that needs its tail otherwise reads to the end of a
1093        // download nobody is waiting for any more, and skipping through a
1094        // queue that is still caching would leave one doing so per skip.
1095        let status = {
1096            let downloading = self.stream_status_fn(id);
1097            let state = self.shared_state.clone();
1098            Arc::new(move || {
1099                if state.is_cursor(id) {
1100                    downloading()
1101                } else {
1102                    streaming::StreamStatus::Failed
1103                }
1104            }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
1105        };
1106
1107        let spawned = thread::Builder::new()
1108            .name("koan-stream-probe".into())
1109            .spawn(move || {
1110                // `wait` says whether a read may sit at the write head for
1111                // more of the download. The first attempt must not: a
1112                // container that goes looking for its tail would wait for the
1113                // whole transfer, and failing at once is how that is detected.
1114                // The second has no length to go looking with, so whatever it
1115                // still wants is in front of it and worth waiting for.
1116                let attempt = |mode, wait: bool| {
1117                    let open = if wait {
1118                        streaming::PartialFileSource::open(
1119                            &path,
1120                            bytes_written.clone(),
1121                            total,
1122                            status.clone(),
1123                            mode,
1124                        )
1125                    } else {
1126                        streaming::PartialFileSource::open_for_probe(
1127                            &path,
1128                            bytes_written.clone(),
1129                            total,
1130                            status.clone(),
1131                            mode,
1132                        )
1133                    };
1134                    open.map_err(buffer::DecodeError::Io).and_then(|source| {
1135                        let mss = symphonia::core::io::MediaSourceStream::new(
1136                            Box::new(source),
1137                            Default::default(),
1138                        );
1139                        buffer::probe_source(mss, &hint)
1140                    })
1141                };
1142
1143                // Ask for the whole description first. Neither attempt waits at
1144                // the write head, so a container that needs bytes which have
1145                // not arrived fails here rather than reading the transfer out.
1146                let info = match attempt(streaming::ProbeMode::Full, false) {
1147                    Ok(info) => Some((info, streaming::ProbeMode::Full)),
1148                    Err(e) => {
1149                        // Try again claiming no length. Ogg goes looking for its
1150                        // last page only when told there is one to find; without
1151                        // it the track opens now and plays, at the price of
1152                        // seeking and of the duration that page carries. Both
1153                        // come back when the download lands.
1154                        log::info!(
1155                            "stream probe: {} needs more than has arrived ({}), opening without a length",
1156                            path.display(),
1157                            e
1158                        );
1159                        let lengthless = lengthless_mode_for(&path);
1160                        attempt(lengthless, true)
1161                            .ok()
1162                            .map(|info| (info, lengthless))
1163                    }
1164                };
1165
1166                match info {
1167                    Some((info, mode)) => {
1168                        tx.send(PlayerCommand::StreamProbed {
1169                            id,
1170                            info: Box::new(info),
1171                            mode,
1172                        })
1173                        .ok();
1174                    }
1175                    // Not a failure of the track: it plays from disk once the
1176                    // download lands, and the cursor is still parked on it.
1177                    None => log::info!(
1178                        "stream probe: {} cannot start early, waiting for the download",
1179                        path.display()
1180                    ),
1181                }
1182            });
1183
1184        if let Err(e) = spawned {
1185            log::warn!("stream probe: could not spawn for {:?}: {}", id, e);
1186        }
1187    }
1188
1189    /// A probe finished. Start the track if it is still the one wanted and
1190    /// nothing has started it in the meantime.
1191    fn stream_probed(
1192        &mut self,
1193        id: QueueItemId,
1194        info: buffer::StreamInfo,
1195        mode: streaming::ProbeMode,
1196    ) {
1197        // Moved on, or already open — the download landed first, or the user
1198        // asked for something else.
1199        let Some(waiting) = self.waiting().filter(|w| self.may_stream(w, id)) else {
1200            return;
1201        };
1202        if !self.shared_state.is_cursor(id) {
1203            return;
1204        }
1205
1206        match self.shared_state.item_playback_source(id) {
1207            // The download landed while probing: play it as an ordinary file.
1208            Some(PlaybackSource::Ready(path)) => {
1209                if let Err(e) = self.start_playback(id, &path, 0, waiting.start) {
1210                    log::error!("stream probe: playback failed: {}", e);
1211                }
1212            }
1213            Some(PlaybackSource::Streaming {
1214                path,
1215                bytes_written,
1216                total,
1217            }) => {
1218                let source = StreamSource {
1219                    path,
1220                    bytes_written,
1221                    total,
1222                    mode,
1223                };
1224                if let Err(e) =
1225                    self.open_session(id, Source::Stream(source), Some(info), 0, waiting.start)
1226                {
1227                    log::error!("stream probe: streaming playback failed: {}", e);
1228                }
1229            }
1230            None => {}
1231        }
1232    }
1233
1234    /// What the streaming source asks per read to know whether the download is
1235    /// still going. Asked each time rather than passed once: a transfer can
1236    /// land, or die, at any point during playback.
1237    fn stream_status_fn(
1238        &self,
1239        id: QueueItemId,
1240    ) -> Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync> {
1241        // The item's own state, which a transfer's end writes before it wakes
1242        // anyone: asked twice per read of a file still arriving, so nothing
1243        // else is looked up.
1244        let state = self.shared_state.clone();
1245        Arc::new(move || match state.item_state(id) {
1246            Some(ItemState::Ready) => streaming::StreamStatus::Complete,
1247            Some(ItemState::Failed(_)) => streaming::StreamStatus::Failed,
1248            _ => streaming::StreamStatus::Downloading,
1249        })
1250    }
1251
1252    /// Seek within the current track, preserving pause state.
1253    ///
1254    /// A track still downloading is seekable only as far as its bytes reach, so
1255    /// the target is clamped to `seekable_ms` and the restart goes back through
1256    /// the streaming path — reopening a partial file as a plain file would
1257    /// decode whatever happens to be on disk and end the track early.
1258    pub fn seek(&mut self, position_ms: u64) {
1259        let Some(id) = self.session().map(|s| s.track.id) else {
1260            return;
1261        };
1262        // Stop just short of the end rather than falling into the next track.
1263        let seekable = self.shared_state.seekable_ms();
1264        if seekable == 0 {
1265            // Nothing of this track can be reached yet — a partial container
1266            // that has not said what it is. Restarting it at zero is not what
1267            // anyone asked for, so the seek is simply declined.
1268            log::debug!("seek declined: {:?} is not seekable yet", id);
1269            return;
1270        }
1271        let ceiling = seekable.min(
1272            self.shared_state
1273                .duration_ms()
1274                .saturating_sub(SEEK_END_GUARD_MS),
1275        );
1276        let clamped = position_ms.min(ceiling);
1277
1278        // A renderer seeks the track it holds; reloading it would let the
1279        // top of the track be heard.
1280        if self.seek_on_renderer(clamped) {
1281            return;
1282        }
1283        if let Err(e) = self.restart_current(clamped) {
1284            log::error!("seek failed: {}", e);
1285        }
1286    }
1287
1288    /// Restart what is playing at `position_ms`, preserving pause state.
1289    ///
1290    /// What a seek does, and what switching output device does, and what going
1291    /// back from the first track does. All three restart the same track, so all
1292    /// three resolve the source the same way: from the queue item, never from
1293    /// the session's path, which names the `.part` file for a track that was still
1294    /// downloading when it started and is not renamed when the download lands.
1295    fn restart_current(&mut self, position_ms: u64) -> Result<(), PlayerError> {
1296        let Some(session) = self.session() else {
1297            return Ok(());
1298        };
1299        let (info, start) = (session.track.clone(), session.run);
1300
1301        match self.shared_state.item_playback_source(info.id) {
1302            // A renderer takes whole files only: wait for this one to land,
1303            // keeping its place and whether it was playing.
1304            Some(PlaybackSource::Streaming { .. }) if self.renderer.is_some() => {
1305                self.park(info.id, position_ms, start);
1306                return Ok(());
1307            }
1308            Some(PlaybackSource::Streaming {
1309                path,
1310                bytes_written,
1311                total,
1312            }) => {
1313                // No probe: what is playing already said what this is, and
1314                // reading an Ogg's last page to learn it again would mean
1315                // waiting for the rest of the download.
1316                let known = buffer::StreamInfo {
1317                    codec: info.codec.clone(),
1318                    sample_rate: info.sample_rate,
1319                    channels: info.channels,
1320                    bit_depth: info.bit_depth,
1321                    bitrate_kbps: info.bitrate_kbps,
1322                    duration_ms: info.duration_ms,
1323                };
1324                let source = StreamSource {
1325                    path,
1326                    bytes_written,
1327                    total,
1328                    mode: self.stream_mode,
1329                };
1330                self.open_session(
1331                    info.id,
1332                    Source::Stream(source),
1333                    Some(known),
1334                    position_ms,
1335                    start,
1336                )?;
1337            }
1338            Some(PlaybackSource::Ready(path)) => {
1339                self.start_playback(info.id, &path, position_ms, start)?;
1340            }
1341            None => return Ok(()),
1342        }
1343
1344        self.report(match start {
1345            Run::Playing => PlaybackReportState::Playing,
1346            Run::Paused => PlaybackReportState::Paused,
1347        });
1348        Ok(())
1349    }
1350
1351    /// Skip to next track in playlist.
1352    pub fn next_track(&mut self) {
1353        match self.shared_state.advance_cursor_loadable() {
1354            Some(id) => self.play(id),
1355            None => {
1356                log::info!("no more tracks in playlist");
1357                self.stop_playback_and_clear_state();
1358            }
1359        }
1360    }
1361
1362    /// Go back to previous track.
1363    pub fn prev_track(&mut self) {
1364        match self.shared_state.retreat_cursor() {
1365            Some((id, _)) => self.play(id),
1366            None => {
1367                // No previous track — restart current from the beginning.
1368                if let Err(e) = self.restart_current(0) {
1369                    log::error!("restart failed: {}", e);
1370                }
1371            }
1372        }
1373    }
1374
1375    /// Pause playback, fading out if the config asks for it.
1376    ///
1377    /// A fade leaves the unit running until it reaches silence;
1378    /// `update_playback_state` stops it from there. A track still on its way
1379    /// opens paused when it arrives.
1380    pub fn pause(&mut self) {
1381        let fade = crate::config::Config::cached().playback.fade_on_pause;
1382        match &mut self.transport {
1383            Transport::Idle => return,
1384            Transport::Waiting(waiting) => {
1385                waiting.start = Run::Paused;
1386                return;
1387            }
1388            Transport::Loaded(session) => {
1389                if let Output::Local(local) = &session.output {
1390                    if fade {
1391                        local.engine.fade_out();
1392                    } else if let Err(e) = local.engine.stop() {
1393                        log::error!("pause failed: {}", e);
1394                        return;
1395                    }
1396                }
1397                session.run = Run::Paused;
1398            }
1399        }
1400        self.pause_renderer();
1401        self.report(PlaybackReportState::Paused);
1402    }
1403
1404    /// Resume playback. Fades back in if the pause faded out.
1405    ///
1406    /// With nothing loaded — a session restored stopped, or a start that
1407    /// failed — there is nothing to resume, and play starts the track under
1408    /// the cursor instead of doing nothing. A track still on its way keeps the
1409    /// position it was waiting to open at.
1410    pub fn resume(&mut self) {
1411        self.answer_silence();
1412        let session = match &mut self.transport {
1413            Transport::Idle => {
1414                if let Some(id) = self.shared_state.cursor() {
1415                    self.play(id);
1416                }
1417                return;
1418            }
1419            Transport::Waiting(waiting) => {
1420                let waiting = *waiting;
1421                self.cue(waiting.id, waiting.position_ms, Run::Playing);
1422                return;
1423            }
1424            Transport::Loaded(session) => session,
1425        };
1426        let Output::Local(local) = &session.output else {
1427            self.resume_renderer();
1428            return;
1429        };
1430        let engine = &local.engine;
1431        let resumed = if engine.is_running() || engine.is_silent() {
1432            self.lead_in_ends = None;
1433            engine.fade_in()
1434        } else {
1435            engine.start()
1436        };
1437        if let Err(e) = resumed {
1438            log::error!("resume failed: {}", e);
1439            return;
1440        }
1441        session.run = Run::Playing;
1442        self.wake_analyzer();
1443        self.report(PlaybackReportState::Playing);
1444    }
1445
1446    /// Tell the analyser there is about to be something to hear.
1447    ///
1448    /// It parks when nothing is playing and nothing is reading, and the one
1449    /// thing it cannot be signalled from is the play head — that counter is
1450    /// written by the audio render callback, which may never take a lock. So
1451    /// the player says so instead, on the two edges where silence ends.
1452    fn wake_analyzer(&self) {
1453        self.viz_snapshot.wake();
1454    }
1455
1456    /// Stop playback and clear playlist.
1457    pub fn stop(&mut self) {
1458        self.shared_state.clear_playlist();
1459        self.stop_playback_and_clear_state();
1460    }
1461
1462    /// Stop the audio engine and decode thread without touching shared state.
1463    ///
1464    /// Output stops first, then the decode thread is joined, then the engine
1465    /// drops: tearing CoreAudio down under a live producer is the end-of-queue
1466    /// crash (#89).
1467    fn stop_engine(&mut self) {
1468        let playback = match std::mem::replace(&mut self.transport, Transport::Idle) {
1469            Transport::Loaded(playback) => playback,
1470            other => {
1471                self.transport = other;
1472                return;
1473            }
1474        };
1475        self.session += 1;
1476        self.bank_listening();
1477        let local = match playback.output {
1478            Output::Local(local) => local,
1479            Output::Renderer(play) => {
1480                self.halt_renderer(*play);
1481                self.answer_silence();
1482                return;
1483            }
1484        };
1485        let Local {
1486            engine,
1487            mut decode_handle,
1488            stream,
1489            ..
1490        } = local;
1491
1492        let _ = engine.stop();
1493        // Stop first, so the failed read the abandon causes reads as a stop
1494        // rather than a bad source to skip past.
1495        decode_handle.signal_stop();
1496        if let Some(stream) = stream {
1497            stream.abandon();
1498        }
1499        decode_handle.stop();
1500        drop(engine);
1501        self.answer_silence();
1502    }
1503
1504    /// Tell whoever waited for the pause where the playhead came to rest.
1505    fn answer_silence(&mut self) {
1506        let position_ms = self.shared_state.position_ms();
1507        for reply in self.silence_waiters.drain(..) {
1508            let _ = reply.send(position_ms);
1509        }
1510    }
1511
1512    /// Full stop: tear down engine + clear all display state.
1513    fn stop_playback_and_clear_state(&mut self) {
1514        self.forget_waiting();
1515        self.report(PlaybackReportState::Stopped);
1516        self.finish_play();
1517        self.stop_engine();
1518        self.timeline.reset();
1519        self.shared_state.set_position_ms(0);
1520    }
1521
1522    /// Remove a track from the playlist. If it was the cursor, move to the
1523    /// track that followed it, as the listener's last request would have it.
1524    ///
1525    /// `remove_item` clears the cursor, and an unset cursor means "start from the
1526    /// top" — so the successor is pinned down by parking the cursor on the removed
1527    /// track's predecessor first. `None` is correct only when it was the first item.
1528    pub fn remove_from_playlist(&mut self, id: QueueItemId) {
1529        let was_cursor = self.shared_state.is_cursor(id);
1530        let resume_after = was_cursor
1531            .then(|| self.shared_state.item_before(id))
1532            .flatten();
1533        self.shared_state.remove_item(id);
1534        if was_cursor {
1535            self.shared_state.set_cursor(resume_after);
1536            let next = self.shared_state.advance_cursor_loadable();
1537            self.carry_on(next, self.intent());
1538        }
1539    }
1540
1541    /// A download finished. A track waiting on it opens; one already
1542    /// streaming from it, playing or paused, re-reads its metadata from the
1543    /// complete file.
1544    ///
1545    /// The item's state is the downloader's to set, before it sends this.
1546    /// Setting it here as well would turn a duplicate entry's `Failed` back to
1547    /// `Ready`, pointing at a `.part` file that was deleted.
1548    pub fn track_ready(&mut self, id: QueueItemId) {
1549        if !self.shared_state.is_cursor(id) {
1550            return;
1551        }
1552
1553        if let Some(waiting) = self.waiting().filter(|w| w.id == id) {
1554            log::info!("track_ready: opening {:?}", id);
1555            self.cue(id, waiting.position_ms, waiting.start);
1556            return;
1557        }
1558
1559        if self.session().is_some_and(|s| s.track.id == id) {
1560            log::info!(
1561                "track_ready: download complete while streaming {:?}, refreshing metadata",
1562                id
1563            );
1564            self.refresh_track_metadata(id);
1565        }
1566    }
1567
1568    /// Enough of a download has landed to stream it. Opens the track if it is
1569    /// waiting to start from the top.
1570    pub fn track_stream_ready(&mut self, id: QueueItemId) {
1571        let Some(waiting) = self.waiting().filter(|w| self.may_stream(w, id)) else {
1572            return;
1573        };
1574        if !self.shared_state.is_cursor(id) {
1575            return;
1576        }
1577
1578        match self.shared_state.item_playback_source(id) {
1579            Some(PlaybackSource::Streaming {
1580                path,
1581                bytes_written,
1582                total,
1583            }) => {
1584                log::info!("track_stream_ready: probing partial file for {:?}", id);
1585                self.probe_stream_for_playback(id, &path, bytes_written, total);
1586            }
1587            Some(PlaybackSource::Ready(path)) => {
1588                // Download finished between threshold and now — just play normally.
1589                log::info!(
1590                    "track_stream_ready: track already ready, starting normal playback for {:?}",
1591                    id
1592                );
1593                if let Err(e) = self.start_playback(id, &path, 0, waiting.start) {
1594                    log::error!("track_stream_ready playback failed: {}", e);
1595                }
1596            }
1597            None => {} // Not enough data yet — wait.
1598        }
1599    }
1600
1601    /// Re-read full lofty metadata for a track after its download completes.
1602    /// Called from track_ready() when a streaming track finishes downloading.
1603    /// What the item takes from it is `update_item_metadata`'s call.
1604    fn refresh_track_metadata(&mut self, id: QueueItemId) {
1605        use crate::index::metadata;
1606
1607        let path = match self.shared_state.item_path_if_ready(id) {
1608            Some(p) => p,
1609            None => return,
1610        };
1611
1612        match metadata::read_metadata(&path) {
1613            Ok(meta) => {
1614                self.shared_state.update_item_metadata(
1615                    id,
1616                    meta.title,
1617                    meta.artist,
1618                    meta.album_artist.unwrap_or_default(),
1619                    meta.album,
1620                    meta.duration_ms.map(|d| d as u64),
1621                );
1622
1623                // Re-probe the complete file for accurate duration + stream info.
1624                // The initial probe was done on partial streaming data and may have
1625                // underestimated duration, causing premature seek clamping or wrong
1626                // progress bar display.
1627                //
1628                // The path is taken over at the same time. Playback started
1629                // against the `.part` file and the download's last act is to
1630                // rename it, so what `track_info` holds now names nothing.
1631                if let Transport::Loaded(session) = &mut self.transport
1632                    && session.track.id == id
1633                {
1634                    let current = &mut session.track;
1635                    let probed = buffer::probe_file(&path).ok();
1636                    let duration_ms = probed
1637                        .as_ref()
1638                        .map(|s| s.duration_ms)
1639                        .filter(|d| *d > current.duration_ms)
1640                        .unwrap_or(current.duration_ms);
1641                    if duration_ms != current.duration_ms {
1642                        log::info!(
1643                            "track_ready: duration corrected {}ms → {}ms",
1644                            current.duration_ms,
1645                            duration_ms
1646                        );
1647                    }
1648                    current.duration_ms = duration_ms;
1649                    current.path = path.clone();
1650                }
1651
1652                // Signal UI to re-read cover art and update souvlaki media controls.
1653                self.shared_state.signal_metadata_refresh();
1654                log::info!("track_ready: metadata refreshed for {:?}", id);
1655            }
1656            Err(e) => {
1657                log::warn!("track_ready: metadata refresh failed for {:?}: {}", id, e);
1658            }
1659        }
1660    }
1661
1662    /// A session has opened on `id`. A seek restarts playback of the same
1663    /// track, so the play in flight carries on when it is the same item —
1664    /// otherwise scrubbing around a track would enter it into history once
1665    /// per seek. Everything that ends a play and opens the same item again,
1666    /// as a new one, closes the old play first: `play`, a repeat at the end
1667    /// of a session.
1668    fn on_track_changed(&mut self, id: QueueItemId, position_ms: u64) {
1669        if let Some(f) = self.in_flight.as_mut().filter(|f| f.item == id) {
1670            f.jump(position_ms);
1671            // The session is new, so the play is its first boundary.
1672            f.boundary = 0;
1673            return;
1674        }
1675        self.begin_play(id, position_ms, 0);
1676    }
1677
1678    /// The needle has moved to a new play of `id`, at `boundary` of the
1679    /// session. Close out the outgoing play and write the new one to history
1680    /// straight away, so history reads in play order even for a track that
1681    /// is skipped a moment later.
1682    fn begin_play(&mut self, id: QueueItemId, position_ms: u64, boundary: usize) {
1683        self.finish_play();
1684        let track_id = self.shared_state.item_db_id(id);
1685        let mut flight = InFlight::new(id, track_id, position_ms);
1686        flight.boundary = boundary;
1687        self.in_flight = Some(flight);
1688        if let (Some(track_id), Some(recorder)) = (track_id, self.history.as_ref()) {
1689            recorder.record(PlayEvent::Started {
1690                track_id,
1691                position_ms,
1692            });
1693        }
1694    }
1695
1696    /// Count what has played of the track in flight since it was last
1697    /// counted. Before anything that resets the timeline, which is where the
1698    /// playhead is read from, and before a track is closed out.
1699    fn bank_listening(&mut self) {
1700        // A renderer's playhead is its clock: set only while one plays.
1701        let renderer_at = self
1702            .shared_state
1703            .renderer_clock()
1704            .map(|_| self.shared_state.position_ms());
1705        if let Some(f) = self.in_flight.as_mut()
1706            && let Some(at) = self.timeline.position_in(f.boundary).or(renderer_at)
1707        {
1708            f.advance(at);
1709        }
1710    }
1711
1712    /// Tell the remote server where the track playing now stands. A track
1713    /// starting is reported by `on_track_changed`; this covers what happens
1714    /// to it afterwards.
1715    fn report(&self, state: PlaybackReportState) {
1716        let Some(track_id) = self.in_flight.as_ref().and_then(InFlight::track_id) else {
1717            return;
1718        };
1719        if let Some(recorder) = self.history.as_ref() {
1720            recorder.record(PlayEvent::Playback(PlaybackReport {
1721                track_id,
1722                state,
1723                position_ms: self.shared_state.position_ms(),
1724            }));
1725        }
1726    }
1727
1728    /// Tell history how long the current track was heard for. Returns what was
1729    /// reported, which is how the tests see it.
1730    fn finish_play(&mut self) -> Option<PlayEvent> {
1731        self.bank_listening();
1732        let flight = self.in_flight.take()?;
1733        let event = PlayEvent::Finished {
1734            track_id: flight.track_id()?,
1735            listened_ms: flight.listened_ms(),
1736        };
1737        if let Some(recorder) = self.history.as_ref() {
1738            recorder.record(event);
1739        }
1740        Some(event)
1741    }
1742
1743    /// What happens without a command: the silence after a rate switch
1744    /// running out, a pause's fade reaching silence, and the playhead crossing
1745    /// into the next queued track. Called from the command loop on each wake.
1746    pub fn update_playback_state(&mut self) {
1747        if self.session().is_some()
1748            && self
1749                .lead_in_ends
1750                .is_some_and(|end| std::time::Instant::now() >= end)
1751        {
1752            // The playhead held still through the silence while clients
1753            // counted on from where it was published; this is where they hear
1754            // it move.
1755            self.lead_in_ends = None;
1756            self.shared_state.changed();
1757        }
1758
1759        if let Some(session) = self.session()
1760            && session.run == Run::Paused
1761            && let Some(engine) = session.engine()
1762            && engine.is_running()
1763            && engine.is_silent()
1764        {
1765            if let Err(e) = engine.stop() {
1766                log::error!("stopping after fade failed: {}", e);
1767            }
1768            self.answer_silence();
1769        }
1770
1771        self.follow_playhead();
1772        self.renderer_tick();
1773        self.publish();
1774    }
1775
1776    /// A gapless transition moves the playhead into the next track without
1777    /// anything on this thread asking, so the play is banked, the session's
1778    /// track moved on and the cursor brought along from here. A play is its
1779    /// boundary: an item repeated runs into itself, and that is a new play.
1780    fn follow_playhead(&mut self) {
1781        if self.session().is_none() {
1782            return;
1783        }
1784        let Some(playhead) = self.timeline.playhead() else {
1785            return;
1786        };
1787        let id = playhead.id;
1788        if self
1789            .in_flight
1790            .as_ref()
1791            .is_none_or(|f| f.item != id || f.boundary != playhead.boundary)
1792        {
1793            self.begin_play(id, playhead.position_ms, playhead.boundary);
1794        }
1795        let Transport::Loaded(session) = &mut self.transport else {
1796            return;
1797        };
1798        if session.track.id == id {
1799            return;
1800        }
1801        let Some((id, path, info, _)) = self.timeline.current_playback() else {
1802            return;
1803        };
1804        log::info!("timeline: now playing {:?}", id);
1805        session.track = TrackInfo {
1806            id,
1807            path,
1808            codec: info.codec,
1809            sample_rate: info.sample_rate,
1810            bit_depth: info.bit_depth,
1811            bitrate_kbps: info.bitrate_kbps,
1812            channels: info.channels,
1813            duration_ms: info.duration_ms,
1814        };
1815        self.shared_state.set_cursor(Some(id));
1816    }
1817
1818    /// A download a waiting track needs will never land.
1819    ///
1820    /// That wait has no end, so walk on to the next item that can still load,
1821    /// opening it the way the failed one would have opened, or stop cleanly if
1822    /// there is none. Only a waiting track is affected: one already streaming
1823    /// sees the failure in its reads and ends the decode, which advances the
1824    /// queue.
1825    pub fn track_failed(&mut self, id: QueueItemId) {
1826        let Some(waiting) = self.waiting().filter(|w| w.id == id) else {
1827            return;
1828        };
1829        if !self.shared_state.is_cursor(id) {
1830            return;
1831        }
1832        log::info!("track {:?} cannot load, moving on", id);
1833        let next = self.shared_state.advance_cursor_loadable();
1834        self.carry_on(next, Some(waiting.start));
1835    }
1836
1837    /// Decode thread naturally finished (playlist exhausted or error).
1838    /// Advance to the next playable track; otherwise stop cleanly.
1839    ///
1840    /// Ignored from a session already torn down: a play or seek handled after
1841    /// the message was sent has replaced what it describes.
1842    fn on_decode_finished(&mut self, session: u64) {
1843        if session != self.session || self.renderer_stream_finished() {
1844            return;
1845        }
1846        log::info!("decode finished, checking for next track");
1847        self.track_ended();
1848    }
1849
1850    /// The track under the playhead played out, here or on a renderer. Go on
1851    /// as the play mode says: the same item again under repeat one, the next
1852    /// otherwise, from the top again when the queue repeats.
1853    ///
1854    /// A track that has not finished downloading parks the cursor on it, so its
1855    /// `TrackReady`/`TrackStreamReady` resumes the queue instead of being
1856    /// discarded as "not the cursor".
1857    ///
1858    /// Opening the item that just ended again — repeating it, or a queue of
1859    /// one repeating — closes the play that ended first, so the session that
1860    /// opens is a new play rather than a seek within it. A repeat of something
1861    /// that played for no time at all — a file that opens but decodes nothing
1862    /// — would go round forever, so it stops instead.
1863    fn track_ended(&mut self) {
1864        self.bank_listening();
1865        let ended = self.in_flight.as_ref().map(|f| f.item);
1866        let heard = self.in_flight.as_ref().is_some_and(|f| f.listened_ms() > 0);
1867        if self.mode.repeat != Repeat::Off && !heard {
1868            log::info!("track ended with nothing heard; not repeating it");
1869            self.stop_playback_and_clear_state();
1870            return;
1871        }
1872        let again = (self.mode.repeat == Repeat::One)
1873            .then(|| self.shared_state.cursor())
1874            .flatten()
1875            .filter(|id| {
1876                self.shared_state
1877                    .item_state(*id)
1878                    .is_some_and(|s| !matches!(s, ItemState::Failed(_)))
1879            });
1880        let next = again.or_else(|| self.shared_state.advance_cursor_loadable());
1881        if next.is_some() && next == ended {
1882            self.finish_play();
1883        }
1884        self.carry_on(next, self.intent());
1885    }
1886
1887    /// Restart the session at the playhead if the queue no longer follows the
1888    /// playing track with what the decoder has already queued.
1889    ///
1890    /// The decoder looks ahead a few seconds before a track ends, and what it
1891    /// queued is in the ring. Without this, removing or moving the next track
1892    /// in that window has no effect, and a track inserted to play next is
1893    /// skipped. A lookahead not yet taken needs nothing: the decoder reads the
1894    /// queue when it gets there. One that found the end of the queue counts as
1895    /// taken, so a track added then follows gaplessly instead of after a stop.
1896    fn revoke_stale_lookahead(&mut self) {
1897        if self.session().is_none() {
1898            return;
1899        }
1900        self.update_playback_state();
1901        let Some(playhead) = self.timeline.playhead() else {
1902            return;
1903        };
1904        let position_ms = playhead.position_ms;
1905        let Some(session) = self.session().filter(|s| s.track.id == playhead.id) else {
1906            return;
1907        };
1908        let stale = {
1909            let steps = session.lookahead.lock();
1910            // The steps the playhead has yet to reach: those that open a
1911            // boundary past the one it is in. By boundary, not by the item
1912            // chosen — an item repeated is chosen by every step.
1913            steps
1914                .iter()
1915                .filter(|step| step.boundary > playhead.boundary)
1916                .any(|step| !self.shared_state.still_follows(step))
1917        };
1918        if !stale {
1919            return;
1920        }
1921        log::info!("queue changed under the lookahead, restarting at {position_ms}ms");
1922        if let Err(e) = self.restart_current(position_ms) {
1923            log::error!("restart after a queue edit failed: {}", e);
1924        }
1925    }
1926
1927    /// Snapshot items with their predecessors for an undo of "these were removed".
1928    /// In playlist order, so undo re-inserts each item after a predecessor that
1929    /// is already back in place.
1930    fn snapshot_for_undo(
1931        &self,
1932        ids: &[QueueItemId],
1933    ) -> Vec<(Box<state::PlaylistItem>, Option<QueueItemId>)> {
1934        self.shared_state
1935            .items_before(ids)
1936            .into_iter()
1937            .filter_map(|(id, after)| Some((Box::new(self.shared_state.get_item(id)?), after)))
1938            .collect()
1939    }
1940
1941    /// Route an undo entry to the batch buffer (if batching) or the undo stack.
1942    fn push_undo(&mut self, entry: UndoEntry) {
1943        if let Some(ref mut batch) = self.batch_buffer {
1944            batch.push(entry);
1945        } else {
1946            self.undo_stack.push(entry);
1947        }
1948    }
1949
1950    /// Process a single command.
1951    pub fn process_command(&mut self, cmd: PlayerCommand) {
1952        // Someone wants to hear something, or has picked where: the renderer
1953        // remembered from last time is no longer what to go back to. A cue is
1954        // how a session is restored at launch, which is what the renderer is
1955        // to carry on with, so it is not one of them.
1956        if cmd.asks_to_play() && !matches!(cmd, PlayerCommand::Cue { .. })
1957            || matches!(
1958                cmd,
1959                PlayerCommand::UseRenderer(_)
1960                    | PlayerCommand::SetOutputDevice(_)
1961                    | PlayerCommand::ClearOutputDevice
1962            )
1963        {
1964            self.resume_renderer = false;
1965        }
1966        let edits_queue = matches!(
1967            cmd,
1968            PlayerCommand::AddToPlaylist(_)
1969                | PlayerCommand::InsertInPlaylist { .. }
1970                | PlayerCommand::RemoveFromPlaylist(_)
1971                | PlayerCommand::RemoveFromPlaylistBatch(_)
1972                | PlayerCommand::MoveInPlaylist { .. }
1973                | PlayerCommand::MoveItemsInPlaylist { .. }
1974                | PlayerCommand::ReorderPlaylist(_)
1975                | PlayerCommand::Undo
1976                | PlayerCommand::Redo
1977                | PlayerCommand::SetShuffle(_)
1978                | PlayerCommand::SetRepeat(_)
1979                | PlayerCommand::RestorePlayMode(_)
1980        );
1981        self.apply_command(cmd);
1982        if edits_queue {
1983            self.revoke_stale_lookahead();
1984        }
1985        self.publish();
1986    }
1987
1988    fn apply_command(&mut self, cmd: PlayerCommand) {
1989        match cmd {
1990            PlayerCommand::Play(id) => self.play(id),
1991            PlayerCommand::Cue {
1992                id,
1993                position_ms,
1994                play,
1995            } => self.cue(
1996                id,
1997                position_ms,
1998                if play { Run::Playing } else { Run::Paused },
1999            ),
2000            PlayerCommand::Pause => self.pause(),
2001            PlayerCommand::PauseAndReport(reply) => {
2002                self.pause();
2003                self.silence_waiters.push(reply);
2004                if self
2005                    .session()
2006                    .is_none_or(|p| p.engine().is_none_or(|e| !e.is_running()))
2007                {
2008                    self.answer_silence();
2009                }
2010            }
2011            PlayerCommand::Resume => self.resume(),
2012            PlayerCommand::Stop => self.stop(),
2013            PlayerCommand::Seek(pos) => self.seek(pos),
2014            PlayerCommand::NextTrack => self.next_track(),
2015            PlayerCommand::PrevTrack => self.prev_track(),
2016            PlayerCommand::AddToPlaylist(items) => {
2017                let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
2018                let whole = self.shared_state.is_empty();
2019                self.shared_state.add_items(items);
2020                // Into an empty queue, an add is a queue arriving whole.
2021                if whole
2022                    && self.mode.shuffle
2023                    && let Some(&first) = ids.first()
2024                {
2025                    self.shared_state.shuffle_from(first);
2026                }
2027                self.push_undo(UndoEntry::Added { ids });
2028            }
2029            PlayerCommand::UpdatePaths(updates) => {
2030                self.shared_state.update_paths(&updates);
2031                if let Transport::Loaded(session) = &mut self.transport
2032                    && let Some((_, new_path)) =
2033                        updates.iter().find(|(id, _)| *id == session.track.id)
2034                {
2035                    session.track.path = new_path.clone();
2036                }
2037            }
2038            PlayerCommand::InsertInPlaylist { items, after } => {
2039                let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
2040                self.shared_state.insert_items_after(items, after);
2041                self.push_undo(UndoEntry::Inserted { ids });
2042            }
2043            PlayerCommand::ClearPlaylist => self.clear_playlist(),
2044            PlayerCommand::ReplacePlaylist {
2045                items,
2046                start,
2047                position_ms,
2048                play,
2049            } => {
2050                if items.is_empty() {
2051                    self.clear_playlist();
2052                    return;
2053                }
2054                let start_id = items.get(start).unwrap_or(&items[0]).id;
2055                // Stopped before the swap, as `clear_playlist` does, and the
2056                // swap one change: see `SharedPlayerState::replace_playlist`.
2057                self.stop_playback_and_clear_state();
2058                let (old_items, cursor) = self.shared_state.replace_playlist(items);
2059                if self.mode.shuffle {
2060                    self.shared_state.shuffle_from(start_id);
2061                }
2062                self.push_undo(UndoEntry::Replaced {
2063                    items: old_items,
2064                    cursor,
2065                });
2066                self.cue(
2067                    start_id,
2068                    position_ms,
2069                    if play { Run::Playing } else { Run::Paused },
2070                );
2071            }
2072            PlayerCommand::RemoveFromPlaylist(id) => {
2073                let item = self.shared_state.get_item(id);
2074                let after = self.shared_state.item_before(id);
2075                self.remove_from_playlist(id);
2076                if let Some(item) = item {
2077                    self.push_undo(UndoEntry::Removed {
2078                        items: vec![(Box::new(item), after)],
2079                    });
2080                }
2081            }
2082            PlayerCommand::RemoveFromPlaylistBatch(ids) => {
2083                // Snapshot before removing anything, and resolve the resume point
2084                // once: removing one at a time would restart the engine for every
2085                // deleted track that the cursor lands on along the way.
2086                let items_with_pos = self.snapshot_for_undo(&ids);
2087                let intent = self.intent();
2088                let resume_after = match self.shared_state.cursor() {
2089                    Some(cursor) if ids.contains(&cursor) => {
2090                        Some(self.shared_state.surviving_item_before(cursor, &ids))
2091                    }
2092                    _ => None,
2093                };
2094
2095                self.shared_state.remove_items(&ids);
2096
2097                if let Some(resume_after) = resume_after {
2098                    self.shared_state.set_cursor(resume_after);
2099                    let next = self.shared_state.advance_cursor_loadable();
2100                    self.carry_on(next, intent);
2101                }
2102
2103                if !items_with_pos.is_empty() {
2104                    self.push_undo(UndoEntry::Removed {
2105                        items: items_with_pos,
2106                    });
2107                }
2108            }
2109            PlayerCommand::MoveInPlaylist { id, target, after } => {
2110                let was_after = self.shared_state.item_before(id);
2111                self.shared_state.move_item(id, target, after);
2112                self.push_undo(UndoEntry::Moved { id, was_after });
2113            }
2114            PlayerCommand::MoveItemsInPlaylist { ids, target, after } => {
2115                let entries = self.shared_state.items_before(&ids);
2116                self.shared_state.move_items(&ids, target, after);
2117                self.push_undo(UndoEntry::MovedBatch { entries });
2118            }
2119            PlayerCommand::ReorderPlaylist(order) => {
2120                // Undoable like any other move: undoing it puts the queue back
2121                // and, by doing so, ends the lock — which is the honest result
2122                // of having rearranged the queue by hand.
2123                let entries = self.shared_state.items_before(&order);
2124                self.shared_state.reorder_to(&order);
2125                self.push_undo(UndoEntry::MovedBatch { entries });
2126            }
2127            PlayerCommand::TrackReady(id) => self.track_ready(id),
2128            PlayerCommand::DecodeFinished(session) => self.on_decode_finished(session),
2129            // Nothing to do but wake: the loop works out when the playhead
2130            // reaches the track just queued.
2131            PlayerCommand::TrackQueued => {}
2132            PlayerCommand::TrackStreamReady(id) => self.track_stream_ready(id),
2133            PlayerCommand::StreamProbed { id, info, mode } => self.stream_probed(id, *info, mode),
2134            PlayerCommand::TrackFailed(id) => self.track_failed(id),
2135            PlayerCommand::CacheTracks(ids) => {
2136                if let Some(downloads) = &self.downloads {
2137                    downloads.cache(ids);
2138                }
2139            }
2140            PlayerCommand::Undo => self.execute_undo(),
2141            PlayerCommand::Redo => self.execute_redo(),
2142            PlayerCommand::BeginUndoBatch => {
2143                self.batch_buffer = Some(Vec::new());
2144            }
2145            PlayerCommand::EndUndoBatch => {
2146                if let Some(entries) = self.batch_buffer.take() {
2147                    if entries.len() == 1 {
2148                        // Single entry — push directly, no wrapping.
2149                        self.undo_stack.push(entries.into_iter().next().unwrap());
2150                    } else if !entries.is_empty() {
2151                        self.undo_stack.push(UndoEntry::Batch(entries));
2152                    }
2153                }
2154            }
2155            PlayerCommand::SetOutputDevice(name) => self.set_output_device(name),
2156            PlayerCommand::RestartOutput => {
2157                log::info!("restarting audio output");
2158                self.restart_on_current_track();
2159            }
2160            PlayerCommand::ClearOutputDevice => self.clear_output_device(),
2161            PlayerCommand::ReloadDsp => self.reload_dsp(),
2162            PlayerCommand::UseRenderer(connection) => {
2163                self.remember_renderer(connection.as_deref());
2164                self.use_renderer(connection);
2165            }
2166            PlayerCommand::ResumeRenderer(connection) => {
2167                if std::mem::take(&mut self.resume_renderer) && self.renderer.is_none() {
2168                    log::info!(
2169                        "upnp: back to {}, as last time",
2170                        connection.session.renderer().name
2171                    );
2172                    self.use_renderer(Some(connection));
2173                } else {
2174                    log::info!(
2175                        "upnp: not going back to {}: playback or the output moved first",
2176                        connection.session.renderer().name
2177                    );
2178                }
2179            }
2180            PlayerCommand::SetRendererVolume(volume) => self.set_renderer_volume(volume),
2181            PlayerCommand::Renderer { session, event } => self.on_renderer_event(session, event),
2182            PlayerCommand::SetShuffle(on) => self.set_shuffle(on),
2183            PlayerCommand::SetRepeat(repeat) => self.mode.repeat = repeat,
2184            PlayerCommand::RestorePlayMode(mode) => self.mode = mode,
2185        }
2186    }
2187
2188    /// Turn shuffle on or off, as one undoable step.
2189    ///
2190    /// On, the items after the cursor go in a random order, each keeping a
2191    /// note of where it stood. Off, they go back where they stood, in the
2192    /// places such items occupy now, so an item added meanwhile stays put.
2193    /// The queue is the order things play in either way: the lookahead, the
2194    /// downloads and every remote queue simply follow it.
2195    fn set_shuffle(&mut self, on: bool) {
2196        if self.mode.shuffle == on {
2197            return;
2198        }
2199        let order = self.shared_state.shuffle_order();
2200        if on {
2201            self.shared_state.shuffle_after_cursor();
2202        } else {
2203            self.shared_state.unshuffle();
2204        }
2205        self.push_undo(UndoEntry::Shuffled {
2206            shuffle: self.mode.shuffle,
2207            order,
2208        });
2209        self.mode.shuffle = on;
2210    }
2211
2212    /// Stop, then clear the playlist as one undoable step. Playback and display
2213    /// state go first, without touching the playlist, so the snapshot taken for
2214    /// undo is of the playlist as it was.
2215    fn clear_playlist(&mut self) {
2216        self.stop_playback_and_clear_state();
2217        let (items, cursor) = self.shared_state.snapshot_playlist();
2218        self.shared_state.clear_playlist();
2219        self.push_undo(UndoEntry::Replaced { items, cursor });
2220    }
2221
2222    /// Apply an undo/redo entry: mutate the playlist and return the inverse entry.
2223    fn apply_entry(&mut self, entry: UndoEntry) -> UndoEntry {
2224        match entry {
2225            UndoEntry::Added { ids } | UndoEntry::Inserted { ids } => {
2226                // Undo of "items were added": snapshot them with positions, then remove.
2227                let items_with_pos = self.snapshot_for_undo(&ids);
2228                self.shared_state.remove_items(&ids);
2229                UndoEntry::Removed {
2230                    items: items_with_pos,
2231                }
2232            }
2233            UndoEntry::Removed { items } => {
2234                // Undo of "items were removed": re-insert each at its position.
2235                let mut ids = Vec::with_capacity(items.len());
2236                for (item, after) in items {
2237                    ids.push(item.id);
2238                    self.shared_state.insert_item_at(*item, after);
2239                }
2240                UndoEntry::Added { ids }
2241            }
2242            UndoEntry::Moved { id, was_after } => {
2243                let current_after = self.shared_state.item_before(id);
2244                self.shared_state.move_item_to(id, was_after);
2245                UndoEntry::Moved {
2246                    id,
2247                    was_after: current_after,
2248                }
2249            }
2250            UndoEntry::MovedBatch { entries } => {
2251                let ids: Vec<QueueItemId> = entries.iter().map(|(id, _)| *id).collect();
2252                let current_positions = self.shared_state.items_before(&ids);
2253                self.shared_state.move_items_to(&entries);
2254                UndoEntry::MovedBatch {
2255                    entries: current_positions,
2256                }
2257            }
2258            UndoEntry::Replaced { items, cursor } => {
2259                let (current_items, current_cursor) = self.shared_state.snapshot_playlist();
2260                self.shared_state.restore_playlist(items, cursor);
2261                UndoEntry::Replaced {
2262                    items: current_items,
2263                    cursor: current_cursor,
2264                }
2265            }
2266            UndoEntry::Shuffled { shuffle, order } => {
2267                let inverse = UndoEntry::Shuffled {
2268                    shuffle: self.mode.shuffle,
2269                    order: self.shared_state.shuffle_order(),
2270                };
2271                self.shared_state.restore_shuffle_order(&order);
2272                self.mode.shuffle = shuffle;
2273                inverse
2274            }
2275            UndoEntry::Batch(entries) => {
2276                // Apply entries in reverse order, collect inverses.
2277                let mut inverses: Vec<_> = entries
2278                    .into_iter()
2279                    .rev()
2280                    .map(|e| self.apply_entry(e))
2281                    .collect();
2282                inverses.reverse();
2283                UndoEntry::Batch(inverses)
2284            }
2285        }
2286    }
2287
2288    /// Put playback back in agreement with the playlist.
2289    ///
2290    /// The engine keeps decoding whatever it was on while the playlist changes
2291    /// underneath it, which an undo can turn into a lie: undoing a replace
2292    /// restores the queue but leaves the engine playing a track that queue does
2293    /// not contain. The transport then describes an item nothing can select,
2294    /// and the decode lookahead — which finds the next track by locating the
2295    /// current one — has nothing to follow, so the queue ends at the end of the
2296    /// track instead of carrying on.
2297    ///
2298    /// Done once, after the entry is applied, rather than inside each variant:
2299    /// any undo that takes items away can orphan the engine, not only
2300    /// `Replaced`. A track waiting to open can be orphaned the same way, and
2301    /// would otherwise open when its download lands.
2302    ///
2303    /// Playing carries on playing and paused stays paused, the same as when
2304    /// the playing track is removed: an undo is not a reason to start the
2305    /// music, or to stop it.
2306    fn reconcile_playback(&mut self) {
2307        let orphaned = match &self.transport {
2308            Transport::Idle => false,
2309            Transport::Waiting(waiting) => self.shared_state.get_item(waiting.id).is_none(),
2310            Transport::Loaded(session) => self.shared_state.get_item(session.track.id).is_none(),
2311        };
2312        if !orphaned {
2313            return;
2314        }
2315        // Pick the restored queue up where its cursor says it was, as the
2316        // listener had it. The position is not part of what was snapshotted,
2317        // so the track begins again.
2318        let cursor = self.shared_state.cursor();
2319        self.carry_on(cursor, self.intent());
2320    }
2321
2322    /// Execute an undo operation, pushing the inverse onto the redo stack.
2323    fn execute_undo(&mut self) {
2324        let Some(entry) = self.undo_stack.pop_undo() else {
2325            return;
2326        };
2327        let inverse = self.apply_entry(entry);
2328        self.undo_stack.push_redo(inverse);
2329        self.reconcile_playback();
2330    }
2331
2332    /// Execute a redo operation, pushing the inverse onto the undo stack.
2333    fn execute_redo(&mut self) {
2334        let Some(entry) = self.undo_stack.pop_redo() else {
2335            return;
2336        };
2337        let inverse = self.apply_entry(entry);
2338        self.undo_stack.push_undo_keep_redo(inverse);
2339        self.reconcile_playback();
2340    }
2341
2342    /// Run the command loop. Blocks until the sender is dropped.
2343    ///
2344    /// Asleep until there is something to do: a command, or one of the two
2345    /// things that happen without one — see `next_wake`. A paused or stopped
2346    /// player does not wake at all.
2347    pub fn run(&mut self) {
2348        use crossbeam_channel::RecvTimeoutError;
2349
2350        let rx = self.commands.rx.clone();
2351        loop {
2352            let received = match self.next_wake() {
2353                Some(at) => rx.recv_deadline(at),
2354                None => rx.recv().map_err(|_| RecvTimeoutError::Disconnected),
2355            };
2356            match received {
2357                Ok(cmd) => self.process_command(cmd),
2358                Err(RecvTimeoutError::Timeout) => {}
2359                Err(RecvTimeoutError::Disconnected) => break,
2360            }
2361            self.update_playback_state();
2362        }
2363        self.stop();
2364    }
2365
2366    /// When something changes that no command announces: the playhead
2367    /// reaching the next queued track, the silence after a rate switch
2368    /// running out, or a pause fading to silence.
2369    fn next_wake(&self) -> Option<std::time::Instant> {
2370        let session = self.session()?;
2371        let now = std::time::Instant::now();
2372        let Output::Local(local) = &session.output else {
2373            // A stream's track changes are the timeline's, as they are here.
2374            let next_track = (session.run == Run::Playing && self.streaming_to_renderer())
2375                .then(|| self.timeline.until_next_track())
2376                .flatten()
2377                .map(|left| now + left + BOUNDARY_SLACK);
2378            return match (next_track, self.renderer_deadline()) {
2379                (Some(a), Some(b)) => Some(a.min(b)),
2380                (a, b) => a.or(b),
2381            };
2382        };
2383        match session.run {
2384            Run::Playing => {
2385                let next_track = self
2386                    .timeline
2387                    .until_next_track()
2388                    .map(|left| now + left + BOUNDARY_SLACK);
2389                match (next_track, self.lead_in_ends) {
2390                    (Some(a), Some(b)) => Some(a.min(b)),
2391                    (a, b) => a.or(b),
2392                }
2393            }
2394            Run::Paused if local.engine.is_running() => Some(now + FADE_CHECK),
2395            Run::Paused => None,
2396        }
2397    }
2398
2399    /// Spawn the player on a background thread, returning the shared state,
2400    /// timeline, visualization snapshot, and command sender.
2401    /// Remember the renderer picked as the output, or that this device's own
2402    /// was, for the next launch. Only a choice is remembered: a renderer that
2403    /// drops off the network is gone back to next time.
2404    fn remember_renderer(&self, connection: Option<&crate::upnp::Connection>) {
2405        let renderer = connection.map(|c| c.session.renderer());
2406        if let Err(e) = crate::config::Config::persist(|cfg| {
2407            cfg.playback.renderer = renderer.map(|r| r.udn.clone());
2408            cfg.playback.renderer_name = renderer.map(|r| r.name.clone());
2409        }) {
2410            log::error!("failed to save the output: {e}");
2411        }
2412    }
2413
2414    /// Spawn a player for a process that does not own an output of its own:
2415    /// `koan serve`, `koan mcp`. It plays where it is told and nowhere else.
2416    pub fn spawn() -> (
2417        Arc<SharedPlayerState>,
2418        Arc<PlaybackTimeline>,
2419        Arc<VizSnapshot>,
2420        crossbeam_channel::Sender<PlayerCommand>,
2421    ) {
2422        Self::spawn_with(false)
2423    }
2424
2425    /// Spawn the player for an app someone listens through: the macOS and iOS
2426    /// apps and `koan play`. It goes back to the renderer used last time if
2427    /// that turns up at launch (`upnp::resume`). A headless process must not:
2428    /// the config is the machine's, and it would take the amplifier from the
2429    /// app.
2430    pub fn spawn_for_listening() -> (
2431        Arc<SharedPlayerState>,
2432        Arc<PlaybackTimeline>,
2433        Arc<VizSnapshot>,
2434        crossbeam_channel::Sender<PlayerCommand>,
2435    ) {
2436        Self::spawn_with(true)
2437    }
2438
2439    /// The renderer to go back to at launch: the one used last time, for a
2440    /// player someone listens through.
2441    fn renderer_to_resume(listening: bool) -> Option<String> {
2442        listening
2443            .then(|| crate::config::Config::cached().playback.renderer.clone())
2444            .flatten()
2445    }
2446
2447    fn spawn_with(
2448        listening: bool,
2449    ) -> (
2450        Arc<SharedPlayerState>,
2451        Arc<PlaybackTimeline>,
2452        Arc<VizSnapshot>,
2453        crossbeam_channel::Sender<PlayerCommand>,
2454    ) {
2455        let mut player = Self::new();
2456        player.history = PlayRecorder::spawn();
2457        let state = player.shared_state();
2458        let timeline = player.timeline();
2459        let viz_snapshot = player.viz_snapshot();
2460        let tx = player.command_sender();
2461        // Downloads follow the playlist, so they come with the player rather
2462        // than being something each front end has to remember to ask for.
2463        player.downloads = Some(crate::remote::queue::DownloadQueue::spawn(
2464            tx.clone(),
2465            state.clone(),
2466        ));
2467
2468        if let Some(udn) = Self::renderer_to_resume(listening) {
2469            player.resume_renderer = true;
2470            crate::upnp::resume(udn, &tx);
2471        }
2472
2473        thread::Builder::new()
2474            .name("koan-player".into())
2475            .spawn(move || player.run())
2476            .expect("failed to spawn player thread");
2477
2478        (state, timeline, viz_snapshot, tx)
2479    }
2480}
2481
2482#[cfg(test)]
2483mod tests {
2484    /// `koan serve` and `koan mcp` read the same config as the app, and must
2485    /// not go looking for the app's amplifier.
2486    #[test]
2487    fn a_headless_player_does_not_go_back_to_a_renderer() {
2488        crate::config::isolate_config_for_tests();
2489        crate::config::Config::persist(|cfg| {
2490            cfg.playback.renderer = Some("uuid:headless-test".into());
2491        })
2492        .unwrap();
2493        assert_eq!(Player::renderer_to_resume(false), None);
2494    }
2495
2496    #[test]
2497    fn a_download_in_progress_is_known_by_its_own_extension() {
2498        use std::path::Path;
2499        // `.part` is the transfer's, not the track's.
2500        assert_eq!(
2501            lengthless_mode_for(Path::new("/c/t.m4a.part")),
2502            streaming::ProbeMode::LengthlessWholeEnd
2503        );
2504        assert_eq!(
2505            lengthless_mode_for(Path::new("/c/t.OPUS.part")),
2506            streaming::ProbeMode::LengthlessWholeEnd
2507        );
2508        assert_eq!(
2509            lengthless_mode_for(Path::new("/c/t.flac.part")),
2510            streaming::ProbeMode::Lengthless
2511        );
2512        assert_eq!(
2513            media_extension(Path::new("/c/t.m4a.part")).as_deref(),
2514            Some("m4a")
2515        );
2516        assert_eq!(
2517            media_extension(Path::new("/c/t.mp3")).as_deref(),
2518            Some("mp3")
2519        );
2520    }
2521
2522    use super::*;
2523    use state::PlaylistItem;
2524    use std::path::PathBuf;
2525    use std::sync::atomic::{AtomicU64, Ordering};
2526
2527    fn make_item(title: &str) -> PlaylistItem {
2528        PlaylistItem {
2529            playlist_entry_id: None,
2530            id: QueueItemId::new(),
2531            db_id: None,
2532            path: PathBuf::from(format!("/music/{title}.flac")),
2533            title: title.to_string(),
2534            artist: String::new(),
2535            album_artist: String::new(),
2536            album: String::new(),
2537            year: None,
2538            codec: None,
2539            track_number: None,
2540            disc: None,
2541            duration_ms: None,
2542            state: ItemState::Ready,
2543            pre_shuffle: None,
2544        }
2545    }
2546
2547    pub(super) fn playlist_ids(player: &Player) -> Vec<QueueItemId> {
2548        let (items, _) = player.shared_state.snapshot_playlist();
2549        items.iter().map(|i| i.id).collect()
2550    }
2551
2552    fn playlist_titles(player: &Player) -> Vec<String> {
2553        let (items, _) = player.shared_state.snapshot_playlist();
2554        items.iter().map(|i| i.title.clone()).collect()
2555    }
2556
2557    fn pending_item(title: &str) -> PlaylistItem {
2558        PlaylistItem {
2559            playlist_entry_id: None,
2560            state: ItemState::Pending,
2561            ..make_item(title)
2562        }
2563    }
2564
2565    /// A session playing `id` on `engine`, with no decode thread behind it.
2566    fn test_session(id: QueueItemId, engine: Box<dyn AudioEngineHandle>) -> Session {
2567        Session {
2568            track: TrackInfo {
2569                id,
2570                path: PathBuf::from("/music/t.flac"),
2571                codec: String::new(),
2572                sample_rate: 44_100,
2573                bit_depth: None,
2574                bitrate_kbps: None,
2575                channels: 2,
2576                duration_ms: 1_000,
2577            },
2578            run: Run::Playing,
2579            lookahead: Default::default(),
2580            output: Output::Local(Local {
2581                engine,
2582                decode_handle: buffer::DecodeHandle::new_for_test(Default::default()),
2583                stream: None,
2584                _rate_watch: None,
2585                dsp: None,
2586            }),
2587        }
2588    }
2589
2590    /// Stand in for an engine that is playing `id`. The test items have no
2591    /// files behind them, so `start_playback` can never get far enough to leave
2592    /// this state on its own.
2593    fn pretend_playing(player: &mut Player, id: QueueItemId) {
2594        assert!(
2595            player.shared_state.get_item(id).is_some(),
2596            "item is in the queue"
2597        );
2598        player.stop_engine();
2599        player.transport = Transport::Loaded(test_session(
2600            id,
2601            Box::new(NullEngine {
2602                starts: Default::default(),
2603                running: Default::default(),
2604                lead_in: Default::default(),
2605            }),
2606        ));
2607        player.shared_state.set_cursor(Some(id));
2608        player.publish();
2609    }
2610
2611    fn playing_id(player: &Player) -> Option<QueueItemId> {
2612        player.shared_state.track_info().map(|t| t.id)
2613    }
2614
2615    /// Build `n` ready items, add them, and return their IDs.
2616    fn seed(player: &mut Player, n: usize) -> Vec<QueueItemId> {
2617        let items: Vec<_> = (0..n).map(|i| make_item(&format!("t{i}"))).collect();
2618        let ids = items.iter().map(|i| i.id).collect();
2619        player.process_command(PlayerCommand::AddToPlaylist(items));
2620        ids
2621    }
2622
2623    // --- cursor transitions ---
2624
2625    /// Feed the player a track's worth of playback ticks, as the 50ms poll would.
2626    fn listen(player: &mut Player, from_ms: u64, to_ms: u64) {
2627        let mut at = from_ms;
2628        if let Some(f) = player.in_flight.as_mut() {
2629            f.advance(at); // the position the needle landed on
2630        }
2631        while at < to_ms {
2632            at = (at + 50).min(to_ms);
2633            if let Some(f) = player.in_flight.as_mut() {
2634                f.advance(at);
2635            }
2636        }
2637    }
2638
2639    fn start(player: &mut Player, track_id: i64) -> QueueItemId {
2640        let id = QueueItemId::new();
2641        player.on_track_changed(id, 0);
2642        // The item is not in a playlist here, so there is no db_id to find.
2643        player
2644            .in_flight
2645            .as_mut()
2646            .unwrap()
2647            .track_id_for_test(track_id);
2648        id
2649    }
2650
2651    #[test]
2652    fn a_gapless_transition_closes_the_outgoing_track_and_opens_the_next() {
2653        let mut player = Player::new();
2654        start(&mut player, 11);
2655        listen(&mut player, 0, 200_000);
2656
2657        let b = QueueItemId::new();
2658        player.on_track_changed(b, 0);
2659        let f = player
2660            .in_flight
2661            .as_ref()
2662            .expect("the next track is counting");
2663        assert_eq!(f.item, b);
2664        assert_eq!(f.listened_ms(), 0, "and starts from nothing");
2665    }
2666
2667    #[test]
2668    fn a_track_skipped_seconds_in_is_still_history() {
2669        let mut player = Player::new();
2670        start(&mut player, 7);
2671        listen(&mut player, 0, 2_000);
2672
2673        let event = player
2674            .finish_play()
2675            .expect("putting something on is a thing you did, however briefly");
2676        assert!(matches!(
2677            event,
2678            history::PlayEvent::Finished {
2679                track_id: 7,
2680                listened_ms: 2_000
2681            }
2682        ));
2683    }
2684
2685    #[test]
2686    fn a_track_is_closed_out_once() {
2687        let mut player = Player::new();
2688        start(&mut player, 7);
2689        listen(&mut player, 0, 200_000);
2690
2691        assert!(player.finish_play().is_some());
2692        assert!(player.finish_play().is_none());
2693    }
2694
2695    #[test]
2696    fn seeking_around_a_track_does_not_enter_it_twice() {
2697        let mut player = Player::new();
2698        let id = start(&mut player, 7);
2699        listen(&mut player, 0, 120_000);
2700
2701        // A seek restarts playback of the same item.
2702        player.on_track_changed(id, 30_000);
2703        assert_eq!(
2704            player.in_flight.as_ref().unwrap().listened_ms(),
2705            120_000,
2706            "the seek kept the count rather than restarting it"
2707        );
2708        listen(&mut player, 30_000, 40_000);
2709
2710        let Some(history::PlayEvent::Finished { listened_ms, .. }) = player.finish_play() else {
2711            panic!("still one play");
2712        };
2713        assert_eq!(listened_ms, 130_000);
2714        assert!(player.finish_play().is_none());
2715    }
2716
2717    #[test]
2718    fn a_track_that_is_not_in_the_library_is_not_recorded() {
2719        let mut player = Player::new();
2720        let id = QueueItemId::new();
2721        player.on_track_changed(id, 0);
2722        listen(&mut player, 0, 200_000);
2723        assert!(player.finish_play().is_none());
2724    }
2725
2726    #[test]
2727    fn stopping_closes_out_what_was_heard() {
2728        let mut player = Player::new();
2729        start(&mut player, 7);
2730        listen(&mut player, 0, 150_000);
2731
2732        player.stop_playback_and_clear_state();
2733        assert!(player.in_flight.is_none(), "the stop consumed it");
2734    }
2735
2736    #[test]
2737    fn resume_with_nothing_loaded_plays_the_cursor() {
2738        let dir = tempfile::tempdir().unwrap();
2739        let path = dir.path().join("t.wav");
2740        crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2741
2742        let mut player = Player::new();
2743        player.backend = Box::new(StuckBackend {
2744            rate: 8_000.0,
2745            asked: Default::default(),
2746            starts: Default::default(),
2747        });
2748        let item = PlaylistItem {
2749            db_id: Some(5),
2750            path,
2751            ..make_item("t")
2752        };
2753        let id = item.id;
2754        player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2755        player.shared_state.set_cursor(Some(id));
2756        assert!(player.session().is_none());
2757
2758        player.process_command(PlayerCommand::Resume);
2759        assert!(player.session().is_some(), "the cursor's track started");
2760        assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2761        player.process_command(PlayerCommand::Stop);
2762    }
2763
2764    #[test]
2765    fn a_player_with_nothing_coming_has_nothing_to_wake_for() {
2766        let dir = tempfile::tempdir().unwrap();
2767        let path = dir.path().join("t.wav");
2768        crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2769
2770        let mut player = Player::new();
2771        player.backend = Box::new(StuckBackend {
2772            rate: 8_000.0,
2773            asked: Default::default(),
2774            starts: Default::default(),
2775        });
2776        let item = PlaylistItem {
2777            path,
2778            ..make_item("t")
2779        };
2780        let id = item.id;
2781        player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2782        assert_eq!(player.next_wake(), None, "stopped");
2783
2784        player.process_command(PlayerCommand::Play(id));
2785        assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2786        assert_eq!(
2787            player.next_wake(),
2788            None,
2789            "playing, with no track queued after it"
2790        );
2791
2792        if let Transport::Loaded(session) = &mut player.transport {
2793            session.engine().unwrap().stop().unwrap();
2794            session.run = Run::Paused;
2795        }
2796        assert_eq!(player.next_wake(), None, "paused");
2797        player.process_command(PlayerCommand::Stop);
2798    }
2799
2800    #[test]
2801    fn a_track_cued_paused_never_starts_the_output() {
2802        let dir = tempfile::tempdir().unwrap();
2803        let path = dir.path().join("t.wav");
2804        crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2805
2806        let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2807        let mut player = Player::new();
2808        player.backend = Box::new(StuckBackend {
2809            rate: 8_000.0,
2810            asked: Default::default(),
2811            starts: starts.clone(),
2812        });
2813        let item = PlaylistItem {
2814            path,
2815            ..make_item("t")
2816        };
2817        let id = item.id;
2818        player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2819
2820        player.process_command(PlayerCommand::Cue {
2821            id,
2822            position_ms: 2_000,
2823            play: false,
2824        });
2825        assert!(player.session().is_some(), "loaded");
2826        assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
2827        // The decoder says where the seek landed when it queues the track,
2828        // and the playhead reads that from then on: the start of the packet
2829        // holding 2s, which can be a little before it.
2830        loop {
2831            match player
2832                .commands
2833                .rx
2834                .recv_timeout(std::time::Duration::from_secs(5))
2835            {
2836                Ok(PlayerCommand::TrackQueued) => break,
2837                Ok(_) => {}
2838                Err(e) => panic!("the decoder never queued the track: {e}"),
2839            }
2840        }
2841        let at = player.shared_state.position_ms();
2842        assert!((1_750..=2_000).contains(&at), "cued at {at}ms");
2843        assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
2844
2845        // Seeking while paused reopens the track, and stays quiet too.
2846        player.process_command(PlayerCommand::Seek(4_000));
2847        assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
2848        assert_eq!(starts.load(Ordering::Relaxed), 0);
2849
2850        player.process_command(PlayerCommand::Resume);
2851        assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2852        assert_eq!(starts.load(Ordering::Relaxed), 1);
2853        player.process_command(PlayerCommand::Stop);
2854    }
2855
2856    /// Load `t.wav` as an item whose download has not landed, under a player
2857    /// on the fake output.
2858    fn downloading_wav(dir: &Path) -> (Player, QueueItemId, Arc<std::sync::atomic::AtomicUsize>) {
2859        let path = dir.join("t.wav");
2860        crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2861        let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2862        let mut player = Player::new();
2863        player.backend = Box::new(StuckBackend {
2864            rate: 8_000.0,
2865            asked: Default::default(),
2866            starts: starts.clone(),
2867        });
2868        let item = PlaylistItem {
2869            path,
2870            state: ItemState::Pending,
2871            ..make_item("t")
2872        };
2873        let id = item.id;
2874        player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2875        (player, id, starts)
2876    }
2877
2878    /// The download lands, as the downloader says so: the item first, then
2879    /// the player.
2880    fn land(player: &mut Player, id: QueueItemId) {
2881        player.shared_state.update_item_state(id, ItemState::Ready);
2882        player.process_command(PlayerCommand::TrackReady(id));
2883    }
2884
2885    #[test]
2886    fn track_ready_never_resurrects_an_item_that_failed() {
2887        let mut player = Player::new();
2888        let item = PlaylistItem {
2889            state: ItemState::Failed("gone".into()),
2890            ..make_item("t")
2891        };
2892        let id = item.id;
2893        player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2894
2895        player.process_command(PlayerCommand::TrackReady(id));
2896
2897        assert!(matches!(
2898            player.shared_state.get_item(id).map(|i| i.state),
2899            Some(ItemState::Failed(_))
2900        ));
2901    }
2902
2903    fn await_queued(player: &Player) {
2904        loop {
2905            match player
2906                .commands
2907                .rx
2908                .recv_timeout(std::time::Duration::from_secs(5))
2909            {
2910                Ok(PlayerCommand::TrackQueued) => return,
2911                Ok(_) => {}
2912                Err(e) => panic!("the decoder never queued the track: {e}"),
2913            }
2914        }
2915    }
2916
2917    #[test]
2918    fn a_cue_on_a_downloading_track_opens_at_its_position_once_it_lands() {
2919        let dir = tempfile::tempdir().unwrap();
2920        let (mut player, id, starts) = downloading_wav(dir.path());
2921
2922        player.process_command(PlayerCommand::Cue {
2923            id,
2924            position_ms: 6_000,
2925            play: true,
2926        });
2927        assert!(player.session().is_none(), "nothing opened early");
2928        assert_eq!(player.shared_state.playback_state(), PlaybackState::Stopped);
2929        assert_eq!(starts.load(Ordering::Relaxed), 0);
2930
2931        land(&mut player, id);
2932        assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2933        assert_eq!(player.playback_starts, 1, "opened once, at the position");
2934        assert_eq!(starts.load(Ordering::Relaxed), 1);
2935        await_queued(&player);
2936        let at = player.shared_state.position_ms();
2937        assert!((5_750..=6_000).contains(&at), "opened at {at}ms");
2938        player.process_command(PlayerCommand::Stop);
2939    }
2940
2941    #[test]
2942    fn a_paused_cue_on_a_downloading_track_lands_paused() {
2943        let dir = tempfile::tempdir().unwrap();
2944        let (mut player, id, starts) = downloading_wav(dir.path());
2945
2946        player.process_command(PlayerCommand::Cue {
2947            id,
2948            position_ms: 3_000,
2949            play: false,
2950        });
2951        land(&mut player, id);
2952        assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
2953        assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
2954        await_queued(&player);
2955        let at = player.shared_state.position_ms();
2956        assert!((2_750..=3_000).contains(&at), "cued at {at}ms");
2957        player.process_command(PlayerCommand::Stop);
2958    }
2959
2960    #[test]
2961    fn playing_something_else_forgets_a_waiting_cue() {
2962        let dir = tempfile::tempdir().unwrap();
2963        let (mut player, id, _) = downloading_wav(dir.path());
2964        let other = seed(&mut player, 1)[0];
2965
2966        player.process_command(PlayerCommand::Cue {
2967            id,
2968            position_ms: 6_000,
2969            play: true,
2970        });
2971        player.process_command(PlayerCommand::Play(other));
2972        player.process_command(PlayerCommand::Play(id));
2973        assert!(player.waiting().is_some_and(|w| w.position_ms == 0));
2974        player.process_command(PlayerCommand::Stop);
2975    }
2976
2977    #[test]
2978    fn a_paused_track_whose_download_lands_stays_paused_where_it_was() {
2979        let dir = tempfile::tempdir().unwrap();
2980        let (mut player, id, starts) = downloading_wav(dir.path());
2981        // Loaded paused mid-track, as a stream is when its download lands.
2982        player.shared_state.update_item_state(id, ItemState::Ready);
2983        player.process_command(PlayerCommand::Cue {
2984            id,
2985            position_ms: 3_000,
2986            play: false,
2987        });
2988        await_queued(&player);
2989
2990        player.process_command(PlayerCommand::TrackReady(id));
2991        assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
2992        assert_eq!(player.playback_starts, 1, "not reopened");
2993        assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
2994        player.process_command(PlayerCommand::Stop);
2995    }
2996
2997    #[test]
2998    fn pausing_a_track_on_its_way_opens_it_paused() {
2999        let dir = tempfile::tempdir().unwrap();
3000        let (mut player, id, starts) = downloading_wav(dir.path());
3001
3002        player.process_command(PlayerCommand::Play(id));
3003        assert!(player.shared_state.is_waiting());
3004        assert!(player.shared_state.wants_to_play(), "a toggle pauses it");
3005        assert!(!player.shared_state.is_idle(), "adding tracks leaves it be");
3006        player.process_command(PlayerCommand::Pause);
3007        assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
3008        assert!(!player.shared_state.wants_to_play());
3009
3010        player.shared_state.update_item_state(id, ItemState::Ready);
3011        player.process_command(PlayerCommand::TrackReady(id));
3012        assert!(player.session().is_some(), "loaded");
3013        assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
3014        assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
3015        player.process_command(PlayerCommand::Stop);
3016    }
3017
3018    #[test]
3019    fn a_cue_paused_and_resumed_on_its_way_keeps_its_position() {
3020        let dir = tempfile::tempdir().unwrap();
3021        let (mut player, id, starts) = downloading_wav(dir.path());
3022
3023        player.process_command(PlayerCommand::Cue {
3024            id,
3025            position_ms: 6_000,
3026            play: true,
3027        });
3028        player.process_command(PlayerCommand::Pause);
3029        player.process_command(PlayerCommand::Resume);
3030        assert!(player.session().is_none(), "still on its way");
3031        assert_eq!(player.shared_state.playback_state(), PlaybackState::Stopped);
3032
3033        player.shared_state.update_item_state(id, ItemState::Ready);
3034        player.process_command(PlayerCommand::TrackReady(id));
3035        assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
3036        assert_eq!(starts.load(Ordering::Relaxed), 1);
3037        await_queued(&player);
3038        let at = player.shared_state.position_ms();
3039        assert!((5_750..=6_000).contains(&at), "opened at {at}ms");
3040        player.process_command(PlayerCommand::Stop);
3041    }
3042
3043    #[test]
3044    fn a_restored_cursor_does_not_start_when_its_download_lands() {
3045        let dir = tempfile::tempdir().unwrap();
3046        let (mut player, id, starts) = downloading_wav(dir.path());
3047        player.shared_state.set_cursor(Some(id));
3048
3049        player.shared_state.update_item_state(id, ItemState::Ready);
3050        player.process_command(PlayerCommand::TrackReady(id));
3051        assert!(player.session().is_none());
3052        assert_eq!(starts.load(Ordering::Relaxed), 0);
3053    }
3054
3055    #[test]
3056    fn replacing_the_queue_keeps_a_transfer_both_queues_want() {
3057        // The download queue syncs from its own thread on each playlist
3058        // change. Replaced as a clear then an add, the playlist was empty in
3059        // between, and a sync there let go of every waiter: the transfer for a
3060        // track in both queues was abandoned and started over. The replace
3061        // must be one change, and that change must still want the track.
3062        let pending = |title: &str| PlaylistItem {
3063            db_id: Some(7),
3064            state: ItemState::Pending,
3065            ..make_item(title)
3066        };
3067        let mut player = Player::new();
3068        let old = pending("old");
3069        let old_id = old.id;
3070        player.process_command(PlayerCommand::AddToPlaylist(vec![old]));
3071        let store = player.shared_state.downloads().clone();
3072        store.claim(7, Some(old_id));
3073
3074        let before = player.shared_state.pending_version();
3075        player.process_command(PlayerCommand::ReplacePlaylist {
3076            items: vec![pending("again")],
3077            start: 0,
3078            position_ms: 0,
3079            play: false,
3080        });
3081        assert_eq!(
3082            player.shared_state.pending_version(),
3083            before + 1,
3084            "one change, so no reader can see the playlist between two"
3085        );
3086
3087        // A sync against the one state it published.
3088        store.resync(&player.shared_state.pending_downloads());
3089        assert!(!store.abandoned(7));
3090    }
3091
3092    #[test]
3093    fn a_hand_off_opens_paused_at_its_position_in_one_command() {
3094        let dir = tempfile::tempdir().unwrap();
3095        let path = dir.path().join("t.wav");
3096        crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
3097        let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
3098        let mut player = Player::new();
3099        player.backend = Box::new(StuckBackend {
3100            rate: 8_000.0,
3101            asked: Default::default(),
3102            starts: starts.clone(),
3103        });
3104        seed(&mut player, 2);
3105        let item = PlaylistItem {
3106            path,
3107            ..make_item("t")
3108        };
3109        let id = item.id;
3110
3111        player.process_command(PlayerCommand::ReplacePlaylist {
3112            items: vec![make_item("before"), item],
3113            start: 1,
3114            position_ms: 4_000,
3115            play: false,
3116        });
3117        assert_eq!(
3118            playlist_titles(&player),
3119            vec!["before", "t"],
3120            "replaced, not added to"
3121        );
3122        assert_eq!(player.shared_state.cursor(), Some(id));
3123        assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
3124        assert_eq!(player.playback_starts, 1, "opened once, at the position");
3125        assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
3126        await_queued(&player);
3127        let at = player.shared_state.position_ms();
3128        assert!((3_750..=4_000).contains(&at), "opened at {at}ms");
3129
3130        player.process_command(PlayerCommand::Undo);
3131        assert_eq!(playlist_titles(&player), vec!["t0", "t1"], "one undo step");
3132        player.process_command(PlayerCommand::Stop);
3133    }
3134
3135    #[test]
3136    fn a_decode_end_from_a_replaced_session_is_ignored() {
3137        let dir = tempfile::tempdir().unwrap();
3138        let mut player = Player::new();
3139        player.backend = Box::new(StuckBackend {
3140            rate: 8_000.0,
3141            asked: Default::default(),
3142            starts: Default::default(),
3143        });
3144        let items: Vec<_> = ["a", "b", "c"]
3145            .iter()
3146            .map(|name| {
3147                let path = dir.path().join(format!("{name}.wav"));
3148                crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
3149                PlaylistItem {
3150                    path,
3151                    ..make_item(name)
3152                }
3153            })
3154            .collect();
3155        let ids: Vec<_> = items.iter().map(|i| i.id).collect();
3156        player.process_command(PlayerCommand::AddToPlaylist(items));
3157
3158        player.process_command(PlayerCommand::Play(ids[0]));
3159        let stale = player.session;
3160        player.process_command(PlayerCommand::Play(ids[1]));
3161        player.process_command(PlayerCommand::DecodeFinished(stale));
3162
3163        assert_eq!(player.shared_state.cursor(), Some(ids[1]));
3164        assert_eq!(player.playback_starts, 2);
3165        player.process_command(PlayerCommand::Stop);
3166    }
3167
3168    /// Short tracks of one format under an output that never plays: the
3169    /// decoder queues them all behind the first.
3170    fn queued_wavs(dir: &Path, names: &[&str]) -> (Player, Vec<QueueItemId>) {
3171        let (player, ids) = wavs_playing(dir, names, 1.0);
3172        for _ in &ids {
3173            await_queued(&player);
3174        }
3175        (player, ids)
3176    }
3177
3178    /// `names` as WAVs of `seconds` each, the first playing. Long enough ones
3179    /// fill the ring before the decoder reaches the end of the queue.
3180    fn wavs_playing(dir: &Path, names: &[&str], seconds: f32) -> (Player, Vec<QueueItemId>) {
3181        wavs_in(dir, names, seconds, Repeat::Off)
3182    }
3183
3184    /// `wavs_playing`, under `repeat` from the start. Each has a library id,
3185    /// its index plus one, so history records it.
3186    fn wavs_in(
3187        dir: &Path,
3188        names: &[&str],
3189        seconds: f32,
3190        repeat: Repeat,
3191    ) -> (Player, Vec<QueueItemId>) {
3192        let mut player = Player::new();
3193        player.backend = Box::new(StuckBackend {
3194            rate: 8_000.0,
3195            asked: Default::default(),
3196            starts: Default::default(),
3197        });
3198        let items: Vec<_> = names
3199            .iter()
3200            .zip(1..)
3201            .map(|(name, track)| {
3202                let path = dir.join(format!("{name}.wav"));
3203                crate::test_utils::generate_wav(&path, 8_000, 1, seconds, 16);
3204                PlaylistItem {
3205                    path,
3206                    db_id: Some(track),
3207                    ..make_item(name)
3208                }
3209            })
3210            .collect();
3211        let ids: Vec<_> = items.iter().map(|i| i.id).collect();
3212        player.process_command(PlayerCommand::AddToPlaylist(items));
3213        player.process_command(PlayerCommand::SetRepeat(repeat));
3214        player.process_command(PlayerCommand::Play(ids[0]));
3215        (player, ids)
3216    }
3217
3218    /// What the decoder has queued after the playhead, once it is at least `n`.
3219    fn queued_at_least(player: &Player, n: usize) -> Vec<QueueItemId> {
3220        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
3221        loop {
3222            let queued = player.timeline.queued_after_playhead();
3223            if queued.len() >= n {
3224                return queued;
3225            }
3226            assert!(
3227                std::time::Instant::now() < deadline,
3228                "the decoder queued {queued:?}, never {n}"
3229            );
3230            std::thread::yield_now();
3231        }
3232    }
3233
3234    // --- play modes ---
3235
3236    #[test]
3237    fn repeating_the_queue_runs_the_last_track_into_the_first() {
3238        let dir = tempfile::tempdir().unwrap();
3239        let (mut player, ids) = wavs_in(dir.path(), &["a", "b"], 1.0, Repeat::Queue);
3240
3241        let queued = queued_at_least(&player, 3);
3242        assert_eq!(
3243            queued[..3],
3244            [ids[1], ids[0], ids[1]],
3245            "gapless, round again"
3246        );
3247        let steps = player.session().unwrap().lookahead.lock().clone();
3248        assert!(!steps[0].wrapped);
3249        let wrap = steps.iter().find(|s| s.after == ids[1]).unwrap();
3250        assert_eq!(wrap.next, Some(ids[0]));
3251        assert!(wrap.wrapped);
3252        assert_eq!(player.playback_starts, 1);
3253        player.process_command(PlayerCommand::Stop);
3254    }
3255
3256    #[test]
3257    fn turning_repeat_off_takes_back_a_wrap_already_queued() {
3258        let dir = tempfile::tempdir().unwrap();
3259        let (mut player, ids) = wavs_in(dir.path(), &["a"], 1.0, Repeat::Queue);
3260        assert_eq!(queued_at_least(&player, 1)[0], ids[0]);
3261
3262        player.process_command(PlayerCommand::SetRepeat(Repeat::Off));
3263        assert_eq!(player.playback_starts, 2, "restarted at the playhead");
3264        wait_for_lookahead(&player);
3265        assert!(player.timeline.queued_after_playhead().is_empty());
3266        player.process_command(PlayerCommand::Stop);
3267    }
3268
3269    #[test]
3270    fn repeating_one_track_queues_it_again_and_each_pass_is_a_play() {
3271        let dir = tempfile::tempdir().unwrap();
3272        let (mut player, ids) = wavs_in(dir.path(), &["a", "b"], 1.0, Repeat::One);
3273        let (recorder, events) = history::PlayRecorder::capture();
3274        player.history = Some(recorder);
3275        // The first pass as a play would have it, now there is somewhere to
3276        // record it.
3277        player.in_flight = None;
3278        player.on_track_changed(ids[0], 0);
3279
3280        assert_eq!(queued_at_least(&player, 2)[..2], [ids[0], ids[0]]);
3281        // Into the second pass: one second of 8 kHz mono and a little more.
3282        player
3283            .timeline
3284            .samples_played
3285            .store(8_400, Ordering::Relaxed);
3286        player.update_playback_state();
3287
3288        let flight = player.in_flight.as_ref().unwrap();
3289        assert_eq!((flight.item, flight.boundary), (ids[0], 1));
3290        assert_eq!(player.shared_state.cursor(), Some(ids[0]));
3291        let events: Vec<_> = events.try_iter().collect();
3292        let started = events
3293            .iter()
3294            .filter(|e| matches!(e, PlayEvent::Started { track_id: 1, .. }))
3295            .count();
3296        assert_eq!(started, 2, "two plays: {events:?}");
3297        assert!(
3298            events.contains(&PlayEvent::Finished {
3299                track_id: 1,
3300                listened_ms: 1_000
3301            }),
3302            "the first pass banked whole: {events:?}"
3303        );
3304        player.process_command(PlayerCommand::Stop);
3305    }
3306
3307    #[test]
3308    fn next_moves_on_from_a_track_repeating() {
3309        let dir = tempfile::tempdir().unwrap();
3310        let (mut player, ids) = wavs_in(dir.path(), &["a", "b"], 1.0, Repeat::One);
3311
3312        player.process_command(PlayerCommand::NextTrack);
3313        assert_eq!(playing_id(&player), Some(ids[1]));
3314        player.process_command(PlayerCommand::NextTrack);
3315        assert_eq!(
3316            playing_id(&player),
3317            Some(ids[0]),
3318            "round from the last, as repeating does"
3319        );
3320        player.process_command(PlayerCommand::Stop);
3321    }
3322
3323    #[test]
3324    fn a_repeat_of_nothing_heard_stops_rather_than_going_round() {
3325        let dir = tempfile::tempdir().unwrap();
3326        let (mut player, _) = wavs_in(dir.path(), &["a"], 1.0, Repeat::One);
3327
3328        // Nothing plays through the stuck output, so the session that ends
3329        // has nothing heard of it: a file that decodes to nothing.
3330        player.process_command(PlayerCommand::DecodeFinished(player.session));
3331        assert!(matches!(player.transport, Transport::Idle));
3332    }
3333
3334    #[test]
3335    fn shuffle_on_and_off_again_puts_the_queue_back() {
3336        let mut player = Player::new();
3337        let ids = seed(&mut player, 20);
3338        pretend_playing(&mut player, ids[3]);
3339
3340        player.process_command(PlayerCommand::SetShuffle(true));
3341        let shuffled = playlist_ids(&player);
3342        assert!(player.shared_state.play_mode().shuffle);
3343        assert_eq!(
3344            shuffled[..4],
3345            ids[..4],
3346            "nothing up to the playing track moves"
3347        );
3348        assert_ne!(shuffled, ids);
3349        let mut sorted = shuffled.clone();
3350        sorted.sort_by_key(|id| ids.iter().position(|i| i == id));
3351        assert_eq!(sorted, ids, "the same items");
3352
3353        let extra = make_item("extra");
3354        let extra_id = extra.id;
3355        player.process_command(PlayerCommand::InsertInPlaylist {
3356            items: vec![extra],
3357            after: ids[3],
3358        });
3359        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[10]));
3360        player.process_command(PlayerCommand::SetShuffle(false));
3361
3362        let mut expected = ids.clone();
3363        expected.remove(10);
3364        expected.insert(4, extra_id);
3365        assert_eq!(playlist_ids(&player), expected, "added since stays put");
3366        assert!(!player.shared_state.play_mode().shuffle);
3367        assert!(
3368            player
3369                .shared_state
3370                .shuffle_order()
3371                .iter()
3372                .all(|(_, pre)| pre.is_none())
3373        );
3374    }
3375
3376    #[test]
3377    fn a_queue_replaced_while_shuffled_plays_shuffled_from_its_start() {
3378        let mut player = Player::new();
3379        seed(&mut player, 3);
3380        player.process_command(PlayerCommand::SetShuffle(true));
3381
3382        let items: Vec<_> = (0..20).map(|i| make_item(&format!("n{i}"))).collect();
3383        let given: Vec<_> = items.iter().map(|i| i.id).collect();
3384        player.process_command(PlayerCommand::ReplacePlaylist {
3385            items,
3386            start: 5,
3387            position_ms: 0,
3388            play: false,
3389        });
3390        let shuffled = playlist_ids(&player);
3391        assert_eq!(shuffled[0], given[5], "the start first");
3392        assert_eq!(player.shared_state.cursor(), Some(given[5]));
3393        assert_ne!(shuffled, given);
3394
3395        player.process_command(PlayerCommand::SetShuffle(false));
3396        assert_eq!(
3397            playlist_ids(&player),
3398            given,
3399            "off gives the queue as it came"
3400        );
3401    }
3402
3403    #[test]
3404    fn a_queue_added_to_an_empty_one_while_shuffled_plays_shuffled() {
3405        let mut player = Player::new();
3406        player.process_command(PlayerCommand::SetShuffle(true));
3407        let given = seed(&mut player, 20);
3408        assert_eq!(playlist_ids(&player)[0], given[0]);
3409        assert_ne!(playlist_ids(&player), given);
3410
3411        player.process_command(PlayerCommand::SetShuffle(false));
3412        assert_eq!(playlist_ids(&player), given);
3413    }
3414
3415    #[test]
3416    fn shuffle_is_one_undo_step() {
3417        let mut player = Player::new();
3418        let ids = seed(&mut player, 10);
3419        pretend_playing(&mut player, ids[0]);
3420
3421        player.process_command(PlayerCommand::SetShuffle(true));
3422        let shuffled = playlist_ids(&player);
3423        player.process_command(PlayerCommand::Undo);
3424        assert_eq!(playlist_ids(&player), ids);
3425        assert!(!player.shared_state.play_mode().shuffle);
3426
3427        player.process_command(PlayerCommand::Redo);
3428        assert_eq!(playlist_ids(&player), shuffled);
3429        assert!(player.shared_state.play_mode().shuffle);
3430        player.process_command(PlayerCommand::SetShuffle(false));
3431        assert_eq!(
3432            playlist_ids(&player),
3433            ids,
3434            "the redone shuffle still unwinds"
3435        );
3436    }
3437
3438    #[test]
3439    fn removing_the_last_track_playing_carries_on_from_the_top_when_repeating() {
3440        let mut player = Player::new();
3441        let ids = seed(&mut player, 3);
3442        player.process_command(PlayerCommand::SetRepeat(Repeat::Queue));
3443        player.shared_state.set_cursor(Some(ids[2]));
3444
3445        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
3446        assert_eq!(player.shared_state.cursor(), Some(ids[0]));
3447    }
3448
3449    #[test]
3450    fn removing_a_track_the_decoder_queued_takes_it_back() {
3451        let dir = tempfile::tempdir().unwrap();
3452        let (mut player, ids) = queued_wavs(dir.path(), &["a", "b", "c"]);
3453        assert_eq!(player.timeline.queued_after_playhead(), ids[1..]);
3454
3455        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[1]));
3456        assert_eq!(player.playback_starts, 2, "restarted at the playhead");
3457        assert_eq!(player.shared_state.cursor(), Some(ids[0]));
3458        await_queued(&player);
3459        await_queued(&player);
3460        assert_eq!(player.timeline.queued_after_playhead(), vec![ids[2]]);
3461        player.process_command(PlayerCommand::Stop);
3462    }
3463
3464    #[test]
3465    fn a_track_inserted_to_play_next_is_not_skipped() {
3466        let dir = tempfile::tempdir().unwrap();
3467        let (mut player, ids) = queued_wavs(dir.path(), &["a", "b"]);
3468        let path = dir.path().join("next.wav");
3469        crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
3470        let next = PlaylistItem {
3471            path,
3472            ..make_item("next")
3473        };
3474        let next_id = next.id;
3475
3476        player.process_command(PlayerCommand::InsertInPlaylist {
3477            items: vec![next],
3478            after: ids[0],
3479        });
3480        assert_eq!(player.playback_starts, 2);
3481        for _ in 0..3 {
3482            await_queued(&player);
3483        }
3484        assert_eq!(
3485            player.timeline.queued_after_playhead(),
3486            vec![next_id, ids[1]]
3487        );
3488        player.process_command(PlayerCommand::Stop);
3489    }
3490
3491    #[test]
3492    fn a_track_the_decoder_could_not_open_does_not_make_every_edit_restart() {
3493        let dir = tempfile::tempdir().unwrap();
3494        let mut player = Player::new();
3495        player.backend = Box::new(StuckBackend {
3496            rate: 8_000.0,
3497            asked: Default::default(),
3498            starts: Default::default(),
3499        });
3500        let wav = |name: &str| {
3501            let path = dir.path().join(format!("{name}.wav"));
3502            crate::test_utils::generate_wav(&path, 8_000, 1, 30.0, 16);
3503            PlaylistItem {
3504                path,
3505                ..make_item(name)
3506            }
3507        };
3508        // Ready, but with no file behind it: the decoder skips it.
3509        let items = vec![wav("a"), make_item("missing"), wav("c")];
3510        let ids: Vec<_> = items.iter().map(|i| i.id).collect();
3511        player.process_command(PlayerCommand::AddToPlaylist(items));
3512        player.process_command(PlayerCommand::Play(ids[0]));
3513        await_queued(&player);
3514        await_queued(&player);
3515
3516        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("later")]));
3517        assert_eq!(
3518            player.playback_starts, 1,
3519            "nothing the decoder decided changed"
3520        );
3521        assert_eq!(player.timeline.queued_after_playhead(), vec![ids[2]]);
3522        player.process_command(PlayerCommand::Stop);
3523    }
3524
3525    #[test]
3526    fn a_track_added_after_the_decoder_reached_the_end_follows_gaplessly() {
3527        let dir = tempfile::tempdir().unwrap();
3528        let (mut player, ids) = queued_wavs(dir.path(), &["a"]);
3529        wait_for_lookahead(&player);
3530        let path = dir.path().join("b.wav");
3531        crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
3532        let b = PlaylistItem {
3533            path,
3534            ..make_item("b")
3535        };
3536        let b_id = b.id;
3537
3538        player.process_command(PlayerCommand::AddToPlaylist(vec![b]));
3539        assert_eq!(player.playback_starts, 2, "restarted to queue it");
3540        assert_eq!(player.shared_state.cursor(), Some(ids[0]));
3541        await_queued(&player);
3542        await_queued(&player);
3543        assert_eq!(player.timeline.queued_after_playhead(), vec![b_id]);
3544        player.process_command(PlayerCommand::Stop);
3545    }
3546
3547    #[test]
3548    fn playing_next_a_track_still_downloading_takes_back_the_lookahead() {
3549        let dir = tempfile::tempdir().unwrap();
3550        let (mut player, ids) = queued_wavs(dir.path(), &["a", "b"]);
3551        let remote = pending_item("remote");
3552        let remote_id = remote.id;
3553
3554        player.process_command(PlayerCommand::InsertInPlaylist {
3555            items: vec![remote],
3556            after: ids[0],
3557        });
3558        assert_eq!(player.playback_starts, 2, "b is taken back out of the ring");
3559        await_queued(&player);
3560        wait_for_lookahead(&player);
3561        let session = player.session().unwrap();
3562        let step = session.lookahead.lock().last().cloned().unwrap();
3563        assert_eq!(step.next, Some(remote_id), "the decoder waits for it");
3564        assert!(player.timeline.queued_after_playhead().is_empty());
3565        player.process_command(PlayerCommand::Stop);
3566    }
3567
3568    #[test]
3569    fn gapless_playback_waits_for_a_track_still_downloading() {
3570        let dir = tempfile::tempdir().unwrap();
3571        let mut player = Player::new();
3572        player.backend = Box::new(StuckBackend {
3573            rate: 8_000.0,
3574            asked: Default::default(),
3575            starts: Default::default(),
3576        });
3577        let wav = |name: &str| {
3578            let path = dir.path().join(format!("{name}.wav"));
3579            crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
3580            PlaylistItem {
3581                path,
3582                ..make_item(name)
3583            }
3584        };
3585        let items = vec![wav("a"), pending_item("arriving"), wav("c")];
3586        let ids: Vec<_> = items.iter().map(|i| i.id).collect();
3587        player.process_command(PlayerCommand::AddToPlaylist(items));
3588        player.process_command(PlayerCommand::Play(ids[0]));
3589        await_queued(&player);
3590        wait_for_lookahead(&player);
3591
3592        assert!(
3593            player.timeline.queued_after_playhead().is_empty(),
3594            "c is not queued over the track before it"
3595        );
3596        // The session drains and ends; the advance parks on the track.
3597        player.process_command(PlayerCommand::DecodeFinished(player.session));
3598        assert_eq!(player.shared_state.cursor(), Some(ids[1]));
3599        assert!(player.waiting().is_some());
3600        player.process_command(PlayerCommand::Stop);
3601    }
3602
3603    /// Until the decoder has taken a step past the playing track.
3604    fn wait_for_lookahead(player: &Player) {
3605        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
3606        while player
3607            .session()
3608            .is_some_and(|s| s.lookahead.lock().is_empty())
3609        {
3610            assert!(
3611                std::time::Instant::now() < deadline,
3612                "the decoder never looked ahead"
3613            );
3614            std::thread::yield_now();
3615        }
3616    }
3617
3618    #[test]
3619    fn an_edit_after_what_the_decoder_queued_leaves_playback_alone() {
3620        let dir = tempfile::tempdir().unwrap();
3621        let (mut player, ids) = wavs_playing(dir.path(), &["a", "b"], 30.0);
3622        await_queued(&player);
3623        await_queued(&player);
3624        let path = dir.path().join("later.wav");
3625        crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
3626
3627        player.process_command(PlayerCommand::AddToPlaylist(vec![PlaylistItem {
3628            path,
3629            ..make_item("later")
3630        }]));
3631        assert_eq!(player.playback_starts, 1);
3632        assert_eq!(player.timeline.queued_after_playhead(), vec![ids[1]]);
3633        player.process_command(PlayerCommand::Stop);
3634    }
3635
3636    #[test]
3637    fn undoing_the_add_of_a_track_on_its_way_forgets_it() {
3638        let dir = tempfile::tempdir().unwrap();
3639        let (mut player, _, starts) = downloading_wav(dir.path());
3640        let path = dir.path().join("added.wav");
3641        crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
3642        let added = PlaylistItem {
3643            path,
3644            state: ItemState::Pending,
3645            ..make_item("added")
3646        };
3647        let added_id = added.id;
3648        player.process_command(PlayerCommand::AddToPlaylist(vec![added]));
3649        player.process_command(PlayerCommand::Play(added_id));
3650
3651        player.process_command(PlayerCommand::Undo);
3652        assert!(player.waiting().is_none());
3653        assert_eq!(player.shared_state.playback_state(), PlaybackState::Stopped);
3654
3655        // Its download lands anyway; nothing asked for it any more.
3656        player.process_command(PlayerCommand::Redo);
3657        player
3658            .shared_state
3659            .update_item_state(added_id, ItemState::Ready);
3660        player.process_command(PlayerCommand::TrackReady(added_id));
3661        assert!(player.session().is_none());
3662        assert_eq!(starts.load(Ordering::Relaxed), 0);
3663    }
3664
3665    /// xorshift64: enough randomness to drive the player, reproducible from
3666    /// its seed, and no dependency for it.
3667    pub(super) struct Rng(pub(super) u64);
3668
3669    impl Rng {
3670        pub(super) fn next(&mut self) -> u64 {
3671            self.0 ^= self.0 << 13;
3672            self.0 ^= self.0 >> 7;
3673            self.0 ^= self.0 << 17;
3674            self.0
3675        }
3676
3677        pub(super) fn below(&mut self, n: usize) -> usize {
3678            (self.next() % n.max(1) as u64) as usize
3679        }
3680
3681        pub(super) fn coin(&mut self) -> bool {
3682            self.next() & 1 == 1
3683        }
3684    }
3685
3686    /// Commands a listener can send that ask for sound. Anything else may only
3687    /// keep playing what was already playing or on its way to.
3688    pub(super) fn asks_to_play(cmd: &PlayerCommand) -> bool {
3689        cmd.asks_to_play()
3690    }
3691
3692    /// Check the invariants #679 sets out, as far as they can be seen from
3693    /// outside the audio thread.
3694    pub(super) fn check_invariants(
3695        player: &Player,
3696        wanted_before: bool,
3697        asked: bool,
3698    ) -> Result<(), String> {
3699        let state = &player.shared_state;
3700        let ids = playlist_ids(player);
3701        let listed = |id: QueueItemId| ids.contains(&id);
3702        let mut seen = std::collections::HashSet::new();
3703        if !ids.iter().all(|id| seen.insert(*id)) {
3704            return Err("an item is in the playlist twice".into());
3705        }
3706        let cursor = state.cursor();
3707        if cursor.is_some_and(|c| !listed(c)) {
3708            return Err("the cursor is on an item not in the playlist".into());
3709        }
3710        let info = state.track_info();
3711        match (&player.transport, &info) {
3712            (Transport::Loaded(_), None) => return Err("loaded, but no track published".into()),
3713            (Transport::Loaded(_), Some(info)) => {
3714                if !listed(info.id) {
3715                    return Err("playing an item not in the playlist".into());
3716                }
3717                if cursor != Some(info.id) {
3718                    return Err("the cursor is not on what is playing".into());
3719                }
3720                if !matches!(
3721                    state.playback_state(),
3722                    PlaybackState::Playing | PlaybackState::Paused
3723                ) {
3724                    return Err("loaded, but published as stopped".into());
3725                }
3726            }
3727            (_, Some(_)) => return Err("a track is published with nothing loaded".into()),
3728            (Transport::Waiting(w), None) => {
3729                if cursor != Some(w.id) || !listed(w.id) {
3730                    return Err("waiting on an item that is not the cursor's".into());
3731                }
3732                let expected = match w.start {
3733                    Run::Playing => PlaybackState::Stopped,
3734                    Run::Paused => PlaybackState::Paused,
3735                };
3736                if state.playback_state() != expected {
3737                    return Err("a wait published in the wrong state".into());
3738                }
3739            }
3740            (Transport::Idle, None) => {
3741                if state.playback_state() != PlaybackState::Stopped {
3742                    return Err("idle, but not published as stopped".into());
3743                }
3744            }
3745        }
3746        if state.is_waiting() != player.waiting().is_some() {
3747            return Err("the published wait disagrees with the player".into());
3748        }
3749        if player
3750            .timeline
3751            .queued_after_playhead()
3752            .into_iter()
3753            .any(|id| !listed(id))
3754        {
3755            return Err("the decoder has queued an item no longer in the playlist".into());
3756        }
3757        if !asked && !wanted_before && state.wants_to_play() {
3758            return Err("started without being asked".into());
3759        }
3760        if state.play_mode() != player.mode {
3761            return Err("the published mode disagrees with the player".into());
3762        }
3763        if !player.mode.shuffle && state.shuffle_order().iter().any(|(_, pre)| pre.is_some()) {
3764            return Err("shuffle is off, but an item remembers a place to go back to".into());
3765        }
3766        if player.mode.repeat == Repeat::Off
3767            && let Some(session) = player.session()
3768            && let Some(playhead) = player.timeline.playhead()
3769            && session.track.id == playhead.id
3770            && session
3771                .lookahead
3772                .lock()
3773                .iter()
3774                .filter(|step| step.boundary > playhead.boundary)
3775                .any(|step| step.wrapped || step.next == Some(step.after))
3776        {
3777            return Err("repeat is off, but the decoder has queued a wrap".into());
3778        }
3779        Ok(())
3780    }
3781
3782    #[test]
3783    fn random_use_keeps_the_player_honest() {
3784        let dir = tempfile::tempdir().unwrap();
3785        let path = dir.path().join("t.wav");
3786        crate::test_utils::generate_wav(&path, 8_000, 1, 0.5, 16);
3787
3788        for seed in 1..=40u64 {
3789            let mut rng = Rng(seed.wrapping_mul(0x9E37_79B9_7F4A_7C15) | 1);
3790            let mut player = Player::new();
3791            player.backend = Box::new(StuckBackend {
3792                rate: 8_000.0,
3793                asked: Default::default(),
3794                starts: Default::default(),
3795            });
3796            let fresh = |rng: &mut Rng| {
3797                let state = match rng.below(3) {
3798                    0 => ItemState::Pending,
3799                    _ => ItemState::Ready,
3800                };
3801                PlaylistItem {
3802                    path: path.clone(),
3803                    state,
3804                    ..make_item("t")
3805                }
3806            };
3807            let start: Vec<_> = (0..5).map(|_| fresh(&mut rng)).collect();
3808            player.process_command(PlayerCommand::AddToPlaylist(start));
3809
3810            let mut history = Vec::new();
3811            for step in 0..150 {
3812                let ids = playlist_ids(&player);
3813                let pick = |rng: &mut Rng| ids.get(rng.below(ids.len())).copied();
3814                let cmd = match rng.below(23) {
3815                    0 => pick(&mut rng).map(PlayerCommand::Play),
3816                    1 => pick(&mut rng).map(|id| PlayerCommand::Cue {
3817                        id,
3818                        position_ms: if rng.coin() { 0 } else { 200 },
3819                        play: rng.coin(),
3820                    }),
3821                    2 => Some(PlayerCommand::Pause),
3822                    3 => Some(PlayerCommand::Resume),
3823                    4 => Some(PlayerCommand::NextTrack),
3824                    5 => Some(PlayerCommand::PrevTrack),
3825                    6 => Some(PlayerCommand::Seek(100)),
3826                    7 => pick(&mut rng).map(PlayerCommand::RemoveFromPlaylist),
3827                    8 => Some(PlayerCommand::RemoveFromPlaylistBatch(
3828                        (0..2).filter_map(|_| pick(&mut rng)).collect(),
3829                    )),
3830                    9 => pick(&mut rng).zip(pick(&mut rng)).map(|(id, target)| {
3831                        PlayerCommand::MoveInPlaylist {
3832                            id,
3833                            target,
3834                            after: rng.coin(),
3835                        }
3836                    }),
3837                    10 => pick(&mut rng).map(|after| PlayerCommand::InsertInPlaylist {
3838                        items: vec![fresh(&mut rng)],
3839                        after,
3840                    }),
3841                    11 => Some(PlayerCommand::AddToPlaylist(vec![fresh(&mut rng)])),
3842                    12 => Some(PlayerCommand::Undo),
3843                    13 => Some(PlayerCommand::Redo),
3844                    14 | 15 => {
3845                        let (items, _) = player.shared_state.snapshot_playlist();
3846                        let pending: Vec<_> = items
3847                            .iter()
3848                            .filter(|i| matches!(i.state, ItemState::Pending))
3849                            .map(|i| i.id)
3850                            .collect();
3851                        pending.get(rng.below(pending.len())).map(|&id| {
3852                            if rng.below(4) == 0 {
3853                                player
3854                                    .shared_state
3855                                    .update_item_state(id, ItemState::Failed("gone".into()));
3856                                PlayerCommand::TrackFailed(id)
3857                            } else {
3858                                player.shared_state.update_item_state(id, ItemState::Ready);
3859                                PlayerCommand::TrackReady(id)
3860                            }
3861                        })
3862                    }
3863                    16 => Some(PlayerCommand::DecodeFinished(player.session)),
3864                    17 => Some(PlayerCommand::DecodeFinished(
3865                        player.session.wrapping_sub(1),
3866                    )),
3867                    18 => Some(PlayerCommand::ClearPlaylist),
3868                    20 => Some(PlayerCommand::SetShuffle(rng.coin())),
3869                    21 => Some(PlayerCommand::SetRepeat(
3870                        [Repeat::Off, Repeat::Queue, Repeat::One][rng.below(3)],
3871                    )),
3872                    _ => Some(PlayerCommand::ReplacePlaylist {
3873                        items: (0..3).map(|_| fresh(&mut rng)).collect(),
3874                        start: rng.below(4),
3875                        position_ms: if rng.coin() { 0 } else { 200 },
3876                        play: rng.coin(),
3877                    }),
3878                };
3879                let Some(cmd) = cmd else { continue };
3880                let label = format!("{cmd:?}");
3881                let asked = asks_to_play(&cmd);
3882                let replaced_shuffled = player.mode.shuffle
3883                    && matches!(&cmd, PlayerCommand::ReplacePlaylist { items, .. } if items.len() > 1);
3884                let wanted_before = player.shared_state.wants_to_play();
3885                player.process_command(cmd);
3886                // What the decode threads sent meanwhile, as the loop would
3887                // see it. A full channel would stall a decode thread on its
3888                // way out.
3889                while let Ok(sent) = player.commands.rx.try_recv() {
3890                    player.process_command(sent);
3891                }
3892                player.update_playback_state();
3893                history.push(label);
3894                let broken = check_invariants(&player, wanted_before, asked)
3895                    .err()
3896                    .or_else(|| {
3897                        (replaced_shuffled
3898                            && player
3899                                .shared_state
3900                                .shuffle_order()
3901                                .iter()
3902                                .all(|(_, pre)| pre.is_none()))
3903                        .then(|| "a queue replaced while shuffled plays in order".to_string())
3904                    });
3905                if let Some(broken) = broken {
3906                    let tail = history[history.len().saturating_sub(8)..].join("\n  ");
3907                    panic!("seed {seed}, step {step}: {broken}\nlast commands:\n  {tail}");
3908                }
3909            }
3910            player.process_command(PlayerCommand::Stop);
3911        }
3912    }
3913
3914    #[test]
3915    fn a_pause_reports_where_the_fade_went_silent() {
3916        use std::sync::atomic::AtomicBool;
3917
3918        struct FadingEngine {
3919            running: Arc<AtomicBool>,
3920            silent: Arc<AtomicBool>,
3921        }
3922        impl AudioEngineHandle for FadingEngine {
3923            fn start(&self) -> Result<(), BackendError> {
3924                self.running.store(true, Ordering::Relaxed);
3925                Ok(())
3926            }
3927            fn stop(&self) -> Result<(), BackendError> {
3928                self.running.store(false, Ordering::Relaxed);
3929                Ok(())
3930            }
3931            fn is_running(&self) -> bool {
3932                self.running.load(Ordering::Relaxed)
3933            }
3934            fn fade_out(&self) {}
3935            fn fade_in(&self) -> Result<(), BackendError> {
3936                Ok(())
3937            }
3938            fn is_silent(&self) -> bool {
3939                self.silent.load(Ordering::Relaxed)
3940            }
3941        }
3942
3943        let running = Arc::new(AtomicBool::new(true));
3944        let silent = Arc::new(AtomicBool::new(false));
3945        let mut player = Player::new();
3946        player.transport = Transport::Loaded(test_session(
3947            QueueItemId::new(),
3948            Box::new(FadingEngine {
3949                running: running.clone(),
3950                silent: silent.clone(),
3951            }),
3952        ));
3953        player
3954            .shared_state
3955            .set_playback_state(PlaybackState::Playing);
3956        player.shared_state.set_position_ms(5_000);
3957
3958        let (reply, answer) = crossbeam_channel::bounded(1);
3959        player.process_command(PlayerCommand::PauseAndReport(reply));
3960
3961        if crate::config::Config::cached().playback.fade_on_pause {
3962            assert!(answer.try_recv().is_err(), "not while the fade is audible");
3963            // The fade plays on, and the playhead with it.
3964            player.shared_state.set_position_ms(5_150);
3965            player.update_playback_state();
3966            assert!(answer.try_recv().is_err());
3967            silent.store(true, Ordering::Relaxed);
3968            player.update_playback_state();
3969            assert!(!running.load(Ordering::Relaxed));
3970            assert_eq!(answer.try_recv().unwrap(), 5_150);
3971        } else {
3972            assert_eq!(answer.try_recv().unwrap(), 5_000);
3973        }
3974    }
3975
3976    #[test]
3977    fn the_server_hears_each_turn_playback_takes() {
3978        use PlaybackReportState::{Paused, Playing, Stopped};
3979        use history::PlaybackReport;
3980
3981        let dir = tempfile::tempdir().unwrap();
3982        let path = dir.path().join("t.wav");
3983        crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
3984
3985        let mut player = Player::new();
3986        player.backend = Box::new(StuckBackend {
3987            rate: 8_000.0,
3988            asked: Default::default(),
3989            starts: Default::default(),
3990        });
3991        let (recorder, events) = PlayRecorder::capture();
3992        player.history = Some(recorder);
3993
3994        let item = PlaylistItem {
3995            db_id: Some(5),
3996            path,
3997            ..make_item("t")
3998        };
3999        let id = item.id;
4000        player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
4001        player.process_command(PlayerCommand::Play(id));
4002        player.process_command(PlayerCommand::Pause);
4003        player.process_command(PlayerCommand::Seek(4_000));
4004        player.process_command(PlayerCommand::Resume);
4005        player.process_command(PlayerCommand::Stop);
4006
4007        let report = |state, position_ms| {
4008            PlayEvent::Playback(PlaybackReport {
4009                track_id: 5,
4010                state,
4011                position_ms,
4012            })
4013        };
4014        assert_eq!(
4015            events.try_iter().collect::<Vec<_>>(),
4016            vec![
4017                PlayEvent::Started {
4018                    track_id: 5,
4019                    position_ms: 0
4020                },
4021                report(Paused, 0),
4022                report(Paused, 4_000),
4023                report(Playing, 4_000),
4024                report(Stopped, 4_000),
4025                PlayEvent::Finished {
4026                    track_id: 5,
4027                    listened_ms: 0
4028                },
4029            ]
4030        );
4031    }
4032
4033    #[test]
4034    fn removing_the_playing_track_resumes_at_its_successor() {
4035        let mut player = Player::new();
4036        let ids = seed(&mut player, 5);
4037        pretend_playing(&mut player, ids[2]);
4038
4039        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
4040
4041        assert_eq!(
4042            player.shared_state.cursor(),
4043            Some(ids[3]),
4044            "playback must continue at the next track, not restart the queue"
4045        );
4046        assert_eq!(player.playback_starts, 1);
4047    }
4048
4049    #[test]
4050    fn removing_the_paused_track_moves_on_paused() {
4051        let mut player = Player::new();
4052        let ids = seed(&mut player, 3);
4053        pretend_playing(&mut player, ids[1]);
4054        if let Transport::Loaded(session) = &mut player.transport {
4055            session.run = Run::Paused;
4056        }
4057
4058        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[1]));
4059        assert_eq!(player.shared_state.cursor(), Some(ids[2]));
4060        assert!(!player.shared_state.wants_to_play(), "still paused");
4061    }
4062
4063    #[test]
4064    fn removing_the_cursor_with_nothing_loaded_starts_nothing() {
4065        let mut player = Player::new();
4066        let ids = seed(&mut player, 3);
4067        player.shared_state.set_cursor(Some(ids[1]));
4068
4069        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[1]));
4070        assert_eq!(player.shared_state.cursor(), Some(ids[2]));
4071        assert_eq!(player.playback_starts, 0);
4072    }
4073
4074    #[test]
4075    fn removing_the_first_playing_track_resumes_at_the_new_first() {
4076        let mut player = Player::new();
4077        let ids = seed(&mut player, 3);
4078        pretend_playing(&mut player, ids[0]);
4079
4080        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[0]));
4081
4082        assert_eq!(player.shared_state.cursor(), Some(ids[1]));
4083    }
4084
4085    #[test]
4086    fn next_track_parks_on_a_track_that_has_not_downloaded_yet() {
4087        let mut player = Player::new();
4088        let playing = make_item("playing");
4089        let waiting = pending_item("waiting");
4090        let later = make_item("later");
4091        let (playing_id, waiting_id) = (playing.id, waiting.id);
4092        player.process_command(PlayerCommand::AddToPlaylist(vec![playing, waiting, later]));
4093        pretend_playing(&mut player, playing_id);
4094
4095        player.process_command(PlayerCommand::DecodeFinished(player.session));
4096
4097        assert_eq!(
4098            player.shared_state.cursor(),
4099            Some(waiting_id),
4100            "the cursor parks on the track being fetched"
4101        );
4102        assert_eq!(
4103            player.playback_starts, 0,
4104            "nothing to play until its bytes land"
4105        );
4106
4107        // The download completes. Because the cursor is parked here, the
4108        // TrackReady actually reaches the player and the queue resumes.
4109        player
4110            .shared_state
4111            .update_item_state(waiting_id, ItemState::Ready);
4112        player.process_command(PlayerCommand::TrackReady(waiting_id));
4113
4114        assert_eq!(player.playback_starts, 1);
4115        assert_eq!(player.shared_state.cursor(), Some(waiting_id));
4116    }
4117
4118    #[test]
4119    fn a_download_that_cannot_land_moves_the_cursor_on() {
4120        let mut player = Player::new();
4121        let waiting = pending_item("waiting");
4122        let later = make_item("later");
4123        let (waiting_id, later_id) = (waiting.id, later.id);
4124        player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, later]));
4125
4126        player.process_command(PlayerCommand::Play(waiting_id));
4127        assert_eq!(player.playback_starts, 0, "nothing to play yet");
4128
4129        // The download gives up. Ready will never come.
4130        player
4131            .shared_state
4132            .update_item_state(waiting_id, ItemState::Failed("remote unavailable".into()));
4133        player.process_command(PlayerCommand::TrackFailed(waiting_id));
4134
4135        assert_eq!(
4136            player.shared_state.cursor(),
4137            Some(later_id),
4138            "the queue moves past a track that can never load"
4139        );
4140        assert_eq!(player.playback_starts, 1);
4141    }
4142
4143    #[test]
4144    fn a_queue_that_can_never_load_stops_rather_than_waiting() {
4145        let mut player = Player::new();
4146        let first = pending_item("first");
4147        let second = pending_item("second");
4148        let (first_id, second_id) = (first.id, second.id);
4149        player.process_command(PlayerCommand::AddToPlaylist(vec![first, second]));
4150
4151        player.process_command(PlayerCommand::Play(first_id));
4152        for id in [first_id, second_id] {
4153            player
4154                .shared_state
4155                .update_item_state(id, ItemState::Failed("remote unavailable".into()));
4156            player.process_command(PlayerCommand::TrackFailed(id));
4157        }
4158
4159        assert_eq!(player.playback_starts, 0);
4160        assert_eq!(
4161            player.shared_state.playback_state(),
4162            PlaybackState::Stopped,
4163            "a stop the UI can see, not an indefinite wait for TrackReady"
4164        );
4165    }
4166
4167    #[test]
4168    fn a_failure_elsewhere_in_the_queue_leaves_the_cursor_alone() {
4169        let mut player = Player::new();
4170        let waiting = pending_item("waiting");
4171        let other = pending_item("other");
4172        let (waiting_id, other_id) = (waiting.id, other.id);
4173        player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, other]));
4174        player.process_command(PlayerCommand::Play(waiting_id));
4175
4176        player
4177            .shared_state
4178            .update_item_state(other_id, ItemState::Failed("remote unavailable".into()));
4179        player.process_command(PlayerCommand::TrackFailed(other_id));
4180
4181        assert_eq!(
4182            player.shared_state.cursor(),
4183            Some(waiting_id),
4184            "a track still downloading keeps the cursor"
4185        );
4186    }
4187
4188    #[test]
4189    fn batch_delete_containing_the_cursor_restarts_the_engine_once() {
4190        let mut player = Player::new();
4191        let ids = seed(&mut player, 5);
4192        pretend_playing(&mut player, ids[2]);
4193
4194        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![
4195            ids[1], ids[2], ids[3],
4196        ]));
4197
4198        assert_eq!(playlist_titles(&player), vec!["t0", "t4"]);
4199        assert_eq!(player.shared_state.cursor(), Some(ids[4]));
4200        assert_eq!(
4201            player.playback_starts, 1,
4202            "one resume for the whole selection, not one per deleted track"
4203        );
4204    }
4205
4206    #[test]
4207    fn batch_delete_below_the_cursor_leaves_playback_alone() {
4208        let mut player = Player::new();
4209        let ids = seed(&mut player, 4);
4210        player.shared_state.set_cursor(Some(ids[0]));
4211
4212        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![ids[2], ids[3]]));
4213
4214        assert_eq!(player.shared_state.cursor(), Some(ids[0]));
4215        assert_eq!(player.playback_starts, 0);
4216    }
4217
4218    #[test]
4219    fn undo_of_a_batch_delete_restores_the_original_order() {
4220        // The TUI collects a selection from a HashSet, so the IDs arrive in
4221        // arbitrary order — scrambled here so a snapshot that trusts that order
4222        // re-inserts C before B and lands it at the end of the playlist.
4223        let mut player = Player::new();
4224        let items = vec![
4225            make_item("A"),
4226            make_item("B"),
4227            make_item("C"),
4228            make_item("D"),
4229        ];
4230        let (b_id, c_id) = (items[1].id, items[2].id);
4231        player.process_command(PlayerCommand::AddToPlaylist(items));
4232
4233        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![c_id, b_id]));
4234        assert_eq!(playlist_titles(&player), vec!["A", "D"]);
4235
4236        player.process_command(PlayerCommand::Undo);
4237        assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
4238    }
4239
4240    // --- AddToPlaylist undo/redo ---
4241
4242    #[test]
4243    fn undo_add_removes_items() {
4244        let mut player = Player::new();
4245        let items = vec![make_item("A"), make_item("B")];
4246        let ids: Vec<_> = items.iter().map(|i| i.id).collect();
4247
4248        player.process_command(PlayerCommand::AddToPlaylist(items));
4249        assert_eq!(playlist_ids(&player), ids);
4250        assert!(player.undo_stack().can_undo());
4251
4252        player.process_command(PlayerCommand::Undo);
4253        assert!(playlist_ids(&player).is_empty());
4254        assert!(player.undo_stack().can_redo());
4255    }
4256
4257    #[test]
4258    fn redo_add_restores_items() {
4259        let mut player = Player::new();
4260        let items = vec![make_item("A"), make_item("B")];
4261
4262        player.process_command(PlayerCommand::AddToPlaylist(items));
4263        player.process_command(PlayerCommand::Undo);
4264        assert!(playlist_ids(&player).is_empty());
4265
4266        player.process_command(PlayerCommand::Redo);
4267        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
4268    }
4269
4270    // --- RemoveFromPlaylist undo/redo ---
4271
4272    #[test]
4273    fn undo_remove_restores_item_at_position() {
4274        let mut player = Player::new();
4275        let items = vec![make_item("A"), make_item("B"), make_item("C")];
4276        let b_id = items[1].id;
4277
4278        player.process_command(PlayerCommand::AddToPlaylist(items));
4279        player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
4280        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
4281
4282        player.process_command(PlayerCommand::Undo);
4283        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4284    }
4285
4286    #[test]
4287    fn undo_remove_first_item() {
4288        let mut player = Player::new();
4289        let items = vec![make_item("A"), make_item("B")];
4290        let a_id = items[0].id;
4291
4292        player.process_command(PlayerCommand::AddToPlaylist(items));
4293        player.process_command(PlayerCommand::RemoveFromPlaylist(a_id));
4294        assert_eq!(playlist_titles(&player), vec!["B"]);
4295
4296        player.process_command(PlayerCommand::Undo);
4297        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
4298    }
4299
4300    #[test]
4301    fn undo_batch_remove_restores_all() {
4302        let mut player = Player::new();
4303        let items = vec![
4304            make_item("A"),
4305            make_item("B"),
4306            make_item("C"),
4307            make_item("D"),
4308        ];
4309        let b_id = items[1].id;
4310        let c_id = items[2].id;
4311
4312        player.process_command(PlayerCommand::AddToPlaylist(items));
4313        let version_before = player.shared_state.playlist_version();
4314        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![b_id, c_id]));
4315        assert_eq!(playlist_titles(&player), vec!["A", "D"]);
4316        // One bump for the whole batch. Bumping per item is what made clearing
4317        // a large queue crawl, and every bump wakes every client watching.
4318        assert_eq!(
4319            player.shared_state.playlist_version(),
4320            version_before + 1,
4321            "batch removal must bump the playlist version exactly once"
4322        );
4323
4324        // Single undo restores both
4325        player.process_command(PlayerCommand::Undo);
4326        assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
4327    }
4328
4329    #[test]
4330    fn redo_batch_remove() {
4331        let mut player = Player::new();
4332        let items = vec![make_item("A"), make_item("B"), make_item("C")];
4333        let a_id = items[0].id;
4334        let b_id = items[1].id;
4335
4336        player.process_command(PlayerCommand::AddToPlaylist(items));
4337        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![a_id, b_id]));
4338        player.process_command(PlayerCommand::Undo);
4339        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4340
4341        player.process_command(PlayerCommand::Redo);
4342        assert_eq!(playlist_titles(&player), vec!["C"]);
4343    }
4344
4345    #[test]
4346    fn redo_remove() {
4347        let mut player = Player::new();
4348        let items = vec![make_item("A"), make_item("B"), make_item("C")];
4349        let b_id = items[1].id;
4350
4351        player.process_command(PlayerCommand::AddToPlaylist(items));
4352        player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
4353        player.process_command(PlayerCommand::Undo);
4354        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4355
4356        player.process_command(PlayerCommand::Redo);
4357        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
4358    }
4359
4360    // --- InsertInPlaylist undo/redo ---
4361
4362    #[test]
4363    fn undo_insert_removes_inserted_items() {
4364        let mut player = Player::new();
4365        let items = vec![make_item("A"), make_item("C")];
4366        let a_id = items[0].id;
4367
4368        player.process_command(PlayerCommand::AddToPlaylist(items));
4369
4370        let inserted = vec![make_item("B")];
4371        player.process_command(PlayerCommand::InsertInPlaylist {
4372            items: inserted,
4373            after: a_id,
4374        });
4375        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4376
4377        player.process_command(PlayerCommand::Undo);
4378        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
4379    }
4380
4381    // --- MoveInPlaylist undo/redo ---
4382
4383    #[test]
4384    fn undo_move_restores_position() {
4385        let mut player = Player::new();
4386        let items = vec![make_item("A"), make_item("B"), make_item("C")];
4387        let a_id = items[0].id;
4388        let c_id = items[2].id;
4389
4390        player.process_command(PlayerCommand::AddToPlaylist(items));
4391
4392        // Move A after C: [B, C, A]
4393        player.process_command(PlayerCommand::MoveInPlaylist {
4394            id: a_id,
4395            target: c_id,
4396            after: true,
4397        });
4398        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
4399
4400        player.process_command(PlayerCommand::Undo);
4401        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4402    }
4403
4404    #[test]
4405    fn redo_move() {
4406        let mut player = Player::new();
4407        let items = vec![make_item("A"), make_item("B"), make_item("C")];
4408        let a_id = items[0].id;
4409        let c_id = items[2].id;
4410
4411        player.process_command(PlayerCommand::AddToPlaylist(items));
4412        player.process_command(PlayerCommand::MoveInPlaylist {
4413            id: a_id,
4414            target: c_id,
4415            after: true,
4416        });
4417        player.process_command(PlayerCommand::Undo);
4418        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4419
4420        player.process_command(PlayerCommand::Redo);
4421        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
4422    }
4423
4424    // --- MoveItemsInPlaylist (batch) undo/redo ---
4425
4426    #[test]
4427    fn undo_batch_move() {
4428        let mut player = Player::new();
4429        let items = vec![
4430            make_item("A"),
4431            make_item("B"),
4432            make_item("C"),
4433            make_item("D"),
4434        ];
4435        let a_id = items[0].id;
4436        let b_id = items[1].id;
4437        let d_id = items[3].id;
4438
4439        player.process_command(PlayerCommand::AddToPlaylist(items));
4440
4441        // Move A,B after D: [C, D, A, B]
4442        player.process_command(PlayerCommand::MoveItemsInPlaylist {
4443            ids: vec![a_id, b_id],
4444            target: d_id,
4445            after: true,
4446        });
4447        assert_eq!(playlist_titles(&player), vec!["C", "D", "A", "B"]);
4448
4449        player.process_command(PlayerCommand::Undo);
4450        assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
4451    }
4452
4453    // --- ClearPlaylist undo/redo ---
4454
4455    #[test]
4456    fn undo_clear_restores_playlist() {
4457        let mut player = Player::new();
4458        let items = vec![make_item("A"), make_item("B"), make_item("C")];
4459
4460        player.process_command(PlayerCommand::AddToPlaylist(items));
4461        player.process_command(PlayerCommand::ClearPlaylist);
4462        assert!(playlist_ids(&player).is_empty());
4463
4464        player.process_command(PlayerCommand::Undo);
4465        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4466    }
4467
4468    /// The bug: replacing the queue starts the new track, and undoing restored
4469    /// the old queue while leaving the engine on a track that queue no longer
4470    /// contains — a transport describing a row nobody can see, and a decode
4471    /// lookahead with nothing to follow.
4472    #[test]
4473    fn undoing_a_replace_does_not_leave_the_engine_on_an_orphaned_track() {
4474        let mut player = Player::new();
4475        let original = seed(&mut player, 3);
4476        player.shared_state.set_cursor(Some(original[0]));
4477        pretend_playing(&mut player, original[0]);
4478
4479        let replacement = vec![make_item("something else")];
4480        let orphan = replacement[0].id;
4481        player.process_command(PlayerCommand::ReplacePlaylist {
4482            items: replacement,
4483            start: 0,
4484            position_ms: 0,
4485            play: true,
4486        });
4487        // What `play()` would have left behind if the file existed.
4488        pretend_playing(&mut player, orphan);
4489
4490        player.process_command(PlayerCommand::Undo);
4491
4492        assert_eq!(playlist_ids(&player), original, "the queue comes back");
4493        assert!(
4494            player.shared_state.get_item(orphan).is_none(),
4495            "and the replacement is gone from it"
4496        );
4497        assert!(
4498            playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
4499            "so nothing may still be playing out of it"
4500        );
4501    }
4502
4503    /// The same orphaning, reached by undoing an add rather than a replace.
4504    #[test]
4505    fn undoing_an_add_does_not_leave_the_engine_on_a_removed_track() {
4506        let mut player = Player::new();
4507        seed(&mut player, 2);
4508        let added = seed(&mut player, 1);
4509        pretend_playing(&mut player, added[0]);
4510
4511        player.process_command(PlayerCommand::Undo);
4512
4513        assert!(player.shared_state.get_item(added[0]).is_none());
4514        assert!(
4515            playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
4516            "the engine cannot be left on the item the undo removed"
4517        );
4518    }
4519
4520    /// An undo that leaves the playing item where it is must not restart it.
4521    #[test]
4522    fn undoing_a_move_leaves_playback_alone() {
4523        let mut player = Player::new();
4524        let ids = seed(&mut player, 3);
4525        player.shared_state.set_cursor(Some(ids[0]));
4526        pretend_playing(&mut player, ids[0]);
4527        let starts = player.playback_starts;
4528
4529        player.process_command(PlayerCommand::MoveInPlaylist {
4530            id: ids[2],
4531            target: ids[0],
4532            after: false,
4533        });
4534        player.process_command(PlayerCommand::Undo);
4535
4536        assert_eq!(playlist_ids(&player), ids);
4537        assert_eq!(playing_id(&player), Some(ids[0]), "still on the same track");
4538        assert_eq!(player.playback_starts, starts, "and not restarted");
4539    }
4540
4541    #[test]
4542    fn redo_clear() {
4543        let mut player = Player::new();
4544        let items = vec![make_item("A"), make_item("B")];
4545
4546        player.process_command(PlayerCommand::AddToPlaylist(items));
4547        player.process_command(PlayerCommand::ClearPlaylist);
4548        player.process_command(PlayerCommand::Undo);
4549        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
4550
4551        player.process_command(PlayerCommand::Redo);
4552        assert!(playlist_ids(&player).is_empty());
4553    }
4554
4555    // --- Multi-step undo/redo ---
4556
4557    #[test]
4558    fn multiple_undos_in_sequence() {
4559        let mut player = Player::new();
4560
4561        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("A")]));
4562        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
4563        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("C")]));
4564        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4565
4566        player.process_command(PlayerCommand::Undo);
4567        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
4568
4569        player.process_command(PlayerCommand::Undo);
4570        assert_eq!(playlist_titles(&player), vec!["A"]);
4571
4572        player.process_command(PlayerCommand::Undo);
4573        assert!(playlist_ids(&player).is_empty());
4574    }
4575
4576    #[test]
4577    fn undo_redo_undo_cycle() {
4578        let mut player = Player::new();
4579        let items = vec![make_item("A"), make_item("B")];
4580
4581        player.process_command(PlayerCommand::AddToPlaylist(items));
4582        player.process_command(PlayerCommand::Undo);
4583        assert!(playlist_ids(&player).is_empty());
4584
4585        player.process_command(PlayerCommand::Redo);
4586        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
4587
4588        player.process_command(PlayerCommand::Undo);
4589        assert!(playlist_ids(&player).is_empty());
4590    }
4591
4592    #[test]
4593    fn new_action_clears_redo_stack() {
4594        let mut player = Player::new();
4595        let items = vec![make_item("A")];
4596
4597        player.process_command(PlayerCommand::AddToPlaylist(items));
4598        player.process_command(PlayerCommand::Undo);
4599        assert!(player.undo_stack().can_redo());
4600
4601        // New action should clear redo
4602        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
4603        assert!(!player.undo_stack().can_redo());
4604    }
4605
4606    #[test]
4607    fn undo_on_empty_stack_is_noop() {
4608        let mut player = Player::new();
4609        player.process_command(PlayerCommand::Undo);
4610        assert!(playlist_ids(&player).is_empty());
4611    }
4612
4613    #[test]
4614    fn redo_on_empty_stack_is_noop() {
4615        let mut player = Player::new();
4616        player.process_command(PlayerCommand::Redo);
4617        assert!(playlist_ids(&player).is_empty());
4618    }
4619
4620    // --- Non-undoable commands don't push entries ---
4621
4622    #[test]
4623    fn playback_commands_not_undoable() {
4624        let mut player = Player::new();
4625        player.process_command(PlayerCommand::Pause);
4626        player.process_command(PlayerCommand::Resume);
4627        player.process_command(PlayerCommand::NextTrack);
4628        player.process_command(PlayerCommand::PrevTrack);
4629        assert!(!player.undo_stack().can_undo());
4630    }
4631
4632    #[test]
4633    fn update_paths_not_undoable() {
4634        let mut player = Player::new();
4635        let items = vec![make_item("A")];
4636        let id = items[0].id;
4637        player.process_command(PlayerCommand::AddToPlaylist(items));
4638
4639        let undo_count = player.undo_stack().undo_len();
4640        player.process_command(PlayerCommand::UpdatePaths(vec![(
4641            id,
4642            PathBuf::from("/new/path.flac"),
4643        )]));
4644        assert_eq!(player.undo_stack().undo_len(), undo_count);
4645    }
4646
4647    // --- Complex scenarios ---
4648
4649    #[test]
4650    fn add_remove_undo_undo_produces_original() {
4651        let mut player = Player::new();
4652        let items = vec![make_item("A"), make_item("B"), make_item("C")];
4653        let b_id = items[1].id;
4654        let original_titles = vec!["A", "B", "C"];
4655
4656        player.process_command(PlayerCommand::AddToPlaylist(items));
4657        player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
4658        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
4659
4660        // Undo remove → back to A, B, C
4661        player.process_command(PlayerCommand::Undo);
4662        assert_eq!(playlist_titles(&player), original_titles);
4663
4664        // Undo add → empty
4665        player.process_command(PlayerCommand::Undo);
4666        assert!(playlist_ids(&player).is_empty());
4667    }
4668
4669    #[test]
4670    fn interleaved_adds_and_moves_undo() {
4671        let mut player = Player::new();
4672        let items = vec![make_item("A"), make_item("B"), make_item("C")];
4673        let a_id = items[0].id;
4674        let c_id = items[2].id;
4675
4676        player.process_command(PlayerCommand::AddToPlaylist(items));
4677
4678        // Move A after C: [B, C, A]
4679        player.process_command(PlayerCommand::MoveInPlaylist {
4680            id: a_id,
4681            target: c_id,
4682            after: true,
4683        });
4684        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
4685
4686        // Add D: [B, C, A, D]
4687        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("D")]));
4688        assert_eq!(playlist_titles(&player), vec!["B", "C", "A", "D"]);
4689
4690        // Undo add D: [B, C, A]
4691        player.process_command(PlayerCommand::Undo);
4692        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
4693
4694        // Undo move: [A, B, C]
4695        player.process_command(PlayerCommand::Undo);
4696        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4697    }
4698
4699    /// Regression test for GitHub #89: AudioEngine must be dropped synchronously
4700    /// in stop_engine() before the caller changes sample rates. If the engine is
4701    /// dropped on a background thread, CoreAudio's internal buffer list can be
4702    /// freed while AudioUnitUninitialize is still tearing it down → crash.
4703    #[test]
4704    fn stop_engine_drops_engine_synchronously() {
4705        use std::sync::atomic::{AtomicBool, Ordering};
4706
4707        struct MockEngine {
4708            dropped: Arc<AtomicBool>,
4709        }
4710        impl AudioEngineHandle for MockEngine {
4711            fn start(&self) -> Result<(), BackendError> {
4712                Ok(())
4713            }
4714            fn stop(&self) -> Result<(), BackendError> {
4715                Ok(())
4716            }
4717            fn is_running(&self) -> bool {
4718                false
4719            }
4720            fn fade_out(&self) {}
4721            fn fade_in(&self) -> Result<(), BackendError> {
4722                Ok(())
4723            }
4724            fn is_silent(&self) -> bool {
4725                false
4726            }
4727        }
4728        impl Drop for MockEngine {
4729            fn drop(&mut self) {
4730                self.dropped.store(true, Ordering::SeqCst);
4731            }
4732        }
4733
4734        let dropped = Arc::new(AtomicBool::new(false));
4735
4736        let mut player = Player::new();
4737        player.transport = Transport::Loaded(test_session(
4738            QueueItemId::new(),
4739            Box::new(MockEngine {
4740                dropped: dropped.clone(),
4741            }),
4742        ));
4743
4744        player.stop_engine();
4745
4746        // The engine must already be dropped when stop_engine returns.
4747        // If this fails, the engine was moved to a background thread — the
4748        // exact race condition that causes the #89 crash.
4749        assert!(
4750            dropped.load(Ordering::SeqCst),
4751            "AudioEngine must be dropped synchronously in stop_engine (GitHub #89)"
4752        );
4753    }
4754
4755    #[test]
4756    fn abandoning_a_stream_wakes_a_reader_parked_at_the_write_head() {
4757        let live = LiveStream {
4758            feed: crate::remote::downloads::ByteFeed::new(),
4759            abandoned: Default::default(),
4760        };
4761        let feed = live.feed.clone();
4762        let started = std::time::Instant::now();
4763        let reader = thread::spawn(move || {
4764            feed.wait_past(
4765                0,
4766                std::time::Instant::now() + std::time::Duration::from_secs(30),
4767            )
4768        });
4769        thread::sleep(std::time::Duration::from_millis(50));
4770        live.abandon();
4771        reader.join().unwrap();
4772
4773        assert!(live.abandoned.load(std::sync::atomic::Ordering::Acquire));
4774        assert!(started.elapsed() < std::time::Duration::from_secs(5));
4775    }
4776
4777    // --- Engine format matches the decoded PCM ---
4778
4779    /// Backend pinned to one sample rate that refuses every switch, recording
4780    /// the format the engine is asked for.
4781    pub(super) struct StuckBackend {
4782        pub(super) rate: f64,
4783        pub(super) asked: Arc<std::sync::Mutex<Option<(f64, u32)>>>,
4784        /// How many times an engine it made was started.
4785        pub(super) starts: Arc<std::sync::atomic::AtomicUsize>,
4786    }
4787
4788    struct NullEngine {
4789        starts: Arc<std::sync::atomic::AtomicUsize>,
4790        running: std::sync::atomic::AtomicBool,
4791        lead_in: Arc<AtomicU64>,
4792    }
4793    impl AudioEngineHandle for NullEngine {
4794        fn start(&self) -> Result<(), BackendError> {
4795            self.starts.fetch_add(1, Ordering::Relaxed);
4796            self.running.store(true, Ordering::Relaxed);
4797            Ok(())
4798        }
4799        fn stop(&self) -> Result<(), BackendError> {
4800            self.running.store(false, Ordering::Relaxed);
4801            Ok(())
4802        }
4803        fn is_running(&self) -> bool {
4804            self.running.load(Ordering::Relaxed)
4805        }
4806        fn fade_out(&self) {}
4807        fn fade_in(&self) -> Result<(), BackendError> {
4808            Ok(())
4809        }
4810        fn is_silent(&self) -> bool {
4811            false
4812        }
4813        fn lead_in(&self, frames: u64) {
4814            self.lead_in.store(frames, Ordering::Relaxed);
4815        }
4816    }
4817
4818    impl AudioBackend for StuckBackend {
4819        fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
4820            Ok(vec![self.default_device()?])
4821        }
4822        fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
4823            Ok(backend::DeviceInfo {
4824                name: "Stuck DAC".into(),
4825                sample_rates: vec![self.rate],
4826                platform_id: 0,
4827                kind: Default::default(),
4828            })
4829        }
4830        fn supported_sample_rates(
4831            &self,
4832            _device: &backend::DeviceInfo,
4833        ) -> Result<Vec<f64>, BackendError> {
4834            Ok(vec![self.rate])
4835        }
4836        fn get_device_sample_rate(
4837            &self,
4838            _device: &backend::DeviceInfo,
4839        ) -> Result<f64, BackendError> {
4840            Ok(self.rate)
4841        }
4842        fn set_device_sample_rate(
4843            &self,
4844            _device: &backend::DeviceInfo,
4845            rate: f64,
4846        ) -> Result<f64, BackendError> {
4847            Err(BackendError::UnsupportedSampleRate(rate))
4848        }
4849        fn create_engine(
4850            &self,
4851            _device: &backend::DeviceInfo,
4852            sample_rate: f64,
4853            channels: u32,
4854            _consumer: rtrb::Consumer<f32>,
4855            _samples_played: Arc<AtomicU64>,
4856        ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
4857            *self.asked.lock().unwrap() = Some((sample_rate, channels));
4858            Ok(Box::new(NullEngine {
4859                starts: self.starts.clone(),
4860                running: Default::default(),
4861                lead_in: Default::default(),
4862            }))
4863        }
4864    }
4865
4866    /// Hands the test every ring buffer an engine is made with, so what
4867    /// reaches the device can be read back.
4868    struct CaptureBackend {
4869        consumers: Arc<std::sync::Mutex<Vec<rtrb::Consumer<f32>>>>,
4870    }
4871
4872    impl AudioBackend for CaptureBackend {
4873        fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
4874            Ok(vec![self.default_device()?])
4875        }
4876        fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
4877            Ok(backend::DeviceInfo {
4878                name: "Capture DAC".into(),
4879                sample_rates: vec![44100.0],
4880                platform_id: 0,
4881                kind: Default::default(),
4882            })
4883        }
4884        fn supported_sample_rates(
4885            &self,
4886            _device: &backend::DeviceInfo,
4887        ) -> Result<Vec<f64>, BackendError> {
4888            Ok(vec![44100.0])
4889        }
4890        fn get_device_sample_rate(
4891            &self,
4892            _device: &backend::DeviceInfo,
4893        ) -> Result<f64, BackendError> {
4894            Ok(44100.0)
4895        }
4896        fn set_device_sample_rate(
4897            &self,
4898            _device: &backend::DeviceInfo,
4899            rate: f64,
4900        ) -> Result<f64, BackendError> {
4901            Ok(rate)
4902        }
4903        fn create_engine(
4904            &self,
4905            _device: &backend::DeviceInfo,
4906            _sample_rate: f64,
4907            _channels: u32,
4908            consumer: rtrb::Consumer<f32>,
4909            _samples_played: Arc<AtomicU64>,
4910        ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
4911            self.consumers.lock().unwrap().push(consumer);
4912            Ok(Box::new(NullEngine {
4913                starts: Default::default(),
4914                running: Default::default(),
4915                lead_in: Default::default(),
4916            }))
4917        }
4918    }
4919
4920    /// The loudest sample of a session's first `want` samples.
4921    fn peak_reaching_the_device(
4922        consumers: &std::sync::Mutex<Vec<rtrb::Consumer<f32>>>,
4923        want: usize,
4924    ) -> f32 {
4925        let mut consumer = consumers.lock().unwrap().pop().expect("an engine was made");
4926        let mut peak = 0.0f32;
4927        let mut got = 0;
4928        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
4929        while got < want && std::time::Instant::now() < deadline {
4930            let n = consumer.slots();
4931            if n == 0 {
4932                std::thread::sleep(std::time::Duration::from_millis(2));
4933                continue;
4934            }
4935            let chunk = consumer.read_chunk(n).unwrap();
4936            let (a, b) = chunk.as_slices();
4937            peak = a.iter().chain(b).fold(peak, |p, s| p.max(s.abs()));
4938            got += n;
4939            chunk.commit_all();
4940        }
4941        assert!(got >= want, "only {got} of {want} samples arrived");
4942        peak
4943    }
4944
4945    /// A file on disk and a download still landing open through the same
4946    /// session, so both are processed. The streaming path is the one that
4947    /// used to be its own function, and the easy one to leave behind.
4948    #[test]
4949    fn a_profile_processes_files_and_streams_alike() {
4950        let dir = tempfile::tempdir().unwrap();
4951        let tone = dir.path().join("tone.wav");
4952        crate::test_utils::generate_wav_tone(&tone, 44100, 440.0, 0.5);
4953        let want = 22050;
4954
4955        let consumers = Arc::new(std::sync::Mutex::new(Vec::new()));
4956        let mut player = Player::new();
4957        player.backend = Box::new(CaptureBackend {
4958            consumers: consumers.clone(),
4959        });
4960        let mut item = make_item("tone");
4961        item.path = tone.clone();
4962        let id = item.id;
4963        player.shared_state.add_items(vec![item]);
4964        player.shared_state.set_cursor(Some(id));
4965
4966        let stream = || {
4967            let feed = crate::remote::downloads::ByteFeed::new();
4968            feed.set(std::fs::metadata(&tone).unwrap().len());
4969            Source::Stream(StreamSource {
4970                path: tone.clone(),
4971                bytes_written: feed,
4972                total: 0,
4973                mode: streaming::ProbeMode::Full,
4974            })
4975        };
4976        let peaks = |player: &mut Player| {
4977            let mut out = Vec::new();
4978            for source in [Source::File(tone.clone()), stream()] {
4979                player
4980                    .try_open_session(id, source, None, 0, Run::Playing)
4981                    .unwrap();
4982                player.publish();
4983                out.push((
4984                    peak_reaching_the_device(&consumers, want),
4985                    player.shared_state.dsp().map(|d| d.profile),
4986                ));
4987                player.stop_engine();
4988            }
4989            out
4990        };
4991
4992        let untouched = peaks(&mut player);
4993
4994        player.dsp_override = Some(Arc::new(
4995            crate::audio::dsp::Setup::new(vec![], vec![]).with_preamp(-6.0206),
4996        ));
4997        let processed = peaks(&mut player);
4998
4999        for (kind, ((before, none), (after, half))) in ["file", "stream"]
5000            .iter()
5001            .zip(untouched.iter().zip(&processed))
5002        {
5003            assert_eq!(
5004                none, &None,
5005                "{kind}: nothing is published without a profile"
5006            );
5007            assert_eq!(half.as_deref(), Some("test"), "{kind}: the badge names it");
5008            assert!(*before > 0.1, "{kind}: the tone reached the device");
5009            assert!(
5010                (after / before - 0.5).abs() < 0.01,
5011                "{kind}: {before} → {after}, not halved"
5012            );
5013        }
5014    }
5015
5016    fn engine_format_for(source_rate: u32, channels: u16, device_rate: f64) -> (f64, u32) {
5017        let asked = Arc::new(std::sync::Mutex::new(None));
5018        let mut player = Player::new();
5019        player.backend = Box::new(StuckBackend {
5020            rate: device_rate,
5021            asked: asked.clone(),
5022            starts: Default::default(),
5023        });
5024
5025        let info = buffer::StreamInfo {
5026            codec: "MP3".into(),
5027            sample_rate: source_rate,
5028            channels,
5029            bit_depth: Some(16),
5030            bitrate_kbps: None,
5031            duration_ms: 1000,
5032        };
5033        let (_producer, consumer) = rtrb::RingBuffer::new(16);
5034        player
5035            .create_engine_for(&info, consumer)
5036            .expect("engine creation should succeed");
5037        let asked = *asked.lock().unwrap();
5038        asked.expect("engine was never created")
5039    }
5040
5041    #[test]
5042    fn engine_uses_source_rate_when_device_refuses_switch() {
5043        // MPEG-2 MP3 rates are routinely rejected by output devices. The engine
5044        // must still be told the rate the PCM actually is.
5045        assert_eq!(engine_format_for(22050, 2, 48000.0), (22050.0, 2));
5046        assert_eq!(engine_format_for(32000, 2, 44100.0), (32000.0, 2));
5047    }
5048
5049    #[test]
5050    fn engine_uses_source_channel_count() {
5051        assert_eq!(engine_format_for(44100, 1, 44100.0), (44100.0, 1));
5052    }
5053
5054    /// The rate the device settled at, as the front ends read it.
5055    fn settled_rate_for(source_rate: u32, device_rate: f64) -> Option<u32> {
5056        let mut player = Player::new();
5057        player.backend = Box::new(StuckBackend {
5058            rate: device_rate,
5059            asked: Arc::new(std::sync::Mutex::new(None)),
5060            starts: Default::default(),
5061        });
5062        let state = player.shared_state.clone();
5063
5064        let info = buffer::StreamInfo {
5065            codec: "MP3".into(),
5066            sample_rate: source_rate,
5067            channels: 2,
5068            bit_depth: Some(16),
5069            bitrate_kbps: None,
5070            duration_ms: 1000,
5071        };
5072        let (_producer, consumer) = rtrb::RingBuffer::new(16);
5073        player
5074            .create_engine_for(&info, consumer)
5075            .expect("engine creation should succeed");
5076        state.output_sample_rate()
5077    }
5078
5079    #[test]
5080    fn settled_device_rate_reaches_the_shared_state() {
5081        // A device that refuses the switch is being fed resampled audio, and
5082        // that is the case the front ends have to be able to see. Before this
5083        // the comparison happened once, in a log line.
5084        assert_eq!(settled_rate_for(22050, 48000.0), Some(48000));
5085        // No switch needed, so nothing resampled: the two rates agree.
5086        assert_eq!(settled_rate_for(44100, 44100.0), Some(44100));
5087    }
5088
5089    /// A device that takes its time reclocking, as real hardware does.
5090    struct SlowBackend {
5091        observed: Arc<std::sync::Mutex<Vec<Option<u32>>>>,
5092        state: Arc<SharedPlayerState>,
5093        lead_in: Arc<AtomicU64>,
5094    }
5095
5096    impl AudioBackend for SlowBackend {
5097        fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
5098            Ok(vec![self.default_device()?])
5099        }
5100        fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
5101            Ok(backend::DeviceInfo {
5102                name: "Slow DAC".into(),
5103                sample_rates: vec![44100.0, 48000.0],
5104                platform_id: 0,
5105                kind: Default::default(),
5106            })
5107        }
5108        fn supported_sample_rates(
5109            &self,
5110            _device: &backend::DeviceInfo,
5111        ) -> Result<Vec<f64>, BackendError> {
5112            Ok(vec![44100.0, 48000.0])
5113        }
5114        fn get_device_sample_rate(
5115            &self,
5116            _device: &backend::DeviceInfo,
5117        ) -> Result<f64, BackendError> {
5118            Ok(48000.0)
5119        }
5120        fn set_device_sample_rate(
5121            &self,
5122            _device: &backend::DeviceInfo,
5123            rate: f64,
5124        ) -> Result<f64, BackendError> {
5125            // What a front end polling mid-switch would see.
5126            self.observed
5127                .lock()
5128                .unwrap()
5129                .push(self.state.output_sample_rate());
5130            Ok(rate)
5131        }
5132        fn create_engine(
5133            &self,
5134            _device: &backend::DeviceInfo,
5135            _sample_rate: f64,
5136            _channels: u32,
5137            _consumer: rtrb::Consumer<f32>,
5138            _samples_played: Arc<AtomicU64>,
5139        ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
5140            Ok(Box::new(NullEngine {
5141                starts: Default::default(),
5142                running: Default::default(),
5143                lead_in: self.lead_in.clone(),
5144            }))
5145        }
5146    }
5147
5148    #[test]
5149    fn the_previous_rate_is_not_published_while_the_device_reclocks() {
5150        // A 48 kHz track followed by a 44.1 kHz one: for as long as the switch
5151        // takes — the better part of a second on USB — the new track's info is
5152        // published against the old track's output rate. A front end polling in
5153        // that window used to latch "44.1 → 48" and, since nothing about the
5154        // codec or the source rate changed afterwards, never let go of it.
5155        let mut player = Player::new();
5156        let state = player.shared_state.clone();
5157        state.set_output_sample_rate(48000);
5158
5159        let observed = Arc::new(std::sync::Mutex::new(Vec::new()));
5160        player.backend = Box::new(SlowBackend {
5161            observed: observed.clone(),
5162            state: state.clone(),
5163            lead_in: Default::default(),
5164        });
5165
5166        let info = buffer::StreamInfo {
5167            codec: "FLAC".into(),
5168            sample_rate: 44100,
5169            channels: 2,
5170            bit_depth: Some(16),
5171            bitrate_kbps: None,
5172            duration_ms: 1000,
5173        };
5174        let (_producer, consumer) = rtrb::RingBuffer::new(16);
5175        player
5176            .create_engine_for(&info, consumer)
5177            .expect("engine creation should succeed");
5178
5179        assert_eq!(
5180            *observed.lock().unwrap(),
5181            vec![None],
5182            "mid-switch the output rate must read as unknown, not as the last track's"
5183        );
5184        assert_eq!(state.output_sample_rate(), Some(44100));
5185    }
5186
5187    /// The lead-in an engine was given, for a track at `source_rate` on a
5188    /// device sitting at 48 kHz.
5189    fn lead_in_for(source_rate: u32) -> u64 {
5190        let mut player = Player::new();
5191        let lead_in = Arc::new(AtomicU64::new(0));
5192        player.backend = Box::new(SlowBackend {
5193            observed: Default::default(),
5194            state: player.shared_state.clone(),
5195            lead_in: lead_in.clone(),
5196        });
5197        let info = buffer::StreamInfo {
5198            codec: "FLAC".into(),
5199            sample_rate: source_rate,
5200            channels: 2,
5201            bit_depth: Some(16),
5202            bitrate_kbps: None,
5203            duration_ms: 1000,
5204        };
5205        let (_producer, consumer) = rtrb::RingBuffer::new(16);
5206        player
5207            .create_engine_for(&info, consumer)
5208            .expect("engine creation should succeed");
5209        lead_in.load(Ordering::Relaxed)
5210    }
5211
5212    #[test]
5213    fn an_engine_made_inside_the_silence_keeps_the_rest_of_it() {
5214        // A seek straight after a rate switch: the device is still relocking,
5215        // though the new engine finds its rate already matching.
5216        let mut player = Player::new();
5217        let lead_in = Arc::new(AtomicU64::new(0));
5218        player.backend = Box::new(SlowBackend {
5219            observed: Default::default(),
5220            state: player.shared_state.clone(),
5221            lead_in: lead_in.clone(),
5222        });
5223        let info = |sample_rate| buffer::StreamInfo {
5224            codec: "FLAC".into(),
5225            sample_rate,
5226            channels: 2,
5227            bit_depth: Some(16),
5228            bitrate_kbps: None,
5229            duration_ms: 1000,
5230        };
5231        let (_p, consumer) = rtrb::RingBuffer::new(16);
5232        player.create_engine_for(&info(44100), consumer).unwrap();
5233        let first = lead_in.swap(0, Ordering::Relaxed);
5234        assert!(first > 0);
5235
5236        let (_p, consumer) = rtrb::RingBuffer::new(16);
5237        player.create_engine_for(&info(48000), consumer).unwrap();
5238        let carried = lead_in.load(Ordering::Relaxed);
5239        // Frames at the new engine's rate: what is left of the same second.
5240        assert!(
5241            carried > 0 && carried < first * 48000 / 44100,
5242            "carried {carried} of {first}"
5243        );
5244    }
5245
5246    #[test]
5247    fn silence_leads_in_only_after_the_device_changed_rate() {
5248        assert!(
5249            lead_in_for(44100) > 0,
5250            "the device is relocking, so the start of the track would be lost"
5251        );
5252        assert_eq!(lead_in_for(48000), 0, "no switch, nothing to wait for");
5253    }
5254
5255    /// Backend that hands its rate-change callback back to the test.
5256    struct WatchedBackend {
5257        inner: StuckBackend,
5258        #[allow(clippy::type_complexity)]
5259        captured: Arc<std::sync::Mutex<Option<Box<dyn Fn(f64) + Send + Sync>>>>,
5260    }
5261
5262    struct NullWatch;
5263    impl backend::SampleRateWatch for NullWatch {}
5264
5265    impl AudioBackend for WatchedBackend {
5266        fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
5267            self.inner.list_devices()
5268        }
5269        fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
5270            self.inner.default_device()
5271        }
5272        fn supported_sample_rates(
5273            &self,
5274            device: &backend::DeviceInfo,
5275        ) -> Result<Vec<f64>, BackendError> {
5276            self.inner.supported_sample_rates(device)
5277        }
5278        fn get_device_sample_rate(
5279            &self,
5280            device: &backend::DeviceInfo,
5281        ) -> Result<f64, BackendError> {
5282            self.inner.get_device_sample_rate(device)
5283        }
5284        fn set_device_sample_rate(
5285            &self,
5286            device: &backend::DeviceInfo,
5287            rate: f64,
5288        ) -> Result<f64, BackendError> {
5289            self.inner.set_device_sample_rate(device, rate)
5290        }
5291        fn watch_device_sample_rate(
5292            &self,
5293            _device: &backend::DeviceInfo,
5294            on_change: Box<dyn Fn(f64) + Send + Sync>,
5295        ) -> Option<Box<dyn backend::SampleRateWatch>> {
5296            *self.captured.lock().unwrap() = Some(on_change);
5297            Some(Box::new(NullWatch))
5298        }
5299        fn create_engine(
5300            &self,
5301            device: &backend::DeviceInfo,
5302            sample_rate: f64,
5303            channels: u32,
5304            consumer: rtrb::Consumer<f32>,
5305            samples_played: Arc<AtomicU64>,
5306        ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
5307            self.inner
5308                .create_engine(device, sample_rate, channels, consumer, samples_played)
5309        }
5310    }
5311
5312    #[test]
5313    fn external_rate_change_reaches_the_shared_state() {
5314        // The device is shared. Another client moving the rate mid-track used
5315        // to leave the front ends asserting bit-perfection while the HAL
5316        // resampled underneath them.
5317        let captured = Arc::new(std::sync::Mutex::new(None));
5318        let mut player = Player::new();
5319        player.backend = Box::new(WatchedBackend {
5320            inner: StuckBackend {
5321                rate: 44100.0,
5322                asked: Arc::new(std::sync::Mutex::new(None)),
5323                starts: Default::default(),
5324            },
5325            captured: captured.clone(),
5326        });
5327        let state = player.shared_state.clone();
5328
5329        let info = buffer::StreamInfo {
5330            codec: "FLAC".into(),
5331            sample_rate: 44100,
5332            channels: 2,
5333            bit_depth: Some(16),
5334            bitrate_kbps: None,
5335            duration_ms: 1000,
5336        };
5337        let (_producer, consumer) = rtrb::RingBuffer::new(16);
5338        player
5339            .create_engine_for(&info, consumer)
5340            .expect("engine creation should succeed");
5341        assert_eq!(state.output_sample_rate(), Some(44100));
5342
5343        let on_change = captured.lock().unwrap().take().expect("watch registered");
5344        on_change(48000.0);
5345        assert_eq!(state.output_sample_rate(), Some(48000));
5346    }
5347}