Skip to main content

koan_core/player/
mod.rs

1pub mod commands;
2pub mod state;
3pub mod undo;
4
5use std::path::Path;
6use std::sync::Arc;
7use std::sync::atomic::{AtomicU64, Ordering};
8use std::thread;
9
10use thiserror::Error;
11
12use crate::audio::{
13    analyzer::VizAnalyzer,
14    backend::{self, AudioBackend, AudioEngineHandle, BackendError},
15    buffer, streaming,
16    viz::{VizBuffer, VizSnapshot},
17};
18use buffer::PlaybackTimeline;
19use commands::{CommandChannel, PlayerCommand};
20use state::{LoadState, PlaybackSource, PlaybackState, QueueItemId, SharedPlayerState, TrackInfo};
21use undo::{UndoEntry, UndoStack};
22
23/// Ring buffer size in samples. ~1s at 192kHz stereo.
24pub(crate) const RING_BUFFER_SIZE: usize = 192_000 * 2;
25
26#[derive(Debug, Error)]
27pub enum PlayerError {
28    #[error("backend error: {0}")]
29    Backend(#[from] BackendError),
30    #[error("decode error: {0}")]
31    Decode(#[from] buffer::DecodeError),
32}
33
34/// The player controller. Owns the audio pipeline and processes commands.
35pub struct Player {
36    shared_state: Arc<SharedPlayerState>,
37    commands: CommandChannel,
38    active_playback: Option<ActivePlayback>,
39    timeline: Arc<PlaybackTimeline>,
40    viz_buffer: Arc<VizBuffer>,
41    viz_snapshot: Arc<VizSnapshot>,
42    /// Background FFT analysis thread. Held for its lifetime; dropped on Player drop.
43    _viz_analyzer: VizAnalyzer,
44    undo_stack: UndoStack,
45    /// When Some, undo entries are collected into this buffer instead of pushed
46    /// directly onto the undo stack. Flushed on EndUndoBatch.
47    batch_buffer: Option<Vec<UndoEntry>>,
48    /// Configured output device name. None = system default.
49    output_device_name: Option<String>,
50    /// Platform audio backend (CoreAudio on macOS, cpal on Linux).
51    backend: Box<dyn AudioBackend>,
52    /// Debounce: timestamp of last NextTrack/PrevTrack to suppress key repeat.
53    last_skip: std::time::Instant,
54    /// Playback sessions started — lets tests assert how many engine restarts
55    /// an operation costs.
56    #[cfg(test)]
57    playback_starts: usize,
58}
59
60/// Holds the resources for an active playback session.
61struct ActivePlayback {
62    engine: Box<dyn AudioEngineHandle>,
63    decode_handle: buffer::DecodeHandle,
64}
65
66impl Default for Player {
67    fn default() -> Self {
68        Self::new()
69    }
70}
71
72impl Player {
73    pub fn new() -> Self {
74        let viz_buffer = VizBuffer::new();
75        let viz_snapshot = VizSnapshot::new();
76        let timeline = PlaybackTimeline::new();
77        let cfg = crate::config::Config::load_or_default();
78        let viz_analyzer = VizAnalyzer::spawn_with_snapshot(
79            Arc::clone(&viz_buffer),
80            &cfg.visualizer,
81            Arc::clone(&viz_snapshot),
82            timeline.samples_played_counter(),
83        );
84
85        Self {
86            shared_state: SharedPlayerState::new(),
87            commands: CommandChannel::new(),
88            active_playback: None,
89            timeline,
90            viz_buffer,
91            viz_snapshot,
92            _viz_analyzer: viz_analyzer,
93            undo_stack: UndoStack::new(),
94            batch_buffer: None,
95            output_device_name: cfg.playback.output_device.clone(),
96            backend: crate::audio::platform_backend(),
97            last_skip: std::time::Instant::now(),
98            #[cfg(test)]
99            playback_starts: 0,
100        }
101    }
102
103    /// Get a clone of the shared state for UI reads.
104    pub fn shared_state(&self) -> Arc<SharedPlayerState> {
105        self.shared_state.clone()
106    }
107
108    /// Get the playback timeline for UI reads.
109    pub fn timeline(&self) -> Arc<PlaybackTimeline> {
110        self.timeline.clone()
111    }
112
113    /// Get the visualization buffer for the TUI.
114    pub fn viz_buffer(&self) -> Arc<VizBuffer> {
115        self.viz_buffer.clone()
116    }
117
118    /// Get the shared analysis snapshot for the TUI.
119    /// The analysis thread writes here; the UI thread reads a clone each frame.
120    pub fn viz_snapshot(&self) -> Arc<VizSnapshot> {
121        self.viz_snapshot.clone()
122    }
123
124    /// Access undo stack (for tests and UI state queries).
125    pub fn undo_stack(&self) -> &UndoStack {
126        &self.undo_stack
127    }
128
129    /// Create an audio engine for a stream, switching the output device to the
130    /// source rate first so output is bit-perfect.
131    ///
132    /// The engine is always configured with the source's own rate and channel
133    /// count — the format the decode thread writes into the ring buffer. A
134    /// device that cannot take the requested rate (MPEG-2/2.5 MP3 rates are
135    /// commonly refused) resamples instead of playing at the wrong speed.
136    fn create_engine_for(
137        &self,
138        info: &buffer::StreamInfo,
139        consumer: rtrb::Consumer<f32>,
140    ) -> Result<Box<dyn AudioEngineHandle>, PlayerError> {
141        let device = self.resolve_device()?;
142        let device_rate = self.backend.get_device_sample_rate(&device)?;
143        let source_rate = info.sample_rate as f64;
144
145        if (device_rate - source_rate).abs() > 0.1 {
146            log::info!(
147                "switching device sample rate: {}Hz → {}Hz",
148                device_rate,
149                source_rate
150            );
151            let settled = match self.backend.set_device_sample_rate(&device, source_rate) {
152                Ok(rate) => rate,
153                Err(e) => {
154                    log::warn!("failed to set device sample rate: {}", e);
155                    device_rate
156                }
157            };
158            if (settled - source_rate).abs() > 0.1 {
159                log::warn!(
160                    "device stayed at {}Hz (wanted {}Hz) — output is resampled, not bit-perfect",
161                    settled,
162                    source_rate
163                );
164            }
165        }
166
167        Ok(self.backend.create_engine(
168            &device,
169            source_rate,
170            info.channels as u32,
171            consumer,
172            self.timeline.samples_played_counter(),
173        )?)
174    }
175
176    /// Resolve the output device: use configured device name if set,
177    /// falling back to system default if not set or if the named device is unavailable.
178    fn resolve_device(&self) -> Result<backend::DeviceInfo, PlayerError> {
179        if let Some(ref name) = self.output_device_name {
180            match self.backend.list_devices() {
181                Ok(devices) => {
182                    if let Some(dev) = devices.into_iter().find(|d| d.name == *name) {
183                        return Ok(dev);
184                    }
185                    log::warn!(
186                        "configured output device '{}' not found, falling back to default",
187                        name,
188                    );
189                }
190                Err(e) => {
191                    log::warn!("failed to list devices while resolving '{}': {}", name, e);
192                }
193            }
194        }
195        Ok(self.backend.default_device()?)
196    }
197
198    /// Switch the output device. Persists to config and restarts the engine
199    /// on the current track if playing.
200    pub fn set_output_device(&mut self, name: String) {
201        log::info!("switching output device to: {}", name);
202        self.output_device_name = Some(name.clone());
203
204        // Persist to config.toml (not the merged config — avoids leaking secrets).
205        if let Err(e) = crate::config::Config::update_base(|cfg| {
206            cfg.playback.output_device = Some(name);
207        }) {
208            log::error!("failed to save output device config: {}", e);
209        }
210
211        self.restart_on_current_track();
212    }
213
214    /// Clear the configured output device, reverting to system default.
215    pub fn clear_output_device(&mut self) {
216        log::info!("reverting to system default output device");
217        self.output_device_name = None;
218
219        if let Err(e) = crate::config::Config::update_base(|cfg| {
220            cfg.playback.output_device = None;
221        }) {
222            log::error!("failed to save output device config: {}", e);
223        }
224
225        self.restart_on_current_track();
226    }
227
228    /// If a track is currently playing or paused, restart playback at the
229    /// current position (e.g. after switching output devices). Preserves pause state.
230    fn restart_on_current_track(&mut self) {
231        if let Some(info) = self.shared_state.track_info() {
232            let was_paused = self.shared_state.playback_state() == PlaybackState::Paused;
233            let position_ms = self.shared_state.position_ms();
234            if let Err(e) = self.start_playback(info.id, &info.path, position_ms) {
235                log::error!("failed to restart playback on device switch: {}", e);
236                return;
237            }
238            if was_paused {
239                self.pause();
240            }
241        }
242    }
243
244    /// Get the current output device name (if configured).
245    pub fn output_device_name(&self) -> Option<&str> {
246        self.output_device_name.as_deref()
247    }
248
249    /// Get a command sender for the UI layer.
250    pub fn command_sender(&self) -> crossbeam_channel::Sender<PlayerCommand> {
251        self.commands.tx.clone()
252    }
253
254    /// Play a specific item in the playlist by ID.
255    /// Sets cursor, starts playback if Ready or streaming-ready, otherwise waits for TrackReady.
256    pub fn play(&mut self, id: QueueItemId) {
257        self.shared_state.set_cursor(Some(id));
258
259        match self.shared_state.item_playback_source(id) {
260            Some(PlaybackSource::Ready(path)) => {
261                if let Err(e) = self.start_playback(id, &path, 0) {
262                    log::error!("play failed: {}", e);
263                }
264            }
265            Some(PlaybackSource::Streaming {
266                path,
267                bytes_written,
268                total,
269            }) => {
270                if let Err(e) = self.start_streaming_playback(id, &path, bytes_written, total) {
271                    // The cursor stays here, so TrackReady starts it once the
272                    // whole file has landed.
273                    log::error!("streaming play failed, waiting for full download: {}", e);
274                }
275            }
276            None => {
277                // Item not ready — stop current playback, wait for TrackReady.
278                self.stop_engine();
279                self.shared_state.set_playback_state(PlaybackState::Stopped);
280                log::info!("play: item {:?} not ready, waiting for TrackReady", id);
281            }
282        }
283    }
284
285    /// Internal: start playback of a file.
286    ///
287    /// A failure leaves the player cleanly stopped. Displaying a track that no
288    /// engine is playing freezes the position and makes the transport lie.
289    fn start_playback(
290        &mut self,
291        id: QueueItemId,
292        path: &Path,
293        seek_ms: u64,
294    ) -> Result<(), PlayerError> {
295        #[cfg(test)]
296        {
297            self.playback_starts += 1;
298        }
299        let result = self.open_playback(id, path, seek_ms);
300        if result.is_err() {
301            self.stop_playback_and_clear_state();
302        }
303        result
304    }
305
306    fn open_playback(
307        &mut self,
308        id: QueueItemId,
309        path: &Path,
310        seek_ms: u64,
311    ) -> Result<(), PlayerError> {
312        self.stop_engine();
313
314        let info = buffer::probe_file(path)?;
315
316        // Set track_info + position immediately so the UI never sees a gap.
317        // For seeks, this keeps the bar at the target position instead of
318        // flashing to 0 while the new timeline spins up.
319        self.shared_state.set_track_info(Some(TrackInfo {
320            id,
321            path: path.to_path_buf(),
322            codec: info.codec.clone(),
323            sample_rate: info.sample_rate,
324            bit_depth: info.bit_depth,
325            bitrate_kbps: info.bitrate_kbps,
326            channels: info.channels,
327            duration_ms: info.duration_ms,
328        }));
329        self.shared_state.set_position_ms(seek_ms);
330        log::info!(
331            "playing: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
332            path.display(),
333            id,
334            info.codec,
335            info.sample_rate,
336            info.channels,
337            info.duration_ms,
338            if seek_ms > 0 {
339                format!(" @{}ms", seek_ms)
340            } else {
341                String::new()
342            }
343        );
344
345        let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
346
347        // Reset timeline for new playback session and start decode.
348        self.timeline.reset();
349
350        // Gapless lookahead: the decode thread maintains its own cursor
351        // (separate from the UI cursor) so it can look ahead through the
352        // playlist without affecting what the UI shows as "now playing".
353        let advance_state = self.shared_state.clone();
354        let decode_cursor = parking_lot::Mutex::new(Some(id));
355        let next_track = move || {
356            let current = decode_cursor.lock().take()?;
357            let next = advance_state.peek_next_ready_after(current);
358            if let Some((next_id, _)) = &next {
359                let mut guard = decode_cursor.lock();
360                *guard = Some(*next_id);
361            }
362            next
363        };
364
365        // Load ReplayGain config for this playback session.
366        let cfg = crate::config::Config::load_or_default();
367        let rg_mode = cfg.playback.replaygain;
368        let pre_amp_db = cfg.playback.pre_amp_db;
369
370        let finish_tx = self.commands.tx.clone();
371        let (_stream_info, decode_handle) = buffer::start_decode_file(
372            id,
373            path,
374            producer,
375            seek_ms,
376            next_track,
377            self.timeline.clone(),
378            Some(self.viz_buffer.clone()),
379            rg_mode,
380            pre_amp_db,
381            move || {
382                finish_tx.send(PlayerCommand::DecodeFinished).ok();
383            },
384        )?;
385
386        let engine = self.create_engine_for(&info, consumer)?;
387        engine.start()?;
388
389        self.shared_state.set_playback_state(PlaybackState::Playing);
390
391        self.active_playback = Some(ActivePlayback {
392            engine,
393            decode_handle,
394        });
395
396        Ok(())
397    }
398
399    /// Internal: start streaming playback from a partially-downloaded file.
400    ///
401    /// Creates a StreamBuffer and a pump thread that reads from the on-disk partial
402    /// file as bytes become available (tracked via `bytes_written`). The decode thread
403    /// reads from a StreamingSource backed by that buffer, blocking briefly when it
404    /// catches up to the write head.
405    fn start_streaming_playback(
406        &mut self,
407        id: QueueItemId,
408        path: &Path,
409        bytes_written: Arc<AtomicU64>,
410        total: u64,
411    ) -> Result<(), PlayerError> {
412        let result = self.open_streaming_playback(id, path, bytes_written, total);
413        if result.is_err() {
414            self.stop_playback_and_clear_state();
415        }
416        result
417    }
418
419    fn open_streaming_playback(
420        &mut self,
421        id: QueueItemId,
422        path: &Path,
423        bytes_written: Arc<AtomicU64>,
424        total: u64,
425    ) -> Result<(), PlayerError> {
426        self.stop_engine();
427
428        // Create a StreamBuffer with known total length.
429        let stream_buf = streaming::StreamBuffer::new(if total > 0 { Some(total) } else { None });
430
431        // Spawn a pump thread: reads bytes from the on-disk partial file as they
432        // become available (per bytes_written) and pushes them into StreamBuffer.
433        // This bridges the disk-based download with StreamingSource's in-memory design.
434        // The playlist item's path points to the .part file during download, so the
435        // pump opens the correct file. After download completes, the .part is renamed
436        // to the final path and the item path is updated — but the pump's open FD
437        // remains valid (Unix rename semantics).
438        let pump_path = path.to_path_buf();
439        let pump_buf = stream_buf.clone();
440        let pump_written = bytes_written.clone();
441        let pump_state = self.shared_state.clone();
442        thread::Builder::new()
443            .name("koan-stream-pump".into())
444            .spawn(move || {
445                use std::fs::File;
446                use std::io::Read;
447                use std::time::{Duration, Instant};
448
449                /// No new bytes for this long and the download is treated as dead.
450                /// `bytes_written` simply stops advancing when one dies, so without
451                /// a deadline the pump spins and the decode thread parks forever.
452                const STALL_LIMIT: Duration = Duration::from_secs(30);
453
454                let mut file = match File::open(&pump_path) {
455                    Ok(f) => f,
456                    Err(e) => {
457                        log::error!("stream pump: failed to open {}: {}", pump_path.display(), e);
458                        pump_buf.fail();
459                        return;
460                    }
461                };
462                let mut buf = vec![0u8; 65536];
463                let mut offset: u64 = 0;
464                let mut last_progress = Instant::now();
465                loop {
466                    // Nothing left to read the bytes, so nothing left to write them for.
467                    if pump_buf.is_abandoned() {
468                        return;
469                    }
470                    match pump_state.item_load_state(id) {
471                        Some(LoadState::Failed(e)) => {
472                            log::warn!("stream pump: download of {:?} failed: {}", id, e);
473                            pump_buf.fail();
474                            return;
475                        }
476                        // The download landed: drain to EOF rather than trusting
477                        // `total`, which is 0 for a chunked transfer.
478                        Some(LoadState::Ready) => match file.read(&mut buf) {
479                            Ok(0) => break,
480                            Ok(n) => {
481                                pump_buf.push(&buf[..n]);
482                                offset += n as u64;
483                                continue;
484                            }
485                            Err(e) => {
486                                log::warn!("stream pump read error: {}", e);
487                                pump_buf.fail();
488                                return;
489                            }
490                        },
491                        _ => {}
492                    }
493
494                    let available = pump_written.load(Ordering::Acquire);
495                    if offset >= available {
496                        if total > 0 && available >= total {
497                            break; // Download complete.
498                        }
499                        if last_progress.elapsed() >= STALL_LIMIT {
500                            log::warn!(
501                                "stream pump: no data for {}s, abandoning {}",
502                                STALL_LIMIT.as_secs(),
503                                pump_path.display()
504                            );
505                            pump_buf.fail();
506                            return;
507                        }
508                        thread::sleep(Duration::from_millis(10));
509                        continue;
510                    }
511                    let to_read = ((available - offset) as usize).min(buf.len());
512                    match file.read(&mut buf[..to_read]) {
513                        Ok(0) => {
514                            // File data may lag behind bytes_written (OS buffer flush timing).
515                            // Only treat as true EOF if we've pumped all expected data.
516                            if total > 0 && offset >= total {
517                                break;
518                            }
519                            let latest = pump_written.load(Ordering::Acquire);
520                            if total > 0 && latest >= total && offset >= latest {
521                                break;
522                            }
523                            // Data not yet visible on disk — back off and retry.
524                            thread::sleep(Duration::from_millis(1));
525                            continue;
526                        }
527                        Ok(n) => {
528                            pump_buf.push(&buf[..n]);
529                            offset += n as u64;
530                            last_progress = Instant::now();
531                        }
532                        Err(e) => {
533                            log::warn!("stream pump read error: {}", e);
534                            pump_buf.fail();
535                            return;
536                        }
537                    }
538                }
539                pump_buf.finish();
540            })
541            .map_err(|e| PlayerError::Decode(buffer::DecodeError::Io(e)))?;
542
543        // Probe via a streaming reader — blocks (via condvar) until enough header data arrives.
544        let probe_reader = stream_buf.reader();
545        let probe_hint = {
546            let mut h = symphonia::core::formats::probe::Hint::new();
547            if let Some(ext) = path.extension().and_then(|e| e.to_str()) {
548                h.with_extension(ext);
549            }
550            h
551        };
552        let probe_mss =
553            symphonia::core::io::MediaSourceStream::new(Box::new(probe_reader), Default::default());
554        let info = buffer::probe_source(probe_mss, &probe_hint)?;
555
556        self.shared_state.set_track_info(Some(TrackInfo {
557            id,
558            path: path.to_path_buf(),
559            codec: info.codec.clone(),
560            sample_rate: info.sample_rate,
561            bit_depth: info.bit_depth,
562            bitrate_kbps: info.bitrate_kbps,
563            channels: info.channels,
564            duration_ms: info.duration_ms,
565        }));
566        self.shared_state.set_position_ms(0);
567        log::info!(
568            "streaming: {} ({:?}) — {} {}Hz/{}ch, {}ms",
569            path.display(),
570            id,
571            info.codec,
572            info.sample_rate,
573            info.channels,
574            info.duration_ms,
575        );
576
577        let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
578
579        self.timeline.reset();
580
581        // Gapless lookahead after streaming: next track uses normal file path.
582        let advance_state = self.shared_state.clone();
583        let decode_cursor = parking_lot::Mutex::new(Some(id));
584        let next_track = move || {
585            let current = decode_cursor.lock().take()?;
586            let next = advance_state.peek_next_ready_after(current);
587            if let Some((next_id, _)) = &next {
588                let mut guard = decode_cursor.lock();
589                *guard = Some(*next_id);
590            }
591            next
592        };
593
594        // Decode using a fresh StreamingSource reader — reads from the StreamBuffer
595        // that the pump thread feeds. The decode thread blocks when it catches up to
596        // the write head, resuming as more data arrives.
597        // Build a SourceEntry using a fresh StreamingSource reader for the decode thread.
598        let decode_reader = stream_buf.reader();
599        let path_buf = path.to_path_buf();
600        let ext = path
601            .extension()
602            .and_then(|e| e.to_str())
603            .unwrap_or("")
604            .to_string();
605        let mut decode_hint = symphonia::core::formats::probe::Hint::new();
606        if !ext.is_empty() {
607            decode_hint.with_extension(&ext);
608        }
609        let first = buffer::SourceEntry {
610            id,
611            path: path_buf,
612            hint: decode_hint,
613            make_mss: Box::new(move || {
614                Ok(symphonia::core::io::MediaSourceStream::new(
615                    Box::new(decode_reader),
616                    Default::default(),
617                ))
618            }),
619        };
620
621        // Load ReplayGain config for this streaming session.
622        let cfg = crate::config::Config::load_or_default();
623        let rg_mode = cfg.playback.replaygain;
624        let pre_amp_db = cfg.playback.pre_amp_db;
625
626        let finish_tx = self.commands.tx.clone();
627        let (_stream_info, decode_handle) = buffer::start_decode(
628            first,
629            producer,
630            0,
631            move || {
632                let (next_id, next_path) = next_track()?;
633                Some(buffer::SourceEntry::from_file(next_id, next_path))
634            },
635            self.timeline.clone(),
636            Some(self.viz_buffer.clone()),
637            rg_mode,
638            pre_amp_db,
639            move || {
640                finish_tx.send(PlayerCommand::DecodeFinished).ok();
641            },
642        )?;
643
644        let engine = self.create_engine_for(&info, consumer)?;
645        engine.start()?;
646
647        self.shared_state.set_playback_state(PlaybackState::Playing);
648
649        self.active_playback = Some(ActivePlayback {
650            engine,
651            decode_handle,
652        });
653
654        Ok(())
655    }
656
657    /// Seek within the current track. Clamps to just before the end to avoid
658    /// accidentally skipping. Preserves pause state.
659    pub fn seek(&mut self, position_ms: u64) {
660        let info = match self.shared_state.track_info() {
661            Some(info) => info,
662            None => return,
663        };
664        let id = info.id;
665        let path = info.path.clone();
666        let duration = info.duration_ms;
667
668        // Clamp to just before the end so we don't skip to the next track.
669        let mut clamped = if duration > 0 {
670            position_ms.min(duration.saturating_sub(500))
671        } else {
672            position_ms
673        };
674
675        // Clamp to downloaded portion if streaming to prevent seeking into
676        // data that hasn't arrived yet.
677        if let Some(dl_frac) = self.shared_state.current_download_fraction() {
678            let max_ms = (dl_frac * duration as f64) as u64;
679            if max_ms > 5_000 {
680                clamped = clamped.min(max_ms - 5_000);
681            }
682        }
683
684        let was_paused = self.shared_state.playback_state() == PlaybackState::Paused;
685
686        if let Err(e) = self.start_playback(id, &path, clamped) {
687            log::error!("seek failed: {}", e);
688            return;
689        }
690
691        if was_paused {
692            self.pause();
693        }
694    }
695
696    /// Skip to next track in playlist.
697    pub fn next_track(&mut self) {
698        match self.shared_state.advance_cursor_loadable() {
699            Some(id) => self.play(id),
700            None => {
701                log::info!("no more tracks in playlist");
702                self.stop_playback_and_clear_state();
703            }
704        }
705    }
706
707    /// Go back to previous track.
708    pub fn prev_track(&mut self) {
709        match self.shared_state.retreat_cursor() {
710            Some((id, path)) => {
711                if matches!(path.try_exists(), Ok(true)) {
712                    if let Err(e) = self.start_playback(id, &path, 0) {
713                        log::error!("prev track failed: {}", e);
714                    }
715                } else {
716                    log::warn!("prev track path doesn't exist: {}", path.display());
717                }
718            }
719            None => {
720                // No previous track — restart current from the beginning.
721                if let Some(info) = self.shared_state.track_info()
722                    && let Err(e) = self.start_playback(info.id, &info.path, 0)
723                {
724                    log::error!("restart failed: {}", e);
725                }
726            }
727        }
728    }
729
730    /// Pause playback.
731    pub fn pause(&mut self) {
732        if let Some(ref playback) = self.active_playback {
733            if let Err(e) = playback.engine.stop() {
734                log::error!("pause failed: {}", e);
735                return;
736            }
737            self.shared_state.set_playback_state(PlaybackState::Paused);
738        }
739    }
740
741    /// Resume playback.
742    pub fn resume(&mut self) {
743        if let Some(ref playback) = self.active_playback {
744            if let Err(e) = playback.engine.start() {
745                log::error!("resume failed: {}", e);
746                return;
747            }
748            self.shared_state.set_playback_state(PlaybackState::Playing);
749        }
750    }
751
752    /// Stop playback and clear playlist.
753    pub fn stop(&mut self) {
754        self.shared_state.clear_playlist();
755        self.stop_playback_and_clear_state();
756    }
757
758    /// Stop the audio engine and decode thread without touching shared state.
759    ///
760    /// The engine is stopped synchronously (silence begins immediately), but
761    /// the heavy teardown (decode thread join + AudioUnit dispose) is moved to
762    /// a background thread so the player command loop never blocks — preventing
763    /// UI freezes when CoreAudio or the decode thread is slow to shut down.
764    fn stop_engine(&mut self) {
765        if let Some(playback) = self.active_playback.take() {
766            // Stop audio output immediately.
767            let _ = playback.engine.stop();
768            // Signal decode thread to exit (non-blocking).
769            playback.decode_handle.signal_stop();
770
771            // Drop the audio engine synchronously. AudioUnitUninitialize must
772            // complete before any device sample rate change, otherwise CoreAudio's
773            // internal ExtendedAudioBufferList can be freed mid-teardown (crash
774            // in caulk::alloc::tiered_allocator::deallocate — GitHub #89).
775            let ActivePlayback {
776                engine,
777                decode_handle,
778            } = playback;
779            drop(engine);
780
781            // The decode handle join is the heavy part (waits for the decode
782            // thread to exit) — run it on a background thread so we don't block
783            // the player command loop.
784            thread::Builder::new()
785                .name("koan-cleanup".into())
786                .spawn(move || drop(decode_handle))
787                .ok();
788        }
789    }
790
791    /// Full stop: tear down engine + clear all display state.
792    fn stop_playback_and_clear_state(&mut self) {
793        self.stop_engine();
794        self.timeline.reset();
795        self.shared_state.set_playback_state(PlaybackState::Stopped);
796        self.shared_state.set_position_ms(0);
797        self.shared_state.set_track_info(None);
798    }
799
800    /// Remove a track from the playlist. If it was the cursor, resume at the
801    /// track that followed it.
802    ///
803    /// `remove_item` clears the cursor, and an unset cursor means "start from the
804    /// top" — so the successor is pinned down by parking the cursor on the removed
805    /// track's predecessor first. `None` is correct only when it was the first item.
806    pub fn remove_from_playlist(&mut self, id: QueueItemId) {
807        let was_cursor = self.shared_state.is_cursor(id);
808        let resume_after = was_cursor
809            .then(|| self.shared_state.item_before(id))
810            .flatten();
811        self.shared_state.remove_item(id);
812        if was_cursor {
813            self.shared_state.set_cursor(resume_after);
814            self.next_track();
815        }
816    }
817
818    /// A download finished — if cursor is waiting on this item, start playback.
819    /// If already streaming this item, trigger progressive metadata enhancement.
820    pub fn track_ready(&mut self, id: QueueItemId) {
821        // Mark as Ready (download thread already did this, but be safe).
822        self.shared_state.update_load_state(id, LoadState::Ready);
823
824        if !self.shared_state.is_cursor(id) {
825            return;
826        }
827
828        let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
829        let current_track_id = self.shared_state.track_info().map(|t| t.id);
830
831        if is_playing && current_track_id == Some(id) {
832            // Already streaming this track — download just finished.
833            // Trigger progressive enhancement: re-read full lofty metadata and update state.
834            log::info!(
835                "track_ready: download complete while streaming {:?}, refreshing metadata",
836                id
837            );
838            self.refresh_track_metadata(id);
839            return;
840        }
841
842        // Cursor is on this item but not yet playing — start playback now.
843        if !is_playing && let Some(path) = self.shared_state.item_path_if_ready(id) {
844            log::info!("track_ready: starting playback for {:?}", id);
845            if let Err(e) = self.start_playback(id, &path, 0) {
846                log::error!("track_ready playback failed: {}", e);
847            }
848        }
849    }
850
851    /// Called when enough data has been buffered for streaming playback.
852    /// If the cursor is waiting on this track and nothing is playing, start streaming.
853    pub fn track_stream_ready(&mut self, id: QueueItemId) {
854        if !self.shared_state.is_cursor(id) {
855            return;
856        }
857
858        let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
859        if is_playing {
860            return; // Already playing something — don't interrupt.
861        }
862
863        match self.shared_state.item_playback_source(id) {
864            Some(PlaybackSource::Streaming {
865                path,
866                bytes_written,
867                total,
868            }) => {
869                log::info!(
870                    "track_stream_ready: starting streaming playback for {:?}",
871                    id
872                );
873                if let Err(e) = self.start_streaming_playback(id, &path, bytes_written, total) {
874                    log::error!("track_stream_ready streaming failed: {}", e);
875                }
876            }
877            Some(PlaybackSource::Ready(path)) => {
878                // Download finished between threshold and now — just play normally.
879                log::info!(
880                    "track_stream_ready: track already ready, starting normal playback for {:?}",
881                    id
882                );
883                if let Err(e) = self.start_playback(id, &path, 0) {
884                    log::error!("track_stream_ready playback failed: {}", e);
885                }
886            }
887            None => {} // Not enough data yet — wait.
888        }
889    }
890
891    /// Re-read full lofty metadata for a track after its download completes.
892    /// Updates the playlist item's tags and track_info with complete metadata.
893    /// Called from track_ready() when a streaming track finishes downloading.
894    fn refresh_track_metadata(&mut self, id: QueueItemId) {
895        use crate::index::metadata;
896
897        let path = match self.shared_state.item_path_if_ready(id) {
898            Some(p) => p,
899            None => return,
900        };
901
902        match metadata::read_metadata(&path) {
903            Ok(meta) => {
904                // Update playlist item with full lofty tags (title, artist, album, duration).
905                self.shared_state.update_item_metadata(
906                    id,
907                    meta.title,
908                    meta.artist,
909                    meta.album_artist.unwrap_or_default(),
910                    meta.album,
911                    meta.duration_ms.map(|d| d as u64),
912                );
913
914                // Re-probe the complete file for accurate duration + stream info.
915                // The initial probe was done on partial streaming data and may have
916                // underestimated duration, causing premature seek clamping or wrong
917                // progress bar display.
918                if let Ok(stream_info) = buffer::probe_file(&path)
919                    && let Some(current) = self.shared_state.track_info()
920                    && current.id == id
921                    && stream_info.duration_ms > current.duration_ms
922                {
923                    log::info!(
924                        "track_ready: duration corrected {}ms → {}ms",
925                        current.duration_ms,
926                        stream_info.duration_ms
927                    );
928                    self.shared_state.set_track_info(Some(TrackInfo {
929                        duration_ms: stream_info.duration_ms,
930                        ..current
931                    }));
932                }
933
934                // Signal UI to re-read cover art and update souvlaki media controls.
935                self.shared_state.signal_metadata_refresh();
936                log::info!("track_ready: metadata refreshed for {:?}", id);
937            }
938            Err(e) => {
939                log::warn!("track_ready: metadata refresh failed for {:?}: {}", id, e);
940            }
941        }
942    }
943
944    /// Poll the timeline and update shared state with current track/position.
945    /// Called from the command loop on each tick.
946    pub fn update_playback_state(&self) {
947        if self.active_playback.is_none() {
948            return;
949        }
950
951        if let Some((id, path, info, position_ms)) = self.timeline.current_playback() {
952            self.shared_state.set_position_ms(position_ms);
953
954            // Update track_info + cursor if the timeline shows a different track
955            // (gapless transition happened).
956            let current_id = self.shared_state.track_info().map(|t| t.id);
957            if current_id != Some(id) {
958                log::info!("timeline: now playing {:?}", id);
959                self.shared_state.set_track_info(Some(TrackInfo {
960                    id,
961                    path,
962                    codec: info.codec,
963                    sample_rate: info.sample_rate,
964                    bit_depth: info.bit_depth,
965                    bitrate_kbps: info.bitrate_kbps,
966                    channels: info.channels,
967                    duration_ms: info.duration_ms,
968                }));
969                self.shared_state.set_cursor(Some(id));
970            }
971        }
972    }
973
974    /// Decode thread naturally finished (playlist exhausted or error).
975    /// Advance to the next playable track; otherwise stop cleanly.
976    ///
977    /// A track that has not finished downloading parks the cursor on it, so its
978    /// `TrackReady`/`TrackStreamReady` resumes the queue instead of being
979    /// discarded as "not the cursor".
980    fn on_decode_finished(&mut self) {
981        log::info!("decode finished, checking for next track");
982        match self.shared_state.advance_cursor_loadable() {
983            Some(id) => self.play(id),
984            None => {
985                log::info!("no more tracks — stopping");
986                self.stop_playback_and_clear_state();
987            }
988        }
989    }
990
991    /// Snapshot items with their predecessors for an undo of "these were removed".
992    /// In playlist order, so undo re-inserts each item after a predecessor that
993    /// is already back in place.
994    fn snapshot_for_undo(
995        &self,
996        ids: &[QueueItemId],
997    ) -> Vec<(Box<state::PlaylistItem>, Option<QueueItemId>)> {
998        self.shared_state
999            .items_before(ids)
1000            .into_iter()
1001            .filter_map(|(id, after)| Some((Box::new(self.shared_state.get_item(id)?), after)))
1002            .collect()
1003    }
1004
1005    /// Route an undo entry to the batch buffer (if batching) or the undo stack.
1006    fn push_undo(&mut self, entry: UndoEntry) {
1007        if let Some(ref mut batch) = self.batch_buffer {
1008            batch.push(entry);
1009        } else {
1010            self.undo_stack.push(entry);
1011        }
1012    }
1013
1014    /// Process a single command.
1015    pub fn process_command(&mut self, cmd: PlayerCommand) {
1016        match cmd {
1017            PlayerCommand::Play(id) => self.play(id),
1018            PlayerCommand::Pause => self.pause(),
1019            PlayerCommand::Resume => self.resume(),
1020            PlayerCommand::Stop => self.stop(),
1021            PlayerCommand::Seek(pos) => self.seek(pos),
1022            PlayerCommand::NextTrack => {
1023                // Debounce: suppress key repeat from terminal (150ms window).
1024                let now = std::time::Instant::now();
1025                if now.duration_since(self.last_skip).as_millis() >= 150 {
1026                    self.last_skip = now;
1027                    self.next_track();
1028                }
1029            }
1030            PlayerCommand::PrevTrack => {
1031                let now = std::time::Instant::now();
1032                if now.duration_since(self.last_skip).as_millis() >= 150 {
1033                    self.last_skip = now;
1034                    self.prev_track();
1035                }
1036            }
1037            PlayerCommand::AddToPlaylist(items) => {
1038                let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1039                self.shared_state.add_items(items);
1040                self.push_undo(UndoEntry::Added { ids });
1041            }
1042            PlayerCommand::UpdatePaths(updates) => {
1043                self.shared_state.update_paths(&updates);
1044                if let Some(info) = self.shared_state.track_info()
1045                    && let Some((_, new_path)) = updates.iter().find(|(id, _)| *id == info.id)
1046                {
1047                    self.shared_state.set_track_info(Some(TrackInfo {
1048                        path: new_path.clone(),
1049                        ..info
1050                    }));
1051                }
1052            }
1053            PlayerCommand::InsertInPlaylist { items, after } => {
1054                let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1055                self.shared_state.insert_items_after(items, after);
1056                self.push_undo(UndoEntry::Inserted { ids });
1057            }
1058            PlayerCommand::ClearPlaylist => {
1059                // Stop engine + clear display state WITHOUT touching the playlist,
1060                // then snapshot, then clear. This avoids the race where stop()
1061                // would clear the playlist before we capture it for undo.
1062                self.stop_playback_and_clear_state();
1063                let (items, cursor) = self.shared_state.snapshot_playlist();
1064                self.shared_state.clear_playlist();
1065                self.push_undo(UndoEntry::Replaced { items, cursor });
1066            }
1067            PlayerCommand::RemoveFromPlaylist(id) => {
1068                let item = self.shared_state.get_item(id);
1069                let after = self.shared_state.item_before(id);
1070                self.remove_from_playlist(id);
1071                if let Some(item) = item {
1072                    self.push_undo(UndoEntry::Removed {
1073                        items: vec![(Box::new(item), after)],
1074                    });
1075                }
1076            }
1077            PlayerCommand::RemoveFromPlaylistBatch(ids) => {
1078                // Snapshot before removing anything, and resolve the resume point
1079                // once: removing one at a time would restart the engine for every
1080                // deleted track that the cursor lands on along the way.
1081                let items_with_pos = self.snapshot_for_undo(&ids);
1082                let resume_after = match self.shared_state.cursor() {
1083                    Some(cursor) if ids.contains(&cursor) => {
1084                        Some(self.shared_state.surviving_item_before(cursor, &ids))
1085                    }
1086                    _ => None,
1087                };
1088
1089                self.shared_state.remove_items(&ids);
1090
1091                if let Some(resume_after) = resume_after {
1092                    self.shared_state.set_cursor(resume_after);
1093                    self.next_track();
1094                }
1095
1096                if !items_with_pos.is_empty() {
1097                    self.push_undo(UndoEntry::Removed {
1098                        items: items_with_pos,
1099                    });
1100                }
1101            }
1102            PlayerCommand::MoveInPlaylist { id, target, after } => {
1103                let was_after = self.shared_state.item_before(id);
1104                self.shared_state.move_item(id, target, after);
1105                self.push_undo(UndoEntry::Moved { id, was_after });
1106            }
1107            PlayerCommand::MoveItemsInPlaylist { ids, target, after } => {
1108                let entries = self.shared_state.items_before(&ids);
1109                self.shared_state.move_items(&ids, target, after);
1110                self.push_undo(UndoEntry::MovedBatch { entries });
1111            }
1112            PlayerCommand::TrackReady(id) => self.track_ready(id),
1113            PlayerCommand::DecodeFinished => self.on_decode_finished(),
1114            PlayerCommand::TrackStreamReady(id) => self.track_stream_ready(id),
1115            PlayerCommand::Undo => self.execute_undo(),
1116            PlayerCommand::Redo => self.execute_redo(),
1117            PlayerCommand::BeginUndoBatch => {
1118                self.batch_buffer = Some(Vec::new());
1119            }
1120            PlayerCommand::EndUndoBatch => {
1121                if let Some(entries) = self.batch_buffer.take() {
1122                    if entries.len() == 1 {
1123                        // Single entry — push directly, no wrapping.
1124                        self.undo_stack.push(entries.into_iter().next().unwrap());
1125                    } else if !entries.is_empty() {
1126                        self.undo_stack.push(UndoEntry::Batch(entries));
1127                    }
1128                }
1129            }
1130            PlayerCommand::SetOutputDevice(name) => self.set_output_device(name),
1131            PlayerCommand::ClearOutputDevice => self.clear_output_device(),
1132        }
1133    }
1134
1135    /// Apply an undo/redo entry: mutate the playlist and return the inverse entry.
1136    fn apply_entry(&mut self, entry: UndoEntry) -> Option<UndoEntry> {
1137        match entry {
1138            UndoEntry::Added { ids } => {
1139                // Undo of "items were added": snapshot them with positions, then remove.
1140                let items_with_pos = self.snapshot_for_undo(&ids);
1141                self.shared_state.remove_items(&ids);
1142                Some(UndoEntry::Removed {
1143                    items: items_with_pos,
1144                })
1145            }
1146            UndoEntry::Removed { items } => {
1147                // Undo of "items were removed": re-insert each at its position.
1148                let mut ids = Vec::with_capacity(items.len());
1149                for (item, after) in items {
1150                    ids.push(item.id);
1151                    self.shared_state.insert_item_at(*item, after);
1152                }
1153                Some(UndoEntry::Added { ids })
1154            }
1155            UndoEntry::Inserted { ids } => {
1156                // Same as Added — snapshot positions, remove items.
1157                let items_with_pos = self.snapshot_for_undo(&ids);
1158                self.shared_state.remove_items(&ids);
1159                Some(UndoEntry::Removed {
1160                    items: items_with_pos,
1161                })
1162            }
1163            UndoEntry::Moved { id, was_after } => {
1164                let current_after = self.shared_state.item_before(id);
1165                self.shared_state.move_item_to(id, was_after);
1166                Some(UndoEntry::Moved {
1167                    id,
1168                    was_after: current_after,
1169                })
1170            }
1171            UndoEntry::MovedBatch { entries } => {
1172                let ids: Vec<QueueItemId> = entries.iter().map(|(id, _)| *id).collect();
1173                let current_positions = self.shared_state.items_before(&ids);
1174                self.shared_state.move_items_to(&entries);
1175                Some(UndoEntry::MovedBatch {
1176                    entries: current_positions,
1177                })
1178            }
1179            UndoEntry::Replaced { items, cursor } => {
1180                let (current_items, current_cursor) = self.shared_state.snapshot_playlist();
1181                self.shared_state.restore_playlist(items, cursor);
1182                Some(UndoEntry::Replaced {
1183                    items: current_items,
1184                    cursor: current_cursor,
1185                })
1186            }
1187            UndoEntry::Batch(entries) => {
1188                // Apply entries in reverse order, collect inverses.
1189                let mut inverses = Vec::with_capacity(entries.len());
1190                for entry in entries.into_iter().rev() {
1191                    if let Some(inverse) = self.apply_entry(entry) {
1192                        inverses.push(inverse);
1193                    }
1194                }
1195                inverses.reverse();
1196                Some(UndoEntry::Batch(inverses))
1197            }
1198        }
1199    }
1200
1201    /// Execute an undo operation, pushing the inverse onto the redo stack.
1202    fn execute_undo(&mut self) {
1203        let Some(entry) = self.undo_stack.pop_undo() else {
1204            return;
1205        };
1206        if let Some(inverse) = self.apply_entry(entry) {
1207            self.undo_stack.push_redo(inverse);
1208        }
1209    }
1210
1211    /// Execute a redo operation, pushing the inverse onto the undo stack.
1212    fn execute_redo(&mut self) {
1213        let Some(entry) = self.undo_stack.pop_redo() else {
1214            return;
1215        };
1216        if let Some(inverse) = self.apply_entry(entry) {
1217            self.undo_stack.push_undo_keep_redo(inverse);
1218        }
1219    }
1220
1221    /// Run the command loop. Blocks until the sender is dropped.
1222    pub fn run(&mut self) {
1223        use std::time::Duration;
1224
1225        let rx = self.commands.rx.clone();
1226        loop {
1227            // Poll with timeout so we update position even without commands.
1228            match rx.recv_timeout(Duration::from_millis(50)) {
1229                Ok(cmd) => self.process_command(cmd),
1230                Err(crossbeam_channel::RecvTimeoutError::Timeout) => {}
1231                Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break,
1232            }
1233            self.update_playback_state();
1234        }
1235        self.stop();
1236    }
1237
1238    /// Spawn the player on a background thread, returning the shared state,
1239    /// timeline, visualization snapshot, and command sender.
1240    pub fn spawn() -> (
1241        Arc<SharedPlayerState>,
1242        Arc<PlaybackTimeline>,
1243        Arc<VizSnapshot>,
1244        crossbeam_channel::Sender<PlayerCommand>,
1245    ) {
1246        let mut player = Self::new();
1247        let state = player.shared_state();
1248        let timeline = player.timeline();
1249        let viz_snapshot = player.viz_snapshot();
1250        let tx = player.command_sender();
1251
1252        thread::Builder::new()
1253            .name("koan-player".into())
1254            .spawn(move || player.run())
1255            .expect("failed to spawn player thread");
1256
1257        (state, timeline, viz_snapshot, tx)
1258    }
1259}
1260
1261#[cfg(test)]
1262mod tests {
1263    use super::*;
1264    use state::PlaylistItem;
1265    use std::path::PathBuf;
1266
1267    fn make_item(title: &str) -> PlaylistItem {
1268        PlaylistItem {
1269            id: QueueItemId::new(),
1270            db_id: None,
1271            path: PathBuf::from(format!("/music/{title}.flac")),
1272            title: title.to_string(),
1273            artist: String::new(),
1274            album_artist: String::new(),
1275            album: String::new(),
1276            year: None,
1277            codec: None,
1278            track_number: None,
1279            disc: None,
1280            duration_ms: None,
1281            load_state: LoadState::Ready,
1282        }
1283    }
1284
1285    fn playlist_ids(player: &Player) -> Vec<QueueItemId> {
1286        let (items, _) = player.shared_state.snapshot_playlist();
1287        items.iter().map(|i| i.id).collect()
1288    }
1289
1290    fn playlist_titles(player: &Player) -> Vec<String> {
1291        let (items, _) = player.shared_state.snapshot_playlist();
1292        items.iter().map(|i| i.title.clone()).collect()
1293    }
1294
1295    fn pending_item(title: &str) -> PlaylistItem {
1296        PlaylistItem {
1297            load_state: LoadState::Pending,
1298            ..make_item(title)
1299        }
1300    }
1301
1302    /// Build `n` ready items, add them, and return their IDs.
1303    fn seed(player: &mut Player, n: usize) -> Vec<QueueItemId> {
1304        let items: Vec<_> = (0..n).map(|i| make_item(&format!("t{i}"))).collect();
1305        let ids = items.iter().map(|i| i.id).collect();
1306        player.process_command(PlayerCommand::AddToPlaylist(items));
1307        ids
1308    }
1309
1310    // --- cursor transitions ---
1311
1312    #[test]
1313    fn removing_the_playing_track_resumes_at_its_successor() {
1314        let mut player = Player::new();
1315        let ids = seed(&mut player, 5);
1316        player.shared_state.set_cursor(Some(ids[2]));
1317
1318        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
1319
1320        assert_eq!(
1321            player.shared_state.cursor(),
1322            Some(ids[3]),
1323            "playback must continue at the next track, not restart the queue"
1324        );
1325        assert_eq!(player.playback_starts, 1);
1326    }
1327
1328    #[test]
1329    fn removing_the_first_playing_track_resumes_at_the_new_first() {
1330        let mut player = Player::new();
1331        let ids = seed(&mut player, 3);
1332        player.shared_state.set_cursor(Some(ids[0]));
1333
1334        player.process_command(PlayerCommand::RemoveFromPlaylist(ids[0]));
1335
1336        assert_eq!(player.shared_state.cursor(), Some(ids[1]));
1337    }
1338
1339    #[test]
1340    fn next_track_parks_on_a_track_that_has_not_downloaded_yet() {
1341        let mut player = Player::new();
1342        let playing = make_item("playing");
1343        let waiting = pending_item("waiting");
1344        let later = make_item("later");
1345        let (playing_id, waiting_id) = (playing.id, waiting.id);
1346        player.process_command(PlayerCommand::AddToPlaylist(vec![playing, waiting, later]));
1347        player.shared_state.set_cursor(Some(playing_id));
1348
1349        player.process_command(PlayerCommand::DecodeFinished);
1350
1351        assert_eq!(
1352            player.shared_state.cursor(),
1353            Some(waiting_id),
1354            "the cursor parks on the track being fetched"
1355        );
1356        assert_eq!(
1357            player.playback_starts, 0,
1358            "nothing to play until its bytes land"
1359        );
1360
1361        // The download completes. Because the cursor is parked here, the
1362        // TrackReady actually reaches the player and the queue resumes.
1363        player
1364            .shared_state
1365            .update_load_state(waiting_id, LoadState::Ready);
1366        player.process_command(PlayerCommand::TrackReady(waiting_id));
1367
1368        assert_eq!(player.playback_starts, 1);
1369        assert_eq!(player.shared_state.cursor(), Some(waiting_id));
1370    }
1371
1372    #[test]
1373    fn batch_delete_containing_the_cursor_restarts_the_engine_once() {
1374        let mut player = Player::new();
1375        let ids = seed(&mut player, 5);
1376        player.shared_state.set_cursor(Some(ids[2]));
1377
1378        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![
1379            ids[1], ids[2], ids[3],
1380        ]));
1381
1382        assert_eq!(playlist_titles(&player), vec!["t0", "t4"]);
1383        assert_eq!(player.shared_state.cursor(), Some(ids[4]));
1384        assert_eq!(
1385            player.playback_starts, 1,
1386            "one resume for the whole selection, not one per deleted track"
1387        );
1388    }
1389
1390    #[test]
1391    fn batch_delete_below_the_cursor_leaves_playback_alone() {
1392        let mut player = Player::new();
1393        let ids = seed(&mut player, 4);
1394        player.shared_state.set_cursor(Some(ids[0]));
1395
1396        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![ids[2], ids[3]]));
1397
1398        assert_eq!(player.shared_state.cursor(), Some(ids[0]));
1399        assert_eq!(player.playback_starts, 0);
1400    }
1401
1402    #[test]
1403    fn undo_of_a_batch_delete_restores_the_original_order() {
1404        // The TUI collects a selection from a HashSet, so the IDs arrive in
1405        // arbitrary order — scrambled here so a snapshot that trusts that order
1406        // re-inserts C before B and lands it at the end of the playlist.
1407        let mut player = Player::new();
1408        let items = vec![
1409            make_item("A"),
1410            make_item("B"),
1411            make_item("C"),
1412            make_item("D"),
1413        ];
1414        let (b_id, c_id) = (items[1].id, items[2].id);
1415        player.process_command(PlayerCommand::AddToPlaylist(items));
1416
1417        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![c_id, b_id]));
1418        assert_eq!(playlist_titles(&player), vec!["A", "D"]);
1419
1420        player.process_command(PlayerCommand::Undo);
1421        assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
1422    }
1423
1424    // --- AddToPlaylist undo/redo ---
1425
1426    #[test]
1427    fn undo_add_removes_items() {
1428        let mut player = Player::new();
1429        let items = vec![make_item("A"), make_item("B")];
1430        let ids: Vec<_> = items.iter().map(|i| i.id).collect();
1431
1432        player.process_command(PlayerCommand::AddToPlaylist(items));
1433        assert_eq!(playlist_ids(&player), ids);
1434        assert!(player.undo_stack().can_undo());
1435
1436        player.process_command(PlayerCommand::Undo);
1437        assert!(playlist_ids(&player).is_empty());
1438        assert!(player.undo_stack().can_redo());
1439    }
1440
1441    #[test]
1442    fn redo_add_restores_items() {
1443        let mut player = Player::new();
1444        let items = vec![make_item("A"), make_item("B")];
1445
1446        player.process_command(PlayerCommand::AddToPlaylist(items));
1447        player.process_command(PlayerCommand::Undo);
1448        assert!(playlist_ids(&player).is_empty());
1449
1450        player.process_command(PlayerCommand::Redo);
1451        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
1452    }
1453
1454    // --- RemoveFromPlaylist undo/redo ---
1455
1456    #[test]
1457    fn undo_remove_restores_item_at_position() {
1458        let mut player = Player::new();
1459        let items = vec![make_item("A"), make_item("B"), make_item("C")];
1460        let b_id = items[1].id;
1461
1462        player.process_command(PlayerCommand::AddToPlaylist(items));
1463        player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
1464        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
1465
1466        player.process_command(PlayerCommand::Undo);
1467        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
1468    }
1469
1470    #[test]
1471    fn undo_remove_first_item() {
1472        let mut player = Player::new();
1473        let items = vec![make_item("A"), make_item("B")];
1474        let a_id = items[0].id;
1475
1476        player.process_command(PlayerCommand::AddToPlaylist(items));
1477        player.process_command(PlayerCommand::RemoveFromPlaylist(a_id));
1478        assert_eq!(playlist_titles(&player), vec!["B"]);
1479
1480        player.process_command(PlayerCommand::Undo);
1481        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
1482    }
1483
1484    #[test]
1485    fn undo_batch_remove_restores_all() {
1486        let mut player = Player::new();
1487        let items = vec![
1488            make_item("A"),
1489            make_item("B"),
1490            make_item("C"),
1491            make_item("D"),
1492        ];
1493        let b_id = items[1].id;
1494        let c_id = items[2].id;
1495
1496        player.process_command(PlayerCommand::AddToPlaylist(items));
1497        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![b_id, c_id]));
1498        assert_eq!(playlist_titles(&player), vec!["A", "D"]);
1499
1500        // Single undo restores both
1501        player.process_command(PlayerCommand::Undo);
1502        assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
1503    }
1504
1505    #[test]
1506    fn redo_batch_remove() {
1507        let mut player = Player::new();
1508        let items = vec![make_item("A"), make_item("B"), make_item("C")];
1509        let a_id = items[0].id;
1510        let b_id = items[1].id;
1511
1512        player.process_command(PlayerCommand::AddToPlaylist(items));
1513        player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![a_id, b_id]));
1514        player.process_command(PlayerCommand::Undo);
1515        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
1516
1517        player.process_command(PlayerCommand::Redo);
1518        assert_eq!(playlist_titles(&player), vec!["C"]);
1519    }
1520
1521    #[test]
1522    fn redo_remove() {
1523        let mut player = Player::new();
1524        let items = vec![make_item("A"), make_item("B"), make_item("C")];
1525        let b_id = items[1].id;
1526
1527        player.process_command(PlayerCommand::AddToPlaylist(items));
1528        player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
1529        player.process_command(PlayerCommand::Undo);
1530        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
1531
1532        player.process_command(PlayerCommand::Redo);
1533        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
1534    }
1535
1536    // --- InsertInPlaylist undo/redo ---
1537
1538    #[test]
1539    fn undo_insert_removes_inserted_items() {
1540        let mut player = Player::new();
1541        let items = vec![make_item("A"), make_item("C")];
1542        let a_id = items[0].id;
1543
1544        player.process_command(PlayerCommand::AddToPlaylist(items));
1545
1546        let inserted = vec![make_item("B")];
1547        player.process_command(PlayerCommand::InsertInPlaylist {
1548            items: inserted,
1549            after: a_id,
1550        });
1551        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
1552
1553        player.process_command(PlayerCommand::Undo);
1554        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
1555    }
1556
1557    // --- MoveInPlaylist undo/redo ---
1558
1559    #[test]
1560    fn undo_move_restores_position() {
1561        let mut player = Player::new();
1562        let items = vec![make_item("A"), make_item("B"), make_item("C")];
1563        let a_id = items[0].id;
1564        let c_id = items[2].id;
1565
1566        player.process_command(PlayerCommand::AddToPlaylist(items));
1567
1568        // Move A after C: [B, C, A]
1569        player.process_command(PlayerCommand::MoveInPlaylist {
1570            id: a_id,
1571            target: c_id,
1572            after: true,
1573        });
1574        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
1575
1576        player.process_command(PlayerCommand::Undo);
1577        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
1578    }
1579
1580    #[test]
1581    fn redo_move() {
1582        let mut player = Player::new();
1583        let items = vec![make_item("A"), make_item("B"), make_item("C")];
1584        let a_id = items[0].id;
1585        let c_id = items[2].id;
1586
1587        player.process_command(PlayerCommand::AddToPlaylist(items));
1588        player.process_command(PlayerCommand::MoveInPlaylist {
1589            id: a_id,
1590            target: c_id,
1591            after: true,
1592        });
1593        player.process_command(PlayerCommand::Undo);
1594        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
1595
1596        player.process_command(PlayerCommand::Redo);
1597        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
1598    }
1599
1600    // --- MoveItemsInPlaylist (batch) undo/redo ---
1601
1602    #[test]
1603    fn undo_batch_move() {
1604        let mut player = Player::new();
1605        let items = vec![
1606            make_item("A"),
1607            make_item("B"),
1608            make_item("C"),
1609            make_item("D"),
1610        ];
1611        let a_id = items[0].id;
1612        let b_id = items[1].id;
1613        let d_id = items[3].id;
1614
1615        player.process_command(PlayerCommand::AddToPlaylist(items));
1616
1617        // Move A,B after D: [C, D, A, B]
1618        player.process_command(PlayerCommand::MoveItemsInPlaylist {
1619            ids: vec![a_id, b_id],
1620            target: d_id,
1621            after: true,
1622        });
1623        assert_eq!(playlist_titles(&player), vec!["C", "D", "A", "B"]);
1624
1625        player.process_command(PlayerCommand::Undo);
1626        assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
1627    }
1628
1629    // --- ClearPlaylist undo/redo ---
1630
1631    #[test]
1632    fn undo_clear_restores_playlist() {
1633        let mut player = Player::new();
1634        let items = vec![make_item("A"), make_item("B"), make_item("C")];
1635
1636        player.process_command(PlayerCommand::AddToPlaylist(items));
1637        player.process_command(PlayerCommand::ClearPlaylist);
1638        assert!(playlist_ids(&player).is_empty());
1639
1640        player.process_command(PlayerCommand::Undo);
1641        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
1642    }
1643
1644    #[test]
1645    fn redo_clear() {
1646        let mut player = Player::new();
1647        let items = vec![make_item("A"), make_item("B")];
1648
1649        player.process_command(PlayerCommand::AddToPlaylist(items));
1650        player.process_command(PlayerCommand::ClearPlaylist);
1651        player.process_command(PlayerCommand::Undo);
1652        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
1653
1654        player.process_command(PlayerCommand::Redo);
1655        assert!(playlist_ids(&player).is_empty());
1656    }
1657
1658    // --- Multi-step undo/redo ---
1659
1660    #[test]
1661    fn multiple_undos_in_sequence() {
1662        let mut player = Player::new();
1663
1664        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("A")]));
1665        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
1666        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("C")]));
1667        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
1668
1669        player.process_command(PlayerCommand::Undo);
1670        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
1671
1672        player.process_command(PlayerCommand::Undo);
1673        assert_eq!(playlist_titles(&player), vec!["A"]);
1674
1675        player.process_command(PlayerCommand::Undo);
1676        assert!(playlist_ids(&player).is_empty());
1677    }
1678
1679    #[test]
1680    fn undo_redo_undo_cycle() {
1681        let mut player = Player::new();
1682        let items = vec![make_item("A"), make_item("B")];
1683
1684        player.process_command(PlayerCommand::AddToPlaylist(items));
1685        player.process_command(PlayerCommand::Undo);
1686        assert!(playlist_ids(&player).is_empty());
1687
1688        player.process_command(PlayerCommand::Redo);
1689        assert_eq!(playlist_titles(&player), vec!["A", "B"]);
1690
1691        player.process_command(PlayerCommand::Undo);
1692        assert!(playlist_ids(&player).is_empty());
1693    }
1694
1695    #[test]
1696    fn new_action_clears_redo_stack() {
1697        let mut player = Player::new();
1698        let items = vec![make_item("A")];
1699
1700        player.process_command(PlayerCommand::AddToPlaylist(items));
1701        player.process_command(PlayerCommand::Undo);
1702        assert!(player.undo_stack().can_redo());
1703
1704        // New action should clear redo
1705        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
1706        assert!(!player.undo_stack().can_redo());
1707    }
1708
1709    #[test]
1710    fn undo_on_empty_stack_is_noop() {
1711        let mut player = Player::new();
1712        player.process_command(PlayerCommand::Undo);
1713        assert!(playlist_ids(&player).is_empty());
1714    }
1715
1716    #[test]
1717    fn redo_on_empty_stack_is_noop() {
1718        let mut player = Player::new();
1719        player.process_command(PlayerCommand::Redo);
1720        assert!(playlist_ids(&player).is_empty());
1721    }
1722
1723    // --- Non-undoable commands don't push entries ---
1724
1725    #[test]
1726    fn playback_commands_not_undoable() {
1727        let mut player = Player::new();
1728        player.process_command(PlayerCommand::Pause);
1729        player.process_command(PlayerCommand::Resume);
1730        player.process_command(PlayerCommand::NextTrack);
1731        player.process_command(PlayerCommand::PrevTrack);
1732        assert!(!player.undo_stack().can_undo());
1733    }
1734
1735    #[test]
1736    fn update_paths_not_undoable() {
1737        let mut player = Player::new();
1738        let items = vec![make_item("A")];
1739        let id = items[0].id;
1740        player.process_command(PlayerCommand::AddToPlaylist(items));
1741
1742        let undo_count = player.undo_stack().undo_len();
1743        player.process_command(PlayerCommand::UpdatePaths(vec![(
1744            id,
1745            PathBuf::from("/new/path.flac"),
1746        )]));
1747        assert_eq!(player.undo_stack().undo_len(), undo_count);
1748    }
1749
1750    // --- Complex scenarios ---
1751
1752    #[test]
1753    fn add_remove_undo_undo_produces_original() {
1754        let mut player = Player::new();
1755        let items = vec![make_item("A"), make_item("B"), make_item("C")];
1756        let b_id = items[1].id;
1757        let original_titles = vec!["A", "B", "C"];
1758
1759        player.process_command(PlayerCommand::AddToPlaylist(items));
1760        player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
1761        assert_eq!(playlist_titles(&player), vec!["A", "C"]);
1762
1763        // Undo remove → back to A, B, C
1764        player.process_command(PlayerCommand::Undo);
1765        assert_eq!(playlist_titles(&player), original_titles);
1766
1767        // Undo add → empty
1768        player.process_command(PlayerCommand::Undo);
1769        assert!(playlist_ids(&player).is_empty());
1770    }
1771
1772    #[test]
1773    fn interleaved_adds_and_moves_undo() {
1774        let mut player = Player::new();
1775        let items = vec![make_item("A"), make_item("B"), make_item("C")];
1776        let a_id = items[0].id;
1777        let c_id = items[2].id;
1778
1779        player.process_command(PlayerCommand::AddToPlaylist(items));
1780
1781        // Move A after C: [B, C, A]
1782        player.process_command(PlayerCommand::MoveInPlaylist {
1783            id: a_id,
1784            target: c_id,
1785            after: true,
1786        });
1787        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
1788
1789        // Add D: [B, C, A, D]
1790        player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("D")]));
1791        assert_eq!(playlist_titles(&player), vec!["B", "C", "A", "D"]);
1792
1793        // Undo add D: [B, C, A]
1794        player.process_command(PlayerCommand::Undo);
1795        assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
1796
1797        // Undo move: [A, B, C]
1798        player.process_command(PlayerCommand::Undo);
1799        assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
1800    }
1801
1802    /// Regression test for GitHub #89: AudioEngine must be dropped synchronously
1803    /// in stop_engine() before the caller changes sample rates. If the engine is
1804    /// dropped on a background thread, CoreAudio's internal buffer list can be
1805    /// freed while AudioUnitUninitialize is still tearing it down → crash.
1806    #[test]
1807    fn stop_engine_drops_engine_synchronously() {
1808        use std::sync::atomic::AtomicBool;
1809
1810        struct MockEngine {
1811            dropped: Arc<AtomicBool>,
1812        }
1813        impl AudioEngineHandle for MockEngine {
1814            fn start(&self) -> Result<(), BackendError> {
1815                Ok(())
1816            }
1817            fn stop(&self) -> Result<(), BackendError> {
1818                Ok(())
1819            }
1820            fn is_running(&self) -> bool {
1821                false
1822            }
1823        }
1824        impl Drop for MockEngine {
1825            fn drop(&mut self) {
1826                self.dropped.store(true, Ordering::SeqCst);
1827            }
1828        }
1829
1830        let dropped = Arc::new(AtomicBool::new(false));
1831
1832        // Build a minimal decode handle that won't block.
1833        let stop_flag = Arc::new(AtomicBool::new(false));
1834        let decode_handle = buffer::DecodeHandle::new_for_test(stop_flag);
1835
1836        let mut player = Player::new();
1837        player.active_playback = Some(ActivePlayback {
1838            engine: Box::new(MockEngine {
1839                dropped: dropped.clone(),
1840            }),
1841            decode_handle,
1842        });
1843
1844        player.stop_engine();
1845
1846        // The engine must already be dropped when stop_engine returns.
1847        // If this fails, the engine was moved to a background thread — the
1848        // exact race condition that causes the #89 crash.
1849        assert!(
1850            dropped.load(Ordering::SeqCst),
1851            "AudioEngine must be dropped synchronously in stop_engine (GitHub #89)"
1852        );
1853    }
1854
1855    // --- Engine format matches the decoded PCM ---
1856
1857    /// Backend pinned to one sample rate that refuses every switch, recording
1858    /// the format the engine is asked for.
1859    struct StuckBackend {
1860        rate: f64,
1861        asked: Arc<std::sync::Mutex<Option<(f64, u32)>>>,
1862    }
1863
1864    struct NullEngine;
1865    impl AudioEngineHandle for NullEngine {
1866        fn start(&self) -> Result<(), BackendError> {
1867            Ok(())
1868        }
1869        fn stop(&self) -> Result<(), BackendError> {
1870            Ok(())
1871        }
1872        fn is_running(&self) -> bool {
1873            false
1874        }
1875    }
1876
1877    impl AudioBackend for StuckBackend {
1878        fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
1879            Ok(vec![self.default_device()?])
1880        }
1881        fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
1882            Ok(backend::DeviceInfo {
1883                name: "Stuck DAC".into(),
1884                sample_rates: vec![self.rate],
1885                platform_id: 0,
1886            })
1887        }
1888        fn supported_sample_rates(
1889            &self,
1890            _device: &backend::DeviceInfo,
1891        ) -> Result<Vec<f64>, BackendError> {
1892            Ok(vec![self.rate])
1893        }
1894        fn get_device_sample_rate(
1895            &self,
1896            _device: &backend::DeviceInfo,
1897        ) -> Result<f64, BackendError> {
1898            Ok(self.rate)
1899        }
1900        fn set_device_sample_rate(
1901            &self,
1902            _device: &backend::DeviceInfo,
1903            rate: f64,
1904        ) -> Result<f64, BackendError> {
1905            Err(BackendError::UnsupportedSampleRate(rate))
1906        }
1907        fn create_engine(
1908            &self,
1909            _device: &backend::DeviceInfo,
1910            sample_rate: f64,
1911            channels: u32,
1912            _consumer: rtrb::Consumer<f32>,
1913            _samples_played: Arc<AtomicU64>,
1914        ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
1915            *self.asked.lock().unwrap() = Some((sample_rate, channels));
1916            Ok(Box::new(NullEngine))
1917        }
1918    }
1919
1920    fn engine_format_for(source_rate: u32, channels: u16, device_rate: f64) -> (f64, u32) {
1921        let asked = Arc::new(std::sync::Mutex::new(None));
1922        let mut player = Player::new();
1923        player.backend = Box::new(StuckBackend {
1924            rate: device_rate,
1925            asked: asked.clone(),
1926        });
1927
1928        let info = buffer::StreamInfo {
1929            codec: "MP3".into(),
1930            sample_rate: source_rate,
1931            channels,
1932            bit_depth: Some(16),
1933            bitrate_kbps: None,
1934            duration_ms: 1000,
1935        };
1936        let (_producer, consumer) = rtrb::RingBuffer::new(16);
1937        player
1938            .create_engine_for(&info, consumer)
1939            .expect("engine creation should succeed");
1940        let asked = *asked.lock().unwrap();
1941        asked.expect("engine was never created")
1942    }
1943
1944    #[test]
1945    fn engine_uses_source_rate_when_device_refuses_switch() {
1946        // MPEG-2 MP3 rates are routinely rejected by output devices. The engine
1947        // must still be told the rate the PCM actually is.
1948        assert_eq!(engine_format_for(22050, 2, 48000.0), (22050.0, 2));
1949        assert_eq!(engine_format_for(32000, 2, 44100.0), (32000.0, 2));
1950    }
1951
1952    #[test]
1953    fn engine_uses_source_channel_count() {
1954        assert_eq!(engine_format_for(44100, 1, 44100.0), (44100.0, 1));
1955    }
1956}